diff --git a/lib/delayed_job/threaded_worker.rb b/lib/delayed_job/threaded_worker.rb index 3f514daaf7b..28ebc36a398 100644 --- a/lib/delayed_job/threaded_worker.rb +++ b/lib/delayed_job/threaded_worker.rb @@ -1,5 +1,18 @@ module Delayed class ThreadedWorker < Delayed::Worker + # Plugin callbacks register against the base Delayed::Worker.lifecycle, so + # delegate to the base class to avoid rebuilding a separate lifecycle here + # (which would leak a callback per ThreadedWorker.new). Use superclass, not + # the Delayed::Worker constant, which tests may reassign to this subclass and + # cause infinite recursion. + def self.setup_lifecycle + superclass.setup_lifecycle + end + + def self.lifecycle + superclass.lifecycle + end + def initialize(options={}) super @num_threads = options[:num_threads] diff --git a/spec/unit/lib/delayed_job/threaded_worker_spec.rb b/spec/unit/lib/delayed_job/threaded_worker_spec.rb index f425a1b38a3..285bc5ecfbf 100644 --- a/spec/unit/lib/delayed_job/threaded_worker_spec.rb +++ b/spec/unit/lib/delayed_job/threaded_worker_spec.rb @@ -22,6 +22,16 @@ worker = Delayed::ThreadedWorker.new({ num_threads: 2 }) expect(worker.instance_variable_get(:@grace_period_seconds)).to eq(30) end + + it 'does not accumulate enqueue lifecycle callbacks across instantiations' do + before_enqueue_callbacks = lambda do + Delayed::Worker.lifecycle.instance_variable_get(:@callbacks)[:enqueue].instance_variable_get(:@before).size + end + + baseline = before_enqueue_callbacks.call + 5.times { Delayed::ThreadedWorker.new(options) } + expect(before_enqueue_callbacks.call).to eq(baseline) + end end describe '#start' do