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
30 changes: 29 additions & 1 deletion lib/solid_queue/dispatcher.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
33 changes: 26 additions & 7 deletions lib/solid_queue/dispatcher/maintenance.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down
36 changes: 36 additions & 0 deletions lib/solid_queue/processes/registrable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down
34 changes: 34 additions & 0 deletions lib/solid_queue/supervisor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand Down
43 changes: 43 additions & 0 deletions lib/solid_queue/supervisor/maintenance.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
Loading
Loading