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
5 changes: 3 additions & 2 deletions lib/async/job/processor/redis/processing_list.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
32 changes: 27 additions & 5 deletions lib/async/job/processor/redis/server.rb
Original file line number Diff line number Diff line change
Expand Up @@ -28,15 +28,19 @@ 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
@client = client
@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")
Expand All @@ -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
Expand Down Expand Up @@ -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.
#
Expand Down
10 changes: 10 additions & 0 deletions test/async/job/processor/processing_list.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 72 additions & 0 deletions test/async/job/processor/server.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading