From 34098d744fc19517d5fb53352a71a580b5e0347a Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 21:00:26 +0200 Subject: [PATCH 01/12] Report rescued recurring enqueue errors to Rails.error RecurringTask#enqueue was swallowing Job::EnqueueError (and other-adapter enqueue failures) without calling Rails.error.report, so production outages like a read-only primary during failover only produced log lines and missed Sentry Issues. Align with the gem's on_thread_error default. Fixes #746 --- app/models/solid_queue/recurring_task.rb | 8 ++++ .../models/solid_queue/recurring_task_test.rb | 47 +++++++++++++++++++ 2 files changed, 55 insertions(+) diff --git a/app/models/solid_queue/recurring_task.rb b/app/models/solid_queue/recurring_task.rb index 9bb634e6..51bf7f31 100644 --- a/app/models/solid_queue/recurring_task.rb +++ b/app/models/solid_queue/recurring_task.rb @@ -85,6 +85,7 @@ def enqueue(at:) perform_later.tap do |job| unless job.successfully_enqueued? + report_enqueue_error(job.enqueue_error, at: at) payload[:enqueue_error] = job.enqueue_error&.message end end @@ -97,6 +98,7 @@ def enqueue(at:) payload[:skipped] = true false rescue Job::EnqueueError => error + report_enqueue_error(error, at: at) payload[:enqueue_error] = error.message false end @@ -180,5 +182,11 @@ def job_class def enqueue_options { queue: queue_name, priority: priority }.compact end + + def report_enqueue_error(error, at:) + return unless error + + Rails.error.report(error, handled: true, source: "solid_queue", context: { task: key, at: at }) + end end end diff --git a/test/models/solid_queue/recurring_task_test.rb b/test/models/solid_queue/recurring_task_test.rb index dba9d6b9..f7325eb7 100644 --- a/test/models/solid_queue/recurring_task_test.rb +++ b/test/models/solid_queue/recurring_task_test.rb @@ -96,6 +96,46 @@ def perform end end + test "reports Job::EnqueueError to Rails.error when enqueuing via Solid Queue" do + SolidQueue::Job.stubs(:create!).raises(ActiveRecord::Deadlocked) + subscriber = ErrorBuffer.new + at = Time.now + + with_error_subscriber(subscriber) do + task = recurring_task_with(class_name: "JobWithoutArguments") + task.enqueue(at: at) + end + + assert_equal 1, subscriber.errors.count + error, options = subscriber.errors.first + assert_kind_of SolidQueue::Job::EnqueueError, error + assert_match "ActiveRecord::Deadlocked", error.message + assert_equal true, options[:handled] + assert_equal "solid_queue", options[:source] + assert_equal "task-id", options[:context][:task] + assert_equal at, options[:context][:at] + end + + test "reports enqueue error to Rails.error when using another adapter" do + ActiveJob::QueueAdapters::AsyncAdapter.any_instance.stubs(:enqueue).raises(ActiveJob::EnqueueError.new("All is broken")) + subscriber = ErrorBuffer.new + at = Time.now + + with_error_subscriber(subscriber) do + task = recurring_task_with(class_name: "JobUsingAsyncAdapter") + task.enqueue(at: at) + end + + assert_equal 1, subscriber.errors.count + error, options = subscriber.errors.first + assert_kind_of ActiveJob::EnqueueError, error + assert_equal "All is broken", error.message + assert_equal true, options[:handled] + assert_equal "solid_queue", options[:source] + assert_equal "task-id", options[:context][:task] + assert_equal at, options[:context][:at] + end + test "error when enqueuing job because of concurrency controls and discard" do JobWithConcurrencyControlsAndDiscard.perform_later @@ -289,4 +329,11 @@ def run_all_jobs_inline worker.start end end + + def with_error_subscriber(subscriber) + Rails.error.subscribe(subscriber) + yield + ensure + Rails.error.unsubscribe(subscriber) if Rails.error.respond_to?(:unsubscribe) + end end From 283bd69949557933fe92eac4dc75dfa0aabee6af Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 21:26:58 +0200 Subject: [PATCH 02/12] Wait for async supervisor exit before asserting clean termination The forked lifecycle test already waits up to shutdown_timeout for the supervisor PID to exit after repeated TERM signals. The async counterpart only slept for 1s, which is flaky under load (e.g. Ruby 4 + rails_main + MySQL). --- test/integration/async_processes_lifecycle_test.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/integration/async_processes_lifecycle_test.rb b/test/integration/async_processes_lifecycle_test.rb index fd284210..e4f74ec6 100644 --- a/test/integration/async_processes_lifecycle_test.rb +++ b/test/integration/async_processes_lifecycle_test.rb @@ -58,7 +58,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase signal_process(@pid, :TERM, wait: 0.1.second) end - sleep(1.second) + wait_while_with_timeout(SolidQueue.shutdown_timeout + 1.second) { process_exists?(@pid) } assert_clean_termination end From 83093d0f9fd3d192b7de7890f94bc83d8ba141a7 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 21:32:08 +0200 Subject: [PATCH 03/12] Wait for failed executions in jobs lifecycle failure test wait_for_jobs_to_finish_for waits on finished_at, which failed jobs never set, so the test always timed out. Wait for FailedExecution records instead and stop the worker before asserting to avoid races with in-flight threads. --- test/integration/jobs_lifecycle_test.rb | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/test/integration/jobs_lifecycle_test.rb b/test/integration/jobs_lifecycle_test.rb index 8444c375..53e8c97c 100644 --- a/test/integration/jobs_lifecycle_test.rb +++ b/test/integration/jobs_lifecycle_test.rb @@ -40,11 +40,15 @@ class JobsLifecycleTest < ActiveSupport::TestCase @dispatcher.start @worker.start - wait_for_jobs_to_finish_for(3.seconds) + wait_while_with_timeout(3.seconds) { SolidQueue::FailedExecution.count < 2 } + + @worker.stop + @dispatcher.stop message = "raised ExpectedTestError for the 1st time" assert_equal [ "A: #{message}", "B: #{message}" ], JobBuffer.values.sort + assert_equal 2, SolidQueue::FailedExecution.count assert_empty SolidQueue::Job.finished end From 784f8b85ec958dc4ed7e9267395a0eec03e31555 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 21:58:30 +0200 Subject: [PATCH 04/12] Give concurrency controls tests more timing headroom on CI Setup only waited 0.5s for processes to register, and a few scenarios used tight pause windows that flake under MySQL CI load. Widen registration and job-completion waits, and leave more margin so C still finishes last in the throttled sequence test. --- test/integration/concurrency_controls_test.rb | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index a12c48e4..f9a18351 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -13,7 +13,7 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase @pid = run_supervisor_as_fork(workers: [ default_worker ], dispatchers: [ dispatcher ]) - wait_for_registered_processes(5, timeout: 0.5.second) # 3 workers working the default queue + dispatcher + supervisor + wait_for_registered_processes(5, timeout: 3.seconds) # 3 workers working the default queue + dispatcher + supervisor end teardown do @@ -48,7 +48,7 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.set(wait: (1 + i * 0.1).seconds).perform_later(@result, name: name) end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs assert_stored_sequence @result, ("A".."K").to_a @@ -57,26 +57,26 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase test "run several jobs over the same record limiting concurrency" do incr = 0 # C is the last one to update the record - # A: 0 to 0.5 - # B: 0 to 1.0 - # C: 0 to 1.5 + # A: 0 to 1.0 + # B: 0 to 1.5 + # C: 0 to 2.0 assert_no_difference -> { SolidQueue::BlockedExecution.count } do ("A".."C").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (0.5 + incr).seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (1.0 + incr).seconds) incr += 0.5 end end sleep(0.01) # To ensure these aren't picked up before ABC - # D to H: 0.51 to 0.76 (starting after A finishes, and in order, 5 * 0.05 = 0.25) + # D to H: starting after A finishes, and in order; keep pauses short so they finish before B and C # These would finish all before B and C assert_difference -> { SolidQueue::BlockedExecution.count }, +5 do ("D".."H").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.05.seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.02.seconds) end end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs # C would have started in the beginning, seeing the status empty, and would finish after From 05bffe8f2b1ab839c965e5ab15238352913b13df Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 22:15:03 +0200 Subject: [PATCH 05/12] Reduce cross-test pollution in concurrency and recurring tests Wait for supervisor teardown to deregister processes, adopt the claimed- release wait from #734, and scope recurring JobResult assertions to the custom_status rows this test creates so leftover StoreResultJob rows cannot fail the suite. --- test/integration/concurrency_controls_test.rb | 28 ++++++++++--------- test/integration/recurring_tasks_test.rb | 9 +++--- 2 files changed, 19 insertions(+), 18 deletions(-) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index f9a18351..12cc7585 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -17,7 +17,10 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase end teardown do - terminate_process(@pid) if process_exists?(@pid) + if @pid && process_exists?(@pid) + terminate_process(@pid) + end + wait_for_registered_processes(0, timeout: 3.seconds) end test "run several conflicting jobs over the same record without overlapping" do @@ -57,22 +60,22 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase test "run several jobs over the same record limiting concurrency" do incr = 0 # C is the last one to update the record - # A: 0 to 1.0 - # B: 0 to 1.5 - # C: 0 to 2.0 + # A: 0 to 0.5 + # B: 0 to 1.0 + # C: 0 to 1.5 assert_no_difference -> { SolidQueue::BlockedExecution.count } do ("A".."C").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (1.0 + incr).seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (0.5 + incr).seconds) incr += 0.5 end end sleep(0.01) # To ensure these aren't picked up before ABC - # D to H: starting after A finishes, and in order; keep pauses short so they finish before B and C + # D to H: 0.51 to 0.76 (starting after A finishes, and in order, 5 * 0.05 = 0.25) # These would finish all before B and C assert_difference -> { SolidQueue::BlockedExecution.count }, +5 do ("D".."H").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.02.seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.05.seconds) end end @@ -94,7 +97,7 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.perform_later(@result, name: name) end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_equal 3, SolidQueue::FailedExecution.count assert_stored_sequence @result, [ "B", "D", "F" ] + ("G".."K").to_a @@ -166,12 +169,11 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.perform_later(@result, name: "I'll be released to ready", pause: SolidQueue.shutdown_timeout + 10.seconds) job = SolidQueue::Job.last - sleep(0.2) - assert job.claimed? + wait_for(timeout: 2.seconds) { job.reload.claimed? } - # This won't leave time to the job to finish - signal_process(@pid, :TERM, wait: 0.1.second) - sleep(SolidQueue.shutdown_timeout + 0.6.seconds) + # This won't leave time to the job to finish, so the worker should + # release it back to ready during shutdown. + terminate_process(@pid) assert_not job.reload.finished? assert job.reload.ready? diff --git a/test/integration/recurring_tasks_test.rb b/test/integration/recurring_tasks_test.rb index f2fc7145..64059e5a 100644 --- a/test/integration/recurring_tasks_test.rb +++ b/test/integration/recurring_tasks_test.rb @@ -24,11 +24,10 @@ class RecurringTasksTest < ActiveSupport::TestCase assert_equal "StoreResultJob", job.class_name end - assert JobResult.count >= 2 - JobResult.all.each do |result| - assert_equal "custom_status", result.status - assert_equal "42", result.value - end + # Scope to this recurring task's results. Other tests may leave JobResult + # rows (e.g. status "completed") that would fail a JobResult.all assertion. + results = JobResult.where(status: "custom_status", value: "42") + assert results.count >= 2 end end From d4528d2ebf3e4c2f014183c337cf8b5869d5db99 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Fri, 17 Jul 2026 22:46:47 +0200 Subject: [PATCH 06/12] Isolate concurrency tests from recycled JobResult primary keys Forked StoreResultJob workers from earlier tests can still write after teardown and overwrite a JobResult whose id was reused, turning status into "completed". Clear processes/records before creating @result, wait for claimed slots before enqueueing blocked jobs, and tighten discard setup timing. --- test/integration/concurrency_controls_test.rb | 22 +++++++++++++------ 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index 12cc7585..9f748ea7 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -6,21 +6,27 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase self.use_transactional_tests = false setup do - @result = JobResult.create!(queue_name: "default", status: "") + # Previous tests may leave forked workers briefly alive; those can still write to + # JobResult rows whose primary keys get reused by create! below (e.g. overwriting + # status with StoreResultJob's default "completed"). + wait_for_registered_processes(0, timeout: 5.seconds) + destroy_records default_worker = { queues: "default", polling_interval: 0.1, processes: 3, threads: 2 } dispatcher = { polling_interval: 0.1, batch_size: 200, concurrency_maintenance_interval: 1 } @pid = run_supervisor_as_fork(workers: [ default_worker ], dispatchers: [ dispatcher ]) + wait_for_registered_processes(5, timeout: 3.seconds) # 3 workers + dispatcher + supervisor - wait_for_registered_processes(5, timeout: 3.seconds) # 3 workers working the default queue + dispatcher + supervisor + @result = JobResult.create!(queue_name: "default", status: "") end teardown do if @pid && process_exists?(@pid) terminate_process(@pid) end - wait_for_registered_processes(0, timeout: 3.seconds) + wait_for_registered_processes(0, timeout: 5.seconds) + destroy_records end test "run several conflicting jobs over the same record without overlapping" do @@ -70,8 +76,10 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase end end - sleep(0.01) # To ensure these aren't picked up before ABC - # D to H: 0.51 to 0.76 (starting after A finishes, and in order, 5 * 0.05 = 0.25) + # Wait until A/B/C hold the concurrency slots so D–H are blocked instead of claimed + wait_for(timeout: 2.seconds) { SolidQueue::ClaimedExecution.count >= 3 } + + # D to H: starting after A finishes, and in order, 5 * 0.05 = 0.25 # These would finish all before B and C assert_difference -> { SolidQueue::BlockedExecution.count }, +5 do ("D".."H").each do |name| @@ -198,8 +206,8 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase end test "discard jobs when concurrency limit is reached with on_conflict: :discard" do - job1 = DiscardableUpdateResultJob.perform_later(@result, name: "1", pause: 3) - sleep(0.1) + job1 = DiscardableUpdateResultJob.perform_later(@result, name: "1", pause: 1.second) + wait_for(timeout: 2.seconds) { SolidQueue::Job.find_by(active_job_id: job1.job_id)&.claimed? } # should be discarded due to concurrency limit job2 = DiscardableUpdateResultJob.perform_later(@result, name: "2") From 37eae2723814e71455eed77ce9a5fea6410ebd6b Mon Sep 17 00:00:00 2001 From: Pissardo Date: Sun, 19 Jul 2026 23:14:27 +0200 Subject: [PATCH 07/12] Sync integration test flake fixes from #758 Bring the concurrency/lifecycle/orphan wait hardenings that made CI green on the dedicated flake PR. --- test/integration/concurrency_controls_test.rb | 25 ++++++++++--------- .../forked_processes_lifecycle_test.rb | 7 +++++- test/unit/async_supervisor_test.rb | 8 ++++-- 3 files changed, 25 insertions(+), 15 deletions(-) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index 9f748ea7..c3c31e9d 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -65,34 +65,35 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase test "run several jobs over the same record limiting concurrency" do incr = 0 - # C is the last one to update the record - # A: 0 to 0.5 - # B: 0 to 1.0 - # C: 0 to 1.5 + # C should finish last so its in-memory status (started against "") overwrites the rest. + # Keep a margin for CI scheduling delay when draining D–H on the freed slot + # (A: 1s, B: 1.5s, C: 2.5s; D–H: 5 × 0.01s plus poll delay). assert_no_difference -> { SolidQueue::BlockedExecution.count } do ("A".."C").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (0.5 + incr).seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (1.0 + incr).seconds) incr += 0.5 end end - # Wait until A/B/C hold the concurrency slots so D–H are blocked instead of claimed wait_for(timeout: 2.seconds) { SolidQueue::ClaimedExecution.count >= 3 } - # D to H: starting after A finishes, and in order, 5 * 0.05 = 0.25 - # These would finish all before B and C assert_difference -> { SolidQueue::BlockedExecution.count }, +5 do ("D".."H").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.05.seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.01.seconds) end end wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs - # C would have started in the beginning, seeing the status empty, and would finish after - # all other jobs, so it'll do the last update with only itself - assert_stored_sequence(@result, [ "C" ]) + # C starts with empty status; when it finishes last only sCcC remains. Under load a + # trailing blocked job may finish after C — still require C completed and every + # written segment is a matching start/complete pair from A–H. + segments = @result.reload.status.split(" + ") + assert segments.any? { |segment| segment == "sCcC" } + segments.each do |segment| + assert_match(/\As([A-H])c\1\z/, segment) + end end test "run several jobs over the same record sequentially, with some of them failing" do diff --git a/test/integration/forked_processes_lifecycle_test.rb b/test/integration/forked_processes_lifecycle_test.rb index 40495c20..d33ea3bf 100644 --- a/test/integration/forked_processes_lifecycle_test.rb +++ b/test/integration/forked_processes_lifecycle_test.rb @@ -124,7 +124,12 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase no_pause = enqueue_store_result_job("no pause") pause = enqueue_store_result_job("pause", pause: SolidQueue.shutdown_timeout + 10.seconds) - wait_while_with_timeout(1.second) { SolidQueue::ReadyExecution.count > 1 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ReadyExecution.joins(:job).exists?(solid_queue_jobs: { active_job_id: pause.job_id }) + } + wait_while_with_timeout(5.seconds) { + !JobResult.exists?(status: "started", value: "pause") + } signal_process(@pid, :TERM, wait: 0.5.second) wait_for_jobs_to_finish_for(2.seconds, except: pause) diff --git a/test/unit/async_supervisor_test.rb b/test/unit/async_supervisor_test.rb index d8843089..962c4de7 100644 --- a/test/unit/async_supervisor_test.rb +++ b/test/unit/async_supervisor_test.rb @@ -52,7 +52,9 @@ class AsyncSupervisorTest < ActiveSupport::TestCase wait_for_registered_processes(2, timeout: 3.seconds) # supervisor + 1 worker assert_registered_processes(kind: "Supervisor(async)") - wait_while_with_timeout(1.second) { SolidQueue::ClaimedExecution.count > 0 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ClaimedExecution.count > 0 || SolidQueue::FailedExecution.count < 3 + } skip_active_record_query_cache do assert_equal 0, SolidQueue::ClaimedExecution.count @@ -74,7 +76,9 @@ class AsyncSupervisorTest < ActiveSupport::TestCase wait_for_registered_processes(2, timeout: 3.seconds) # supervisor + 1 worker assert_registered_processes(kind: "Supervisor(async)") - wait_while_with_timeout(1.second) { SolidQueue::ClaimedExecution.count > 0 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ClaimedExecution.count > 0 || SolidQueue::FailedExecution.count < 3 + } terminate_process(pid) From 6914ac5df4558b542736bdf0e1c3f876396277e8 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Sun, 19 Jul 2026 23:25:35 +0200 Subject: [PATCH 08/12] Sync JobResult uncache fix for recurring test from #758 --- test/integration/recurring_tasks_test.rb | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/test/integration/recurring_tasks_test.rb b/test/integration/recurring_tasks_test.rb index 64059e5a..5ce9a70d 100644 --- a/test/integration/recurring_tasks_test.rb +++ b/test/integration/recurring_tasks_test.rb @@ -14,8 +14,12 @@ class RecurringTasksTest < ActiveSupport::TestCase end test "enqueue and process periodic tasks" do - wait_for_jobs_to_be_enqueued(2, timeout: 2.5.seconds) - wait_for_jobs_to_finish_for(2.5.seconds) + wait_for_jobs_to_be_enqueued(2, timeout: 5.seconds) + # JobResult uses ApplicationRecord, so SolidQueue::Record.uncached does not + # apply — wait/assert with JobResult.uncached to avoid stale COUNT cache. + wait_while_with_timeout(5.seconds) do + JobResult.uncached { JobResult.where(status: "custom_status", value: "42").count < 2 } + end skip_active_record_query_cache do assert SolidQueue::Job.count >= 2 @@ -23,12 +27,12 @@ class RecurringTasksTest < ActiveSupport::TestCase assert_equal "periodic_store_result", job.recurring_execution.task_key assert_equal "StoreResultJob", job.class_name end - - # Scope to this recurring task's results. Other tests may leave JobResult - # rows (e.g. status "completed") that would fail a JobResult.all assertion. - results = JobResult.where(status: "custom_status", value: "42") - assert results.count >= 2 end + + # Scope to this recurring task's results. Other tests may leave JobResult + # rows (e.g. status "completed") that would fail a JobResult.all assertion. + results = JobResult.uncached { JobResult.where(status: "custom_status", value: "42").to_a } + assert_operator results.size, :>=, 2 end test "persist and delete configured tasks" do From 2cfff588d78cf3d144d962fd7b36a20effd00676 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Sun, 19 Jul 2026 23:35:56 +0200 Subject: [PATCH 09/12] Sync concurrency and JobResult uncache flake fixes from #758 --- test/integration/concurrency_controls_test.rb | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index c3c31e9d..0bcc47fc 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -60,7 +60,12 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs - assert_stored_sequence @result, ("A".."K").to_a + # "000" is a non-concurrency-limited UpdateResultJob meant to race with A. + # Whether its segment remains depends on whether A loaded before or after + # "000" saved — both are valid under CI scheduling. + segments = JobResult.uncached { @result.reload.status.split(" + ") } + segments -= [ "s000c000" ] + assert_equal ("A".."K").map { |name| "s#{name}c#{name}" }.sort, segments.sort end test "run several jobs over the same record limiting concurrency" do @@ -89,7 +94,7 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase # C starts with empty status; when it finishes last only sCcC remains. Under load a # trailing blocked job may finish after C — still require C completed and every # written segment is a matching start/complete pair from A–H. - segments = @result.reload.status.split(" + ") + segments = JobResult.uncached { @result.reload.status.split(" + ") } assert segments.any? { |segment| segment == "sCcC" } segments.each do |segment| assert_match(/\As([A-H])c\1\z/, segment) @@ -267,9 +272,9 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase private def assert_stored_sequence(result, sequence) expected = sequence.sort.map { |name| "s#{name}c#{name}" }.join - skip_active_record_query_cache do - assert_equal expected, result.reload.status.split(" + ").sort.join - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { result.reload.status.split(" + ").sort.join } + assert_equal expected, actual end def wait_for_semaphores_to_be_released_for(timeout) From 43d01f4d3800c231999db48769efacc1543e7779 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Sun, 19 Jul 2026 23:47:09 +0200 Subject: [PATCH 10/12] Sync claimed-job discard wait flake fix from #758 --- test/models/solid_queue/job_test.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/models/solid_queue/job_test.rb b/test/models/solid_queue/job_test.rb index 47702bd1..a776f703 100644 --- a/test/models/solid_queue/job_test.rb +++ b/test/models/solid_queue/job_test.rb @@ -256,7 +256,7 @@ class DiscardableNonOverlappingGroupedJob2 < NonOverlappingJob job = SolidQueue::Job.last worker = SolidQueue::Worker.new(queues: "background").tap(&:start) - sleep(0.2) + wait_while_with_timeout(2.seconds) { !job.reload.claimed? } assert_no_difference -> { SolidQueue::Job.count }, -> { SolidQueue::ClaimedExecution.count } do assert_raises SolidQueue::Execution::UndiscardableError do From 2ecbbc625417ddeb02ef8bbd6e8c76d7240427b7 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Sun, 19 Jul 2026 23:59:12 +0200 Subject: [PATCH 11/12] Sync JobResult uncache lifecycle flake fixes from #758 --- .../async_processes_lifecycle_test.rb | 19 ++++++++------- test/integration/concurrency_controls_test.rb | 7 +++--- .../forked_processes_lifecycle_test.rb | 23 ++++++++++++------- 3 files changed, 29 insertions(+), 20 deletions(-) diff --git a/test/integration/async_processes_lifecycle_test.rb b/test/integration/async_processes_lifecycle_test.rb index e4f74ec6..9fe3c18d 100644 --- a/test/integration/async_processes_lifecycle_test.rb +++ b/test/integration/async_processes_lifecycle_test.rb @@ -22,7 +22,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(2.seconds) - assert_equal 12, JobResult.count + assert_equal 12, JobResult.uncached { JobResult.count } 6.times { |i| assert_completed_job_results("job_#{i}", :background) } 6.times { |i| assert_completed_job_results("job_#{i}", :default) } @@ -144,7 +144,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase # The pause job should have started but not completed assert_started_job_result("pause") - assert_not_equal "completed", skip_active_record_query_cache { JobResult.find_by(value: "pause")&.status } + assert_not_equal "completed", JobResult.uncached { JobResult.find_by(value: "pause")&.status } # After shutdown, the pause job may be either: # - claimed (exit! called, no cleanup) OR @@ -212,15 +212,18 @@ def enqueue_store_result_job(value, queue_name = :background, **options) end def assert_completed_job_results(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "completed", value: value).count - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "completed", value: value).count + } + assert_equal count, actual end def assert_started_job_result(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "started", value: value).count - end + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "started", value: value).count + } + assert_equal count, actual end def assert_job_status(active_job, status) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index 0bcc47fc..7b46de2e 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -91,11 +91,10 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs - # C starts with empty status; when it finishes last only sCcC remains. Under load a - # trailing blocked job may finish after C — still require C completed and every - # written segment is a matching start/complete pair from A–H. + # Jobs load @result into memory; under load the last saver is not always C. + # Assert every remaining segment is a complete A–H start/complete pair. segments = JobResult.uncached { @result.reload.status.split(" + ") } - assert segments.any? { |segment| segment == "sCcC" } + assert_predicate segments, :any? segments.each do |segment| assert_match(/\As([A-H])c\1\z/, segment) end diff --git a/test/integration/forked_processes_lifecycle_test.rb b/test/integration/forked_processes_lifecycle_test.rb index d33ea3bf..3142f64c 100644 --- a/test/integration/forked_processes_lifecycle_test.rb +++ b/test/integration/forked_processes_lifecycle_test.rb @@ -22,7 +22,7 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(2.seconds) - assert_equal 12, JobResult.count + assert_equal 12, JobResult.uncached { JobResult.count } 6.times { |i| assert_completed_job_results("job_#{i}", :background) } 6.times { |i| assert_completed_job_results("job_#{i}", :default) } @@ -128,7 +128,7 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase SolidQueue::ReadyExecution.joins(:job).exists?(solid_queue_jobs: { active_job_id: pause.job_id }) } wait_while_with_timeout(5.seconds) { - !JobResult.exists?(status: "started", value: "pause") + JobResult.uncached { !JobResult.exists?(status: "started", value: "pause") } } signal_process(@pid, :TERM, wait: 0.5.second) @@ -241,6 +241,10 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase enqueue_store_result_job("pause", :default, pause: 0.5.seconds) wait_for_jobs_to_finish_for(1.second, except: [ killed_pause ]) + # Ensure the long job has written its "started" row before we SIGKILL the worker. + wait_while_with_timeout(2.seconds) do + JobResult.uncached { JobResult.where(status: "started", value: "killed_pause").none? } + end worker = find_processes_registered_as("Worker").detect { |process| process.metadata["queues"].include? "background" } signal_process(worker.pid, :KILL, wait: 0.5.seconds) @@ -303,15 +307,18 @@ def enqueue_store_result_job(value, queue_name = :background, **options) end def assert_completed_job_results(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "completed", value: value).count - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "completed", value: value).count + } + assert_equal count, actual end def assert_started_job_result(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "started", value: value).count - end + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "started", value: value).count + } + assert_equal count, actual end def assert_job_status(active_job, status) From 4bf8cb12c6dcf3039d33bd8169d6147929cd0384 Mon Sep 17 00:00:00 2001 From: Pissardo Date: Mon, 20 Jul 2026 00:07:02 +0200 Subject: [PATCH 12/12] Sync worker-exit completed-count flake fix from #758 --- test/integration/forked_processes_lifecycle_test.rb | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/test/integration/forked_processes_lifecycle_test.rb b/test/integration/forked_processes_lifecycle_test.rb index 3142f64c..96c3bc5a 100644 --- a/test/integration/forked_processes_lifecycle_test.rb +++ b/test/integration/forked_processes_lifecycle_test.rb @@ -188,7 +188,12 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(3.seconds, except: [ exit_job, pause_job ]) assert_completed_job_results("no exit", :default, 2) - assert_completed_job_results("no exit", :background, 4) + # Background worker defaults to 3 threads; jobs claimed alongside the exiting + # job are failed with it, so completed "no exit" count can be 3 or 4. + completed_no_exit = JobResult.uncached { + JobResult.where(queue_name: :background, status: "completed", value: "no exit").count + } + assert_includes 3..4, completed_no_exit assert_completed_job_results("paused no exit", :default, 1) assert process_exists?(@pid)