From f1e50aca952c791b52f4276f7fa9b6c8837af16b Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 10 Sep 2026 12:22:11 +0000 Subject: [PATCH 1/3] refactor(functions): deslop cron scheduler and move tests Flatten refreshRoutes, drop redundant cron prep/active checks, and move unit tests into src/__tests__. Revert unrelated 2FA and pnpm allowBuilds drive-bys from the cron PR. Co-authored-by: Christina Papadogianni --- modules/authentication/src/handlers/twoFa.ts | 2 +- modules/functions/package.json | 2 +- modules/functions/src/Functions.ts | 1 - .../cron.utils.test.ts | 2 +- .../utils.background.test.ts | 2 +- .../src/controllers/cronQueue.controller.ts | 105 +++++++-------- .../src/controllers/function.controller.ts | 127 ++++++++---------- modules/functions/src/migrations/index.ts | 1 - modules/functions/tsconfig.json | 9 +- .../src/admin/routes/ToggleTwoFa.route.ts | 2 +- packages/core/src/admin/utils/auth.ts | 2 +- pnpm-workspace.yaml | 2 - 12 files changed, 123 insertions(+), 134 deletions(-) rename modules/functions/src/{controllers => __tests__}/cron.utils.test.ts (99%) rename modules/functions/src/{controllers => __tests__}/utils.background.test.ts (99%) diff --git a/modules/authentication/src/handlers/twoFa.ts b/modules/authentication/src/handlers/twoFa.ts index 3972cbe66..fa7fa0c77 100644 --- a/modules/authentication/src/handlers/twoFa.ts +++ b/modules/authentication/src/handlers/twoFa.ts @@ -535,7 +535,7 @@ export class TwoFa implements IAuthenticationStrategy { } private async enableAuthenticator2Fa(user: User): Promise { - const secret = await node2fa.generateSecret({ + const secret = node2fa.generateSecret({ //to do: add logic for app name insertion name: 'Conduit', // add another string when mail is not available diff --git a/modules/functions/package.json b/modules/functions/package.json index 358ac4baa..26ec46610 100644 --- a/modules/functions/package.json +++ b/modules/functions/package.json @@ -25,7 +25,7 @@ "build:bundle": "rimraf bundle && node ../../libraries/service-bundle/dist/cli.js generate-manifest && tsup && node ../../libraries/service-bundle/dist/cli.js copy-assets && node ../../libraries/service-bundle/dist/cli.js generate-lockfile", "prepare": "npm run build", "generateTypes": "sh build.sh", - "test": "tsc -p tsconfig.test.json && node --test dist-test/controllers/*.test.js", + "test": "tsc -p tsconfig.test.json && node --test dist-test/__tests__/*.test.js", "build:docker": "docker build -t ghcr.io/conduitplatform/functions:latest -f ./Dockerfile ../../ && docker push ghcr.io/conduitplatform/functions:latest" }, "directories": { diff --git a/modules/functions/src/Functions.ts b/modules/functions/src/Functions.ts index 6398d3da0..7cfa48072 100644 --- a/modules/functions/src/Functions.ts +++ b/modules/functions/src/Functions.ts @@ -76,7 +76,6 @@ export default class Functions extends ManagedModule { if (moduleName !== 'router' || !serving || this.isRunning) return; this.isRunning = true; this.functionsController = new FunctionController(this.grpcServer, this.grpcSdk); - CronQueueController.getInstance(this.grpcSdk); this.adminRouter = new AdminHandlers( this.grpcServer, this.grpcSdk, diff --git a/modules/functions/src/controllers/cron.utils.test.ts b/modules/functions/src/__tests__/cron.utils.test.ts similarity index 99% rename from modules/functions/src/controllers/cron.utils.test.ts rename to modules/functions/src/__tests__/cron.utils.test.ts index 56b424267..31c1b3a26 100644 --- a/modules/functions/src/controllers/cron.utils.test.ts +++ b/modules/functions/src/__tests__/cron.utils.test.ts @@ -8,7 +8,7 @@ import { parseCronJobFunctionId, planCronSync, validateCronPattern, -} from './cron.utils.js'; +} from '../controllers/cron.utils.js'; describe('cron.utils', () => { describe('parseCronJobFunctionId', () => { diff --git a/modules/functions/src/controllers/utils.background.test.ts b/modules/functions/src/__tests__/utils.background.test.ts similarity index 99% rename from modules/functions/src/controllers/utils.background.test.ts rename to modules/functions/src/__tests__/utils.background.test.ts index 7b7f5a110..5855db4c4 100644 --- a/modules/functions/src/controllers/utils.background.test.ts +++ b/modules/functions/src/__tests__/utils.background.test.ts @@ -3,7 +3,7 @@ import assert from 'node:assert/strict'; import { ConduitGrpcSdk } from '@conduitplatform/grpc-sdk'; import { FunctionExecutions } from '../models/index.js'; import type { Functions } from '../models/index.js'; -import { compileFunctionCode, executeBackgroundFunction } from './utils.js'; +import { compileFunctionCode, executeBackgroundFunction } from '../controllers/utils.js'; const originalGetInstance = FunctionExecutions.getInstance; diff --git a/modules/functions/src/controllers/cronQueue.controller.ts b/modules/functions/src/controllers/cronQueue.controller.ts index 5736faa2a..6f284bef2 100644 --- a/modules/functions/src/controllers/cronQueue.controller.ts +++ b/modules/functions/src/controllers/cronQueue.controller.ts @@ -61,6 +61,40 @@ export class CronQueueController { return this.cronQueue.getRepeatableJobs(); } + private async executeCronJob(job: Job<{ functionId: string }>): Promise { + const func = await Functions.getInstance().findOne( + { _id: job.data.functionId }, + { readPreference: 'primary' }, + ); + if (!func || func.functionType !== 'cron') { + return; + } + const cronPattern = getCronPatternFromInputs(func.inputs); + if (!cronPattern) { + ConduitGrpcSdk.Logger.warn( + `Cron function ${func.name} (${func._id}) has no pattern; skipping tick`, + ); + return; + } + const compiled = + this.compiledFunctions.get(func._id) ?? compileFunctionCode(func.functionCode); + const scheduledAt = new Date().toISOString(); + ConduitGrpcSdk.Logger.log( + `Cron tick for ${func.name} (${cronPattern}) at ${scheduledAt}`, + ); + await executeBackgroundFunction( + func, + { + scheduledAt, + cronPattern, + trigger: 'cron', + }, + compiled, + this.grpcSdk, + ); + ConduitGrpcSdk.Logger.log(`Cron execution completed for ${func.name}`); + } + async ensureWorker(lockDuration: number): Promise { if (!this.shouldSchedule()) { return this.cronWorker; @@ -75,39 +109,7 @@ export class CronQueueController { } this.cronWorker = new Worker( CRON_QUEUE_NAME, - async (job: Job<{ functionId: string }>) => { - const func = await Functions.getInstance().findOne( - { _id: job.data.functionId }, - { readPreference: 'primary' }, - ); - if (!func || func.functionType !== 'cron') { - return; - } - const cronPattern = getCronPatternFromInputs(func.inputs); - if (!cronPattern) { - ConduitGrpcSdk.Logger.warn( - `Cron function ${func.name} (${func._id}) has no pattern; skipping tick`, - ); - return; - } - const compiled = - this.compiledFunctions.get(func._id) ?? compileFunctionCode(func.functionCode); - const scheduledAt = new Date().toISOString(); - ConduitGrpcSdk.Logger.log( - `Cron tick for ${func.name} (${cronPattern}) at ${scheduledAt}`, - ); - await executeBackgroundFunction( - func, - { - scheduledAt, - cronPattern, - trigger: 'cron', - }, - compiled, - this.grpcSdk, - ); - ConduitGrpcSdk.Logger.log(`Cron execution completed for ${func.name}`); - }, + job => this.executeCronJob(job), { concurrency: 1, lockDuration, @@ -220,25 +222,21 @@ export class CronQueueController { } } + private lockDurationFromTimeouts(timeouts: Array): number { + const maxTimeout = timeouts.reduce( + (max, timeout) => Math.max(max, timeout ?? DEFAULT_FUNCTION_TIMEOUT_MS), + DEFAULT_FUNCTION_TIMEOUT_MS, + ); + return maxTimeout + LOCK_BUFFER_MS; + } + private async lockDurationForCronFunctions(): Promise { type CronTimeout = Pick; const cronDocs = (await Functions.getInstance().findMany( { functionType: 'cron' }, { select: 'timeout', readPreference: 'primary' }, )) as CronTimeout[]; - const maxTimeout = cronDocs.reduce( - (max, func) => Math.max(max, func.timeout ?? DEFAULT_FUNCTION_TIMEOUT_MS), - DEFAULT_FUNCTION_TIMEOUT_MS, - ); - return maxTimeout + LOCK_BUFFER_MS; - } - - private lockDurationFromFunctions(cronFunctions: Functions[]): number { - const maxTimeout = cronFunctions.reduce( - (max, func) => Math.max(max, func.timeout ?? DEFAULT_FUNCTION_TIMEOUT_MS), - DEFAULT_FUNCTION_TIMEOUT_MS, - ); - return maxTimeout + LOCK_BUFFER_MS; + return this.lockDurationFromTimeouts(cronDocs.map(func => func.timeout)); } private async reconcileCronJobs(): Promise { @@ -270,9 +268,6 @@ export class CronQueueController { try { if (item.existingKey) { await this.cronQueue.removeRepeatableByKey(item.existingKey); - updated += 1; - } else { - registered += 1; } await this.cronQueue.add( CRON_JOB_NAME, @@ -284,23 +279,23 @@ export class CronQueueController { removeOnFail: { age: 24 * 3600 }, }, ); + if (item.existingKey) { + updated += 1; + } else { + registered += 1; + } } catch (err) { ConduitGrpcSdk.Logger.error( `Failed to schedule cron job ${item.jobId}: ${(err as Error).message}`, ); errors += 1; - if (item.existingKey) { - updated -= 1; - } else { - registered -= 1; - } } } ConduitGrpcSdk.Logger.log( `Cron sync complete: registered=${registered}, updated=${updated}, unchanged=${plan.unchangedJobIds.length}, removed=${removed}, skipped=${plan.skipped.length}, errors=${errors}`, ); - return this.lockDurationFromFunctions(cronFunctions); + return this.lockDurationFromTimeouts(cronFunctions.map(func => func.timeout)); } private setupWorkerEventHandlers(worker: Worker): void { diff --git a/modules/functions/src/controllers/function.controller.ts b/modules/functions/src/controllers/function.controller.ts index 2b999b363..5f151e0dd 100644 --- a/modules/functions/src/controllers/function.controller.ts +++ b/modules/functions/src/controllers/function.controller.ts @@ -6,12 +6,12 @@ import { ConduitSocketOptions, } from '@conduitplatform/grpc-sdk'; import { + ConfigController, GrpcServer, RequestHandlers, RoutingManager, SocketEventHandler, } from '@conduitplatform/module-tools'; -import { ConfigController } from '@conduitplatform/module-tools'; import { Functions } from '../models/index.js'; import { CronQueueController } from './cronQueue.controller.js'; @@ -33,8 +33,10 @@ type Middleware = { handler: RequestHandlers; }; +type FunctionRoute = Route | Socket | Middleware; + export class FunctionController { - private functionRoutes: (Route | Socket | Middleware)[] = []; + private functionRoutes: FunctionRoute[] = []; private readonly compiledCronFunctions = new Map(); private _routingManager: RoutingManager; @@ -54,77 +56,66 @@ export class FunctionController { }); } - refreshRoutes() { - return Functions.getInstance() - .findMany({}, { readPreference: 'primary' }) - .then(async r => { - if (!r || r.length == 0) { - ConduitGrpcSdk.Logger.log('No functions to register'); - } - this.functionRoutes = []; - this.compiledCronFunctions.clear(); + async refreshRoutes() { + try { + const functions = await Functions.getInstance().findMany( + {}, + { readPreference: 'primary' }, + ); + if (!functions || functions.length === 0) { + ConduitGrpcSdk.Logger.log('No functions to register'); + } + this.functionRoutes = []; + this.compiledCronFunctions.clear(); - for (const func of r) { - try { - if (func.functionType === 'cron') { - try { - this.compiledCronFunctions.set(func._id, tryPrepareCronFunction(func)); - } catch (err) { - ConduitGrpcSdk.Logger.error( - `Failed to prepare cron function ${func.name} (${func._id})`, - ); - ConduitGrpcSdk.Logger.error(err as Error); - } - continue; - } - const route = createFunctionRoute(func, this.grpcSdk); - if (route) { - this.functionRoutes.push(route as any); - } - } catch (err) { - ConduitGrpcSdk.Logger.error( - `Failed to process function ${func.name} (${func._id}); skipping`, - ); - ConduitGrpcSdk.Logger.error(err as Error); + for (const func of functions) { + try { + if (func.functionType === 'cron') { + this.compiledCronFunctions.set(func._id, tryPrepareCronFunction(func)); + continue; } - } - this._routingManager.clear(); - this.functionRoutes.forEach(route => { - if ((route as Socket).events) { - this._routingManager.socket( - (route as Socket).input, - (route as Socket).events, - ); - } else if (!(route as Middleware).hasOwnProperty('returnType')) { - this._routingManager.middleware( - (route as Middleware).input, - (route as Middleware).handler, - ); - } else { - this._routingManager.route( - (route as Route).input, - (route as Route).returnType, - (route as Route).handler, - ); - } - }); - await this._routingManager.registerRoutes(); - ConduitGrpcSdk.Logger.log('Refreshed routes'); - - if (ConfigController.getInstance().config.active) { - const cronQueue = CronQueueController.getInstance(this.grpcSdk); - cronQueue.setCompiledFunctions(this.compiledCronFunctions); - if (ConfigController.getInstance().config.active) { - await cronQueue.syncCronJobs(); + const route = createFunctionRoute(func, this.grpcSdk); + if (route) { + this.functionRoutes.push(route as FunctionRoute); } + } catch (err) { + ConduitGrpcSdk.Logger.error( + `Failed to process function ${func.name} (${func._id}); skipping`, + ); + ConduitGrpcSdk.Logger.error(err as Error); + } + } + this._routingManager.clear(); + this.functionRoutes.forEach(route => { + if ((route as Socket).events) { + this._routingManager.socket((route as Socket).input, (route as Socket).events); + } else if (!(route as Middleware).hasOwnProperty('returnType')) { + this._routingManager.middleware( + (route as Middleware).input, + (route as Middleware).handler, + ); + } else { + this._routingManager.route( + (route as Route).input, + (route as Route).returnType, + (route as Route).handler, + ); } - }) - .catch((err: Error) => { - ConduitGrpcSdk.Logger.error( - 'Something went wrong when loading functions to the router', - ); - ConduitGrpcSdk.Logger.error(err); }); + await this._routingManager.registerRoutes(); + ConduitGrpcSdk.Logger.log('Refreshed routes'); + + if (ConfigController.getInstance().config.active) { + const cronQueue = CronQueueController.getInstance(this.grpcSdk); + cronQueue.setCompiledFunctions(this.compiledCronFunctions); + await cronQueue.syncCronJobs(); + } + } catch (err) { + ConduitGrpcSdk.Logger.error( + 'Something went wrong when loading functions to the router', + ); + ConduitGrpcSdk.Logger.error(err as Error); + } } refreshEndpoints(): void { diff --git a/modules/functions/src/migrations/index.ts b/modules/functions/src/migrations/index.ts index 41085ec5f..430da9a4b 100644 --- a/modules/functions/src/migrations/index.ts +++ b/modules/functions/src/migrations/index.ts @@ -3,7 +3,6 @@ import { Functions } from '../models/index.js'; import { getCronPatternFromInputs } from '../controllers/cron.utils.js'; export async function runMigrations(_grpcSdk: ConduitGrpcSdk) { - void _grpcSdk; const cronFunctions = await Functions.getInstance().findMany({ functionType: 'cron', }); diff --git a/modules/functions/tsconfig.json b/modules/functions/tsconfig.json index 9c119b48f..0df2b7cf2 100644 --- a/modules/functions/tsconfig.json +++ b/modules/functions/tsconfig.json @@ -67,5 +67,12 @@ "forceConsistentCasingInFileNames": true /* Disallow inconsistently-cased references to the same file. */ }, "include": ["src/**/*"], - "exclude": ["node_modules", "dist", "bundle", "tsup.config.ts", "src/**/*.test.ts"] + "exclude": [ + "node_modules", + "dist", + "bundle", + "tsup.config.ts", + "src/**/*.test.ts", + "src/__tests__" + ] } diff --git a/packages/core/src/admin/routes/ToggleTwoFa.route.ts b/packages/core/src/admin/routes/ToggleTwoFa.route.ts index 0a9d9d761..0e0337e4c 100644 --- a/packages/core/src/admin/routes/ToggleTwoFa.route.ts +++ b/packages/core/src/admin/routes/ToggleTwoFa.route.ts @@ -37,7 +37,7 @@ export function toggleTwoFaRoute() { return '2FA already enabled'; } - const secret = await generateSecret({ + const secret = generateSecret({ name: 'Conduit', account: admin.username, }); diff --git a/packages/core/src/admin/utils/auth.ts b/packages/core/src/admin/utils/auth.ts index 841ab25e1..57f0f9eae 100644 --- a/packages/core/src/admin/utils/auth.ts +++ b/packages/core/src/admin/utils/auth.ts @@ -51,6 +51,6 @@ export async function verify2Fa(admin: Admin, code: string) { return signToken({ id: admin._id }, tokenSecret, tokenExpirationTime); } -export async function generateSecret(options?: { name: string; account: string }) { +export function generateSecret(options?: { name: string; account: string }) { return twoFactor.generateSecret(options); } diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 450315c1b..3c971f2c3 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -22,8 +22,6 @@ overrides: allowBuilds: '@apollo/protobufjs': true '@firebase/util': true - '@parcel/watcher': true - '@scarf/scarf': false bcrypt: true core-js: true esbuild: true From 6ec8220c444a549563b1c9a4d50026e7e2a8b229 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 10 Sep 2026 12:24:27 +0000 Subject: [PATCH 2/3] fix(functions): type lock-duration helper without reduce inference TypeScript 6 infers Array.reduce as possibly undefined; use an explicit loop so tsc --noEmit stays clean. Co-authored-by: Christina Papadogianni --- modules/functions/src/controllers/cronQueue.controller.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/modules/functions/src/controllers/cronQueue.controller.ts b/modules/functions/src/controllers/cronQueue.controller.ts index 6f284bef2..18bd33fe7 100644 --- a/modules/functions/src/controllers/cronQueue.controller.ts +++ b/modules/functions/src/controllers/cronQueue.controller.ts @@ -223,10 +223,10 @@ export class CronQueueController { } private lockDurationFromTimeouts(timeouts: Array): number { - const maxTimeout = timeouts.reduce( - (max, timeout) => Math.max(max, timeout ?? DEFAULT_FUNCTION_TIMEOUT_MS), - DEFAULT_FUNCTION_TIMEOUT_MS, - ); + let maxTimeout = DEFAULT_FUNCTION_TIMEOUT_MS; + for (const timeout of timeouts) { + maxTimeout = Math.max(maxTimeout, timeout ?? DEFAULT_FUNCTION_TIMEOUT_MS); + } return maxTimeout + LOCK_BUFFER_MS; } From c70ed1d3b06d9fcc65ba40fcbe4c1aca64abda2b Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 10 Sep 2026 12:24:36 +0000 Subject: [PATCH 3/3] style(functions): format cronQueue controller with prettier Co-authored-by: Christina Papadogianni --- .../src/controllers/cronQueue.controller.ts | 20 ++++++++----------- 1 file changed, 8 insertions(+), 12 deletions(-) diff --git a/modules/functions/src/controllers/cronQueue.controller.ts b/modules/functions/src/controllers/cronQueue.controller.ts index 18bd33fe7..417e239c4 100644 --- a/modules/functions/src/controllers/cronQueue.controller.ts +++ b/modules/functions/src/controllers/cronQueue.controller.ts @@ -107,18 +107,14 @@ export class CronQueueController { this.cronWorker = undefined; this.workerLockDuration = 0; } - this.cronWorker = new Worker( - CRON_QUEUE_NAME, - job => this.executeCronJob(job), - { - concurrency: 1, - lockDuration, - maxStalledCount: 0, - removeOnComplete: { age: 3600, count: 1000 }, - removeOnFail: { age: 24 * 3600 }, - connection: this.redisConnection, - }, - ); + this.cronWorker = new Worker(CRON_QUEUE_NAME, job => this.executeCronJob(job), { + concurrency: 1, + lockDuration, + maxStalledCount: 0, + removeOnComplete: { age: 3600, count: 1000 }, + removeOnFail: { age: 24 * 3600 }, + connection: this.redisConnection, + }); this.workerLockDuration = lockDuration; this.setupWorkerEventHandlers(this.cronWorker); if (!this.shouldSchedule()) {