Skip to content
Merged
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
7 changes: 2 additions & 5 deletions lib/solid_queue/processes/runnable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ def run_in_mode(&block)
case
when running_as_fork?
@boot_guard = BootGuards::ForkGuard.new
fork(&block).tap { @boot_guard.start }
create_fork(&block).tap { @boot_guard.start }
when running_async?
@boot_guard = BootGuards::NullGuard.new
@thread = create_thread(&block)
Expand All @@ -64,10 +64,7 @@ def run_in_mode(&block)
def boot
SolidQueue.instrument(:start_process, process: self) do
run_callbacks(:boot) do
if running_as_fork?
register_signal_handlers
set_procline
end
set_procline if running_as_fork?
end
end

Expand Down
7 changes: 7 additions & 0 deletions lib/solid_queue/processes/supervised.rb
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,13 @@ def supervised?
supervisor.present?
end

def create_fork(&block)
fork do
register_signal_handlers
block.call
end
end

def register_signal_handlers
%w[ INT TERM ].each do |signal|
trap(signal) do
Expand Down
45 changes: 29 additions & 16 deletions lib/solid_queue/supervisor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,13 @@ def start
run_start_hooks

start_processes
launch_maintenance_task

supervise
if stopped?
shutdown
else
launch_maintenance_task
supervise
end
end

def stop
Expand All @@ -65,27 +69,33 @@ def boot
end

def start_processes
configuration.configured_processes.each { |configured_process| start_process(configured_process) }
configuration.configured_processes.each do |configured_process|
# Honour signals that arrive during boot or start hooks: a queued TERM
# stops us here, before starting children, instead of in #supervise,
# after all of them have been started
break if time_to_stop?

start_process(configured_process)
end
end

def supervise
loop do
break if stopped?

if standalone?
set_procline
process_signal_queue
end

unless stopped?
check_and_replace_terminated_processes
interruptible_sleep(1.second)
end
until time_to_stop?
set_procline
check_and_replace_terminated_processes
interruptible_sleep(1.second)
end
ensure
shutdown
end

# Process any signals queued while we were busy and report whether
# we've been asked to stop
def time_to_stop?
process_signal_queue
stopped?
end

def start_process(configured_process)
process_instance = configured_process.instantiate.tap do |instance|
instance.supervised_by process
Expand Down Expand Up @@ -139,7 +149,10 @@ def shutdown
end

def set_procline
procline "supervising #{configured_processes.keys.join(", ")}"
# Embedded supervisors don't own their process's title
if standalone?
procline "supervising #{configured_processes.keys.join(", ")}"
end
end

def sync_std_streams
Expand Down
3 changes: 3 additions & 0 deletions lib/solid_queue/supervisor/signals.rb
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@ def restore_default_signal_handlers
end

def process_signal_queue
# Embedded supervisors don't own their process's signals
return unless standalone?

while signal = signal_queue.shift
handle_signal(signal)
end
Expand Down
82 changes: 82 additions & 0 deletions test/integration/supervisor_boot_signal_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
# frozen_string_literal: true

require "test_helper"

class SupervisorBootSignalTest < ActiveSupport::TestCase
self.use_transactional_tests = false

setup do
@marker_path = Rails.root.join("tmp/solid_queue_start_hook_#{SecureRandom.hex(8)}")
@release_path = Rails.root.join("tmp/solid_queue_release_start_hook_#{SecureRandom.hex(8)}")
end

teardown do
File.delete(@marker_path) if File.exist?(@marker_path)
File.delete(@release_path) if File.exist?(@release_path)
end

# Regression for https://github.com/rails/solid_queue/issues/755:
# TERM received while the supervisor is still running start hooks used to sit
# in the signal queue until after workers were forked. Those workers then
# inherited the supervisor's enqueue-only trap and kept claiming jobs until
# SIGKILL.
test "TERM during supervisor start hooks exits without starting workers" do
resetting_hooks do
marker_path = @marker_path
release_path = @release_path

SolidQueue.on_start do
File.write(marker_path, Process.pid.to_s)
Timeout.timeout(5) { sleep 0.05 until File.exist?(release_path) }
end

SolidQueue.on_worker_start do
JobResult.create!(queue_name: "background", status: "hook_called", value: "worker_started")
end

3.times { |i| StoreResultJob.set(queue: :background).perform_later("should_not_run_#{i}") }

pid = run_supervisor_as_fork(
workers: [ { queues: :background, threads: 1, polling_interval: 0.1 } ],
dispatchers: []
)

wait_while_with_timeout!(5.seconds) { !File.exist?(marker_path) }
Process.kill(:TERM, pid)
File.write(release_path, "1")

wait_for_process_termination_with_timeout(pid, timeout: SolidQueue.shutdown_timeout + 5.seconds)

skip_active_record_query_cache do
assert_equal 0, JobResult.where(value: "worker_started").count
assert_equal 0, JobResult.where(status: "completed").count
assert_equal 0, SolidQueue::ClaimedExecution.count
assert_equal 3, SolidQueue::ReadyExecution.count
assert SolidQueue::Process.none?
end
end
end

private
def resetting_hooks
reset_hooks(SolidQueue::Supervisor) do
reset_hooks(SolidQueue::Worker) do
yield
end
end
end

def reset_hooks(process)
exit_hooks = process.lifecycle_hooks[:exit]
start_hooks = process.lifecycle_hooks[:start]
stop_hooks = process.lifecycle_hooks[:stop]
process.lifecycle_hooks[:exit] = []
process.lifecycle_hooks[:start] = []
process.lifecycle_hooks[:stop] = []
yield
ensure
process.lifecycle_hooks[:exit] = exit_hooks
process.lifecycle_hooks[:start] = start_hooks
process.lifecycle_hooks[:stop] = stop_hooks
end
end
Loading