From dbab4f9969d9e18f4440b8805d3824cd64a6d3b7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Galisteo?= Date: Sun, 4 Oct 2026 19:43:59 +0200 Subject: [PATCH] Acquire concurrency locks by key when dispatching several jobs Dispatching a set of jobs with concurrency controls acquires each job's semaphore one by one, inside the transaction that dispatches the whole set, and in whatever order the jobs come in. Two dispatchers whose batches share concurrency keys can then acquire them in opposite order and deadlock, and so can bulk enqueues. The batch that's aborted stays scheduled, so no job is lost and the limits hold, but the dispatcher that hits the error is delayed. Go through the jobs by concurrency key instead, so everyone acquires the locks they share in the same order. Within a key, jobs keep their id order, so the same job as before gets the lock first. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01Acq8BRNz9v7w9mKYLYRdCi --- app/models/solid_queue/job/executable.rb | 6 +- .../concurrent_dispatching_test.rb | 86 +++++++++++++++++++ test/models/solid_queue/job_test.rb | 15 ++++ 3 files changed, 106 insertions(+), 1 deletion(-) create mode 100644 test/integration/concurrent_dispatching_test.rb diff --git a/app/models/solid_queue/job/executable.rb b/app/models/solid_queue/job/executable.rb index 75a5d2111..4d0706c3b 100644 --- a/app/models/solid_queue/job/executable.rb +++ b/app/models/solid_queue/job/executable.rb @@ -39,8 +39,12 @@ def dispatch_all_at_once(jobs) ReadyExecution.create_all_from_jobs jobs end + # Each job acquires its concurrency lock here, and holds it until the + # transaction dispatching the whole set ends. Going through them by + # concurrency key makes everyone dispatching at the same time acquire + # the locks they share in the same order, so they can't deadlock. def dispatch_all_one_by_one(jobs) - jobs.each(&:dispatch) + jobs.sort_by { |job| [ job.concurrency_key, job.id ] }.each(&:dispatch) end def successfully_dispatched(jobs) diff --git a/test/integration/concurrent_dispatching_test.rb b/test/integration/concurrent_dispatching_test.rb new file mode 100644 index 000000000..cc97c96ea --- /dev/null +++ b/test/integration/concurrent_dispatching_test.rb @@ -0,0 +1,86 @@ +# frozen_string_literal: true + +require "test_helper" + +class ConcurrentDispatchingTest < ActiveSupport::TestCase + self.use_transactional_tests = false + + class NonOverlappingJob < ApplicationJob + limits_concurrency key: ->(key) { key } + + def perform(key) + end + end + + setup do + # SQLite allows a single writer at a time, so two dispatchers never hold + # locks at once. + skip "SQLite serializes all writes" if SolidQueue::Record.connection_pool.db_config.adapter == "sqlite3" + + @original_wait = SolidQueue::Semaphore.method(:wait) + end + + teardown do + SolidQueue::Semaphore.define_singleton_method(:wait, @original_wait) if @original_wait + end + + test "dispatchers with jobs for the same concurrency keys in different order don't deadlock" do + # By job id: one, two, two, one. In batches of two, a dispatcher gets the + # first two jobs and another one, skipping those, the other two. + jobs = %w[ one two two one ].map do |key| + NonOverlappingJob.set(wait: 1.minute).perform_later(key) + SolidQueue::Job.last + end + travel_to 2.minutes.from_now + + errors = dispatch_in_two_batches_at_once + + assert_empty errors + assert_equal 0, SolidQueue::ScheduledExecution.count + + # The first job for each key gets the lock + assert_equal [ :ready, :ready, :blocked, :blocked ], jobs.map { |job| job.reload.status } + end + + private + # Two dispatchers, each with a batch of two jobs. Each one, after acquiring + # its first lock, gives the other the chance to acquire its first one too, + # so that both hold a lock when they go for the second. + def dispatch_in_two_batches_at_once + acquired_first = [ Concurrent::Event.new, Concurrent::Event.new ] + pause_after_first_lock(acquired_first) + + dispatchers = 2.times.map do |index| + Thread.new do + Thread.current[:dispatcher] = index + # The second one starts once the first has its batch and its first lock + acquired_first[0].wait(5.seconds) if index == 1 + + SolidQueue::Record.connection_pool.with_connection do + SolidQueue::ScheduledExecution.dispatch_next_batch(2) + end + + nil + rescue => error + error + end + end + + dispatchers.map(&:value).compact + end + + def pause_after_first_lock(acquired_first) + original_wait = @original_wait + + SolidQueue::Semaphore.define_singleton_method(:wait) do |job| + original_wait.call(job).tap do + index = Thread.current[:dispatcher] + + if index && !acquired_first[index].set? + acquired_first[index].set + acquired_first[1 - index].wait(2.seconds) + end + end + end + end +end diff --git a/test/models/solid_queue/job_test.rb b/test/models/solid_queue/job_test.rb index 7e9c1568f..f7db31261 100644 --- a/test/models/solid_queue/job_test.rb +++ b/test/models/solid_queue/job_test.rb @@ -227,6 +227,21 @@ class DiscardableNonOverlappingGroupedJob2 < NonOverlappingJob assert_not not_enqueued.successfully_enqueued? end + test "enqueue jobs in bulk acquiring their concurrency locks by key" do + results = 3.times.map { JobResult.create!(queue_name: "default") } + + acquired = [] + SolidQueue::Semaphore.stubs(:wait).with { |job| acquired << [ job.concurrency_key, job.id ] }.returns(true) + + # Two jobs for each key, with the keys in descending order + active_jobs = (results.reverse + results.reverse).map { |result| NonOverlappingJob.new(result) } + ActiveJob.perform_all_later(active_jobs) + + # By key and, within a key, in the order they were enqueued + assert_equal 6, acquired.size + assert_equal acquired.sort, acquired + end + test "discard ready job" do AddToBufferJob.perform_later(1) job = SolidQueue::Job.last