diff --git a/lib/solid_queue.rb b/lib/solid_queue.rb index 0f7a40d16..d8c963286 100644 --- a/lib/solid_queue.rb +++ b/lib/solid_queue.rb @@ -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" @@ -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 diff --git a/lib/solid_queue/engine.rb b/lib/solid_queue/engine.rb index 718a6696a..482dc1acb 100644 --- a/lib/solid_queue/engine.rb +++ b/lib/solid_queue/engine.rb @@ -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 diff --git a/lib/solid_queue/pidfile.rb b/lib/solid_queue/pidfile.rb new file mode 100644 index 000000000..14ac2a090 --- /dev/null +++ b/lib/solid_queue/pidfile.rb @@ -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 diff --git a/lib/solid_queue/supervisor.rb b/lib/solid_queue/supervisor.rb index dc4ded998..8efb70981 100644 --- a/lib/solid_queue/supervisor.rb +++ b/lib/solid_queue/supervisor.rb @@ -18,6 +18,7 @@ def initialize(runners) end def start + setup_pidfile register_signal_handlers start_process_prune @@ -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 @@ -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 diff --git a/test/integration/processes_lifecycle_test.rb b/test/integration/processes_lifecycle_test.rb index 6ce9b0ced..9a5138c4e 100644 --- a/test/integration/processes_lifecycle_test.rb +++ b/test/integration/processes_lifecycle_test.rb @@ -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 @@ -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 diff --git a/test/test_helper.rb b/test/test_helper.rb index c342b824f..ab9cbde06 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -18,6 +18,7 @@ class ActiveSupport::TestCase teardown do JobBuffer.clear + File.delete(SolidQueue.supervisor_pidfile) if File.exist?(SolidQueue.supervisor_pidfile) end private @@ -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 diff --git a/test/unit/supervisor_test.rb b/test/unit/supervisor_test.rb new file mode 100644 index 000000000..0aa7beaf3 --- /dev/null +++ b/test/unit/supervisor_test.rb @@ -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