Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions context/redis-queue.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
8 changes: 8 additions & 0 deletions guides/redis-queue/readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
63 changes: 61 additions & 2 deletions lib/async/job/processor/redis/delayed_jobs.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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])
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
10 changes: 8 additions & 2 deletions lib/async/job/processor/redis/server.rb
Original file line number Diff line number Diff line change
Expand Up @@ -28,15 +28,17 @@ 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
@client = client
@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")
Expand Down Expand Up @@ -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
Expand Down
129 changes: 129 additions & 0 deletions test/async/job/processor/delayed_jobs.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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