diff --git a/context/getting-started.md b/context/getting-started.md index 29d590d..8f695f3 100644 --- a/context/getting-started.md +++ b/context/getting-started.md @@ -38,4 +38,23 @@ Async do server.stop end -``` \ No newline at end of file +``` + +## 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. diff --git a/context/redis-queue.md b/context/redis-queue.md index 5a85b15..8703fa1 100644 --- a/context/redis-queue.md +++ b/context/redis-queue.md @@ -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. diff --git a/guides/getting-started/readme.md b/guides/getting-started/readme.md index 29d590d..8f695f3 100644 --- a/guides/getting-started/readme.md +++ b/guides/getting-started/readme.md @@ -38,4 +38,23 @@ Async do server.stop end -``` \ No newline at end of file +``` + +## 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. diff --git a/guides/redis-queue/readme.md b/guides/redis-queue/readme.md index 5a85b15..8703fa1 100644 --- a/guides/redis-queue/readme.md +++ b/guides/redis-queue/readme.md @@ -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. diff --git a/lib/async/job/processor/redis/server.rb b/lib/async/job/processor/redis/server.rb index d3ead22..46a9a20 100644 --- a/lib/async/job/processor/redis/server.rb +++ b/lib/async/job/processor/redis/server.rb @@ -6,6 +6,7 @@ require "async/idler" require "async/job/coder" require "async/job/processor/generic" +require "async/semaphore" require "securerandom" @@ -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) @@ -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. @@ -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 @@ -73,7 +81,7 @@ 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 @@ -81,6 +89,7 @@ def start # Stop the server and all background processing tasks. def stop @task&.stop + @task = nil super end @@ -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" diff --git a/readme.md b/readme.md index 5f1600d..a5684a6 100644 --- a/readme.md +++ b/readme.md @@ -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) diff --git a/test/async/job/processor/server.rb b/test/async/job/processor/server.rb index a56075b..f948a3a 100644 --- a/test/async/job/processor/server.rb +++ b/test/async/job/processor/server.rb @@ -5,6 +5,7 @@ require "async" require "async/redis" +require "async/semaphore" require "sus/fixtures/async/reactor_context" require "sus/fixtures/console" @@ -37,6 +38,17 @@ expect(buffer.pop).to be == job end + it "can dequeue a job synchronously" do + synchronous_server = subject.new(buffer, prefix: "#{prefix}:synchronous") + synchronous_server.call(job) + + synchronous_server.__send__(:dequeue) + + expect(buffer.pop).to be == job + ensure + synchronous_server&.instance_variable_get(:@client)&.close + end + with "delayed job" do it "can schedule a job and have it processed after a delay" do now = Time.now @@ -80,8 +92,278 @@ expect(server.status_string).to be == "R=0 D=0 P=0/0" server.call(job) + sleep 0.1 # Allow some time for the job to be processed. expect(server.status_string).to be == "R=0 D=0 P=0/1" end end + + with "concurrency limit" do + let(:server) {subject.new(slow_delegate, prefix:, resolution: 1, parent:)} + let(:parent) {nil} + + # A delegate that yields long enough to observe concurrent processing: + let(:slow_delegate) do + Class.new do + attr :started + attr :cancelled + + def start + end + + def stop + end + + def call(_job) + @started = true + sleep 5 + ensure + @cancelled = true + end + end.new + end + + it "preserves default concurrent processing without concurrent fetches" do + 4.times do |i| + server.call({"data" => "job #{i}"}) + end + + Async::Task.current.with_timeout(2) do + sleep(0.01) until server.status_string.match?(/P=4\//) + end + + expect(server.status_string).to be =~ /P=4\// + end + + with "Async::Task" do + let(:parent) {Async::Task.current} + let(:fetch_attempts) {[]} + let(:server) do + subject.new(slow_delegate, prefix:, resolution: 1, parent:).tap do |server| + processing_list = server.instance_variable_get(:@processing_list) + attempts = fetch_attempts + + processing_list.define_singleton_method(:fetch) do + attempts << true + sleep + end + end + end + + it "keeps only one blocking fetch in flight" do + sleep 0.01 + + expect(fetch_attempts).to have_attributes(size: be == 1) + end + end + + with "Async::Semaphore" do + let(:parent) {Async::Semaphore.new(2)} + + it "uses a transient dispatcher" do + dispatcher = server.instance_variable_get(:@task) + + expect(dispatcher).to be(:transient?) + end + + it "can limit concurrent job processing to 2" do + 4.times do |i| + server.call({"data" => "job #{i}"}) + end + + Async::Task.current.with_timeout(2) do + sleep(0.01) until server.status_string.match?(/P=2\//) + end + + expect(server.status_string).to be =~ /P=2\// + end + + with "a completed job" do + let(:server) {subject.new(buffer, prefix:, resolution: 1, parent:)} + + it "releases the worker permit after processing" do + server.call(job) + expect(buffer.pop).to be == job + + Async::Task.current.with_timeout(2) do + sleep(0.01) until parent.count == 1 + end + + expect(parent.count).to be == 1 + end + end + + it "cancels in-flight workers when stopped" do + server.call(job) + + Async::Task.current.with_timeout(2) do + sleep(0.01) until slow_delegate.started + end + + expect(parent.count).to be > 0 + server.stop + + Async::Task.current.with_timeout(2) do + sleep(0.01) until parent.count == 0 + end + + expect(slow_delegate.cancelled).to be == true + expect(parent.count).to be == 0 + end + + with "a fetch failure after dispatch" do + let(:fetch_attempts) {[]} + let(:server) do + subject.new(slow_delegate, prefix:, resolution: 1, parent:).tap do |server| + processing_list = server.instance_variable_get(:@processing_list) + fetch = processing_list.method(:fetch) + attempts = fetch_attempts + delegate = slow_delegate + + processing_list.define_singleton_method(:fetch) do + attempts << true + if attempts.size > 1 + sleep(0.01) until delegate.started + raise "Redis unavailable" + end + + fetch.call + end + end + end + + it "retains workers for cancellation by stop" do + server.call(job) + + dispatcher = nil + Async::Task.current.with_timeout(2) do + loop do + dispatcher = server.instance_variable_get(:@task) + break if dispatcher&.failed? + + sleep(0.01) + end + end + + expect(dispatcher).to be(:children?) + expect(parent.count).to be == 1 + + server.stop + + Async::Task.current.with_timeout(2) do + sleep(0.01) until parent.count == 0 + end + + expect(slow_delegate.cancelled).to be == true + expect(server.instance_variable_get(:@task)).to be_nil + end + end + + with "an idle queue" do + let(:fetch_attempts) {[]} + let(:server) do + subject.new(slow_delegate, prefix:, resolution: 1, parent:).tap do |server| + processing_list = server.instance_variable_get(:@processing_list) + attempts = fetch_attempts + + processing_list.define_singleton_method(:fetch) do + attempts << true + sleep + end + end + end + + it "keeps only one blocking fetch in flight" do + sleep 0.01 + + expect(fetch_attempts).to have_attributes(size: be == 1) + expect(parent.count).to be == 1 + end + end + + it "preserves a bounded Redis connection for callers" do + client = Async::Redis::Client.new(limit: 2) + buffer = Async::Job::Buffer.new + bounded_server = Async::Job::Processor::Redis::Server.new( + buffer, + client, + prefix: "#{prefix}:bounded", + parent: Async::Semaphore.new(2), + ) + bounded_server.start! + + bounded_server.call(job) + expect(buffer.pop).to be == job + ensure + bounded_server&.stop + client&.close + end + + with "a stopped heartbeat" do + let(:fetch_started) {[]} + let(:release_fetch) {[]} + let(:retried) {[]} + let(:server) do + subject.new(slow_delegate, prefix:, resolution: 1, parent:).tap do |server| + processing_list = server.instance_variable_get(:@processing_list) + started = fetch_started + release = release_fetch + retried_jobs = retried + + processing_list.define_singleton_method(:fetch) do + started << true + sleep(0.01) until release.any? + "fetched-job" + end + + processing_list.define_singleton_method(:retry) do |id| + retried_jobs << id + end + end + end + + it "does not process a fetched job" do + Async::Task.current.with_timeout(2) do + sleep(0.01) until fetch_started.any? + end + + server.instance_variable_get(:@processing_task).stop + release_fetch << true + + Async::Task.current.with_timeout(2) do + sleep(0.01) until retried.any? + end + + expect(retried).to be == ["fetched-job"] + expect(slow_delegate.started).to be_nil + expect(parent.count).to be == 0 + end + end + + with "a Redis outage" do + let(:parent) {Async::Semaphore.new(1)} + let(:fetch_attempts) {[]} + let(:server) do + subject.new(slow_delegate, prefix:, resolution: 1, parent:).tap do |server| + processing_list = server.instance_variable_get(:@processing_list) + attempts = fetch_attempts + + processing_list.define_singleton_method(:fetch) do + attempts << Process.clock_gettime(Process::CLOCK_MONOTONIC) + raise "Redis unavailable" + end + end + end + + it "fails closed without retrying" do + Async::Task.current.with_timeout(2) do + sleep(0.01) until server.instance_variable_get(:@task).nil? + end + + expect(fetch_attempts).to have_attributes(size: be == 1) + expect(parent.count).to be == 0 + end + end + end + end end