From 90cb8b1c90d62d3ff436bee8e561086278cbc0a1 Mon Sep 17 00:00:00 2001 From: Julik Tarkhanov Date: Thu, 24 Sep 2026 15:33:28 +0200 Subject: [PATCH 1/4] Stop processes whose heartbeats stop returning #778 stops a process whose heartbeats keep failing, but both of its paths hang off a raise. A heartbeat can also never return at all: Process#heartbeat does a real round-trip, and on a half-open socket to a database that stopped answering, the thread blocks in the driver indefinitely. Nothing is raised, so presumed_dead? is never reached. Concurrent::TimerTask does not help either, because it reschedules the task and notifies its observers only after the task returns - so the heartbeat thread goes quiet permanently, having neither succeeded nor failed. Its timeout_interval is a no-op that warns it was never implementable. So record when the heartbeat last returned, in an ensure so that success and failure both count, and have a second timer stop the process once that goes older than the alive threshold. The watchdog reads nothing but memory, so it cannot block the way the heartbeat it watches can. This only covers the case where the run loop is still able to act on being unregistered. A process whose run loop is itself blocked needs something harsher, which I have deliberately left out of this change. --- lib/solid_queue/processes/registrable.rb | 36 ++++++++++++++++++++++++ test/unit/worker_test.rb | 30 ++++++++++++++++++++ 2 files changed, 66 insertions(+) diff --git a/lib/solid_queue/processes/registrable.rb b/lib/solid_queue/processes/registrable.rb index 08b2750c6..c5e938d1f 100644 --- a/lib/solid_queue/processes/registrable.rb +++ b/lib/solid_queue/processes/registrable.rb @@ -38,6 +38,8 @@ def registered? end def launch_heartbeat + @last_heartbeat_returned_at = Concurrent::AtomicReference.new(SolidQueue::Timer.monotonic_time_now) + @heartbeat_task = Concurrent::TimerTask.new(execution_interval: SolidQueue.process_heartbeat_interval) do wrap_in_app_executor { heartbeat } end @@ -47,10 +49,42 @@ def launch_heartbeat end @heartbeat_task.execute + + launch_heartbeat_watchdog + end + + # A heartbeat that raises is handled below, but one that never returns is not. + # Concurrent::TimerTask reschedules a task and notifies its observers only once + # the task has completed, so a heartbeat blocked on an unresponsive database -- + # a half-open socket to a connection pooler that stopped serving, say -- stalls + # the heartbeat thread silently and indefinitely, and `presumed_dead?` is never + # reached because nothing was raised. + # + # The watchdog's check reads nothing but memory, so it cannot block the + # same way; the stop it sets off can still run blocking code (stop hooks, + # for one), but only after the stall has been noticed. + def launch_heartbeat_watchdog + @heartbeat_watchdog_task = Concurrent::TimerTask.new(execution_interval: SolidQueue.process_heartbeat_interval) do + stop_to_be_replaced if heartbeat_stalled? + end + + @heartbeat_watchdog_task.add_observer do |_, _, error| + handle_thread_error(error) if error + end + + @heartbeat_watchdog_task.execute end def stop_heartbeat @heartbeat_task&.shutdown + @heartbeat_watchdog_task&.shutdown + end + + # Whether the heartbeat has stopped returning altogether. Note this asks only + # whether the call came back, not whether it succeeded: a heartbeat that keeps + # raising is `presumed_dead?`'s business. + def heartbeat_stalled? + SolidQueue::Timer.monotonic_time_now - @last_heartbeat_returned_at.get > SolidQueue.process_alive_threshold end def heartbeat @@ -64,6 +98,8 @@ def heartbeat # registration is still there stop_to_be_replaced if presumed_dead? raise error + ensure + @last_heartbeat_returned_at.set(SolidQueue::Timer.monotonic_time_now) end # Whether this process's registration is prunable: if the last heartbeat that diff --git a/test/unit/worker_test.rb b/test/unit/worker_test.rb index 77f41edf5..068ee76fb 100644 --- a/test/unit/worker_test.rb +++ b/test/unit/worker_test.rb @@ -271,6 +271,36 @@ class WorkerTest < ActiveSupport::TestCase SolidQueue.process_alive_threshold = old_alive_threshold end + test "terminate when heartbeats stop returning for longer than the alive threshold" do + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.2.seconds + old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 1.second + + # Unlike a heartbeat that raises, this one never comes back at all, so the + # heartbeat TimerTask is never rescheduled and never notifies its observers + unblock_heartbeats = Concurrent::Event.new + SolidQueue::Process.class_eval do + alias_method :heartbeat_without_blocking, :heartbeat + define_method(:heartbeat) { unblock_heartbeats.wait } + end + + @worker.start + wait_for_registered_processes(1, timeout: 1.second) + + assert_not @worker.pool.shutdown? + + wait_while_with_timeout(3) { !@worker.pool.shutdown? } + assert @worker.pool.shutdown? + ensure + unblock_heartbeats.set + SolidQueue::Process.class_eval do + remove_method :heartbeat + alias_method :heartbeat, :heartbeat_without_blocking + remove_method :heartbeat_without_blocking + end + SolidQueue.process_heartbeat_interval = old_heartbeat_interval + SolidQueue.process_alive_threshold = old_alive_threshold + end + test "sleeps `10.minutes` if at capacity" do 3.times { |i| StoreResultJob.perform_later(i, pause: 5.seconds) } From bd96a3db140b3caf15e9449aacf7e0ad25135e18 Mon Sep 17 00:00:00 2001 From: Julik Tarkhanov Date: Thu, 24 Sep 2026 15:41:53 +0200 Subject: [PATCH 2/4] Stop a supervisor whose maintenance stops returning Same defect as the heartbeat watchdog, one level up. Pruning is how supervisors notice dead processes, and launch_maintenance_task runs it in a Concurrent::TimerTask, so a prune blocked on an unresponsive database stops this supervisor pruning ever again - the mechanism meant to notice dead processes is built from the same material as the processes it watches, and fails at exactly the moment they do. Track when maintenance last returned and stop the supervisor once that goes past STALL_FACTOR times the alive threshold, which leaves a full missed cycle of slack since the task's own interval is the alive threshold. Stopping rather than arranging replacement, because a supervisor cannot replace itself - its run loop breaks on stopped?, so whatever runs it gets to start a new one. The way out must not depend on the database that just stopped answering, and must not abandon the supervised processes. So the stalled stop drops the local registration first, turning shutdown's deregister callback into a no-op instead of a round-trip that would block exactly like the prune did - the stale row is left for another supervisor to prune. And a standalone supervisor stops through its own signal pipeline rather than a bare stop, so that handle_signal pairs stop with terminate_gracefully on the supervise thread and the forks get TERMed instead of orphaned; an embedded supervisor never drains its signal queue, but it already terminates its threads from an after_shutdown hook, so a plain stop is enough there. Supervisors need this separately from the heartbeat watchdog in Registrable: stop_to_be_replaced only unregisters and wakes the loop, and unlike Runnable#shutting_down?, Supervisor#supervise never checks registered?. --- lib/solid_queue/supervisor/maintenance.rb | 69 ++++++++++++++++++ test/unit/async_supervisor_test.rb | 86 +++++++++++++++++++++++ test/unit/fork_supervisor_test.rb | 41 +++++++++++ 3 files changed, 196 insertions(+) diff --git a/lib/solid_queue/supervisor/maintenance.rb b/lib/solid_queue/supervisor/maintenance.rb index 23e321f47..ce8674dc6 100644 --- a/lib/solid_queue/supervisor/maintenance.rb +++ b/lib/solid_queue/supervisor/maintenance.rb @@ -6,8 +6,14 @@ module Supervisor::Maintenance after_boot :fail_orphaned_executions end + # How far past its own interval a maintenance run is allowed to go before we + # treat it as stalled rather than slow. One full missed cycle of slack. + STALL_FACTOR = 2 + private def launch_maintenance_task + @last_maintenance_returned_at = Concurrent::AtomicReference.new(SolidQueue::Timer.monotonic_time_now) + @maintenance_task = Concurrent::TimerTask.new(run_now: true, execution_interval: SolidQueue.process_alive_threshold) do prune_dead_processes end @@ -17,14 +23,77 @@ def launch_maintenance_task end @maintenance_task.execute + + launch_maintenance_watchdog + end + + # Pruning is how a supervisor notices dead processes, and it runs in a + # Concurrent::TimerTask, which reschedules only once its task returns. A + # prune blocked on an unresponsive database therefore stops this supervisor + # pruning ever again, silently -- the mechanism meant to notice dead + # processes is built from the same material as the processes it watches. + # + # A supervisor cannot replace itself, so stop instead, and leave it to + # whatever runs this supervisor to start a new one. + # + # The watchdog ticks at the heartbeat interval - not that heartbeats are + # involved, but any working configuration already keeps that cadence well + # inside the alive threshold, so a stall is noticed at most one tick + # after it crosses the line. + def launch_maintenance_watchdog + @maintenance_watchdog_task = Concurrent::TimerTask.new(execution_interval: SolidQueue.process_heartbeat_interval) do + stop_stalled_supervisor if maintenance_stalled? + end + + @maintenance_watchdog_task.add_observer do |_, _, error| + handle_thread_error(error) if error + end + + @maintenance_watchdog_task.execute end def stop_maintenance_task @maintenance_task&.shutdown + @maintenance_watchdog_task&.shutdown + end + + # Whether maintenance has stopped returning altogether, as opposed to + # failing: a prune that raises is reported by the task's observer. + def maintenance_stalled? + SolidQueue::Timer.monotonic_time_now - @last_maintenance_returned_at.get > + STALL_FACTOR * SolidQueue.process_alive_threshold + end + + # The database is presumed unresponsive, so the way out must not depend + # on it. Dropping the registration locally first makes shutdown's + # deregister callback a no-op instead of a round-trip that would block + # exactly like the prune did; the stale row is left for another + # supervisor to prune. + # + # A standalone supervisor stops through the signal pipeline, so that + # handle_signal pairs stop with terminate_gracefully on the supervise + # thread: a bare stop would exit without ever signalling the forks, whose + # only other way of noticing is polling their parent pid once per run + # loop iteration. An embedded supervisor never drains its signal queue, + # but it already terminates its threads from an after_shutdown hook, so a + # plain stop is enough there. + def stop_stalled_supervisor + return if stopped? + + self.process = nil + + if standalone? + signal_queue << :TERM + interrupt + else + stop + end end def prune_dead_processes wrap_in_app_executor { SolidQueue::Process.prune(excluding: process) } + ensure + @last_maintenance_returned_at.set(SolidQueue::Timer.monotonic_time_now) end def fail_orphaned_executions diff --git a/test/unit/async_supervisor_test.rb b/test/unit/async_supervisor_test.rb index f2a70f0ff..37a5b5387 100644 --- a/test/unit/async_supervisor_test.rb +++ b/test/unit/async_supervisor_test.rb @@ -15,6 +15,92 @@ class AsyncSupervisorTest < ActiveSupport::TestCase assert_no_registered_processes end + test "stop when maintenance stops returning for longer than the stall threshold" do + # The heartbeat interval must sit well below the alive threshold: a healthy + # gap between heartbeat returns is already one interval plus a database + # round-trip, so equal settings would make the heartbeat watchdog replace + # perfectly healthy children mid-test + old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.3.seconds + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + # A prune that never returns, so the maintenance TimerTask is never + # rescheduled and this supervisor would silently stop pruning forever + unblock_prunes = Concurrent::Event.new + SolidQueue::Supervisor::Maintenance.module_eval do + alias_method :prune_dead_processes_without_blocking, :prune_dead_processes + define_method(:prune_dead_processes) { unblock_prunes.wait } + end + + supervisor = run_supervisor_as_thread + wait_while_with_timeout(5) { !supervisor.send(:stopped?) } + + assert supervisor.send(:stopped?) + + # The children deregister on the way down; the supervisor's own row cannot + # be removed -- the prune is what is wedged -- and stays for another + # supervisor to prune + wait_for_registered_processes(1, timeout: 3.seconds) + assert_registered_processes(kind: "Supervisor(async)") + ensure + unblock_prunes.set + SolidQueue::Supervisor::Maintenance.module_eval do + if private_method_defined?(:prune_dead_processes_without_blocking) + remove_method :prune_dead_processes + alias_method :prune_dead_processes, :prune_dead_processes_without_blocking + remove_method :prune_dead_processes_without_blocking + end + end + SolidQueue.process_alive_threshold = old_alive_threshold if old_alive_threshold + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + supervisor&.stop + end + + test "complete shutdown when maintenance stalls and the database is unresponsive" do + old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.3.seconds + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + # The same unresponsive database that wedges the prune wedges the + # deregister on the way out, so shutdown must not attempt it: a supervisor + # that detects the stall but blocks in its own shutdown never lets whatever + # runs it start a replacement + unblock_stalled_calls = Concurrent::Event.new + SolidQueue::Supervisor::Maintenance.module_eval do + alias_method :prune_dead_processes_without_blocking, :prune_dead_processes + define_method(:prune_dead_processes) { unblock_stalled_calls.wait } + end + SolidQueue::Process.class_eval do + alias_method :deregister_without_blocking, :deregister + define_method(:deregister) do |pruned: false| + kind.start_with?("Supervisor") ? unblock_stalled_calls.wait : deregister_without_blocking(pruned: pruned) + end + end + + supervisor = run_supervisor_as_thread + supervise_thread = supervisor.instance_variable_get(:@thread) + wait_while_with_timeout(5) { supervise_thread.alive? } + + assert_not supervise_thread.alive? + ensure + unblock_stalled_calls.set + SolidQueue::Supervisor::Maintenance.module_eval do + if private_method_defined?(:prune_dead_processes_without_blocking) + remove_method :prune_dead_processes + alias_method :prune_dead_processes, :prune_dead_processes_without_blocking + remove_method :prune_dead_processes_without_blocking + end + end + SolidQueue::Process.class_eval do + if method_defined?(:deregister_without_blocking) + remove_method :deregister + alias_method :deregister, :deregister_without_blocking + remove_method :deregister_without_blocking + end + end + SolidQueue.process_alive_threshold = old_alive_threshold if old_alive_threshold + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + supervisor&.stop + end + test "start standalone" do pid = run_supervisor_as_fork(mode: :async) wait_for_registered_processes(4, timeout: 5.seconds) # supervisor + dispatcher + 2 workers diff --git a/test/unit/fork_supervisor_test.rb b/test/unit/fork_supervisor_test.rb index 0d8092278..fd2bb27d0 100644 --- a/test/unit/fork_supervisor_test.rb +++ b/test/unit/fork_supervisor_test.rb @@ -273,6 +273,47 @@ def register assert_equal "RuntimeError", failed.exception_class end + test "terminate forks when maintenance stalls instead of abandoning them" do + old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.3.seconds + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + # A prune that never returns, so the maintenance watchdog stops this + # supervisor; the supervisor fork inherits the stub + SolidQueue::Supervisor::Maintenance.module_eval do + alias_method :prune_dead_processes_without_blocking, :prune_dead_processes + define_method(:prune_dead_processes) { sleep 30 } + end + + # A fork stalled in boot has no run loop to notice anything with: only a + # signal reaches it. A supervisor that stops without signalling its forks + # leaves it orphaned until its own stall runs its course. + run_stalled_supervisor(startup_delay: 60.seconds) + wait_while_with_timeout(3) { startup_pids.empty? } + stalled_fork_pid = startup_pids.first + + # The fork must go because it was signalled, so check it before the + # supervisor: this cannot pass just because the supervisor exited + wait_while_with_timeout(5) { process_exists?(stalled_fork_pid) } + assert_not process_exists?(stalled_fork_pid) + + wait_for_process_termination_with_timeout(@stalled_supervisor_pid, timeout: 5) + ensure + SolidQueue::Supervisor::Maintenance.module_eval do + if private_method_defined?(:prune_dead_processes_without_blocking) + remove_method :prune_dead_processes + alias_method :prune_dead_processes, :prune_dead_processes_without_blocking + remove_method :prune_dead_processes_without_blocking + end + end + SolidQueue.process_alive_threshold = old_alive_threshold if old_alive_threshold + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + + startup_pids.each do |leftover_pid| + ::Process.kill(:KILL, leftover_pid) if process_exists?(leftover_pid) + rescue Errno::ESRCH + end + end + test "replace only the fork that does not finish booting" do SolidQueue.fork_boot_timeout = 0.2.seconds run_stalled_supervisor(startup_delay: 60.seconds, with_healthy_worker: true) From c3eabccb2566722f43d4de3664027fdde7e1aff7 Mon Sep 17 00:00:00 2001 From: Julik Tarkhanov Date: Mon, 28 Sep 2026 17:04:09 +0200 Subject: [PATCH 3/4] Make the supervisor honour its own deregistration A supervisor includes Registrable like any other process, so a pruned registration or a tripped heartbeat watchdog reaches stop_to_be_replaced on it too -- but Registrable's version only deregisters locally and wakes the run loop, counting on the supervisor to notice. Nothing supervises a supervisor, and #supervise breaks only on stopped?, never looking at the registration, so on a supervisor both mechanisms were inert: it would run on as a zombie after the rest of the system had declared it dead and failed its children's claimed executions. Override stop_to_be_replaced on the supervisor to stop it instead, through the same route the maintenance watchdog already uses: the signal pipeline when standalone, so the forks are terminated rather than abandoned, and a plain stop when embedded. The maintenance watchdog now funnels into the same override. Once the registration is dropped locally, prune_dead_processes runs with excluding: nil, which excludes nothing: the supervisor may then prune its own stale row. That is benign -- it is exactly what another supervisor would do to that row -- and only reachable while already stopping, so it is left alone. --- lib/solid_queue/supervisor.rb | 34 ++++++++++++ lib/solid_queue/supervisor/maintenance.rb | 28 +--------- test/unit/async_supervisor_test.rb | 64 +++++++++++++++++++++++ test/unit/fork_supervisor_test.rb | 19 +++++++ 4 files changed, 118 insertions(+), 27 deletions(-) diff --git a/lib/solid_queue/supervisor.rb b/lib/solid_queue/supervisor.rb index e6c76c564..1084ef019 100644 --- a/lib/solid_queue/supervisor.rb +++ b/lib/solid_queue/supervisor.rb @@ -148,6 +148,40 @@ def shutdown end end + # Deregister locally and stop. Registrable's version deregisters locally + # and wakes the run loop so the supervisor replaces this process, but + # nothing supervises a supervisor and #supervise never looks at the + # registration: stopping, and leaving a replacement to whatever runs this + # supervisor, is all it can do. Reached when this supervisor's + # registration is gone -- pruned by another supervisor -- and when its + # heartbeat or maintenance watchdog trips. + # + # The database may be unresponsive -- that is what trips the watchdogs -- + # so the way out must not depend on it. Dropping the registration locally + # first makes shutdown's deregister callback a no-op instead of a + # round-trip that could block exactly like the stalled task did; a stale + # row is left for another supervisor to prune. + # + # A standalone supervisor stops through the signal pipeline, so that + # handle_signal pairs stop with terminate_gracefully on the supervise + # thread: a bare stop would exit without ever signalling the forks, whose + # only other way of noticing is polling their parent pid once per run + # loop iteration. An embedded supervisor never drains its signal queue, + # but it already terminates its threads from an after_shutdown hook, so a + # plain stop is enough there. + def stop_to_be_replaced + return if stopped? + + self.process = nil + + if standalone? + signal_queue << :TERM + interrupt + else + stop + end + end + def set_procline # Embedded supervisors don't own their process's title if standalone? diff --git a/lib/solid_queue/supervisor/maintenance.rb b/lib/solid_queue/supervisor/maintenance.rb index ce8674dc6..ecad8f9ad 100644 --- a/lib/solid_queue/supervisor/maintenance.rb +++ b/lib/solid_queue/supervisor/maintenance.rb @@ -42,7 +42,7 @@ def launch_maintenance_task # after it crosses the line. def launch_maintenance_watchdog @maintenance_watchdog_task = Concurrent::TimerTask.new(execution_interval: SolidQueue.process_heartbeat_interval) do - stop_stalled_supervisor if maintenance_stalled? + stop_to_be_replaced if maintenance_stalled? end @maintenance_watchdog_task.add_observer do |_, _, error| @@ -64,32 +64,6 @@ def maintenance_stalled? STALL_FACTOR * SolidQueue.process_alive_threshold end - # The database is presumed unresponsive, so the way out must not depend - # on it. Dropping the registration locally first makes shutdown's - # deregister callback a no-op instead of a round-trip that would block - # exactly like the prune did; the stale row is left for another - # supervisor to prune. - # - # A standalone supervisor stops through the signal pipeline, so that - # handle_signal pairs stop with terminate_gracefully on the supervise - # thread: a bare stop would exit without ever signalling the forks, whose - # only other way of noticing is polling their parent pid once per run - # loop iteration. An embedded supervisor never drains its signal queue, - # but it already terminates its threads from an after_shutdown hook, so a - # plain stop is enough there. - def stop_stalled_supervisor - return if stopped? - - self.process = nil - - if standalone? - signal_queue << :TERM - interrupt - else - stop - end - end - def prune_dead_processes wrap_in_app_executor { SolidQueue::Process.prune(excluding: process) } ensure diff --git a/test/unit/async_supervisor_test.rb b/test/unit/async_supervisor_test.rb index 37a5b5387..ece1bc921 100644 --- a/test/unit/async_supervisor_test.rb +++ b/test/unit/async_supervisor_test.rb @@ -15,6 +15,70 @@ class AsyncSupervisorTest < ActiveSupport::TestCase assert_no_registered_processes end + test "stop when its registration has been pruned" do + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + supervisor = run_supervisor_as_thread + wait_for_registered_processes(4, timeout: 3.seconds) + + # Simulate another supervisor pruning this one's registration + find_processes_registered_as("Supervisor(async)").first.delete + + # The next heartbeat finds the registration gone: the rest of the system + # already considers this supervisor dead, so it must stop rather than run on + wait_while_with_timeout(3) { !supervisor.send(:stopped?) } + + assert supervisor.send(:stopped?) + + # The children deregister on the way down, and this supervisor's own row is + # already gone + wait_for_registered_processes(0, timeout: 3.seconds) + assert_no_registered_processes + ensure + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + supervisor&.stop + end + + test "stop when heartbeats stop returning for longer than the alive threshold" do + old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.3.seconds + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + # Only the supervisor's heartbeat blocks; the children stay healthy, so the + # supervisor stopping can only mean its own heartbeat watchdog acted + unblock_heartbeats = Concurrent::Event.new + SolidQueue::Process.class_eval do + alias_method :heartbeat_without_blocking, :heartbeat + define_method(:heartbeat) do + kind.start_with?("Supervisor") ? unblock_heartbeats.wait : heartbeat_without_blocking + end + end + + supervisor = run_supervisor_as_thread + wait_while_with_timeout(5) { !supervisor.send(:stopped?) } + + assert supervisor.send(:stopped?) + + # The children deregister on the way down. The supervisor's own stale row + # may or may not be left: once the local registration is dropped, a last + # prune can remove it, which is just what another supervisor would do + wait_while_with_timeout(3) { SolidQueue::Process.where.not(kind: "Supervisor(async)").any? } + skip_active_record_query_cache do + assert_empty SolidQueue::Process.where.not(kind: "Supervisor(async)") + end + ensure + unblock_heartbeats.set + SolidQueue::Process.class_eval do + if method_defined?(:heartbeat_without_blocking) + remove_method :heartbeat + alias_method :heartbeat, :heartbeat_without_blocking + remove_method :heartbeat_without_blocking + end + end + SolidQueue.process_alive_threshold = old_alive_threshold if old_alive_threshold + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + supervisor&.stop + end + test "stop when maintenance stops returning for longer than the stall threshold" do # The heartbeat interval must sit well below the alive threshold: a healthy # gap between heartbeat returns is already one interval plus a database diff --git a/test/unit/fork_supervisor_test.rb b/test/unit/fork_supervisor_test.rb index fd2bb27d0..6dc2d56f0 100644 --- a/test/unit/fork_supervisor_test.rb +++ b/test/unit/fork_supervisor_test.rb @@ -273,6 +273,25 @@ def register assert_equal "RuntimeError", failed.exception_class end + test "terminate forks when the supervisor's registration is pruned" do + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + pid = run_supervisor_as_fork + wait_for_registered_processes(4, timeout: 3.seconds) + + # Simulate another supervisor pruning this one's registration + find_processes_registered_as("Supervisor(fork)").first.delete + + # The next heartbeat finds the registration gone; the supervisor must stop + # through its signal pipeline, so its forks are terminated, not abandoned + wait_for_process_termination_with_timeout(pid, timeout: 5) + + assert_no_registered_processes + ensure + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + terminate_process(pid) if pid && process_exists?(pid) + end + test "terminate forks when maintenance stalls instead of abandoning them" do old_alive_threshold, SolidQueue.process_alive_threshold = SolidQueue.process_alive_threshold, 0.3.seconds old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds From fab72c39585dad2299c34ad3f2568dd938a3c70b Mon Sep 17 00:00:00 2001 From: Julik Tarkhanov Date: Mon, 28 Sep 2026 17:05:54 +0200 Subject: [PATCH 4/4] Stop a dispatcher whose maintenance stops returning The dispatcher's maintenance -- semaphore expiry, unblocking blocked executions, sweeping stalled batches -- runs in a Concurrent::TimerTask, which reschedules only once its task returns. A run blocked on an unresponsive database therefore stops maintenance ever running again, silently: nothing is raised for the task's observer to report, and concurrency-limited jobs stall system-wide while the dispatcher keeps polling as if healthy. Record when a maintenance run last returned, and watch it from a timer that reads nothing but memory, so it cannot block the same way. When the gap exceeds twice the maintenance interval, the dispatcher stops to be replaced: its supervisor starts a fresh dispatcher, maintenance task and all, exactly as it would for a stalled heartbeat. --- lib/solid_queue/dispatcher.rb | 30 +++++++++++++++++++- lib/solid_queue/dispatcher/maintenance.rb | 33 +++++++++++++++++----- test/unit/dispatcher_test.rb | 34 +++++++++++++++++++++++ 3 files changed, 89 insertions(+), 8 deletions(-) diff --git a/lib/solid_queue/dispatcher.rb b/lib/solid_queue/dispatcher.rb index 112399ef5..59fb01074 100644 --- a/lib/solid_queue/dispatcher.rb +++ b/lib/solid_queue/dispatcher.rb @@ -46,11 +46,39 @@ def dispatch_next_batch end def start_maintenance - maintenance&.start + return unless maintenance + + maintenance.start + launch_maintenance_watchdog + end + + # A stalled maintenance run means semaphores are no longer expired and + # blocked executions no longer unblocked, so concurrency-limited jobs + # stall system-wide -- silently, because a run that never returns raises + # nothing for the task's observer to report. The watchdog's check reads + # nothing but memory, so it cannot block the same way, and stopping to + # be replaced hands the work to a fresh dispatcher with a fresh + # maintenance task. + # + # The watchdog ticks at the heartbeat interval, like the supervisor's + # maintenance watchdog: not that heartbeats are involved, it is just the + # liveness-checking cadence SolidQueue already has, short against any + # sane stall threshold. + def launch_maintenance_watchdog + @maintenance_watchdog_task = Concurrent::TimerTask.new(execution_interval: SolidQueue.process_heartbeat_interval) do + stop_to_be_replaced if maintenance.stalled? + end + + @maintenance_watchdog_task.add_observer do |_, _, error| + handle_thread_error(error) if error + end + + @maintenance_watchdog_task.execute end def stop_maintenance maintenance&.stop + @maintenance_watchdog_task&.shutdown end def all_work_completed? diff --git a/lib/solid_queue/dispatcher/maintenance.rb b/lib/solid_queue/dispatcher/maintenance.rb index 596dcb35b..b87a178e5 100644 --- a/lib/solid_queue/dispatcher/maintenance.rb +++ b/lib/solid_queue/dispatcher/maintenance.rb @@ -25,15 +25,14 @@ def metadata { concurrency_maintenance_interval: (interval if concurrency?), batch_maintenance: batches? } end + # How far past its own interval a maintenance run is allowed to go before we + # treat it as stalled rather than slow. One full missed cycle of slack. + STALL_FACTOR = 2 + def start - @maintenance_task = Concurrent::TimerTask.new(run_now: true, execution_interval: interval) do - if concurrency? - expire_semaphores - unblock_blocked_executions - end + @last_run_returned_at = Concurrent::AtomicReference.new(SolidQueue::Timer.monotonic_time_now) - sweep_stalled_batches if batches? - end + @maintenance_task = Concurrent::TimerTask.new(run_now: true, execution_interval: interval) { run } @maintenance_task.add_observer do |_, _, error| handle_thread_error(error) if error @@ -42,11 +41,31 @@ def start @maintenance_task.execute end + # Whether maintenance has stopped returning altogether, as opposed to + # failing: a run that raises is reported by the task's observer and gets + # rescheduled, but one blocked on an unresponsive database never returns to + # its TimerTask, which reschedules only once its task completes, so no + # maintenance would run ever again. + def stalled? + SolidQueue::Timer.monotonic_time_now - @last_run_returned_at.get > STALL_FACTOR * interval + end + def stop @maintenance_task&.shutdown end private + def run + if concurrency? + expire_semaphores + unblock_blocked_executions + end + + sweep_stalled_batches if batches? + ensure + @last_run_returned_at.set(SolidQueue::Timer.monotonic_time_now) + end + def expire_semaphores wrap_in_app_executor do Semaphore.expired.in_batches(of: batch_size, &:delete_all) diff --git a/test/unit/dispatcher_test.rb b/test/unit/dispatcher_test.rb index 5c26f8287..f99090c4a 100644 --- a/test/unit/dispatcher_test.rb +++ b/test/unit/dispatcher_test.rb @@ -65,6 +65,40 @@ class DispatcherTest < ActiveSupport::TestCase end end + test "stop to be replaced when maintenance stops returning for longer than the stall threshold" do + old_heartbeat_interval, SolidQueue.process_heartbeat_interval = SolidQueue.process_heartbeat_interval, 0.1.seconds + + # A maintenance run that never returns, so its TimerTask is never + # rescheduled and semaphores would never be expired again + unblock_maintenance = Concurrent::Event.new + SolidQueue::Dispatcher::Maintenance.class_eval do + alias_method :run_without_blocking, :run + define_method(:run) { unblock_maintenance.wait } + end + + dispatcher = SolidQueue::Dispatcher.new(polling_interval: 0.1, batch_size: 10, concurrency_maintenance_interval: 0.5) + dispatcher.start + wait_for_registered_processes(1, timeout: 1.second) + + assert dispatcher.alive? + + # The dispatcher deregisters locally and stops its run loop, so a + # supervisor would replace it, maintenance task and all + wait_while_with_timeout(3) { dispatcher.alive? } + assert_not dispatcher.alive? + ensure + unblock_maintenance.set + SolidQueue::Dispatcher::Maintenance.class_eval do + if private_method_defined?(:run_without_blocking) || method_defined?(:run_without_blocking) + remove_method :run + alias_method :run, :run_without_blocking + remove_method :run_without_blocking + end + end + SolidQueue.process_heartbeat_interval = old_heartbeat_interval if old_heartbeat_interval + dispatcher&.stop + end + test "ConcurrencyMaintenance remains constructible with its original signature" do maintenance = SolidQueue::Dispatcher::ConcurrencyMaintenance.new(600, 100)