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
10 changes: 8 additions & 2 deletions app/models/solid_queue/blocked_execution.rb
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,12 @@ def releasable(concurrency_keys)
def release
SolidQueue.instrument(:release_blocked, job_id: job.id, concurrency_key: concurrency_key, released: false) do |payload|
transaction do
if acquire_concurrency_lock
if job.job_class.nil?
# The job's class no longer resolves (renamed/removed between deploys).
# Destroy the orphan row so the dispatcher stops retrying it forever.
destroy!
payload[:orphaned] = true
elsif acquire_concurrency_lock
promote_to_ready
destroy!

Expand All @@ -58,7 +63,8 @@ def release

private
def set_expires_at
self.expires_at = job.concurrency_duration.from_now
duration = job.job_class ? job.concurrency_duration : SolidQueue.default_concurrency_control_period
self.expires_at = duration.from_now
end

def acquire_concurrency_lock
Expand Down
8 changes: 4 additions & 4 deletions app/models/solid_queue/job/concurrency_controls.rb
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@ def blocked?
blocked_execution.present?
end

def job_class
@job_class ||= class_name.safe_constantize
end

private
def concurrency_on_conflict
job_class.concurrency_on_conflict.to_s.inquiry
Expand Down Expand Up @@ -66,10 +70,6 @@ def release_next_blocked_job
BlockedExecution.release_one(concurrency_key)
end

def job_class
@job_class ||= class_name.safe_constantize
end

def execution
super || blocked_execution
end
Expand Down
44 changes: 44 additions & 0 deletions test/models/solid_queue/blocked_execution_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
require "test_helper"

class SolidQueue::BlockedExecutionTest < ActiveSupport::TestCase
self.use_transactional_tests = false

class NonOverlappingJob < ApplicationJob
limits_concurrency key: ->(job_result, **) { job_result }

def perform(job_result)
end
end

setup do
@result = JobResult.create!(queue_name: "default")
end

teardown do
SolidQueue::Job.destroy_all
SolidQueue::Semaphore.delete_all
JobResult.delete_all
end

test "release destroys the blocked row when the job class no longer resolves" do
# Enqueue and consume the semaphore so the next job blocks.
NonOverlappingJob.perform_later(@result)
blocking_job = SolidQueue::Job.last
NonOverlappingJob.perform_later(@result)
blocked_job = SolidQueue::Job.last
blocked = blocked_job.blocked_execution
assert blocked, "expected the second job to be blocked"

# Simulate the class being renamed/removed between deploys
blocked_job.update_columns(class_name: "GoneJob")

assert_difference -> { SolidQueue::BlockedExecution.count }, -1 do
assert_nothing_raised do
blocked.reload.release
end
end

# No ready execution was promoted — the orphan row was just cleaned up.
assert_nil SolidQueue::ReadyExecution.find_by(job_id: blocked_job.id)
end
end
Loading