From a5e106d53f96b2b95177cb6faf2774dc3c6e4949 Mon Sep 17 00:00:00 2001 From: Konstantinos Kopanidis Date: Sun, 21 Jun 2026 19:32:00 +0300 Subject: [PATCH] feat(functions): implement cron scheduling with BullMQ Cron functions now register repeatable BullMQ jobs instead of erroneously subscribing to the event bus, with pattern validation and execution logging. --- modules/functions/package.json | 3 + modules/functions/src/Functions.ts | 9 + modules/functions/src/admin/index.ts | 19 +- .../functions/src/controllers/cron.utils.ts | 47 +++++ .../src/controllers/cronQueue.controller.ts | 181 ++++++++++++++++++ .../src/controllers/function.controller.ts | 41 ++-- modules/functions/src/controllers/utils.ts | 57 ++++-- .../src/interfaces/IWebInputs.interface.ts | 2 + modules/functions/src/migrations/index.ts | 25 ++- pnpm-lock.yaml | 9 + 10 files changed, 362 insertions(+), 31 deletions(-) create mode 100644 modules/functions/src/controllers/cron.utils.ts create mode 100644 modules/functions/src/controllers/cronQueue.controller.ts diff --git a/modules/functions/package.json b/modules/functions/package.json index eaf0a6df9..b1f5d19bd 100644 --- a/modules/functions/package.json +++ b/modules/functions/package.json @@ -38,8 +38,11 @@ "@grpc/grpc-js": "^1.14.4", "@grpc/proto-loader": "^0.8.1", "axios": "^1.18.0", + "bullmq": "^5.79.0", "convict": "^6.2.5", + "cron-parser": "^4.9.0", "escape-string-regexp": "^5.0.0", + "ioredis": "^5.11.1", "lodash-es": "^4.18.1" }, "devDependencies": { diff --git a/modules/functions/src/Functions.ts b/modules/functions/src/Functions.ts index 60443f75d..3eb73be53 100644 --- a/modules/functions/src/Functions.ts +++ b/modules/functions/src/Functions.ts @@ -8,6 +8,8 @@ import { AdminHandlers } from './admin/index.js'; import * as models from './models/index.js'; import AppConfigSchema, { Config } from './config/index.js'; import { FunctionController } from './controllers/function.controller.js'; +import { CronQueueController } from './controllers/cronQueue.controller.js'; +import { runMigrations } from './migrations/index.js'; import { ConfigController, ManagedModule, @@ -37,6 +39,7 @@ export default class Functions extends ManagedModule { await this.awaitPeersFromManifest(); this.database = this.grpcSdk.database!; await this.registerSchemas(); + await runMigrations(this.grpcSdk); } async onConfig() { @@ -46,6 +49,11 @@ export default class Functions extends ManagedModule { this.isRunning = false; this.updateHealth(HealthCheckStatus.NOT_SERVING); this.trustModelNoticeLoggedForActiveSession = false; + try { + await CronQueueController.getInstance().drainCronQueue(); + } catch { + // Queue was never initialized while inactive. + } } else if (!this.trustModelNoticeLoggedForActiveSession) { ConduitGrpcSdk.Logger.log( 'Functions module active: user-defined code runs with full server privileges. Only deploy code you trust.', @@ -65,6 +73,7 @@ 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/admin/index.ts b/modules/functions/src/admin/index.ts index 7d3a5f45d..df700fd1b 100644 --- a/modules/functions/src/admin/index.ts +++ b/modules/functions/src/admin/index.ts @@ -21,6 +21,11 @@ import { status } from '@grpc/grpc-js'; import { FunctionExecutions, Functions } from '../models/index.js'; import { FunctionController } from '../controllers/function.controller.js'; import { compileUserFunctionScript } from '../sandbox/functionSandbox.js'; +import { + getCronPatternFromInputs, + normalizeCronInputs, + validateCronPattern, +} from '../controllers/cron.utils.js'; import escapeStringRegexp from 'escape-string-regexp'; @@ -48,11 +53,15 @@ export class AdminHandlers { } catch { throw new GrpcError(status.INVALID_ARGUMENT, 'Invalid function code'); } + let normalizedInputs = inputs; + if (functionType === 'cron') { + normalizedInputs = normalizeCronInputs(inputs); + } const query = { name, functionType, functionCode, - inputs, + inputs: normalizedInputs, returns, timeout: timeout ?? 180000, }; @@ -133,6 +142,14 @@ export class AdminHandlers { returns: returns ?? func.returns, timeout: timeout ?? func.timeout, }; + if (func.functionType === 'cron') { + query.inputs = normalizeCronInputs(query.inputs); + } else if (inputs && getCronPatternFromInputs(inputs)) { + const pattern = getCronPatternFromInputs(inputs); + if (pattern) { + validateCronPattern(pattern); + } + } try { compileUserFunctionScript(query.functionCode); } catch { diff --git a/modules/functions/src/controllers/cron.utils.ts b/modules/functions/src/controllers/cron.utils.ts new file mode 100644 index 000000000..7af9ac06f --- /dev/null +++ b/modules/functions/src/controllers/cron.utils.ts @@ -0,0 +1,47 @@ +import { GrpcError } from '@conduitplatform/grpc-sdk'; +import { status } from '@grpc/grpc-js'; +import { parseExpression } from 'cron-parser'; +import type { IWebInputsInterface } from '../interfaces/IWebInputs.interface.js'; + +export function buildCronJobId(functionId: string): string { + return `cron-${functionId}`; +} + +export function getCronPatternFromInputs( + inputs?: IWebInputsInterface | null, +): string | undefined { + if (!inputs) return undefined; + const pattern = inputs.cronPattern ?? inputs.event; + if (typeof pattern !== 'string') return undefined; + const trimmed = pattern.trim(); + return trimmed.length > 0 ? trimmed : undefined; +} + +export function validateCronPattern(pattern: string): void { + try { + parseExpression(pattern, { tz: 'UTC' }); + } catch { + throw new GrpcError( + status.INVALID_ARGUMENT, + `Invalid cron pattern: "${pattern}". Expected 5-field format: minute hour day month weekday (UTC).`, + ); + } +} + +export function normalizeCronInputs( + inputs: IWebInputsInterface | undefined, +): IWebInputsInterface { + const pattern = getCronPatternFromInputs(inputs); + if (!pattern) { + throw new GrpcError( + status.INVALID_ARGUMENT, + 'Cron pattern is required (inputs.cronPattern or inputs.event)', + ); + } + validateCronPattern(pattern); + return { + ...inputs, + cronPattern: pattern, + event: pattern, + }; +} diff --git a/modules/functions/src/controllers/cronQueue.controller.ts b/modules/functions/src/controllers/cronQueue.controller.ts new file mode 100644 index 000000000..8a1dffd98 --- /dev/null +++ b/modules/functions/src/controllers/cronQueue.controller.ts @@ -0,0 +1,181 @@ +import { Job, Queue, Worker } from 'bullmq'; +import { ConduitGrpcSdk } from '@conduitplatform/grpc-sdk'; +import { Cluster, Redis } from 'ioredis'; +import { Functions } from '../models/index.js'; +import { + buildCronJobId, + getCronPatternFromInputs, + validateCronPattern, +} from './cron.utils.js'; +import type { CompiledUserFunction } from '../sandbox/functionSandbox.js'; +import { compileFunctionCode, executeBackgroundFunction } from './utils.js'; + +const CRON_QUEUE_NAME = 'functions-cron-queue'; +const CRON_JOB_NAME = 'execute-cron'; + +export class CronQueueController { + private static _instance: CronQueueController; + private readonly redisConnection: Redis | Cluster; + private readonly cronQueue: Queue; + private cronWorker?: Worker; + private compiledFunctions = new Map(); + + private constructor(private readonly grpcSdk: ConduitGrpcSdk) { + this.redisConnection = this.grpcSdk.redisManager.getClient(); + this.cronQueue = new Queue(CRON_QUEUE_NAME, { + connection: this.redisConnection, + }); + } + + static getInstance(grpcSdk?: ConduitGrpcSdk): CronQueueController { + if (CronQueueController._instance) { + return CronQueueController._instance; + } + if (!grpcSdk) { + throw new Error('No grpcSdk instance provided!'); + } + CronQueueController._instance = new CronQueueController(grpcSdk); + return CronQueueController._instance; + } + + setCompiledFunctions(compiled: Map): void { + this.compiledFunctions = compiled; + } + + ensureWorker(): Worker { + if (this.cronWorker) { + return this.cronWorker; + } + 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}`); + }, + { + concurrency: 1, + removeOnComplete: { age: 3600, count: 1000 }, + removeOnFail: { age: 24 * 3600 }, + connection: this.redisConnection, + }, + ); + this.setupWorkerEventHandlers(this.cronWorker); + return this.cronWorker; + } + + async syncCronJobs(cronFunctions: Functions[]): Promise { + this.ensureWorker(); + const repeatables = await this.cronQueue.getRepeatableJobs(); + const expectedJobIds = new Set(cronFunctions.map(func => buildCronJobId(func._id))); + + let removed = 0; + for (const repeatable of repeatables) { + if (!repeatable.id || !expectedJobIds.has(repeatable.id)) { + await this.cronQueue.removeRepeatableByKey(repeatable.key); + removed += 1; + } + } + + let registered = 0; + let skipped = 0; + const refreshedRepeatables = await this.cronQueue.getRepeatableJobs(); + for (const func of cronFunctions) { + const pattern = getCronPatternFromInputs(func.inputs); + if (!pattern) { + ConduitGrpcSdk.Logger.warn( + `Cron function ${func.name} (${func._id}) missing pattern; skipping schedule`, + ); + skipped += 1; + continue; + } + try { + validateCronPattern(pattern); + } catch (err) { + ConduitGrpcSdk.Logger.error( + `Cron function ${func.name} (${func._id}) has invalid pattern "${pattern}": ${(err as Error).message}`, + ); + skipped += 1; + continue; + } + + const jobId = buildCronJobId(func._id); + for (const repeatable of refreshedRepeatables) { + if (repeatable.id === jobId) { + await this.cronQueue.removeRepeatableByKey(repeatable.key); + } + } + + await this.cronQueue.add( + CRON_JOB_NAME, + { functionId: func._id }, + { + jobId, + repeat: { pattern, tz: func.inputs?.timezone ?? 'UTC' }, + removeOnComplete: { age: 3600, count: 1000 }, + removeOnFail: { age: 24 * 3600 }, + }, + ); + registered += 1; + } + + ConduitGrpcSdk.Logger.log( + `Cron sync complete: registered=${registered}, removed=${removed}, skipped=${skipped}`, + ); + } + + async drainCronQueue(): Promise { + if (this.cronWorker) { + await this.cronWorker.close(); + this.cronWorker = undefined; + } + await this.cronQueue.drain(); + const repeatables = await this.cronQueue.getRepeatableJobs(); + for (const job of repeatables) { + await this.cronQueue.removeRepeatableByKey(job.key); + } + this.compiledFunctions.clear(); + } + + private setupWorkerEventHandlers(worker: Worker): void { + worker.on('error', (error: Error) => { + ConduitGrpcSdk.Logger.error('Functions cron worker error:'); + ConduitGrpcSdk.Logger.error(error); + }); + worker.on('failed', (job: Job | undefined, error: Error) => { + ConduitGrpcSdk.Logger.error( + job + ? `Cron job failed: ${job.id}, ${error.message}` + : `Cron job error: ${error.message}`, + ); + }); + } +} diff --git a/modules/functions/src/controllers/function.controller.ts b/modules/functions/src/controllers/function.controller.ts index 9d6b8d5e2..6c410257c 100644 --- a/modules/functions/src/controllers/function.controller.ts +++ b/modules/functions/src/controllers/function.controller.ts @@ -11,9 +11,12 @@ import { RoutingManager, SocketEventHandler, } from '@conduitplatform/module-tools'; +import { ConfigController } from '@conduitplatform/module-tools'; import { Functions } from '../models/index.js'; -import { createFunctionRoute } from './utils.js'; +import { CronQueueController } from './cronQueue.controller.js'; +import { compileFunctionCode, createFunctionRoute } from './utils.js'; +import type { CompiledUserFunction } from '../sandbox/functionSandbox.js'; type Socket = { input: ConduitSocketOptions; @@ -32,6 +35,7 @@ type Middleware = { export class FunctionController { private functionRoutes: (Route | Socket | Middleware)[] = []; + private readonly compiledCronFunctions = new Map(); private _routingManager: RoutingManager; @@ -53,13 +57,29 @@ export class FunctionController { refreshRoutes() { return Functions.getInstance() .findMany({}, { readPreference: 'primary' }) - .then(r => { + .then(async r => { if (!r || r.length == 0) { ConduitGrpcSdk.Logger.log('No functions to register'); } this.functionRoutes = []; + this.compiledCronFunctions.clear(); + const cronFunctions: Functions[] = []; r.forEach(func => { + if (func.functionType === 'cron') { + try { + this.compiledCronFunctions.set( + func._id, + compileFunctionCode(func.functionCode), + ); + cronFunctions.push(func); + } catch (err) { + ConduitGrpcSdk.Logger.error( + `Failed to compile cron function ${func.name} (${func._id})`, + ); + ConduitGrpcSdk.Logger.error(err as Error); + } + } const route = createFunctionRoute(func, this.grpcSdk); if (route) { this.functionRoutes.push(route as any); @@ -85,15 +105,14 @@ export class FunctionController { ); } }); - this._routingManager - .registerRoutes() - .then(() => { - ConduitGrpcSdk.Logger.log('Refreshed routes'); - }) - .catch((err: Error) => { - ConduitGrpcSdk.Logger.error('Failed to register routes for module'); - 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(cronFunctions); + } }) .catch((err: Error) => { ConduitGrpcSdk.Logger.error( diff --git a/modules/functions/src/controllers/utils.ts b/modules/functions/src/controllers/utils.ts index f9ede13d3..b8b8f5457 100644 --- a/modules/functions/src/controllers/utils.ts +++ b/modules/functions/src/controllers/utils.ts @@ -15,6 +15,7 @@ import { compileUserFunctionScript, runUserFunctionInSandbox, } from '../sandbox/functionSandbox.js'; +import { getCronPatternFromInputs, validateCronPattern } from './cron.utils.js'; function getOperation(op: string) { switch (op) { @@ -179,6 +180,27 @@ export function createMiddlewareFunction(func: Functions, grpcSdk: ConduitGrpcSd }; } +export async function executeBackgroundFunction( + func: Functions, + requestPayload: unknown, + functionCodeCompiled: CompiledUserFunction, + grpcSdk: ConduitGrpcSdk, +): Promise { + try { + await executeFunction( + func, + { request: requestPayload } as ParsedRouterRequest, + functionCodeCompiled, + func.timeout, + func.name, + grpcSdk, + ); + } catch (err) { + ConduitGrpcSdk.Logger.error(`Execution failed for ${func.name}:`); + ConduitGrpcSdk.Logger.error(err as Error); + } +} + export function createEventFunction(func: Functions, grpcSdk: ConduitGrpcSdk) { if (!func.inputs?.event) { throw new GrpcError(status.INVALID_ARGUMENT, 'Event not found'); @@ -193,30 +215,30 @@ export function createEventFunction(func: Functions, grpcSdk: ConduitGrpcSdk) { parsedData = JSON.parse(data); } - executeFunction( + void executeBackgroundFunction( func, - { request: parsedData }, + parsedData, compiledFunctionCode, - func.timeout, - func.name, grpcSdk, - ) - .then(() => { - ConduitGrpcSdk.Logger.log( - `Executed ${func.name} for event ${func.inputs.event}`, - ); - }) - .catch(err => { - ConduitGrpcSdk.Logger.error( - `Execution failed for ${func.name} for event ${func.inputs.event} with: `, - ); - ConduitGrpcSdk.Logger.error(err); - }); + ).then(() => { + ConduitGrpcSdk.Logger.log(`Executed ${func.name} for event ${func.inputs.event}`); + }); }, func._id, ); } +export function validateCronFunctionConfig(func: Functions): void { + const pattern = getCronPatternFromInputs(func.inputs); + if (!pattern) { + throw new GrpcError( + status.INVALID_ARGUMENT, + 'Cron pattern is required (inputs.cronPattern or inputs.event)', + ); + } + validateCronPattern(pattern); +} + export function createFunctionRoute(func: Functions, grpcSdk: ConduitGrpcSdk) { switch (func.functionType) { case 'request': @@ -224,7 +246,8 @@ export function createFunctionRoute(func: Functions, grpcSdk: ConduitGrpcSdk) { case 'webhook': return createRequestOrWebhookFunction(func, grpcSdk); case 'cron': - //todo + validateCronFunctionConfig(func); + return null; case 'event': createEventFunction(func, grpcSdk); return null; diff --git a/modules/functions/src/interfaces/IWebInputs.interface.ts b/modules/functions/src/interfaces/IWebInputs.interface.ts index d521a6a9a..544b4fda0 100644 --- a/modules/functions/src/interfaces/IWebInputs.interface.ts +++ b/modules/functions/src/interfaces/IWebInputs.interface.ts @@ -7,6 +7,8 @@ import { export interface IWebInputsInterface { method?: 'GET' | 'POST' | 'PUT' | 'DELETE' | 'PATCH'; event?: string; + cronPattern?: string; + timezone?: string; bodyParams?: ConduitModel; urlParams?: ConduitUrlParams; queryParams?: ConduitQueryParams; diff --git a/modules/functions/src/migrations/index.ts b/modules/functions/src/migrations/index.ts index a58e7c052..e9c337f89 100644 --- a/modules/functions/src/migrations/index.ts +++ b/modules/functions/src/migrations/index.ts @@ -1,3 +1,24 @@ -export async function runMigrations() { - // ... +import { ConduitGrpcSdk } from '@conduitplatform/grpc-sdk'; +import { Functions } from '../models/index.js'; +import { getCronPatternFromInputs } from '../controllers/cron.utils.js'; + +export async function runMigrations(grpcSdk: ConduitGrpcSdk) { + const cronFunctions = await Functions.getInstance().findMany({ + functionType: 'cron', + }); + for (const func of cronFunctions) { + const pattern = getCronPatternFromInputs(func.inputs); + if (!pattern || func.inputs?.cronPattern) { + continue; + } + await Functions.getInstance().findByIdAndUpdate(func._id, { + inputs: { + ...func.inputs, + cronPattern: pattern, + }, + }); + ConduitGrpcSdk.Logger.log( + `Migrated cron function ${func.name} (${func._id}) to inputs.cronPattern`, + ); + } } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index d493bd2f1..3b6308f85 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -892,12 +892,21 @@ importers: axios: specifier: ^1.18.0 version: 1.18.0(debug@4.4.3) + bullmq: + specifier: ^5.79.0 + version: 5.79.0 convict: specifier: ^6.2.5 version: 6.2.5 + cron-parser: + specifier: ^4.9.0 + version: 4.9.0 escape-string-regexp: specifier: ^5.0.0 version: 5.0.0 + ioredis: + specifier: 5.11.1 + version: 5.11.1 lodash-es: specifier: ^4.18.1 version: 4.18.1