diff --git a/lib/async/job/processor/redis/processing_list.rb b/lib/async/job/processor/redis/processing_list.rb index 53fafb3..e47a0b9 100644 --- a/lib/async/job/processor/redis/processing_list.rb +++ b/lib/async/job/processor/redis/processing_list.rb @@ -93,9 +93,10 @@ def size # Fetch the next job from the ready queue, moving it to this worker's pending list. # This is a blocking operation that waits until a job is available. - # @returns [String, nil] The job ID, or nil if no job is available. + # @returns [String] The job ID. + # @raises [IOError] If the blocking operation returns without a job. def fetch - @client.brpoplpush(@ready_list.key, @pending_key, 0) + @client.brpoplpush(@ready_list.key, @pending_key, 0) or raise IOError, "Blocking dequeue returned no job!" end # Mark a job as completed, removing it from the pending list and job store. diff --git a/lib/async/job/processor/redis/server.rb b/lib/async/job/processor/redis/server.rb index d3ead22..3e3639b 100644 --- a/lib/async/job/processor/redis/server.rb +++ b/lib/async/job/processor/redis/server.rb @@ -28,8 +28,10 @@ class Server < Generic # @parameter prefix [String] The Redis key prefix for job data. # @parameter coder [Async::Job::Coder] The job serialization codec. # @parameter resolution [Integer] The resolution in seconds for delayed job processing. + # @parameter retry_delay [Numeric] The initial delay before retrying a failed dequeue operation. + # @parameter retry_delay_limit [Numeric] The maximum delay between failed dequeue operations. # @parameter parent [Async::Task] The parent task for background processing. - def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, resolution: 10, parent: nil) + def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, resolution: 10, retry_delay: 1, retry_delay_limit: 30, parent: nil) super(delegate) @id = SecureRandom.uuid @@ -37,6 +39,8 @@ def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, res @prefix = prefix @coder = coder @resolution = resolution + @retry_delay = retry_delay + @retry_delay_limit = retry_delay_limit @job_store = JobStore.new(@client, "#{@prefix}:jobs") @delayed_jobs = DelayedJobs.new(@client, "#{@prefix}:delayed") @@ -56,9 +60,7 @@ def start! @parent.async(transient: true, annotation: self.class.name) do |task| @task = task - while true - self.dequeue(task) - end + self.run(task) ensure @task = nil end @@ -113,7 +115,27 @@ def call(job) end end - protected + protected + + # Run the dequeue loop, retrying transient failures with bounded exponential backoff. + # @parameter parent [Async::Task] The parent task used to process dequeued jobs. + def run(parent) + retry_delay = @retry_delay + + while true + begin + self.dequeue(parent) + retry_delay = @retry_delay + rescue => error + delay = retry_delay * (0.5 + rand * 0.5) + + Console.error(self, "Failed to dequeue job; retrying.", retry_in: delay, exception: error) + sleep(delay) + + retry_delay = [retry_delay * 2, @retry_delay_limit].min + end + end + end # Dequeue a job from the ready list and process it. # diff --git a/test/async/job/processor/processing_list.rb b/test/async/job/processor/processing_list.rb index ab37ac5..2d6e787 100644 --- a/test/async/job/processor/processing_list.rb +++ b/test/async/job/processor/processing_list.rb @@ -51,6 +51,16 @@ pending_job = client.lpop(pending_key) expect(pending_job).to be == job_id end + + it "raises an error if the blocking dequeue returns no job" do + mock(client) do |mock| + mock.replace(:brpoplpush) {nil} + end + + expect do + processing_list.fetch + end.to raise_exception(IOError, message: be(:include?, "Blocking dequeue returned no job")) + end end with "#complete" do diff --git a/test/async/job/processor/server.rb b/test/async/job/processor/server.rb index a56075b..44e56d9 100644 --- a/test/async/job/processor/server.rb +++ b/test/async/job/processor/server.rb @@ -83,5 +83,77 @@ expect(server.status_string).to be == "R=0 D=0 P=0/1" end + + it "formats large counts" do + expect(server.send(:format_count, 1_234)).to be == "1.23K" + expect(server.send(:format_count, 1_234_567)).to be == "1.23M" + end + end +end + +describe Async::Job::Processor::Redis::Server do + include Sus::Fixtures::Console::CapturedLogger + + let(:client) do + Object.new.tap do |client| + def client.script(...) + "script" + end + end + end + + let(:server) {subject.new(nil, client, retry_delay: 1.0, retry_delay_limit: 4.0)} + + with "#run" do + it "retries dequeue failures with bounded exponential backoff" do + attempts = 0 + delays = [] + + mock(server) do |mock| + mock.replace(:dequeue) do |_parent| + attempts += 1 + + throw :finished if attempts > 4 + raise IOError, "Redis connection failed!" + end + + mock.replace(:rand) {1.0} + mock.replace(:sleep) {|delay| delays << delay} + end + + catch(:finished) do + server.send(:run, nil) + end + + expect(delays).to be == [1.0, 2.0, 4.0, 4.0] + expect_console.to have_logged(severity: be(:==, :error), message: be(:include?, "Failed to dequeue job")) + end + + it "resets the retry delay after a successful dequeue" do + attempts = 0 + delays = [] + + mock(server) do |mock| + mock.replace(:dequeue) do |_parent| + attempts += 1 + + case attempts + when 1, 3 + raise IOError, "Redis connection failed!" + when 4 + throw :finished + end + end + + mock.replace(:rand) {1.0} + mock.replace(:sleep) {|delay| delays << delay} + end + + catch(:finished) do + server.send(:run, nil) + end + + expect(delays).to be == [1.0, 1.0] + end end end