fix(server): dispatch each cron trigger exactly once (#22113)
## Problem
App/logic-function crons occasionally fire **twice, ~1 minute apart**.
The most visible symptom is a notification cron sending the same Discord
DM (or channel post) at e.g. `17:00` and again at `17:01`.
## Root cause
`CronTriggerCronJob` runs every minute (`* * * * *`) and re-dispatches
any logic function whose pattern is "due" according to `shouldRunNow`:
```ts
const diff = Math.abs(prevTriggerDate.getTime() - now.getTime());
return diff < rootCronIntervalMs; // 60_000
```
The detection window (`60_000ms`) is **equal to** the 60s tick interval.
So when a root tick drifts across a minute boundary (runs slightly
early/late, or BullMQ fires a catch-up), two adjacent ticks can both see
the *same* trigger as "within the last 60s" and each enqueue a
`LogicFunctionTriggerJob`. The dispatch isn't idempotent, so the
function runs twice.
## Fix
Make dispatch idempotent, keyed on the trigger itself:
- New `getMatchingTriggerTimestamp(pattern, now)` returns the epoch-ms
of the matched trigger (stable regardless of *when* within the window
the root job runs), or `null`. `shouldRunNow` now delegates to it —
behaviour unchanged.
- Before enqueuing, `CronTriggerCronJob` claims a
`logic-function-cron:{workspace}:{function}:{triggerTs}` key in the
`EngineLock` cache. A second tick that resolves to the same trigger
finds the key and skips.
Distinct triggers always have distinct timestamps (hence distinct keys),
so a later legitimate run is never suppressed. The TTL (2 min) only
needs to outlive the detection window.
## Notes
- `WorkflowCronTriggerCronJob` uses the same `shouldRunNow` pattern and
has the same latent double-dispatch; left out of this PR to keep it
focused, but the new helper makes the same guard a small follow-up.
- The cache `get`-then-`set` isn't atomic; for the observed failure mode
(ticks ~1 min apart, sequential) it's reliable. A Redis `SET NX` would
also close the rare concurrent-multi-instance race.
## Test plan
- [x] `should-run-now.utils.spec.ts` extended: two ticks within one
window resolve to the same timestamp; out-of-window and invalid patterns
return `null`. All 8 pass.
- [x] `oxlint --type-aware` + `oxfmt` clean on changed files.
<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/22113?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->
This commit is contained in:
@@ -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 {}
|
||||
+50
@@ -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<boolean> {
|
||||
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,
|
||||
);
|
||||
}
|
||||
}
|
||||
+2
@@ -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,
|
||||
|
||||
+10
-2
@@ -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<WorkspaceEntity>,
|
||||
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;
|
||||
}
|
||||
|
||||
|
||||
+2
@@ -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,
|
||||
],
|
||||
|
||||
+14
-1
@@ -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,
|
||||
|
||||
+18
-3
@@ -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}`,
|
||||
);
|
||||
|
||||
@@ -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,
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -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;
|
||||
}
|
||||
};
|
||||
Reference in New Issue
Block a user