diff --git a/packages/twenty-server/src/engine/core-modules/cron/cron.module.ts b/packages/twenty-server/src/engine/core-modules/cron/cron.module.ts new file mode 100644 index 0000000000..27a1779eea --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/cron/cron.module.ts @@ -0,0 +1,9 @@ +import { Module } from '@nestjs/common'; + +import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; + +@Module({ + providers: [CronTriggerDeduplicationService], + exports: [CronTriggerDeduplicationService], +}) +export class CronModule {} diff --git a/packages/twenty-server/src/engine/core-modules/cron/services/cron-trigger-deduplication.service.ts b/packages/twenty-server/src/engine/core-modules/cron/services/cron-trigger-deduplication.service.ts new file mode 100644 index 0000000000..be78949b6c --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/cron/services/cron-trigger-deduplication.service.ts @@ -0,0 +1,50 @@ +import { Injectable } from '@nestjs/common'; + +import { CronExpressionParser } from 'cron-parser'; + +import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator'; +import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; +import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; + +const ROOT_CRON_INTERVAL_MS = 60_000; +const CRON_DISPATCH_DEDUP_TTL_MS = 2 * 60_000; + +@Injectable() +export class CronTriggerDeduplicationService { + constructor( + @InjectCacheStorage(CacheStorageNamespace.EngineLock) + private readonly cacheStorageService: CacheStorageService, + ) {} + + async shouldDispatch( + keyPrefix: string, + pattern: string, + now: Date, + ): Promise { + let lastTriggerTimestamp: number; + + try { + lastTriggerTimestamp = CronExpressionParser.parse(pattern, { + currentDate: now, + }) + .prev() + .getTime(); + } catch { + return false; + } + + const isDueWithinThisTick = + now.getTime() - lastTriggerTimestamp < ROOT_CRON_INTERVAL_MS; + + if (!isDueWithinThisTick) { + return false; + } + + const dedupKey = `${keyPrefix}:${lastTriggerTimestamp}`; + + return this.cacheStorageService.acquireLock( + dedupKey, + CRON_DISPATCH_DEDUP_TTL_MS, + ); + } +} diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/logic-function-trigger.module.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/logic-function-trigger.module.ts index 9475c3a504..d9ce5ef75c 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/logic-function-trigger.module.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/logic-function-trigger.module.ts @@ -2,6 +2,7 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { TokenModule } from 'src/engine/core-modules/auth/token/token.module'; +import { CronModule } from 'src/engine/core-modules/cron/cron.module'; import { WorkspaceDomainsModule } from 'src/engine/core-modules/domain/workspace-domains/workspace-domains.module'; import { LogicFunctionTriggerJob } from 'src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job'; import { CronTriggerCronCommand } from 'src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.command'; @@ -19,6 +20,7 @@ import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache TokenModule, WorkspaceDomainsModule, WorkspaceCacheModule, + CronModule, ], providers: [ LogicFunctionTriggerJob, diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts index ee22cf8cb0..5fc29782d7 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts @@ -5,6 +5,7 @@ import { isDefined } from 'twenty-shared/utils'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; import { Repository } from 'typeorm'; +import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; @@ -18,7 +19,6 @@ import { LogicFunctionTriggerJobData, } from 'src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job'; import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; -import { shouldRunNow } from 'src/utils/should-run-now.utils'; export const CRON_TRIGGER_CRON_PATTERN = '* * * * *'; @@ -33,6 +33,7 @@ export class CronTriggerCronJob { private readonly workspaceRepository: Repository, private readonly workspaceCacheService: WorkspaceCacheService, private readonly exceptionHandlerService: ExceptionHandlerService, + private readonly cronTriggerDeduplicationService: CronTriggerDeduplicationService, ) {} @Process(CronTriggerCronJob.name) @@ -73,7 +74,14 @@ export class CronTriggerCronJob { continue; } - if (!shouldRunNow(cronSettings.pattern, now)) { + const shouldDispatch = + await this.cronTriggerDeduplicationService.shouldDispatch( + `logic-function-cron:${activeWorkspace.id}:${logicFunction.id}`, + cronSettings.pattern, + now, + ); + + if (!shouldDispatch) { continue; } diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts index 71cc30c885..d39f53dc0e 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts @@ -2,6 +2,7 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { CacheStorageModule } from 'src/engine/core-modules/cache-storage/cache-storage.module'; +import { CronModule } from 'src/engine/core-modules/cron/cron.module'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { WorkspaceDataSourceModule } from 'src/engine/workspace-datasource/workspace-datasource.module'; import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module'; @@ -14,6 +15,7 @@ import { WorkflowDatabaseEventTriggerListener } from 'src/modules/workflow/workf imports: [ TypeOrmModule.forFeature([WorkspaceEntity]), CacheStorageModule, + CronModule, WorkflowCommonModule, WorkspaceDataSourceModule, ], diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts index feed6255bb..2f2b6099b2 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts @@ -2,6 +2,7 @@ import { Test, type TestingModule } from '@nestjs/testing'; import { getDataSourceToken, getRepositoryToken } from '@nestjs/typeorm'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; +import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { WORKFLOW_CRON_TRIGGER_CACHE_KEY } from 'src/modules/workflow/workflow-trigger/automated-trigger/crons/constants/workflow-cron-trigger-cache-key.constant'; @@ -35,6 +36,10 @@ const mockCacheStorageService = { hashSetWithExpire: jest.fn(), }; +const mockCronTriggerDeduplicationService = { + shouldDispatch: jest.fn(), +}; + describe('WorkflowCronTriggerCronJob', () => { let job: WorkflowCronTriggerCronJob; @@ -42,6 +47,7 @@ describe('WorkflowCronTriggerCronJob', () => { jest.clearAllMocks(); jest.useFakeTimers(); jest.setSystemTime(new Date('2026-04-02T15:00:30.000Z')); + mockCronTriggerDeduplicationService.shouldDispatch.mockResolvedValue(true); const module: TestingModule = await Test.createTestingModule({ providers: [ @@ -66,6 +72,10 @@ describe('WorkflowCronTriggerCronJob', () => { provide: CacheStorageNamespace.ModuleWorkflow, useValue: mockCacheStorageService, }, + { + provide: CronTriggerDeduplicationService, + useValue: mockCronTriggerDeduplicationService, + }, ], }).compile(); @@ -117,7 +127,10 @@ describe('WorkflowCronTriggerCronJob', () => { ); }); - it('should not enqueue jobs when cron pattern does not match', async () => { + it('should not enqueue jobs when the trigger is not due', async () => { + mockCronTriggerDeduplicationService.shouldDispatch.mockResolvedValue( + false, + ); mockCacheStorageService.hashGetValues.mockResolvedValue([ JSON.stringify({ workspaceId: WORKSPACE_1, diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts index 2d81137b59..61b3c5a945 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts @@ -9,6 +9,7 @@ import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decora import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; +import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; @@ -26,7 +27,6 @@ import { WorkflowTriggerJob, type WorkflowTriggerJobData, } from 'src/modules/workflow/workflow-trigger/jobs/workflow-trigger.job'; -import { shouldRunNow } from 'src/utils/should-run-now.utils'; export const WORKFLOW_CRON_TRIGGER_CRON_PATTERN = '* * * * *'; @@ -44,6 +44,7 @@ export class WorkflowCronTriggerCronJob { private readonly exceptionHandlerService: ExceptionHandlerService, @InjectCacheStorage(CacheStorageNamespace.ModuleWorkflow) private readonly cacheStorageService: CacheStorageService, + private readonly cronTriggerDeduplicationService: CronTriggerDeduplicationService, ) {} @Process(WorkflowCronTriggerCronJob.name) @@ -82,7 +83,14 @@ export class WorkflowCronTriggerCronJob { continue; } - if (!shouldRunNow(trigger.pattern, now)) { + const shouldDispatch = + await this.cronTriggerDeduplicationService.shouldDispatch( + `workflow-cron:${trigger.workspaceId}:${trigger.workflowId}`, + trigger.pattern, + now, + ); + + if (!shouldDispatch) { continue; } @@ -178,7 +186,14 @@ export class WorkflowCronTriggerCronJob { triggersToCache.push(cachedTrigger); - if (shouldRunNow(settings.pattern, now)) { + const shouldDispatch = + await this.cronTriggerDeduplicationService.shouldDispatch( + `workflow-cron:${workspaceId}:${trigger.workflowId}`, + settings.pattern, + now, + ); + + if (shouldDispatch) { this.logger.log( `Trigger ${trigger.id}: enqueuing WorkflowTriggerJob for workflow ${trigger.workflowId}`, ); diff --git a/packages/twenty-server/src/utils/__test__/should-run-now.utils.spec.ts b/packages/twenty-server/src/utils/__test__/should-run-now.utils.spec.ts deleted file mode 100644 index c1fecc191f..0000000000 --- a/packages/twenty-server/src/utils/__test__/should-run-now.utils.spec.ts +++ /dev/null @@ -1,47 +0,0 @@ -import { shouldRunNow } from 'src/utils/should-run-now.utils'; - -const getNowDate = (hour: string) => { - return new Date(`2025-01-01T${hour}.100Z`); -}; - -describe('shouldRunNow', () => { - it('returns true when now matches cron pattern */1 * * * *', () => { - const cron = '*/1 * * * *'; - - expect(shouldRunNow(cron, getNowDate('10:00:00'))).toBe(true); - }); - - it('returns true with a 50s root cron delay', () => { - const cron = '*/1 * * * *'; - - expect(shouldRunNow(cron, getNowDate('10:00:50'))).toBe(true); - }); - - it('returns true 5 times in a row for a */5 pattern', () => { - const cron = '*/5 * * * *'; // every 5 minutes - - expect(shouldRunNow(cron, getNowDate('09:59:00'))).toBe(false); - expect(shouldRunNow(cron, getNowDate('10:00:00'))).toBe(true); - expect(shouldRunNow(cron, getNowDate('10:01:00'))).toBe(false); - expect(shouldRunNow(cron, getNowDate('10:02:00'))).toBe(false); - expect(shouldRunNow(cron, getNowDate('10:03:00'))).toBe(false); - expect(shouldRunNow(cron, getNowDate('10:04:00'))).toBe(false); - expect(shouldRunNow(cron, getNowDate('10:05:00'))).toBe(true); - expect(shouldRunNow(cron, getNowDate('10:06:00'))).toBe(false); - }); - - it('returns false for invalid cron pattern', () => { - const cron = 'invalid-cron'; - - expect(shouldRunNow(cron, getNowDate('10:00:00'))).toBe(false); - }); - - it('returns false if the next run is outside the interval window (2 minutes)', () => { - const cron = '*/10 * * * *'; // every 10 minutes - const interval2min = 2 * 60_000; - - expect(shouldRunNow(cron, getNowDate('10:06:00'), interval2min)).toBe( - false, - ); - }); -}); diff --git a/packages/twenty-server/src/utils/should-run-now.utils.ts b/packages/twenty-server/src/utils/should-run-now.utils.ts deleted file mode 100644 index 9605678209..0000000000 --- a/packages/twenty-server/src/utils/should-run-now.utils.ts +++ /dev/null @@ -1,20 +0,0 @@ -import { CronExpressionParser } from 'cron-parser'; - -export const shouldRunNow = ( - pattern: string, - now: Date, - rootCronIntervalMs = 60_000, -) => { - try { - const interval = CronExpressionParser.parse(pattern, { - currentDate: now, - }); - - const prevTriggerDate = interval.prev(); - const diff = Math.abs(prevTriggerDate.getTime() - now.getTime()); - - return diff < rootCronIntervalMs; - } catch { - return false; - } -};