From 4ae5c121fca584e5a0345f5ad5445e8a87150f37 Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 17:29:09 +0100 Subject: [PATCH 1/8] add @temporalio/client --- package-lock.json | 209 +++++++++++++++++++++++++++++++++++++++++++++- package.json | 1 + 2 files changed, 209 insertions(+), 1 deletion(-) diff --git a/package-lock.json b/package-lock.json index c3209dd..0eb77d1 100644 --- a/package-lock.json +++ b/package-lock.json @@ -13,6 +13,7 @@ "@discordjs/rest": "^1.7.0", "@notionhq/client": "^2.2.3", "@sentry/node": "^7.50.0", + "@temporalio/client": "^1.11.3", "@togethercrew.dev/db": "^3.0.72", "@togethercrew.dev/tc-messagebroker": "^0.0.50", "@types/express-session": "^1.17.7", @@ -1611,6 +1612,65 @@ "node": ">=14" } }, + "node_modules/@grpc/grpc-js": { + "version": "1.12.2", + "resolved": "https://registry.npmjs.org/@grpc/grpc-js/-/grpc-js-1.12.2.tgz", + "integrity": "sha512-bgxdZmgTrJZX50OjyVwz3+mNEnCTNkh3cIqGPWVNeW9jX6bn1ZkU80uPd+67/ZpIJIjRQ9qaHCjhavyoWYxumg==", + "dependencies": { + "@grpc/proto-loader": "^0.7.13", + "@js-sdsl/ordered-map": "^4.4.2" + }, + "engines": { + "node": ">=12.10.0" + } + }, + "node_modules/@grpc/proto-loader": { + "version": "0.7.13", + "resolved": "https://registry.npmjs.org/@grpc/proto-loader/-/proto-loader-0.7.13.tgz", + "integrity": "sha512-AiXO/bfe9bmxBjxxtYxFAXGZvMaN5s8kO+jBHAJCON8rJoB5YS/D6X7ZNc6XQkuHNmyl4CYaMI1fJ/Gn27RGGw==", + "dependencies": { + "lodash.camelcase": "^4.3.0", + "long": "^5.0.0", + "protobufjs": "^7.2.5", + "yargs": "^17.7.2" + }, + "bin": { + "proto-loader-gen-types": "build/bin/proto-loader-gen-types.js" + }, + "engines": { + "node": ">=6" + } + }, + "node_modules/@grpc/proto-loader/node_modules/cliui": { + "version": "8.0.1", + "resolved": "https://registry.npmjs.org/cliui/-/cliui-8.0.1.tgz", + "integrity": "sha512-BSeNnyus75C4//NQ9gQt1/csTXyo/8Sb+afLAkzAptFuMsod9HFokGNudZpi/oQV73hnVK+sR+5PVRMd+Dr7YQ==", + "dependencies": { + "string-width": "^4.2.0", + "strip-ansi": "^6.0.1", + "wrap-ansi": "^7.0.0" + }, + "engines": { + "node": ">=12" + } + }, + "node_modules/@grpc/proto-loader/node_modules/yargs": { + "version": "17.7.2", + "resolved": "https://registry.npmjs.org/yargs/-/yargs-17.7.2.tgz", + "integrity": "sha512-7dSzzRQ++CKnNI/krKnYRV7JKKPUXMEh61soaHKg9mrWEhzFWhFnxPxGl+69cD1Ou63C13NUPCnmIcrvqCuM6w==", + "dependencies": { + "cliui": "^8.0.1", + "escalade": "^3.1.1", + "get-caller-file": "^2.0.5", + "require-directory": "^2.1.1", + "string-width": "^4.2.3", + "y18n": "^5.0.5", + "yargs-parser": "^21.1.1" + }, + "engines": { + "node": ">=12" + } + }, "node_modules/@hapi/hoek": { "version": "9.3.0", "resolved": "https://registry.npmjs.org/@hapi/hoek/-/hoek-9.3.0.tgz", @@ -2351,6 +2411,15 @@ "resolved": "https://registry.npmjs.org/@jridgewell/sourcemap-codec/-/sourcemap-codec-1.4.14.tgz", "integrity": "sha512-XPSJHWmi394fuUuzDnGz1wiKqWfo1yXecHQMRf2l6hztTO+nPru658AyDngaBe7isIxEkRsPR3FZh+s7iVa4Uw==" }, + "node_modules/@js-sdsl/ordered-map": { + "version": "4.4.2", + "resolved": "https://registry.npmjs.org/@js-sdsl/ordered-map/-/ordered-map-4.4.2.tgz", + "integrity": "sha512-iUKgm52T8HOE/makSxjqoWhe95ZJA1/G1sYsGev2JDKUSS14KAgg1LHb+Ba+IPow0xflbnSkOsZcO08C7w1gYw==", + "funding": { + "type": "opencollective", + "url": "https://opencollective.com/js-sdsl" + } + }, "node_modules/@jsdevtools/ono": { "version": "7.1.3", "resolved": "https://registry.npmjs.org/@jsdevtools/ono/-/ono-7.1.3.tgz", @@ -2424,6 +2493,60 @@ "node": ">=12" } }, + "node_modules/@protobufjs/aspromise": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/@protobufjs/aspromise/-/aspromise-1.1.2.tgz", + "integrity": "sha512-j+gKExEuLmKwvz3OgROXtrJ2UG2x8Ch2YZUxahh+s1F2HZ+wAceUNLkvy6zKCPVRkU++ZWQrdxsUeQXmcg4uoQ==" + }, + "node_modules/@protobufjs/base64": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/@protobufjs/base64/-/base64-1.1.2.tgz", + "integrity": "sha512-AZkcAA5vnN/v4PDqKyMR5lx7hZttPDgClv83E//FMNhR2TMcLUhfRUBHCmSl0oi9zMgDDqRUJkSxO3wm85+XLg==" + }, + "node_modules/@protobufjs/codegen": { + "version": "2.0.4", + "resolved": "https://registry.npmjs.org/@protobufjs/codegen/-/codegen-2.0.4.tgz", + "integrity": "sha512-YyFaikqM5sH0ziFZCN3xDC7zeGaB/d0IUb9CATugHWbd1FRFwWwt4ld4OYMPWu5a3Xe01mGAULCdqhMlPl29Jg==" + }, + "node_modules/@protobufjs/eventemitter": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/@protobufjs/eventemitter/-/eventemitter-1.1.0.tgz", + "integrity": "sha512-j9ednRT81vYJ9OfVuXG6ERSTdEL1xVsNgqpkxMsbIabzSo3goCjDIveeGv5d03om39ML71RdmrGNjG5SReBP/Q==" + }, + "node_modules/@protobufjs/fetch": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/@protobufjs/fetch/-/fetch-1.1.0.tgz", + "integrity": "sha512-lljVXpqXebpsijW71PZaCYeIcE5on1w5DlQy5WH6GLbFryLUrBD4932W/E2BSpfRJWseIL4v/KPgBFxDOIdKpQ==", + "dependencies": { + "@protobufjs/aspromise": "^1.1.1", + "@protobufjs/inquire": "^1.1.0" + } + }, + "node_modules/@protobufjs/float": { + "version": "1.0.2", + "resolved": "https://registry.npmjs.org/@protobufjs/float/-/float-1.0.2.tgz", + "integrity": "sha512-Ddb+kVXlXst9d+R9PfTIxh1EdNkgoRe5tOX6t01f1lYWOvJnSPDBlG241QLzcyPdoNTsblLUdujGSE4RzrTZGQ==" + }, + "node_modules/@protobufjs/inquire": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/@protobufjs/inquire/-/inquire-1.1.0.tgz", + "integrity": "sha512-kdSefcPdruJiFMVSbn801t4vFK7KB/5gd2fYvrxhuJYg8ILrmn9SKSX2tZdV6V+ksulWqS7aXjBcRXl3wHoD9Q==" + }, + "node_modules/@protobufjs/path": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/@protobufjs/path/-/path-1.1.2.tgz", + "integrity": "sha512-6JOcJ5Tm08dOHAbdR3GrvP+yUUfkjG5ePsHYczMFLq3ZmMkAD98cDgcT2iA1lJ9NVwFd4tH/iSSoe44YWkltEA==" + }, + "node_modules/@protobufjs/pool": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/@protobufjs/pool/-/pool-1.1.0.tgz", + "integrity": "sha512-0kELaGSIDBKvcgS4zkjz1PeddatrjYcmMWOlAuAPwAeccUrPHdUqo/J6LiymHHEiJT5NrF1UVwxY14f+fy4WQw==" + }, + "node_modules/@protobufjs/utf8": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/@protobufjs/utf8/-/utf8-1.1.0.tgz", + "integrity": "sha512-Vvn3zZrhQZkkBE8LSuW3em98c0FwgO4nxzv6OdSxPKJIEKY2bGbHn+mhGIPerzI4twdxaP8/0+06HBpwf345Lw==" + }, "node_modules/@sapphire/async-queue": { "version": "1.5.0", "resolved": "https://registry.npmjs.org/@sapphire/async-queue/-/async-queue-1.5.0.tgz", @@ -3165,6 +3288,47 @@ "node": ">=16.0.0" } }, + "node_modules/@temporalio/client": { + "version": "1.11.3", + "resolved": "https://registry.npmjs.org/@temporalio/client/-/client-1.11.3.tgz", + "integrity": "sha512-2x30xAXbUuqelrWe3Vd1FVC0+Z2Cfh6m2W5yDUZBjqTMdNP6qd8nH4S4mceRtZ4TipYSPmaONaiWoAU2VvwEIg==", + "dependencies": { + "@grpc/grpc-js": "^1.10.7", + "@temporalio/common": "1.11.3", + "@temporalio/proto": "1.11.3", + "abort-controller": "^3.0.0", + "long": "^5.2.3", + "uuid": "^9.0.1" + } + }, + "node_modules/@temporalio/common": { + "version": "1.11.3", + "resolved": "https://registry.npmjs.org/@temporalio/common/-/common-1.11.3.tgz", + "integrity": "sha512-dzCrwiE9ox/Q16AjBsUKr4djg1ovYHNCjH36WZadwsemXINRWa5eW53j0WZOlmFF8/CbcHIhiD5N18rZqjiYkg==", + "dependencies": { + "@temporalio/proto": "1.11.3", + "long": "^5.2.3", + "ms": "^3.0.0-canary.1", + "proto3-json-serializer": "^2.0.0" + } + }, + "node_modules/@temporalio/common/node_modules/ms": { + "version": "3.0.0-canary.1", + "resolved": "https://registry.npmjs.org/ms/-/ms-3.0.0-canary.1.tgz", + "integrity": "sha512-kh8ARjh8rMN7Du2igDRO9QJnqCb2xYTJxyQYK7vJJS4TvLLmsbyhiKpSW+t+y26gyOyMd0riphX0GeWKU3ky5g==", + "engines": { + "node": ">=12.13" + } + }, + "node_modules/@temporalio/proto": { + "version": "1.11.3", + "resolved": "https://registry.npmjs.org/@temporalio/proto/-/proto-1.11.3.tgz", + "integrity": "sha512-X+xV75m11BvvS5MljagtYCybRNxpLNJM24eWyfv+uyU4WZSvgQCUh21fY4FOUDJS66DPvO1mefSPu0Nunp1PHg==", + "dependencies": { + "long": "^5.2.3", + "protobufjs": "^7.2.5" + } + }, "node_modules/@textlint/ast-node-types": { "version": "14.0.3", "resolved": "https://registry.npmjs.org/@textlint/ast-node-types/-/ast-node-types-14.0.3.tgz", @@ -9527,6 +9691,11 @@ "resolved": "https://registry.npmjs.org/lodash/-/lodash-4.17.21.tgz", "integrity": "sha512-v2kDEe57lecTulaDIuNTPy3Ry4gLGJ6Z1O3vE1krgXZNrsQ+LFTGHVxVjcXPs17LhbZVGedAJv8XZ1tvj5FvSg==" }, + "node_modules/lodash.camelcase": { + "version": "4.3.0", + "resolved": "https://registry.npmjs.org/lodash.camelcase/-/lodash.camelcase-4.3.0.tgz", + "integrity": "sha512-TwuEnCnxbc3rAvhf/LbG7tJUDzhqXyFnv3dtzLOPgCG/hODL7WFnsbwktkD7yUV0RrreP/l1PALq/YSg6VvjlA==" + }, "node_modules/lodash.defaults": { "version": "4.2.0", "resolved": "https://registry.npmjs.org/lodash.defaults/-/lodash.defaults-4.2.0.tgz", @@ -9579,6 +9748,11 @@ "integrity": "sha512-jttmRe7bRse52OsWIMDLaXxWqRAmtIUccAQ3garviCqJjafXOfNMO0yMfNpdD6zbGaTU0P5Nz7e7gAT6cKmJRw==", "dev": true }, + "node_modules/long": { + "version": "5.2.3", + "resolved": "https://registry.npmjs.org/long/-/long-5.2.3.tgz", + "integrity": "sha512-lcHwpNoggQTObv5apGNCTdJrO69eHOZMi4BNC+rTLER8iHAqGrUVeLh/irVIM7zTw2bOXA8T6uNPeujwOLg/2Q==" + }, "node_modules/longest-streak": { "version": "2.0.4", "resolved": "https://registry.npmjs.org/longest-streak/-/longest-streak-2.0.4.tgz", @@ -11239,6 +11413,40 @@ "node": ">= 6" } }, + "node_modules/proto3-json-serializer": { + "version": "2.0.2", + "resolved": "https://registry.npmjs.org/proto3-json-serializer/-/proto3-json-serializer-2.0.2.tgz", + "integrity": "sha512-SAzp/O4Yh02jGdRc+uIrGoe87dkN/XtwxfZ4ZyafJHymd79ozp5VG5nyZ7ygqPM5+cpLDjjGnYFUkngonyDPOQ==", + "dependencies": { + "protobufjs": "^7.2.5" + }, + "engines": { + "node": ">=14.0.0" + } + }, + "node_modules/protobufjs": { + "version": "7.4.0", + "resolved": "https://registry.npmjs.org/protobufjs/-/protobufjs-7.4.0.tgz", + "integrity": "sha512-mRUWCc3KUU4w1jU8sGxICXH/gNS94DvI1gxqDvBzhj1JpcsimQkYiOJfwsPUykUI5ZaspFbSgmBLER8IrQ3tqw==", + "hasInstallScript": true, + "dependencies": { + "@protobufjs/aspromise": "^1.1.2", + "@protobufjs/base64": "^1.1.2", + "@protobufjs/codegen": "^2.0.4", + "@protobufjs/eventemitter": "^1.1.0", + "@protobufjs/fetch": "^1.1.0", + "@protobufjs/float": "^1.0.2", + "@protobufjs/inquire": "^1.1.0", + "@protobufjs/path": "^1.1.2", + "@protobufjs/pool": "^1.1.0", + "@protobufjs/utf8": "^1.1.0", + "@types/node": ">=13.7.0", + "long": "^5.0.0" + }, + "engines": { + "node": ">=12.0.0" + } + }, "node_modules/proxy-addr": { "version": "2.0.7", "resolved": "https://registry.npmjs.org/proxy-addr/-/proxy-addr-2.0.7.tgz", @@ -13632,7 +13840,6 @@ "version": "21.1.1", "resolved": "https://registry.npmjs.org/yargs-parser/-/yargs-parser-21.1.1.tgz", "integrity": "sha512-tVpsJW7DdjecAiFpbIB1e3qxIQsE6NoPc5/eTdrbbIC4h0LVsWhnoa3g+m2HclBIujHzsxZ4VJVA+GUuc2/LBw==", - "dev": true, "engines": { "node": ">=12" } diff --git a/package.json b/package.json index 70a4471..d8bea98 100644 --- a/package.json +++ b/package.json @@ -27,6 +27,7 @@ "@discordjs/rest": "^1.7.0", "@notionhq/client": "^2.2.3", "@sentry/node": "^7.50.0", + "@temporalio/client": "^1.11.3", "@togethercrew.dev/db": "^3.0.72", "@togethercrew.dev/tc-messagebroker": "^0.0.50", "@types/express-session": "^1.17.7", From 36856cb8a3f0b23e358a4c11486ae390d1caef33 Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 18:23:36 +0100 Subject: [PATCH 2/8] replace discourse extract with temporal --- src/config/index.ts | 8 ++++-- src/services/discourse/core.service.ts | 32 ++++++---------------- src/services/platform.service.ts | 6 ++-- src/services/temporal/core.service.ts | 20 ++++++++++++++ src/services/temporal/discourse.service.ts | 29 ++++++++++++++++++++ 5 files changed, 67 insertions(+), 28 deletions(-) create mode 100644 src/services/temporal/core.service.ts create mode 100644 src/services/temporal/discourse.service.ts diff --git a/src/config/index.ts b/src/config/index.ts index 11db304..930c1cc 100644 --- a/src/config/index.ts +++ b/src/config/index.ts @@ -55,8 +55,9 @@ const envVarsSchema = Joi.object() REDIS_HOST: Joi.string().required().description('Redis host'), REDIS_PORT: Joi.string().required().description('Redis port'), REDIS_PASSWORD: Joi.string().required().description('Reids password').allow(''), - DISCOURSE_EXTRACTION_URL: Joi.string().required().description('Discourse extraction url'), OCI_BACKEND_URL: Joi.string().required().description('Oci Backend url'), + TEMPORAL_URI: Joi.string().required().description('Temporal address'), + TEMPORAL_QUEUE_HEAVY: Joi.string().required().description('Queue for heavy workflows'), }) .unknown(); @@ -156,8 +157,9 @@ export default { session: { secret: envVars.SESSION_SECRET, }, - discourse: { - extractionURL: envVars.DISCOURSE_EXTRACTION_URL, + temporal: { + uri: envVars.TEMPORAL_URI, + heavyQueue: envVars.TEMPORAL_QUEUE_HEAVY, }, ociBackendURL: envVars.OCI_BACKEND_URL, }; diff --git a/src/services/discourse/core.service.ts b/src/services/discourse/core.service.ts index 3d9c905..3dfdea8 100644 --- a/src/services/discourse/core.service.ts +++ b/src/services/discourse/core.service.ts @@ -1,9 +1,8 @@ import parentLogger from '../../config/logger'; import { IAuthAndPlatform } from '../../interfaces'; import categoryService from './category.service'; -import { ApiError, pick, sort } from '../../utils'; -import { Types } from 'mongoose'; -import config from '../../config'; +import { ApiError, pick } from '../../utils'; +import temporalDiscourse from '../temporal/discourse.service' const logger = parentLogger.child({ module: 'DiscourseCoreService' }); async function getPropertyHandler(req: IAuthAndPlatform) { @@ -19,31 +18,18 @@ async function getPropertyHandler(req: IAuthAndPlatform) { * @param {String} platformId * @returns {Promise} */ -async function runDiscourseExtraction(platformId: string): Promise { +async function createDiscourseSchedule(platformId: string, endpoint: string): Promise { try { - const data = { - platform_id: platformId, - }; - logger.debug(data); - const response = await fetch(config.discourse.extractionURL, { - method: 'POST', - body: JSON.stringify(data), - headers: { 'Content-Type': 'application/json' }, - }); - if (response.ok) { - logger.debug(await response.json()); - return; - } else { - const errorResponse = await response.text(); - logger.error({ error: errorResponse }); - } + const schedule = await temporalDiscourse.createSchedule(platformId, endpoint) + console.log(`Started schedule '${schedule.scheduleId}`) + return schedule.scheduleId } catch (error) { - logger.error(error, 'Failed to run discourse extraction discourse'); - throw new ApiError(590, 'Failed to run discourse extraction discourse'); + logger.error(error, 'Failed to create discourse schedule'); + throw new ApiError(590, 'Failed to create discourse schedule'); } } export default { getPropertyHandler, - runDiscourseExtraction, + createDiscourseSchedule, }; diff --git a/src/services/platform.service.ts b/src/services/platform.service.ts index b70dfef..7ff7c1b 100644 --- a/src/services/platform.service.ts +++ b/src/services/platform.service.ts @@ -65,10 +65,12 @@ const getPlatformById = async (id: Types.ObjectId): Promise} platform * @returns {Promise} */ -const callExtractionApp = (platform: HydratedDocument): void => { +const callExtractionApp = async (platform: HydratedDocument): Promise => { switch (platform.name) { case PlatformNames.Discourse: { - discourseService.coreService.runDiscourseExtraction(platform.id as string); + const scheduleId = await discourseService.coreService.createDiscourseSchedule(platform.id as string, platform.metadata?.id as string); + platform.set('metadata.scheduleId', scheduleId) + await platform.save() return; } default: { diff --git a/src/services/temporal/core.service.ts b/src/services/temporal/core.service.ts new file mode 100644 index 0000000..ab313d6 --- /dev/null +++ b/src/services/temporal/core.service.ts @@ -0,0 +1,20 @@ +import { Client, Connection } from '@temporalio/client'; +import config from '../../config'; + +export class TemporalCoreService { + private connection: Connection | undefined + private client: Client | undefined + + constructor() { } + + protected async getClient(): Promise { + if (!this.client) { + if (!this.connection) { + this.connection = await Connection.connect({ address: config.temporal.uri }); + } + this.client = new Client({ connection: this.connection }) + } + return this.client + } + +} \ No newline at end of file diff --git a/src/services/temporal/discourse.service.ts b/src/services/temporal/discourse.service.ts new file mode 100644 index 0000000..e7e7081 --- /dev/null +++ b/src/services/temporal/discourse.service.ts @@ -0,0 +1,29 @@ +import { Client, ScheduleHandle, ScheduleOverlapPolicy } from "@temporalio/client"; +import { TemporalCoreService } from "./core.service"; +import config from "src/config"; + +class TemporalDiscourseService extends TemporalCoreService { + + public async createSchedule(platformId: string, endpoint: string): Promise { + const client: Client = await this.getClient() + + return client.schedule.create({ + action: { + type: 'startWorkflow', + workflowType: 'DiscourseExtractWorkflow', + args: [endpoint], // add platformId + taskQueue: config.temporal.heavyQueue + }, + scheduleId: `discourse/${endpoint}`, + policies: { + catchupWindow: '1 day', + overlap: ScheduleOverlapPolicy.ALLOW_ALL + }, + spec: { + intervals: [{ every: '1d' }] + } + }) + } +} + +export default new TemporalDiscourseService() \ No newline at end of file From 2b57ecd5c841018430dc15e3752601ec02f1dfda Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 19:29:48 +0100 Subject: [PATCH 3/8] code improvements + prettier --- src/services/discourse/core.service.ts | 8 +++---- src/services/platform.service.ts | 9 +++++--- src/services/temporal/core.service.ts | 27 ++++++++++++++++------ src/services/temporal/discourse.service.ts | 25 ++++++++++---------- 4 files changed, 42 insertions(+), 27 deletions(-) diff --git a/src/services/discourse/core.service.ts b/src/services/discourse/core.service.ts index 3dfdea8..47b1a9b 100644 --- a/src/services/discourse/core.service.ts +++ b/src/services/discourse/core.service.ts @@ -2,7 +2,7 @@ import parentLogger from '../../config/logger'; import { IAuthAndPlatform } from '../../interfaces'; import categoryService from './category.service'; import { ApiError, pick } from '../../utils'; -import temporalDiscourse from '../temporal/discourse.service' +import temporalDiscourse from '../temporal/discourse.service'; const logger = parentLogger.child({ module: 'DiscourseCoreService' }); async function getPropertyHandler(req: IAuthAndPlatform) { @@ -20,9 +20,9 @@ async function getPropertyHandler(req: IAuthAndPlatform) { */ async function createDiscourseSchedule(platformId: string, endpoint: string): Promise { try { - const schedule = await temporalDiscourse.createSchedule(platformId, endpoint) - console.log(`Started schedule '${schedule.scheduleId}`) - return schedule.scheduleId + const schedule = await temporalDiscourse.createSchedule(platformId, endpoint); + logger.info(`Started schedule '${schedule.scheduleId}`); + return schedule.scheduleId; } catch (error) { logger.error(error, 'Failed to create discourse schedule'); throw new ApiError(590, 'Failed to create discourse schedule'); diff --git a/src/services/platform.service.ts b/src/services/platform.service.ts index 7ff7c1b..a0dfa9b 100644 --- a/src/services/platform.service.ts +++ b/src/services/platform.service.ts @@ -68,9 +68,12 @@ const getPlatformById = async (id: Types.ObjectId): Promise): Promise => { switch (platform.name) { case PlatformNames.Discourse: { - const scheduleId = await discourseService.coreService.createDiscourseSchedule(platform.id as string, platform.metadata?.id as string); - platform.set('metadata.scheduleId', scheduleId) - await platform.save() + const scheduleId = await discourseService.coreService.createDiscourseSchedule( + platform.id as string, + platform.metadata?.id as string, + ); + platform.set('metadata.scheduleId', scheduleId); + await platform.save(); return; } default: { diff --git a/src/services/temporal/core.service.ts b/src/services/temporal/core.service.ts index ab313d6..fae5035 100644 --- a/src/services/temporal/core.service.ts +++ b/src/services/temporal/core.service.ts @@ -2,19 +2,32 @@ import { Client, Connection } from '@temporalio/client'; import config from '../../config'; export class TemporalCoreService { - private connection: Connection | undefined - private client: Client | undefined + private connection: Connection | undefined; + private client: Client | undefined; - constructor() { } + private async createConnection(): Promise { + try { + return await Connection.connect({ address: config.temporal.uri }); + } catch (error) { + throw new Error(`Failed to connect to Temporal: ${(error as Error).message}`); + } + } protected async getClient(): Promise { if (!this.client) { if (!this.connection) { - this.connection = await Connection.connect({ address: config.temporal.uri }); + this.connection = await this.createConnection(); } - this.client = new Client({ connection: this.connection }) + this.client = new Client({ connection: this.connection }); } - return this.client + return this.client; } -} \ No newline at end of file + public async disconnect(): Promise { + if (this.connection) { + await this.connection.close(); + this.connection = undefined; + this.client = undefined; + } + } +} diff --git a/src/services/temporal/discourse.service.ts b/src/services/temporal/discourse.service.ts index e7e7081..1722461 100644 --- a/src/services/temporal/discourse.service.ts +++ b/src/services/temporal/discourse.service.ts @@ -1,29 +1,28 @@ -import { Client, ScheduleHandle, ScheduleOverlapPolicy } from "@temporalio/client"; -import { TemporalCoreService } from "./core.service"; -import config from "src/config"; +import { Client, ScheduleHandle, ScheduleOverlapPolicy } from '@temporalio/client'; +import { TemporalCoreService } from './core.service'; +import config from 'src/config'; class TemporalDiscourseService extends TemporalCoreService { - public async createSchedule(platformId: string, endpoint: string): Promise { - const client: Client = await this.getClient() + const client: Client = await this.getClient(); return client.schedule.create({ action: { type: 'startWorkflow', workflowType: 'DiscourseExtractWorkflow', - args: [endpoint], // add platformId - taskQueue: config.temporal.heavyQueue + args: [endpoint], // TODO add platformId + taskQueue: config.temporal.heavyQueue, }, - scheduleId: `discourse/${endpoint}`, + scheduleId: `discourse/${encodeURIComponent(endpoint)}`, policies: { catchupWindow: '1 day', - overlap: ScheduleOverlapPolicy.ALLOW_ALL + overlap: ScheduleOverlapPolicy.SKIP, }, spec: { - intervals: [{ every: '1d' }] - } - }) + intervals: [{ every: '1d' }], + }, + }); } } -export default new TemporalDiscourseService() \ No newline at end of file +export default new TemporalDiscourseService(); From e45f4dafb60f08c6dccfe4a2ca2e862027e6a2b9 Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 19:43:15 +0100 Subject: [PATCH 4/8] fix config --- src/services/temporal/discourse.service.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/services/temporal/discourse.service.ts b/src/services/temporal/discourse.service.ts index 1722461..8a070f5 100644 --- a/src/services/temporal/discourse.service.ts +++ b/src/services/temporal/discourse.service.ts @@ -1,6 +1,6 @@ import { Client, ScheduleHandle, ScheduleOverlapPolicy } from '@temporalio/client'; import { TemporalCoreService } from './core.service'; -import config from 'src/config'; +import config from '../../config'; class TemporalDiscourseService extends TemporalCoreService { public async createSchedule(platformId: string, endpoint: string): Promise { From 4fd85fdf4b3f4f0e4769a2fdb59fd05804fcf248 Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 19:49:20 +0100 Subject: [PATCH 5/8] Add error handling around Temporal client operations. --- src/services/temporal/discourse.service.ts | 39 ++++++++++++---------- 1 file changed, 21 insertions(+), 18 deletions(-) diff --git a/src/services/temporal/discourse.service.ts b/src/services/temporal/discourse.service.ts index 8a070f5..85657de 100644 --- a/src/services/temporal/discourse.service.ts +++ b/src/services/temporal/discourse.service.ts @@ -4,24 +4,27 @@ import config from '../../config'; class TemporalDiscourseService extends TemporalCoreService { public async createSchedule(platformId: string, endpoint: string): Promise { - const client: Client = await this.getClient(); - - return client.schedule.create({ - action: { - type: 'startWorkflow', - workflowType: 'DiscourseExtractWorkflow', - args: [endpoint], // TODO add platformId - taskQueue: config.temporal.heavyQueue, - }, - scheduleId: `discourse/${encodeURIComponent(endpoint)}`, - policies: { - catchupWindow: '1 day', - overlap: ScheduleOverlapPolicy.SKIP, - }, - spec: { - intervals: [{ every: '1d' }], - }, - }); + try { + const client: Client = await this.getClient(); + return client.schedule.create({ + action: { + type: 'startWorkflow', + workflowType: 'DiscourseExtractWorkflow', + args: [endpoint], // TODO add platformId + taskQueue: config.temporal.heavyQueue, + }, + scheduleId: `discourse/${encodeURIComponent(endpoint)}`, + policies: { + catchupWindow: '1 day', + overlap: ScheduleOverlapPolicy.SKIP, + }, + spec: { + intervals: [{ every: '1d' }], + }, + }); + } catch (error) { + throw new Error(`Failed to create Temporal schedule: ${(error as Error).message}`); + } } } From 83964cc0cad0e9a9c5eeaf4dcc6b7ff77962d7ce Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 19:49:56 +0100 Subject: [PATCH 6/8] add platformId to args --- src/services/temporal/discourse.service.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/services/temporal/discourse.service.ts b/src/services/temporal/discourse.service.ts index 85657de..5956a4e 100644 --- a/src/services/temporal/discourse.service.ts +++ b/src/services/temporal/discourse.service.ts @@ -10,7 +10,7 @@ class TemporalDiscourseService extends TemporalCoreService { action: { type: 'startWorkflow', workflowType: 'DiscourseExtractWorkflow', - args: [endpoint], // TODO add platformId + args: [endpoint, platformId], taskQueue: config.temporal.heavyQueue, }, scheduleId: `discourse/${encodeURIComponent(endpoint)}`, From 6e06ec36d3dccbe6a6c648694e8cb93d6269a8d0 Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 19:53:27 +0100 Subject: [PATCH 7/8] added ' --- src/services/discourse/core.service.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/services/discourse/core.service.ts b/src/services/discourse/core.service.ts index 47b1a9b..873e50c 100644 --- a/src/services/discourse/core.service.ts +++ b/src/services/discourse/core.service.ts @@ -21,7 +21,7 @@ async function getPropertyHandler(req: IAuthAndPlatform) { async function createDiscourseSchedule(platformId: string, endpoint: string): Promise { try { const schedule = await temporalDiscourse.createSchedule(platformId, endpoint); - logger.info(`Started schedule '${schedule.scheduleId}`); + logger.info(`Started schedule '${schedule.scheduleId}'`); return schedule.scheduleId; } catch (error) { logger.error(error, 'Failed to create discourse schedule'); From 9b1e19272af7053229cc4d7c687c22f709d9867a Mon Sep 17 00:00:00 2001 From: Cyrille Derche Date: Sat, 16 Nov 2024 20:12:30 +0100 Subject: [PATCH 8/8] trigger the first run --- src/services/discourse/core.service.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/src/services/discourse/core.service.ts b/src/services/discourse/core.service.ts index 873e50c..a496fb9 100644 --- a/src/services/discourse/core.service.ts +++ b/src/services/discourse/core.service.ts @@ -22,6 +22,7 @@ async function createDiscourseSchedule(platformId: string, endpoint: string): Pr try { const schedule = await temporalDiscourse.createSchedule(platformId, endpoint); logger.info(`Started schedule '${schedule.scheduleId}'`); + await schedule.trigger(); return schedule.scheduleId; } catch (error) { logger.error(error, 'Failed to create discourse schedule');