Fix workflow cron trigger cache stuck without TTL (#20812)
## 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 <img width="1627" height="721" alt="Screenshot 2026-05-21 at 16 46 13" src="https://github.com/user-attachments/assets/8262dd5f-fbbd-43c9-aede-c0ce5d6a0f59" /> --------- Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+25
@@ -342,6 +342,31 @@ end`;
|
||||
}) as Promise<number>;
|
||||
}
|
||||
|
||||
async hashSetWithExpire({
|
||||
key,
|
||||
field,
|
||||
value,
|
||||
ttlMs,
|
||||
}: {
|
||||
key: string;
|
||||
field: string;
|
||||
value: string;
|
||||
ttlMs: Milliseconds;
|
||||
}): Promise<void> {
|
||||
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,
|
||||
|
||||
+37
-25
@@ -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();
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
+3
-8
@@ -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`);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user