diff --git a/context/redis-queue.md b/context/redis-queue.md index 5a85b15..5cb9c13 100644 --- a/context/redis-queue.md +++ b/context/redis-queue.md @@ -14,6 +14,14 @@ The ready queue is where jobs that are immediately available for processing are The delayed queue holds jobs that are not meant to be executed immediately but at a specified future time. This functionality is crucial for tasks that need to be executed at a later stage, such as scheduled notifications or time-dependent processes. Jobs in the delayed queue are sorted according to their execution time. When possible, they are moved to the ready queue to be executed by the next available worker. This transition is managed through Redis's sorted sets, allowing efficient retrieval and management of timed events. +### Delayed Promotion Recovery + +Redis and promotion errors do not terminate delayed promotion. The promoter retries with exponential backoff, then returns to the configured polling interval after a successful promotion. Task cancellation remains normal lifecycle control. + +Configure retry timing with `ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_INITIAL_RETRY_DELAY` (default `0.25` seconds) and `ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_MAXIMUM_RETRY_DELAY` (default `5` seconds). + +The server logs failures and recovery through `Console`. Pass an optional `delayed_jobs_instrumentation` callable to receive `:failure` events with `error`, `consecutive_failures`, and `retry_in_seconds`, and `:recovered` events with the previous `consecutive_failures`. Callback failures are isolated from the promoter. + ## Processing Queue Once a job is dequeued from the ready queue, it enters the processing queue, signifying that it is currently being executed by a worker. The processing queue is crucial for tracking the progress of jobs and for ensuring that jobs can be retried or recovered in case of worker failure. Each worker emits a heartbeat, and if a worker fails to emit a heartbeat within a specified time, any jobs associated with that worker are automatically moved back to the ready queue for reprocessing. diff --git a/guides/redis-queue/readme.md b/guides/redis-queue/readme.md index 5a85b15..5cb9c13 100644 --- a/guides/redis-queue/readme.md +++ b/guides/redis-queue/readme.md @@ -14,6 +14,14 @@ The ready queue is where jobs that are immediately available for processing are The delayed queue holds jobs that are not meant to be executed immediately but at a specified future time. This functionality is crucial for tasks that need to be executed at a later stage, such as scheduled notifications or time-dependent processes. Jobs in the delayed queue are sorted according to their execution time. When possible, they are moved to the ready queue to be executed by the next available worker. This transition is managed through Redis's sorted sets, allowing efficient retrieval and management of timed events. +### Delayed Promotion Recovery + +Redis and promotion errors do not terminate delayed promotion. The promoter retries with exponential backoff, then returns to the configured polling interval after a successful promotion. Task cancellation remains normal lifecycle control. + +Configure retry timing with `ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_INITIAL_RETRY_DELAY` (default `0.25` seconds) and `ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_MAXIMUM_RETRY_DELAY` (default `5` seconds). + +The server logs failures and recovery through `Console`. Pass an optional `delayed_jobs_instrumentation` callable to receive `:failure` events with `error`, `consecutive_failures`, and `retry_in_seconds`, and `:recovered` events with the previous `consecutive_failures`. Callback failures are isolated from the promoter. + ## Processing Queue Once a job is dequeued from the ready queue, it enters the processing queue, signifying that it is currently being executed by a worker. The processing queue is crucial for tracking the progress of jobs and for ensuring that jobs can be retried or recovered in case of worker failure. Each worker emits a heartbeat, and if a worker fails to emit a heartbeat within a specified time, any jobs associated with that worker are automatically moved back to the ready queue for reprocessing. diff --git a/lib/async/job/processor/redis/delayed_jobs.rb b/lib/async/job/processor/redis/delayed_jobs.rb index d40ff32..83ffcee 100644 --- a/lib/async/job/processor/redis/delayed_jobs.rb +++ b/lib/async/job/processor/redis/delayed_jobs.rb @@ -3,6 +3,8 @@ # Released under the MIT License. # Copyright, 2024-2025, by Samuel Williams. +require "protocol/redis/error" + module Async module Job module Processor @@ -11,6 +13,9 @@ module Redis # Jobs are stored with their execution timestamps and automatically moved # to the ready queue when their scheduled time arrives. class DelayedJobs + INITIAL_RETRY_DELAY = Float(ENV.fetch("ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_INITIAL_RETRY_DELAY", 0.25)) + MAXIMUM_RETRY_DELAY = Float(ENV.fetch("ASYNC_JOB_PROCESSOR_REDIS_DELAYED_JOBS_MAXIMUM_RETRY_DELAY", 5)) + ADD = <<~LUA redis.call('HSET', KEYS[1], ARGV[1], ARGV[2]) redis.call('ZADD', KEYS[2], ARGV[3], ARGV[1]) @@ -45,17 +50,32 @@ def size # @parameter ready_list [ReadyList] The ready list to move jobs to. # @parameter resolution [Integer] The check interval in seconds. # @parameter parent [Async::Task] The parent task to run the background loop in. + # @parameter instrumentation [Interface(:call) | Nil] An optional callback for promoter failure and recovery events. # @returns [Async::Task] The background processing task. - def start(ready_list, resolution: 10, parent: Async::Task.current) + def start(ready_list, resolution: 10, parent: Async::Task.current, instrumentation: nil) parent.async do - while true + consecutive_failures = 0 + + loop do count = move(destination: ready_list.key) + if consecutive_failures > 0 + report_recovery(instrumentation, consecutive_failures) + consecutive_failures = 0 + end + if count > 0 Console.debug(self, "Moved #{count} delayed jobs to ready list.") end sleep(resolution) + rescue Async::Stop + raise + rescue => error + consecutive_failures += 1 + retry_in_seconds = retry_delay(consecutive_failures) + report_failure(instrumentation, error, consecutive_failures, retry_in_seconds) + sleep(retry_in_seconds) end end end @@ -82,6 +102,45 @@ def add(job, timestamp, job_store) # @returns [Integer] The number of jobs moved. def move(destination:, now: Time.now.to_f) @client.evalsha(@move, 2, @key, destination, now) + rescue Protocol::Redis::ServerError => error + raise unless error.message.start_with?("NOSCRIPT") + + @move = @client.script(:load, MOVE) + @client.evalsha(@move, 2, @key, destination, now) + end + + private + + def retry_delay(consecutive_failures) + [INITIAL_RETRY_DELAY * (2 ** (consecutive_failures - 1)), MAXIMUM_RETRY_DELAY].min + end + + def report_failure(instrumentation, error, consecutive_failures, retry_in_seconds) + Console.warn( + self, + "Delayed job promotion failed; retrying in #{retry_in_seconds} seconds.", + error, + consecutive_failures:, + retry_in_seconds:, + ) + rescue + # Logging must not terminate the promoter: + ensure + instrument(instrumentation, :failure, error:, consecutive_failures:, retry_in_seconds:) + end + + def report_recovery(instrumentation, consecutive_failures) + Console.info(self, "Delayed job promotion recovered.", consecutive_failures:) + rescue + # Logging must not terminate the promoter: + ensure + instrument(instrumentation, :recovered, consecutive_failures:) + end + + def instrument(instrumentation, event, **details) + instrumentation&.call(event, **details) + rescue + # Instrumentation must not terminate the promoter: end end end diff --git a/lib/async/job/processor/redis/server.rb b/lib/async/job/processor/redis/server.rb index d3ead22..b9c428f 100644 --- a/lib/async/job/processor/redis/server.rb +++ b/lib/async/job/processor/redis/server.rb @@ -28,8 +28,9 @@ 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 delayed_jobs_instrumentation [Interface(:call) | Nil] An optional callback for delayed promoter events. # @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, delayed_jobs_instrumentation: nil, parent: nil) super(delegate) @id = SecureRandom.uuid @@ -37,6 +38,7 @@ def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, res @prefix = prefix @coder = coder @resolution = resolution + @delayed_jobs_instrumentation = delayed_jobs_instrumentation @job_store = JobStore.new(@client, "#{@prefix}:jobs") @delayed_jobs = DelayedJobs.new(@client, "#{@prefix}:delayed") @@ -70,7 +72,11 @@ def start super # Start the delayed processor, which will move jobs to the ready processor when they are ready: - @delayed_jobs.start(@ready_list, resolution: @resolution) + @delayed_jobs.start( + @ready_list, + resolution: @resolution, + instrumentation: @delayed_jobs_instrumentation, + ) # Start the processing processor, which will move jobs to the ready processor when they are abandoned: @processing_list.start diff --git a/test/async/job/processor/delayed_jobs.rb b/test/async/job/processor/delayed_jobs.rb index e957be0..d30c55b 100644 --- a/test/async/job/processor/delayed_jobs.rb +++ b/test/async/job/processor/delayed_jobs.rb @@ -107,6 +107,48 @@ expect(ready_jobs).to be(:include?, job_id_2) expect(ready_jobs).to be(:include?, job_id_3) end + + it "reloads the move script after Redis flushes its script cache" do + past_time = Time.now - 60 + job_id = delayed_jobs.add(test_job, past_time, job_store) + destination = ready_list.key + client.script(:flush) + + count = delayed_jobs.move(destination:) + + expect(count).to be == 1 + expect(client.lpop(destination)).to be == job_id + end + + it "retries NOSCRIPT only once" do + failing_client = Class.new do + attr :evalsha_count + attr :script_load_count + + def initialize + @evalsha_count = 0 + @script_load_count = 0 + end + + def script(subcommand, source = nil) + @script_load_count += 1 + source + end + + def evalsha(*) + @evalsha_count += 1 + raise Protocol::Redis::ServerError, "NOSCRIPT No matching script." + end + end.new + failing_delayed_jobs = subject.new(failing_client, "#{prefix}:failing") + + expect do + failing_delayed_jobs.move(destination: ready_list.key) + end.to raise_exception(Protocol::Redis::ServerError) + + expect(failing_client.evalsha_count).to be == 2 + expect(failing_client.script_load_count).to be == 3 + end end with "#start" do @@ -148,5 +190,92 @@ message: be(:include?, "Moved 1 delayed jobs to ready list") ) end + + it "retries failed promotions and reports recovery" do + attempts = 0 + events = [] + instrumentation = proc{|event, **details| events << [event, details]} + + delayed_jobs.define_singleton_method(:move) do |destination:| + attempts += 1 + raise "Redis unavailable" if attempts == 1 + + 0 + end + + task = delayed_jobs.start(ready_list, resolution: 60, instrumentation:) + + Async::Task.current.with_timeout(2) do + sleep(0.01) until attempts >= 2 + end + task.stop + + expect(events).to have_attributes(size: be == 2) + expect(events[0]).to have_attributes( + first: be == :failure, + last: have_keys( + error: be_a(RuntimeError), + consecutive_failures: be == 1, + retry_in_seconds: be == 0.25, + ), + ) + expect(events[1]).to be == [:recovered, {consecutive_failures: 1}] + + expect_console.to have_logged( + severity: be == :warn, + message: be(:include?, "Delayed job promotion failed"), + ) + expect_console.to have_logged( + severity: be == :info, + message: be(:include?, "Delayed job promotion recovered"), + ) + ensure + task&.stop + end + + it "continues when instrumentation fails" do + attempts = 0 + instrumentation = proc do + raise "Instrumentation unavailable" + end + + delayed_jobs.define_singleton_method(:move) do |destination:| + attempts += 1 + raise "Redis unavailable" if attempts == 1 + + 0 + end + + task = delayed_jobs.start(ready_list, resolution: 60, instrumentation:) + + Async::Task.current.with_timeout(2) do + sleep(0.01) until attempts >= 2 + end + expect(task).not.to be(:finished?) + ensure + task&.stop + end + + it "caps exponential retry delays" do + expect(delayed_jobs.send(:retry_delay, 1)).to be == 0.25 + expect(delayed_jobs.send(:retry_delay, 2)).to be == 0.5 + expect(delayed_jobs.send(:retry_delay, 10)).to be == 5 + end + + it "does not report task cancellation as a promoter failure" do + events = [] + task = delayed_jobs.start( + ready_list, + resolution: 60, + instrumentation: proc{|event, **details| events << [event, details]}, + ) + + sleep(0.01) + task.stop + + expect(events).to be(:empty?) + ensure + task&.stop + end end end