From 1298feedfd6c0c9ade2664f1116e1ad9c09c6be6 Mon Sep 17 00:00:00 2001 From: Jeff Baxendale Date: Thu, 3 Sep 2026 15:05:02 -0400 Subject: [PATCH] Stop Redis server lifecycle tasks The server creates a dispatcher, a delayed-job promoter, and a heartbeat worker, but previously stopped only its dispatcher. In long-running application processes that can leave background Redis work running after service shutdown. Track the maintenance tasks explicitly, make their starts idempotent, and stop them with the server. This gives every host a complete lifecycle boundary while preserving the existing at-least-once recovery model for pending jobs. Signed-off-by: Jeff Baxendale --- lib/async/job/processor/redis/delayed_jobs.rb | 31 +++++++++++++++++-- .../job/processor/redis/processing_list.rb | 30 ++++++++++++++++-- lib/async/job/processor/redis/server.rb | 4 +++ test/async/job/processor/delayed_jobs.rb | 19 +++++++++++- test/async/job/processor/processing_list.rb | 17 +++++++++- test/async/job/processor/server.rb | 27 ++++++++++++++++ 6 files changed, 122 insertions(+), 6 deletions(-) diff --git a/lib/async/job/processor/redis/delayed_jobs.rb b/lib/async/job/processor/redis/delayed_jobs.rb index d40ff32..14ff0e4 100644 --- a/lib/async/job/processor/redis/delayed_jobs.rb +++ b/lib/async/job/processor/redis/delayed_jobs.rb @@ -34,6 +34,7 @@ def initialize(client, key) @add = @client.script(:load, ADD) @move = @client.script(:load, MOVE) + @task = nil end # @returns [Integer] The number of jobs currently in the delayed queue. @@ -45,9 +46,16 @@ 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. - # @returns [Async::Task] The background processing task. + # @returns [Async::Task | false] The background processing task, or false if already started. def start(ready_list, resolution: 10, parent: Async::Task.current) - parent.async do + return false if @task + + # Reserve ownership before spawning because Async tasks may begin eagerly. + @task = true + + task = parent.async do |task| + @task = task + while true count = move(destination: ready_list.key) @@ -57,7 +65,26 @@ def start(ready_list, resolution: 10, parent: Async::Task.current) sleep(resolution) end + ensure + @task = nil if @task.equal?(task) end + + # A non-greedy parent may return before the child assigns its handle. + @task = task if @task == true + return task + rescue + @task = nil + raise + end + + # Stop the owned delayed job promotion task. + def stop + task = @task + @task = nil + return unless task.respond_to?(:stop) + + task.stop + task.wait if task.alive? end # @attribute [String] The Redis key for this delayed jobs queue. diff --git a/lib/async/job/processor/redis/processing_list.rb b/lib/async/job/processor/redis/processing_list.rb index 53fafb3..63ad53c 100644 --- a/lib/async/job/processor/redis/processing_list.rb +++ b/lib/async/job/processor/redis/processing_list.rb @@ -75,6 +75,7 @@ def initialize(client, key, id, ready_list, job_store) @complete = @client.script(:load, COMPLETE) @complete_count = 0 + @task = nil end # @attribute [String] The base Redis key for this processing list. @@ -133,11 +134,17 @@ def requeue(start_time, delay, factor) # @parameter delay [Integer] The heartbeat update interval in seconds. # @parameter factor [Integer] The heartbeat expiration factor. # @parameter parent [Async::Task] The parent task to run the background loop in. - # @returns [Async::Task] The background processing task. + # @returns [Async::Task | false] The background processing task, or false if already started. def start(delay: 5, factor: 2, parent: Async::Task.current) + return false if @task + + # Reserve ownership before spawning because Async tasks may begin eagerly. + @task = true start_time = Time.now.to_f - parent.async do |task| + task = parent.async do |task| + @task = task + while true task.defer_stop do count = self.requeue(start_time, delay, factor) @@ -149,7 +156,26 @@ def start(delay: 5, factor: 2, parent: Async::Task.current) sleep(delay) end + ensure + @task = nil if @task.equal?(task) end + + # A non-greedy parent may return before the child assigns its handle. + @task = task if @task == true + return task + rescue + @task = nil + raise + end + + # Stop the owned heartbeat and abandoned job recovery task. + def stop + task = @task + @task = nil + return unless task.respond_to?(:stop) + + task.stop + task.wait if task.alive? end end end diff --git a/lib/async/job/processor/redis/server.rb b/lib/async/job/processor/redis/server.rb index d3ead22..7e14759 100644 --- a/lib/async/job/processor/redis/server.rb +++ b/lib/async/job/processor/redis/server.rb @@ -80,7 +80,11 @@ def start # Stop the server and all background processing tasks. def stop + # Stop workers before their heartbeat protection. Pending jobs remain in Redis + # and become recoverable when this server's heartbeat expires naturally. @task&.stop + @processing_list.stop + @delayed_jobs.stop super end diff --git a/test/async/job/processor/delayed_jobs.rb b/test/async/job/processor/delayed_jobs.rb index e957be0..264c4a5 100644 --- a/test/async/job/processor/delayed_jobs.rb +++ b/test/async/job/processor/delayed_jobs.rb @@ -125,7 +125,7 @@ remaining_score = client.zscore(delayed_jobs.key, job_id) expect(remaining_score).to be_nil ensure - task&.stop + delayed_jobs.stop end it "logs debug messages when moving jobs" do @@ -147,6 +147,23 @@ severity: be == :debug, message: be(:include?, "Moved 1 delayed jobs to ready list") ) + ensure + delayed_jobs.stop + end + + it "owns a single background task" do + task = delayed_jobs.start(ready_list, resolution: 1) + + expect(task.alive?).to be == true + expect(delayed_jobs.start(ready_list, resolution: 1)).to be == false + + delayed_jobs.stop + expect(task.finished?).to be == true + + restarted_task = delayed_jobs.start(ready_list, resolution: 1) + expect(restarted_task).not.to be == false + ensure + delayed_jobs.stop end end end diff --git a/test/async/job/processor/processing_list.rb b/test/async/job/processor/processing_list.rb index ab37ac5..04414a4 100644 --- a/test/async/job/processor/processing_list.rb +++ b/test/async/job/processor/processing_list.rb @@ -154,7 +154,22 @@ fetched_job = processing_list.fetch expect(fetched_job).to be == abandoned_job_id ensure - task&.stop + processing_list.stop + end + + it "owns a single background task" do + task = processing_list.start(delay: 1) + + expect(task.alive?).to be == true + expect(processing_list.start(delay: 1)).to be == false + + processing_list.stop + expect(task.finished?).to be == true + + restarted_task = processing_list.start(delay: 1) + expect(restarted_task).not.to be == false + ensure + processing_list.stop end end end diff --git a/test/async/job/processor/server.rb b/test/async/job/processor/server.rb index a56075b..10dcb25 100644 --- a/test/async/job/processor/server.rb +++ b/test/async/job/processor/server.rb @@ -84,4 +84,31 @@ expect(server.status_string).to be == "R=0 D=0 P=0/1" end end + + with "#stop" do + it "stops every server-owned task" do + dispatcher_task = server.instance_variable_get(:@task) + processing_list = server.instance_variable_get(:@processing_list) + processing_task = processing_list.instance_variable_get(:@task) + delayed_jobs = server.instance_variable_get(:@delayed_jobs) + delayed_task = delayed_jobs.instance_variable_get(:@task) + + server.stop + + expect(dispatcher_task.finished?).to be == true + expect(processing_task.finished?).to be == true + expect(delayed_task.finished?).to be == true + + # Shutdown remains safe when both the owner and its parent call it. + server.stop + end + + it "can restart after stopping all server-owned tasks" do + server.stop + server.start + server.call(job) + + expect(buffer.pop).to be == job + end + end end