Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion app/models/solid_queue/job/executable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
86 changes: 86 additions & 0 deletions test/integration/concurrent_dispatching_test.rb
Original file line number Diff line number Diff line change
@@ -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
15 changes: 15 additions & 0 deletions test/models/solid_queue/job_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading