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
3 changes: 3 additions & 0 deletions lib/solid_queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

require "solid_queue/app_executor"
require "solid_queue/interruptible"
require "solid_queue/pidfile"
require "solid_queue/signals"
require "solid_queue/configuration"
require "solid_queue/runner"
Expand All @@ -21,4 +22,6 @@ module SolidQueue
mattr_accessor :process_alive_threshold, default: 5.minutes

mattr_accessor :shutdown_timeout, default: 5.seconds

mattr_accessor :supervisor_pidfile
end
1 change: 1 addition & 0 deletions lib/solid_queue/engine.rb
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ class Engine < ::Rails::Engine
SolidQueue.process_heartbeat_interval = app.config.solid_queue.process_heartbeat_interval || 60.seconds
SolidQueue.process_alive_threshold = app.config.solid_queue.process_alive_threshold || 5.minutes
SolidQueue.shutdown_timeout = app.config.solid_queue.shutdown_timeout || 5.seconds
SolidQueue.supervisor_pidfile = app.config.solid_queue.supervisor_pidfile || app.root.join("tmp", "pids", "solid_queue_supervisor.pid")
end
end

Expand Down
56 changes: 56 additions & 0 deletions lib/solid_queue/pidfile.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
# frozen_string_literal: true

module SolidQueue
class Pidfile
def initialize(path)
@path = path
@pid = ::Process.pid
end

def setup
check_status
write_file
set_at_exit_hook
end

def delete
delete_file
end

private
attr_reader :path, :pid

def check_status
if ::File.exist?(path)
existing_pid = ::File.open(path).read.strip.to_i
existing_pid > 0 && ::Process.kill(0, existing_pid)

already_running!
end
rescue Errno::ESRCH => e
# Process is dead, ignore, just delete the file
delete
rescue Errno::EPERM
already_running!
end

def write_file
::File.open(path, ::File::CREAT | ::File::EXCL | ::File::WRONLY) { |file| file.write(pid.to_s) }
rescue Errno::EEXIST
check_status
retry
end

def set_at_exit_hook
at_exit { delete if ::Process.pid == pid }
end

def delete_file
::File.delete(path) if ::File.exist?(path)
end

def already_running!
abort "A Solid Queue supervisor is already running. Check #{path}"
end
end
end
14 changes: 13 additions & 1 deletion lib/solid_queue/supervisor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ def initialize(runners)
end

def start
setup_pidfile
register_signal_handlers
start_process_prune

Expand All @@ -31,11 +32,18 @@ def start
ensure
stop_process_prune
restore_default_signal_handlers
delete_pidfile
end

private
attr_reader :runners, :forks

def setup_pidfile
@pidfile = if SolidQueue.supervisor_pidfile
Pidfile.new(SolidQueue.supervisor_pidfile).tap(&:setup)
end
end

def start_process_prune
@prune_task = Concurrent::TimerTask.new(run_now: true, execution_interval: SolidQueue.process_alive_threshold) { prune_dead_processes }
@prune_task.execute
Expand Down Expand Up @@ -76,7 +84,11 @@ def quit_runners
end

def stop_process_prune
@prune_task.shutdown
@prune_task&.shutdown
end

def delete_pidfile
@pidfile&.delete
end

def prune_dead_processes
Expand Down
51 changes: 0 additions & 51 deletions test/integration/processes_lifecycle_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -181,12 +181,6 @@ class ProcessLifecycleTest < ActiveSupport::TestCase
end

private
def run_supervisor_as_fork
fork do
SolidQueue::Supervisor.start
end
end

def terminate_supervisor
terminate_process(@pid)
end
Expand All @@ -204,51 +198,6 @@ def assert_clean_termination
assert_no_claimed_jobs
end

def terminate_process(pid, from_parent: true)
signal_process(pid, :TERM)
wait_for_process_with_timeout(pid, from_parent: from_parent)
end

def signal_process(pid, signal, wait: nil)
Thread.new do
sleep(wait) if wait
Process.kill(signal, pid)
end
end

def wait_for_process_with_timeout(pid, timeout: 10, from_parent: true)
Timeout.timeout(timeout) do
if from_parent
Process.waitpid(pid)
assert 0, $?.exitstatus
else
loop do
break unless process_exists?(pid)
sleep(0.1)
end
end
end
rescue Timeout::Error
Process.kill(:KILL, pid)
raise
end

def process_exists?(pid)
Process.getpgid(pid)
true
rescue Errno::ESRCH
false
end

def wait_for_registered_processes(count, timeout: 10.seconds)
Timeout.timeout(timeout) do
while SolidQueue::Process.count < count do
sleep 0.25
end
end
rescue Timeout::Error
end

def assert_registered_processes_for(*queues)
uncached do
registered_queues = SolidQueue::Process.all.map { |process| process.metadata["queue"] }.compact
Expand Down
52 changes: 52 additions & 0 deletions test/test_helper.rb
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
class ActiveSupport::TestCase
teardown do
JobBuffer.clear
File.delete(SolidQueue.supervisor_pidfile) if File.exist?(SolidQueue.supervisor_pidfile)
end

private
Expand All @@ -29,4 +30,55 @@ def wait_for_jobs_to_finish_for(timeout = 10.seconds)
end
rescue Timeout::Error
end

def run_supervisor_as_fork(**options)
fork do
SolidQueue::Supervisor.start(**options)
end
end

def wait_for_registered_processes(count, timeout: 10.seconds)
Timeout.timeout(timeout) do
while SolidQueue::Process.count < count do
sleep 0.25
end
end
rescue Timeout::Error
end

def terminate_process(pid, timeout: 10, signal: :TERM, from_parent: true)
signal_process(pid, signal)
wait_for_process_termination_with_timeout(pid, timeout: timeout, from_parent: from_parent)
end

def wait_for_process_termination_with_timeout(pid, timeout: 10, from_parent: true, exitstatus: 0)
Timeout.timeout(timeout) do
if from_parent
Process.waitpid(pid)
assert exitstatus, $?.exitstatus
else
loop do
break unless process_exists?(pid)
sleep(0.1)
end
end
end
rescue Timeout::Error
signal_process(pid, :KILL)
raise
end

def signal_process(pid, signal, wait: nil)
Thread.new do
sleep(wait) if wait
Process.kill(signal, pid)
end
end

def process_exists?(pid)
Process.getpgid(pid)
true
rescue Errno::ESRCH
false
end
end
61 changes: 61 additions & 0 deletions test/unit/supervisor_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
require "test_helper"

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

setup do
FileUtils.mkdir_p Rails.application.root.join("tmp")
@previous_pidfile = SolidQueue.supervisor_pidfile
@pidfile = Rails.application.root.join("tmp/pidfile_#{SecureRandom.hex}.pid")
SolidQueue.supervisor_pidfile = @pidfile
end

teardown do
SolidQueue.supervisor_pidfile = @previous_pidfile
File.delete(@pidfile) if File.exist?(@pidfile)
end

test "create and delete pidfile" do
assert_not File.exist?(@pidfile)

pid = run_supervisor_as_fork(mode: :all)
wait_for_registered_processes(3)

assert File.exist?(@pidfile)
assert_equal pid, File.read(@pidfile).strip.to_i

terminate_process(pid)

assert_not File.exist?(@pidfile)
end

test "abort if there's already a pidfile for a supervisor" do
File.write(@pidfile, ::Process.pid.to_s)

pid = run_supervisor_as_fork(mode: :all)
wait_for_registered_processes(3)

assert File.exist?(@pidfile)
assert_not_equal pid, File.read(@pidfile).strip.to_i

wait_for_process_termination_with_timeout(pid, exitstatus: 1)
end

test "deletes previous pidfile if the owner is dead" do
pid = run_supervisor_as_fork(mode: :all)
wait_for_registered_processes(3)

terminate_process(pid, signal: :KILL)

assert File.exist?(@pidfile)
assert_equal pid, File.read(@pidfile).strip.to_i

pid = run_supervisor_as_fork(mode: :all)
wait_for_registered_processes(3)

assert File.exist?(@pidfile)
assert_equal pid, File.read(@pidfile).strip.to_i

terminate_process(pid)
end
end