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
21 changes: 20 additions & 1 deletion context/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,4 +38,23 @@ Async do

server.stop
end
```
```

## Bounded processing

Pass `Async::Semaphore` as the parent to bound blocking Redis fetches and job
processing together:

``` ruby
require "async/semaphore"

queue = Async::Job::Builder.build(buffer) do
dequeue Async::Job::Processor::Redis,
parent: Async::Semaphore.new(20)
end
```

Omit `parent`, or pass an `Async::Task`, to retain the compatibility behavior:
one blocking fetch stays in flight while fetched jobs run as children of the
dispatcher. A Redis fetch failure stops that dispatcher so it cannot recover
independently of the processing heartbeat.
16 changes: 16 additions & 0 deletions context/redis-queue.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,3 +17,19 @@ The delayed queue holds jobs that are not meant to be executed immediately but a
## 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.

## Processing concurrency

Without a semaphore `parent`, the server keeps one blocking Redis fetch in
flight and runs each fetched job as a child of the dispatcher. Passing an
`Async::Task` as `parent` preserves this behavior and places the dispatcher
under that task.

Pass `Async::Semaphore` as `parent` to set an explicit bound. The dispatcher
reserves a semaphore slot before the blocking fetch and holds it through
processing, so blocked fetches and executing jobs share the same limit.

A failed blocking fetch terminates the dispatcher instead of retrying in
process, so dequeuing cannot recover independently of the processing heartbeat.
Stopping the server cancels in-flight workers and releases their semaphore
slots.
21 changes: 20 additions & 1 deletion guides/getting-started/readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,4 +38,23 @@ Async do

server.stop
end
```
```

## Bounded processing

Pass `Async::Semaphore` as the parent to bound blocking Redis fetches and job
processing together:

``` ruby
require "async/semaphore"

queue = Async::Job::Builder.build(buffer) do
dequeue Async::Job::Processor::Redis,
parent: Async::Semaphore.new(20)
end
```

Omit `parent`, or pass an `Async::Task`, to retain the compatibility behavior:
one blocking fetch stays in flight while fetched jobs run as children of the
dispatcher. A Redis fetch failure stops that dispatcher so it cannot recover
independently of the processing heartbeat.
16 changes: 16 additions & 0 deletions guides/redis-queue/readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,3 +17,19 @@ The delayed queue holds jobs that are not meant to be executed immediately but a
## 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.

## Processing concurrency

Without a semaphore `parent`, the server keeps one blocking Redis fetch in
flight and runs each fetched job as a child of the dispatcher. Passing an
`Async::Task` as `parent` preserves this behavior and places the dispatcher
under that task.

Pass `Async::Semaphore` as `parent` to set an explicit bound. The dispatcher
reserves a semaphore slot before the blocking fetch and holds it through
processing, so blocked fetches and executing jobs share the same limit.

A failed blocking fetch terminates the dispatcher instead of retrying in
process, so dequeuing cannot recover independently of the processing heartbeat.
Stopping the server cancels in-flight workers and releases their semaphore
slots.
73 changes: 58 additions & 15 deletions lib/async/job/processor/redis/server.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
require "async/idler"
require "async/job/coder"
require "async/job/processor/generic"
require "async/semaphore"

require "securerandom"

Expand All @@ -28,7 +29,7 @@ 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 parent [Async::Task] The parent task for background processing.
# @parameter parent [Async::Task | Async::Semaphore | Nil] An optional lifecycle parent or concurrency limiter for job processing.
def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, resolution: 10, parent: nil)
super(delegate)

Expand All @@ -43,7 +44,11 @@ def initialize(delegate, client, prefix: "async-job", coder: Coder::DEFAULT, res
@ready_list = ReadyList.new(@client, "#{@prefix}:ready")
@processing_list = ProcessingList.new(@client, "#{@prefix}:processing", @id, @ready_list, @job_store)

# Ordinary task parents preserve the original single blocking fetch loop.
@parent = parent || Async::Idler.new
# A semaphore limits both claimed and executing jobs, while the dispatcher
# task remains responsible for worker lifecycle.
@semaphore = parent if parent.is_a?(Async::Semaphore)
end

# Start the job processing loop immediately.
Expand All @@ -53,14 +58,17 @@ def start!

@task = true

@parent.async(transient: true, annotation: self.class.name) do |task|
start_dispatcher do |task|
@task = task

while true
self.dequeue(task)
self.dequeue(task, @semaphore)
end
rescue
dispatcher_failed = true
raise
ensure
@task = nil
@task = nil unless dispatcher_failed && task.children?
end
end

Expand All @@ -73,14 +81,15 @@ def start
@delayed_jobs.start(@ready_list, resolution: @resolution)

# Start the processing processor, which will move jobs to the ready processor when they are abandoned:
@processing_list.start
@processing_task = @processing_list.start

self.start!
end

# Stop the server and all background processing tasks.
def stop
@task&.stop
@task = nil

super
end
Expand Down Expand Up @@ -120,25 +129,59 @@ def call(job)
# If the job fails for any reason, it will be retried.
#
# If you do not desire this behavior, you should catch exceptions in the delegate.
def dequeue(parent)
def dequeue(parent = nil, semaphore = nil)
self.ensure_processing_task_alive!

if semaphore
semaphore.acquire
semaphore_acquired = true
self.ensure_processing_task_alive!
end

_id = @processing_list.fetch
self.ensure_processing_task_alive!

parent.async do
id = _id; _id = nil

job = @coder.load(@job_store.get(id))
@delegate.call(job)
@processing_list.complete(id)
rescue => error
Console.error(self, "Job failed with error!", id: id, exception: error)
@processing_list.retry(id)
id = _id
if parent
parent.async {self.process(id, semaphore)}
else
self.process(id, semaphore)
end
semaphore_acquired = false
_id = nil
ensure
semaphore.release if semaphore_acquired
@processing_list.retry(_id) if _id
end

def process(id, semaphore = nil)
job = @coder.load(@job_store.get(id))
@delegate.call(job)
@processing_list.complete(id)
rescue => error
Console.error(self, "Job failed with error!", id: id, exception: error)
@processing_list.retry(id)
ensure
semaphore&.release
end

private

def start_dispatcher(&block)
if @semaphore
# Keep semaphore permits available for claimed jobs, not the dispatcher.
Async(transient: true, annotation: self.class.name, &block)
else
@parent.async(transient: true, annotation: self.class.name, &block)
end
end

def ensure_processing_task_alive!
if @processing_task && !@processing_task.alive?
raise "Heartbeat task stopped; refusing to dequeue jobs."
end
end

def format_count(value)
if value > 1_000_000
"#{(value/1_000_000.0).round(2)}M"
Expand Down
5 changes: 4 additions & 1 deletion readme.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
# Async::Job::Processor::Redis

Provides an asynchronous job server.
Provides a Redis-backed asynchronous job server with durable ready, delayed,
and processing queues. Processing can be bounded by `Async::Semaphore` while
keeping a single blocking dequeue in flight. A dequeue failure stops the
dispatcher so it cannot recover independently of the processing heartbeat.

[![Development Status](https://github.com/socketry/async-job-processor-redis/workflows/Test/badge.svg)](https://github.com/socketry/async-job-processor-redis/actions?workflow=Test)

Expand Down
Loading