From 1f42551a538f72725c65c0360c6b1116c7736c25 Mon Sep 17 00:00:00 2001 From: Philipp Thun Date: Tue, 11 Aug 2026 18:02:13 +0200 Subject: [PATCH] Prevent ThreadedWorker from leaking enqueue lifecycle callbacks Plugin callbacks register against the base Delayed::Worker.lifecycle, but setup_lifecycle rebuilt the lifecycle on the class it was called on. On the ThreadedWorker subclass that reset a separate lifecycle, so the base one was never reset and each ThreadedWorker.new leaked another before(:enqueue) callback, creating many PollableJobModel rows per enqueue. Production is unaffected (one worker per process). Delegate setup_lifecycle and lifecycle to the base class via superclass, not the Delayed::Worker constant, which tests may reassign to this subclass and cause infinite recursion. --- lib/delayed_job/threaded_worker.rb | 13 +++++++++++++ spec/unit/lib/delayed_job/threaded_worker_spec.rb | 10 ++++++++++ 2 files changed, 23 insertions(+) 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