Skip to content

Commit 821b3c1

Browse files
committed
fix: dangling jobs when batched
1 parent dcc7870 commit 821b3c1

2 files changed

Lines changed: 0 additions & 86 deletions

File tree

src/lua/reserve-batch.lua

Lines changed: 0 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -14,49 +14,6 @@ end
1414

1515
local out = {}
1616

17-
-- STALLED JOB RECOVERY WITH THROTTLING
18-
-- Check for stalled jobs periodically to avoid overhead in hot path
19-
-- This ensures stalled jobs are recovered even in high-load systems where ready queue is never empty
20-
-- Check interval is adaptive: 1/4 of jobTimeout (to check 4x during visibility window), max 5s
21-
local stalledCheckKey = ns .. ":stalled:lastcheck"
22-
local lastCheck = tonumber(redis.call("GET", stalledCheckKey)) or 0
23-
local stalledCheckInterval = math.min(math.floor(vt / 4), 5000)
24-
25-
if (now - lastCheck) >= stalledCheckInterval then
26-
-- Update last check timestamp
27-
redis.call("SET", stalledCheckKey, tostring(now))
28-
29-
-- Check for expired jobs and recover them
30-
local expiredJobs = redis.call("ZRANGEBYSCORE", processingKey, 0, now)
31-
if #expiredJobs > 0 then
32-
for _, jobId in ipairs(expiredJobs) do
33-
local deadlineAt = tonumber(redis.call("ZSCORE", processingKey, jobId))
34-
local gid = redis.call("HGET", ns .. ":job:" .. jobId, "groupId")
35-
if gid and deadlineAt and now > deadlineAt then
36-
local jobKey = ns .. ":job:" .. jobId
37-
local jobScore = redis.call("HGET", jobKey, "score")
38-
if jobScore then
39-
local gZ = ns .. ":g:" .. gid
40-
-- Remove from group active list BEFORE re-adding to group set
41-
-- This prevents the job from blocking the group after recovery
42-
local groupActiveKey = ns .. ":g:" .. gid .. ":active"
43-
redis.call("LREM", groupActiveKey, 1, jobId)
44-
redis.call("ZADD", gZ, tonumber(jobScore), jobId)
45-
-- Reset status so the job is visible as waiting again
46-
redis.call("HSET", jobKey, "status", "waiting")
47-
local head = redis.call("ZRANGE", gZ, 0, 0, "WITHSCORES")
48-
if head and #head >= 2 then
49-
local headScore = tonumber(head[2])
50-
redis.call("ZADD", readyKey, headScore, gid)
51-
end
52-
redis.call("DEL", ns .. ":lock:" .. gid)
53-
redis.call("ZREM", processingKey, jobId)
54-
end
55-
end
56-
end
57-
end
58-
end
59-
6017
-- Pop up to maxBatch groups from ready set (lowest score first)
6118
local groups = redis.call("ZRANGE", readyKey, 0, maxBatch - 1, "WITHSCORES")
6219
if not groups or #groups == 0 then

src/lua/scripts.generated.ts

Lines changed: 0 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -1490,49 +1490,6 @@ end
14901490
14911491
local out = {}
14921492
1493-
-- STALLED JOB RECOVERY WITH THROTTLING
1494-
-- Check for stalled jobs periodically to avoid overhead in hot path
1495-
-- This ensures stalled jobs are recovered even in high-load systems where ready queue is never empty
1496-
-- Check interval is adaptive: 1/4 of jobTimeout (to check 4x during visibility window), max 5s
1497-
local stalledCheckKey = ns .. ":stalled:lastcheck"
1498-
local lastCheck = tonumber(redis.call("GET", stalledCheckKey)) or 0
1499-
local stalledCheckInterval = math.min(math.floor(vt / 4), 5000)
1500-
1501-
if (now - lastCheck) >= stalledCheckInterval then
1502-
-- Update last check timestamp
1503-
redis.call("SET", stalledCheckKey, tostring(now))
1504-
1505-
-- Check for expired jobs and recover them
1506-
local expiredJobs = redis.call("ZRANGEBYSCORE", processingKey, 0, now)
1507-
if #expiredJobs > 0 then
1508-
for _, jobId in ipairs(expiredJobs) do
1509-
local deadlineAt = tonumber(redis.call("ZSCORE", processingKey, jobId))
1510-
local gid = redis.call("HGET", ns .. ":job:" .. jobId, "groupId")
1511-
if gid and deadlineAt and now > deadlineAt then
1512-
local jobKey = ns .. ":job:" .. jobId
1513-
local jobScore = redis.call("HGET", jobKey, "score")
1514-
if jobScore then
1515-
local gZ = ns .. ":g:" .. gid
1516-
-- Remove from group active list BEFORE re-adding to group set
1517-
-- This prevents the job from blocking the group after recovery
1518-
local groupActiveKey = ns .. ":g:" .. gid .. ":active"
1519-
redis.call("LREM", groupActiveKey, 1, jobId)
1520-
redis.call("ZADD", gZ, tonumber(jobScore), jobId)
1521-
-- Reset status so the job is visible as waiting again
1522-
redis.call("HSET", jobKey, "status", "waiting")
1523-
local head = redis.call("ZRANGE", gZ, 0, 0, "WITHSCORES")
1524-
if head and #head >= 2 then
1525-
local headScore = tonumber(head[2])
1526-
redis.call("ZADD", readyKey, headScore, gid)
1527-
end
1528-
redis.call("DEL", ns .. ":lock:" .. gid)
1529-
redis.call("ZREM", processingKey, jobId)
1530-
end
1531-
end
1532-
end
1533-
end
1534-
end
1535-
15361493
-- Pop up to maxBatch groups from ready set (lowest score first)
15371494
local groups = redis.call("ZRANGE", readyKey, 0, maxBatch - 1, "WITHSCORES")
15381495
if not groups or #groups == 0 then

0 commit comments

Comments
 (0)