diff --git a/lib/solid_queue/async_supervisor.rb b/lib/solid_queue/async_supervisor.rb index f6ab3834..c22af014 100644 --- a/lib/solid_queue/async_supervisor.rb +++ b/lib/solid_queue/async_supervisor.rb @@ -38,12 +38,26 @@ def perform_graceful_termination process_instances.values.each(&:stop) Timer.wait_until(SolidQueue.shutdown_timeout, -> { all_processes_terminated? }) + + # The exit! that follows skips the shutdown callbacks, so release the jobs + # of the threads still running now instead of leaving them to be pruned + attempt_to_deregister unless all_processes_terminated? end def perform_immediate_termination exit! end + def attempt_to_deregister + wrap_in_app_executor do + # Not our own copy: it shares the threads' registrations in memory, and + # those of threads that already deregistered are destroyed and frozen + SolidQueue::Process.find_by(id: process_id)&.deregister + end + rescue StandardError => error + handle_thread_error(error) + end + def all_processes_terminated? process_instances.values.none?(&:alive?) end diff --git a/test/integration/async_processes_lifecycle_test.rb b/test/integration/async_processes_lifecycle_test.rb index dc8e4c7e..8185835d 100644 --- a/test/integration/async_processes_lifecycle_test.rb +++ b/test/integration/async_processes_lifecycle_test.rb @@ -173,14 +173,11 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase "registered processes: #{SolidQueue::Process.all.map { |p| { id: p.id, kind: p.kind, pid: p.pid, last_heartbeat_at: p.last_heartbeat_at } }.inspect}" end - # After shutdown, the pause job may be either: - # - claimed (exit! called, no cleanup) OR - # - ready (graceful exit, job released back to queue) - # Both are valid outcomes depending on the timing race between supervisor and worker timeouts. - skip_active_record_query_cache do - job = SolidQueue::Job.find_by(active_job_id: pause.job_id) - assert job.claimed? || job.ready?, "Expected pause job to be claimed or ready, but was neither" - end + # Workers were shutdown without a chance to terminate orderly, but + # since they're linked to the supervisor, the supervisor deregistering + # also deregistered them and released claimed jobs + assert_job_status(pause, :ready) + assert_clean_termination end test "process some jobs that raise errors" do