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/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/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 23e321f47..ecad8f9ad 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,51 @@ 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_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_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 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..ece1bc921 100644 --- a/test/unit/async_supervisor_test.rb +++ b/test/unit/async_supervisor_test.rb @@ -15,6 +15,156 @@ 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 + # 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/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) diff --git a/test/unit/fork_supervisor_test.rb b/test/unit/fork_supervisor_test.rb index 0d8092278..6dc2d56f0 100644 --- a/test/unit/fork_supervisor_test.rb +++ b/test/unit/fork_supervisor_test.rb @@ -273,6 +273,66 @@ 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 + + # 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) 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) }