mirror of
https://github.com/mperham/sidekiq.git
synced 2022-11-09 13:52:34 -05:00
Scheduled and retry jobs add an enqueued_at time when they are put into the queue
This commit is contained in:
parent
a264d2acd6
commit
fabe0e0fed
2 changed files with 29 additions and 23 deletions
|
@ -33,14 +33,15 @@ module Sidekiq
|
|||
# going wrong between the time jobs are popped from the scheduled queue and when
|
||||
# they are pushed onto a work queue and losing the jobs.
|
||||
while message = conn.zrangebyscore(sorted_set, '-inf', now, :limit => [0, 1]).first do
|
||||
msg = Sidekiq.load_json(message)
|
||||
msg = Sidekiq.load_json(message).tap { |msg| msg['enqueued_at'] = Time.now.to_f }
|
||||
|
||||
# Pop item off the queue and add it to the work queue. If the job can't be popped from
|
||||
# the queue, it's because another process already popped it so we can move on to the
|
||||
# next one.
|
||||
if conn.zrem(sorted_set, message)
|
||||
conn.multi do
|
||||
conn.sadd('queues', msg['queue'])
|
||||
conn.lpush("queue:#{msg['queue']}", message)
|
||||
conn.lpush("queue:#{msg['queue']}", Sidekiq.dump_json(msg))
|
||||
end
|
||||
logger.debug { "enqueued #{sorted_set}: #{message}" }
|
||||
end
|
||||
|
|
|
@ -17,30 +17,35 @@ class TestScheduled < Minitest::Test
|
|||
end
|
||||
|
||||
it 'should empty the retry and scheduled queues up to the current time' do
|
||||
Sidekiq.redis do |conn|
|
||||
error_1 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["error_1"], 'queue' => 'queue_1')
|
||||
error_2 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["error_2"], 'queue' => 'queue_2')
|
||||
error_3 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["error_3"], 'queue' => 'queue_3')
|
||||
future_1 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["future_1"], 'queue' => 'queue_4')
|
||||
future_2 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["future_2"], 'queue' => 'queue_5')
|
||||
future_3 = Sidekiq.dump_json('class' => ScheduledWorker.name, 'args' => ["future_3"], 'queue' => 'queue_6')
|
||||
enqueued_time = Time.new(2013, 2, 4)
|
||||
|
||||
conn.zadd("retry", (Time.now - 60).to_f.to_s, error_1)
|
||||
conn.zadd("retry", (Time.now - 50).to_f.to_s, error_2)
|
||||
conn.zadd("retry", (Time.now + 60).to_f.to_s, error_3)
|
||||
conn.zadd("schedule", (Time.now - 60).to_f.to_s, future_1)
|
||||
conn.zadd("schedule", (Time.now - 50).to_f.to_s, future_2)
|
||||
conn.zadd("schedule", (Time.now + 60).to_f.to_s, future_3)
|
||||
Time.stub(:now, enqueued_time) do
|
||||
Sidekiq.redis do |conn|
|
||||
error_1 = { 'class' => ScheduledWorker.name, 'args' => ["error_1"], 'queue' => 'queue_1' }
|
||||
error_2 = { 'class' => ScheduledWorker.name, 'args' => ["error_2"], 'queue' => 'queue_2' }
|
||||
error_3 = { 'class' => ScheduledWorker.name, 'args' => ["error_3"], 'queue' => 'queue_3' }
|
||||
future_1 = { 'class' => ScheduledWorker.name, 'args' => ["future_1"], 'queue' => 'queue_4' }
|
||||
future_2 = { 'class' => ScheduledWorker.name, 'args' => ["future_2"], 'queue' => 'queue_5' }
|
||||
future_3 = { 'class' => ScheduledWorker.name, 'args' => ["future_3"], 'queue' => 'queue_6' }
|
||||
|
||||
poller = Sidekiq::Scheduled::Poller.new
|
||||
poller.poll
|
||||
conn.zadd("retry", (Time.now - 60).to_f.to_s, Sidekiq.dump_json(error_1))
|
||||
conn.zadd("retry", (Time.now - 50).to_f.to_s, Sidekiq.dump_json(error_2))
|
||||
conn.zadd("retry", (Time.now + 60).to_f.to_s, Sidekiq.dump_json(error_3))
|
||||
conn.zadd("schedule", (Time.now - 60).to_f.to_s, Sidekiq.dump_json(future_1))
|
||||
conn.zadd("schedule", (Time.now - 50).to_f.to_s, Sidekiq.dump_json(future_2))
|
||||
conn.zadd("schedule", (Time.now + 60).to_f.to_s, Sidekiq.dump_json(future_3))
|
||||
|
||||
assert_equal [error_1], conn.lrange("queue:queue_1", 0, -1)
|
||||
assert_equal [error_2], conn.lrange("queue:queue_2", 0, -1)
|
||||
assert_equal [error_3], conn.zrange("retry", 0, -1)
|
||||
assert_equal [future_1], conn.lrange("queue:queue_4", 0, -1)
|
||||
assert_equal [future_2], conn.lrange("queue:queue_5", 0, -1)
|
||||
assert_equal [future_3], conn.zrange("schedule", 0, -1)
|
||||
poller = Sidekiq::Scheduled::Poller.new
|
||||
poller.poll
|
||||
|
||||
assert_equal [Sidekiq.dump_json(error_1.merge('enqueued_at' => enqueued_time.to_f))], conn.lrange("queue:queue_1", 0, -1)
|
||||
assert_equal [Sidekiq.dump_json(error_2.merge('enqueued_at' => enqueued_time.to_f))], conn.lrange("queue:queue_2", 0, -1)
|
||||
assert_equal [Sidekiq.dump_json(future_1.merge('enqueued_at' => enqueued_time.to_f))], conn.lrange("queue:queue_4", 0, -1)
|
||||
assert_equal [Sidekiq.dump_json(future_2.merge('enqueued_at' => enqueued_time.to_f))], conn.lrange("queue:queue_5", 0, -1)
|
||||
|
||||
assert_equal [Sidekiq.dump_json(error_3)], conn.zrange("retry", 0, -1)
|
||||
assert_equal [Sidekiq.dump_json(future_3)], conn.zrange("schedule", 0, -1)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
Loading…
Add table
Reference in a new issue