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`); }