From 869f8af4f0c47997256bdfb80eba42e5a1194647 Mon Sep 17 00:00:00 2001 From: Marie <51697796+ijreilly@users.noreply.github.com> Date: Thu, 21 May 2026 18:03:09 +0200 Subject: [PATCH] Fix workflow cron trigger cache stuck without TTL (#20812) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Problem The cron-trigger cache key (`module:workflow:workflow-cron-triggers`) can get stuck without a TTL, silently halting **all** cron-triggered workflows for a whole tenant until the key is manually deleted from Redis. Repro path: 1. Cache miss → DB-scan branch runs. 2. Inner loop writes triggers via `hashSet` (creates the key, **no TTL yet**). 3. Worker crashes / OOMs / gets killed by a deploy between any `hashSet` and the trailing `expire(1h)` call. 4. Key now exists with TTL = `-1` and a partial set of fields. 5. Next tick: `hashGetValues` returns those fields → `cachedValues.length > 0` → **cache-hit branch** → `expire` is never called. 6. Key has no TTL, so it never auto-expires. The DB-scan branch never runs again. New / missing triggers are never picked up. Workflows go silent. Observed in production: cache key with `TTL: no limit` and 121 fields. Deleting the key restored normal behaviour (next tick rebuilt with TTL ~3600). ## Fix Set the TTL right after first value is added ## Monitoring Added a "Cache miss" log count in workflow dashboard, counted among the last 6 hours. Turns green if >= 5 Screenshot 2026-05-21 at 16 46 13 --------- Co-authored-by: Cursor --- .../services/cache-storage.service.ts | 25 ++++++++ .../workflow-cron-trigger-cron.job.spec.ts | 62 +++++++++++-------- .../jobs/workflow-cron-trigger-cron.job.ts | 11 +--- 3 files changed, 65 insertions(+), 33 deletions(-) diff --git a/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts b/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts index 9ed6d0e023..74562bdc36 100644 --- a/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts +++ b/packages/twenty-server/src/engine/core-modules/cache-storage/services/cache-storage.service.ts @@ -342,6 +342,31 @@ end`; }) as Promise; } + async hashSetWithExpire({ + key, + field, + value, + ttlMs, + }: { + key: string; + field: string; + value: string; + ttlMs: Milliseconds; + }): Promise { + if (!this.isRedisCache()) { + throw new Error('hashSetWithExpire is only supported with Redis cache'); + } + + const redisClient = (this.cache as RedisCache).store.client; + const prefixedKey = this.getKey(key); + + await redisClient + .multi() + .hSet(prefixedKey, field, value) + .pExpire(prefixedKey, ttlMs) + .exec(); + } + async hashDelete({ key, field, 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 52abb33330..feed6255bb 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 @@ -32,7 +32,7 @@ const mockExceptionHandlerService = { const mockCacheStorageService = { hashGetValues: jest.fn(), hashSet: jest.fn(), - expire: jest.fn(), + hashSetWithExpire: jest.fn(), }; describe('WorkflowCronTriggerCronJob', () => { @@ -156,6 +156,7 @@ describe('WorkflowCronTriggerCronJob', () => { await job.handle(); expect(mockCacheStorageService.hashSet).not.toHaveBeenCalled(); + expect(mockCacheStorageService.hashSetWithExpire).not.toHaveBeenCalled(); }); }); @@ -175,7 +176,7 @@ describe('WorkflowCronTriggerCronJob', () => { expect(mockCoreDataSource.query).toHaveBeenCalledTimes(3); }); - it('should write each trigger to cache immediately and set TTL', async () => { + it('should use hashSetWithExpire for every trigger so the TTL is always set', async () => { mockCacheStorageService.hashGetValues.mockResolvedValue([]); mockWorkspaceRepository.find.mockResolvedValue([ { id: WORKSPACE_1 }, @@ -202,28 +203,39 @@ describe('WorkflowCronTriggerCronJob', () => { await job.handle(); - expect(mockCacheStorageService.hashSet).toHaveBeenCalledTimes(2); - expect(mockCacheStorageService.hashSet).toHaveBeenCalledWith({ - key: WORKFLOW_CRON_TRIGGER_CACHE_KEY, - field: 'workflow-1', - value: JSON.stringify({ - workspaceId: WORKSPACE_1, - workflowId: 'workflow-1', - pattern: '* * * * *', - }), - }); - expect(mockCacheStorageService.hashSet).toHaveBeenCalledWith({ - key: WORKFLOW_CRON_TRIGGER_CACHE_KEY, - field: 'workflow-2', - value: JSON.stringify({ - workspaceId: WORKSPACE_3, - workflowId: 'workflow-2', - pattern: '* * * * *', - }), - }); - expect(mockCacheStorageService.expire).toHaveBeenCalledWith( - WORKFLOW_CRON_TRIGGER_CACHE_KEY, - WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS, + // hashSet must never be called during a rebuild: any non-atomic write + // could recreate the key without a TTL if it was flushed/evicted + // mid-rebuild, leaving cron workflows permanently stuck. + expect(mockCacheStorageService.hashSet).not.toHaveBeenCalled(); + + expect(mockCacheStorageService.hashSetWithExpire).toHaveBeenCalledTimes( + 2, + ); + expect(mockCacheStorageService.hashSetWithExpire).toHaveBeenNthCalledWith( + 1, + { + key: WORKFLOW_CRON_TRIGGER_CACHE_KEY, + field: 'workflow-1', + value: JSON.stringify({ + workspaceId: WORKSPACE_1, + workflowId: 'workflow-1', + pattern: '* * * * *', + }), + ttlMs: WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS, + }, + ); + expect(mockCacheStorageService.hashSetWithExpire).toHaveBeenNthCalledWith( + 2, + { + key: WORKFLOW_CRON_TRIGGER_CACHE_KEY, + field: 'workflow-2', + value: JSON.stringify({ + workspaceId: WORKSPACE_3, + workflowId: 'workflow-2', + pattern: '* * * * *', + }), + ttlMs: WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS, + }, ); }); @@ -235,7 +247,7 @@ describe('WorkflowCronTriggerCronJob', () => { await job.handle(); expect(mockCacheStorageService.hashSet).not.toHaveBeenCalled(); - expect(mockCacheStorageService.expire).not.toHaveBeenCalled(); + expect(mockCacheStorageService.hashSetWithExpire).not.toHaveBeenCalled(); }); }); 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 85b084cee2..2d81137b59 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 @@ -125,22 +125,17 @@ export class WorkflowCronTriggerCronJob { ); for (const trigger of triggersToCache) { - await this.cacheStorageService.hashSet({ + await this.cacheStorageService.hashSetWithExpire({ key: WORKFLOW_CRON_TRIGGER_CACHE_KEY, field: trigger.workflowId, value: JSON.stringify(trigger), + ttlMs: WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS, }); + triggerCount++; } } - if (triggerCount > 0) { - await this.cacheStorageService.expire( - WORKFLOW_CRON_TRIGGER_CACHE_KEY, - WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS, - ); - } - this.logger.log(`Cache rebuilt with ${triggerCount} cron triggers`); }