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