| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849 |
- --[[
- Function to move job from wait state to active.
- Input:
- opts - token - lock token
- opts - lockDuration
- opts - limiter
- ]]
- -- Includes
- --- @include "addBaseMarkerIfNeeded"
- local function prepareJobForProcessing(keyPrefix, rateLimiterKey, eventStreamKey,
- jobId, processedOn, maxJobs, limiterDuration, markerKey, opts)
- local jobKey = keyPrefix .. jobId
- -- Check if we need to perform rate limiting.
- if maxJobs then
- local jobCounter = tonumber(rcall("INCR", rateLimiterKey))
- if jobCounter == 1 then
- local integerDuration = math.floor(math.abs(limiterDuration))
- rcall("PEXPIRE", rateLimiterKey, integerDuration)
- end
- end
- -- get a lock
- if opts['token'] ~= "0" then
- local lockKey = jobKey .. ':lock'
- rcall("SET", lockKey, opts['token'], "PX", opts['lockDuration'])
- end
- local optionalValues = {}
- if opts['name'] then
- -- Set "processedBy" field to the worker name
- table.insert(optionalValues, "pb")
- table.insert(optionalValues, opts['name'])
- end
- rcall("XADD", eventStreamKey, "*", "event", "active", "jobId", jobId, "prev", "waiting")
- rcall("HMSET", jobKey, "processedOn", processedOn, unpack(optionalValues))
- rcall("HINCRBY", jobKey, "ats", 1)
- addBaseMarkerIfNeeded(markerKey, false)
- -- rate limit delay must be 0 in this case to prevent adding more delay
- -- when job that is moved to active needs to be processed
- return {rcall("HGETALL", jobKey), jobId, 0, 0} -- get job data
- end
|