Skip to content

Commit 9df9aa9

Browse files
committed
fix(redis): revoke cron claims when the next run is updated
A cron claim clears next_run_at and stores a claim token. The worker then computes the next run in JS and finalizes it only if next_run_at is still empty, the cron expression and config revision are unchanged, and the token still matches. updateSchedule({ nextRunAt: null }) writes the same empty next_run_at but kept the token and the revision, so the finalizer could not tell the difference and wrote the next run back. updateSchedule() now removes the claim token whenever it sets next_run_at.
1 parent d560a4a commit 9df9aa9

3 files changed

Lines changed: 68 additions & 4 deletions

File tree

‎.changelog/redis-schedule-due-index.md‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,10 @@ claiming schedules maintain the derived index, while claiming repairs stale entr
1111
and index disagree.
1212

1313
Schedule hash and index writes are atomic, and index consistency is preserved across concurrent
14-
lifecycle changes. Cron finalization now rejects stale calculations when a schedule is deleted or
15-
reconfigured while its next run is being calculated. Paused schedules retain their next run but stay
16-
out of the due index until resumed. Claiming also discards malformed due scores instead of allowing
17-
one corrupt entry to block later schedules.
14+
lifecycle changes. Cron finalization now rejects stale calculations when a schedule is deleted,
15+
reconfigured, or has its next run cleared during the calculation. Paused schedules retain their
16+
next run but stay out of the due index until resumed. Claiming also discards malformed due scores
17+
instead of allowing one corrupt entry to block later schedules.
1818

1919
## Bug Fix
2020

‎src/drivers/redis_scripts.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -575,6 +575,12 @@ ${SCHEDULE_DUE_INDEX_LUA}
575575
redis.call('HSET', schedule_key, field, value)
576576
end
577577
578+
-- An explicit next run supersedes an in-flight cron claim, so the claim's
579+
-- finalization must not overwrite it.
580+
if updates.next_run_at ~= nil then
581+
redis.call('HDEL', schedule_key, 'claim_token')
582+
end
583+
578584
sync_schedule_due_index(schedule_key, due_key, id)
579585
580586
return 1

‎tests/adapter.spec.ts‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1333,6 +1333,64 @@ test.group('Adapter | Redis', (group) => {
13331333
)
13341334
})
13351335

1336+
test('cron finalization does not undo a cleared next run', async ({ assert, cleanup }) => {
1337+
const secondConnection = new Redis({
1338+
host: process.env.REDIS_HOST || 'localhost',
1339+
port: Number.parseInt(process.env.REDIS_PORT || '6379', 10),
1340+
keyPrefix: KEY_PREFIX,
1341+
db: 15,
1342+
})
1343+
cleanup(async () => {
1344+
await secondConnection.quit()
1345+
})
1346+
1347+
const adapter = new RedisAdapter(connection)
1348+
const secondAdapter = new RedisAdapter(secondConnection)
1349+
const id = 'cron-finalize-cleared-next-run'
1350+
1351+
await adapter.upsertSchedule({
1352+
id,
1353+
name: 'CronClearedJob',
1354+
payload: {},
1355+
cronExpression: '* * * * *',
1356+
timezone: 'UTC',
1357+
})
1358+
await adapter.updateSchedule(id, { nextRunAt: new Date(Date.now() - 1_000) })
1359+
1360+
const originalEval = connection.eval.bind(connection)
1361+
let releaseClaim!: () => void
1362+
let claimReturned!: () => void
1363+
const claimReleased = new Promise<void>((resolve) => {
1364+
releaseClaim = resolve
1365+
})
1366+
const claimHasReturned = new Promise<void>((resolve) => {
1367+
claimReturned = resolve
1368+
})
1369+
let gateNextEval = true
1370+
1371+
connection.eval = (async (...args: Parameters<typeof connection.eval>) => {
1372+
const result = await originalEval(...args)
1373+
if (gateNextEval) {
1374+
gateNextEval = false
1375+
claimReturned()
1376+
await claimReleased
1377+
}
1378+
return result
1379+
}) as typeof connection.eval
1380+
1381+
const claim = adapter.claimDueSchedule()
1382+
await claimHasReturned
1383+
await secondAdapter.updateSchedule(id, { nextRunAt: null })
1384+
1385+
releaseClaim()
1386+
await claim
1387+
connection.eval = originalEval
1388+
1389+
const schedule = await adapter.getSchedule(id)
1390+
assert.isNull(schedule!.nextRunAt)
1391+
assert.isNull(await connection.zscore('schedules::due', id))
1392+
})
1393+
13361394
test('stale cron finalization cannot modify a recreated schedule claim', async ({
13371395
assert,
13381396
cleanup,

0 commit comments

Comments
 (0)