local rcall = redis.call
local jobIdKey = KEYS[12]
if rcall("EXISTS", jobIdKey) == 1 then if ARGV[5] == "completed" then
if rcall("SCARD", jobIdKey .. ":dependencies") ~= 0 then
return -4
end
if rcall("ZCARD", jobIdKey .. ":unsuccessful") ~= 0 then
return -9
end
end
local opts = cmsgpack.unpack(ARGV[8])
local token = opts['token']
local errorCode = removeLock(jobIdKey, KEYS[5], token, ARGV[1])
if errorCode < 0 then
return errorCode
end
updateJobFields(jobIdKey, ARGV[9]);
local attempts = opts['attempts']
local maxMetricsSize = opts['maxMetricsSize']
local maxCount = opts['keepJobs']['count']
local maxAge = opts['keepJobs']['age']
local maxLimit = opts['keepJobs']['limit'] or 1000
local jobAttributes = rcall("HMGET", jobIdKey, "parentKey", "parent", "deid")
local parentKey = jobAttributes[1] or ""
local parentId = ""
local parentQueueKey = ""
if jobAttributes[2] then local jsonDecodedParent = cjson.decode(jobAttributes[2])
parentId = jsonDecodedParent['id']
parentQueueKey = jsonDecodedParent['queueKey']
end
local jobId = ARGV[1]
local timestamp = ARGV[2]
local numRemovedElements = rcall("LREM", KEYS[2], -1, jobId)
if (numRemovedElements < 1) then
return -3
end
local eventStreamKey = KEYS[4]
local metaKey = KEYS[9]
trimEvents(metaKey, eventStreamKey)
local prefix = ARGV[7]
removeDeduplicationKeyIfNeededOnFinalization(prefix, jobAttributes[3], jobId)
if parentId == "" and parentKey ~= "" then
parentId = getJobIdFromKey(parentKey)
parentQueueKey = getJobKeyPrefix(parentKey, ":" .. parentId)
end
if parentId ~= "" then
if ARGV[5] == "completed" then
local dependenciesSet = parentKey .. ":dependencies"
if rcall("SREM", dependenciesSet, jobIdKey) == 1 then
updateParentDepsIfNeeded(parentKey, parentQueueKey, dependenciesSet, parentId, jobIdKey, ARGV[4],
timestamp)
end
else
moveChildFromDependenciesIfNeeded(jobAttributes[2], jobIdKey, ARGV[4], timestamp)
end
end
local attemptsMade = rcall("HINCRBY", jobIdKey, "atm", 1)
if maxCount ~= 0 then
local targetSet = KEYS[11]
rcall("ZADD", targetSet, timestamp, jobId)
rcall("HSET", jobIdKey, ARGV[3], ARGV[4], "finishedOn", timestamp)
if ARGV[5] == "failed" then
rcall("HDEL", jobIdKey, "defa")
end
if maxAge ~= nil then
removeJobsByMaxAge(timestamp, maxAge, targetSet, prefix, maxLimit)
end
if maxCount ~= nil and maxCount > 0 then
removeJobsByMaxCount(maxCount, targetSet, prefix)
end
else
removeJobKeys(jobIdKey)
if parentKey ~= "" then
removeParentDependencyKey(jobIdKey, false, parentKey, jobAttributes[3])
end
end
rcall("XADD", eventStreamKey, "*", "event", ARGV[5], "jobId", jobId, ARGV[3], ARGV[4], "prev", "active")
if ARGV[5] == "failed" then
if tonumber(attemptsMade) >= tonumber(attempts) then
rcall("XADD", eventStreamKey, "*", "event", "retries-exhausted", "jobId", jobId, "attemptsMade",
attemptsMade)
end
end
if maxMetricsSize ~= "" then
collectMetrics(KEYS[13], KEYS[13] .. ':data', maxMetricsSize, timestamp)
end
if (ARGV[6] == "1") then
local target, isPausedOrMaxed, rateLimitMax, rateLimitDuration = getTargetQueueList(metaKey, KEYS[2],
KEYS[1], KEYS[8])
local markerKey = KEYS[14]
promoteDelayedJobs(KEYS[7], markerKey, target, KEYS[3], eventStreamKey, prefix, timestamp, KEYS[10],
isPausedOrMaxed)
local maxJobs = tonumber(rateLimitMax or (opts['limiter'] and opts['limiter']['max']))
local expireTime = getRateLimitTTL(maxJobs, KEYS[6])
if expireTime > 0 then
return {0, 0, expireTime, 0}
end
if isPausedOrMaxed then
return {0, 0, 0, 0}
end
local limiterDuration = (opts['limiter'] and opts['limiter']['duration']) or rateLimitDuration
jobId = rcall("RPOPLPUSH", KEYS[1], KEYS[2])
if jobId then
if string.sub(jobId, 1, 2) == "0:" then
rcall("LREM", KEYS[2], 1, jobId)
if jobId == "0:0" then
jobId = moveJobFromPrioritizedToActive(KEYS[3], KEYS[2], KEYS[10])
return prepareJobForProcessing(prefix, KEYS[6], eventStreamKey, jobId, timestamp, maxJobs,
limiterDuration, markerKey, opts)
end
else
return prepareJobForProcessing(prefix, KEYS[6], eventStreamKey, jobId, timestamp, maxJobs,
limiterDuration, markerKey, opts)
end
else
jobId = moveJobFromPrioritizedToActive(KEYS[3], KEYS[2], KEYS[10])
if jobId then
return prepareJobForProcessing(prefix, KEYS[6], eventStreamKey, jobId, timestamp, maxJobs,
limiterDuration, markerKey, opts)
end
end
local nextTimestamp = getNextDelayedTimestamp(KEYS[7])
if nextTimestamp ~= nil then
return {0, 0, 0, nextTimestamp}
end
end
local waitLen = rcall("LLEN", KEYS[1])
if waitLen == 0 then
local activeLen = rcall("LLEN", KEYS[2])
if activeLen == 0 then
local prioritizedLen = rcall("ZCARD", KEYS[3])
if prioritizedLen == 0 then
rcall("XADD", eventStreamKey, "*", "event", "drained")
end
end
end
return 0
else
return -1
end