Skip to content

Commit 85eefc6

Browse files
authored
Merge pull request #5 from starburstlabs/feature/concurrency_max_blocked
AIA-733: Cap the number of blocked executions per concurrency key
2 parents c2cf4a5 + aca950b commit 85eefc6

5 files changed

Lines changed: 137 additions & 3 deletions

File tree

app/models/solid_queue/job/concurrency_controls.rb

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ module ConcurrencyControls
88
included do
99
has_one :blocked_execution
1010

11-
delegate :concurrency_limit, :concurrency_duration, to: :job_class
11+
delegate :concurrency_limit, :concurrency_duration, :concurrency_max_blocked, to: :job_class
1212

1313
before_destroy :unblock_next_blocked_job, if: -> { concurrency_limited? && ready? }
1414
end
@@ -59,7 +59,18 @@ def handle_concurrency_conflict
5959
end
6060

6161
def block
62-
BlockedExecution.create_or_find_by!(job_id: id)
62+
return BlockedExecution.create_or_find_by!(job_id: id) unless concurrency_max_blocked
63+
64+
transaction do
65+
# Lock the semaphore row to serialize concurrent enqueues racing to fill the cap.
66+
Semaphore.lock.find_by(key: concurrency_key)
67+
68+
if BlockedExecution.where(concurrency_key: concurrency_key).count < concurrency_max_blocked
69+
BlockedExecution.create_or_find_by!(job_id: id)
70+
else
71+
destroy
72+
end
73+
end
6374
end
6475

6576
def release_next_blocked_job

lib/active_job/concurrency_controls.rb

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,15 +14,17 @@ module ConcurrencyControls
1414
class_attribute :concurrency_limit
1515
class_attribute :concurrency_duration, default: SolidQueue.default_concurrency_control_period
1616
class_attribute :concurrency_on_conflict, default: :block
17+
class_attribute :concurrency_max_blocked
1718
end
1819

1920
class_methods do
20-
def limits_concurrency(key:, to: 1, group: DEFAULT_CONCURRENCY_GROUP, duration: SolidQueue.default_concurrency_control_period, on_conflict: :block)
21+
def limits_concurrency(key:, to: 1, group: DEFAULT_CONCURRENCY_GROUP, duration: SolidQueue.default_concurrency_control_period, on_conflict: :block, max_blocked: nil)
2122
self.concurrency_key = key
2223
self.concurrency_limit = to
2324
self.concurrency_group = group
2425
self.concurrency_duration = duration
2526
self.concurrency_on_conflict = on_conflict.presence_in(CONCURRENCY_ON_CONFLICT_BEHAVIOUR) || :block
27+
self.concurrency_max_blocked = max_blocked
2628
end
2729
end
2830

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
class BoundedBlockedUpdateResultJob < UpdateResultJob
2+
limits_concurrency key: ->(job_result, **) { job_result }, on_conflict: :block, max_blocked: 1
3+
end

test/integration/concurrency_controls_test.rb

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -235,6 +235,28 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase
235235
assert_stored_sequence(another_result, [ "2" ])
236236
end
237237

238+
test "block up to max_blocked and discard the rest" do
239+
# to: 1, max_blocked: 1 — so 1 runs, 1 blocks, the rest get discarded
240+
job1 = BoundedBlockedUpdateResultJob.perform_later(@result, name: "1", pause: 0.5)
241+
sleep(0.1) # ensure "1" is claimed first
242+
243+
job2 = BoundedBlockedUpdateResultJob.perform_later(@result, name: "2") # blocks under cap
244+
job3 = BoundedBlockedUpdateResultJob.perform_later(@result, name: "3") # cap full, discarded
245+
job4 = BoundedBlockedUpdateResultJob.perform_later(@result, name: "4") # cap full, discarded
246+
247+
wait_for_jobs_to_finish_for(5.seconds)
248+
assert_no_unfinished_jobs
249+
250+
# 1 ran, 2 promoted after 1 finished and ran; 3 and 4 never ran
251+
assert_stored_sequence(@result, [ "1", "2" ])
252+
253+
jobs = SolidQueue::Job.where(active_job_id: [ job1, job2, job3, job4 ].map(&:job_id))
254+
assert_equal 2, jobs.count
255+
assert_equal [ job1.provider_job_id, job2.provider_job_id ].sort, jobs.pluck(:id).sort
256+
assert_nil job3.provider_job_id
257+
assert_nil job4.provider_job_id
258+
end
259+
238260
test "discard on conflict and release semaphore" do
239261
DiscardableUpdateResultJob.perform_later(@result, name: "1", pause: 0.1)
240262
# will be discarded

test/models/solid_queue/job_test.rb

Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,18 @@ class DiscardableNonOverlappingJob < NonOverlappingJob
1414
limits_concurrency key: ->(job_result, **) { job_result }, on_conflict: :discard
1515
end
1616

17+
class BoundedBlockedJob < NonOverlappingJob
18+
limits_concurrency key: ->(job_result, **) { job_result }, on_conflict: :block, max_blocked: 1
19+
end
20+
21+
class BoundedBlockedThrottledJob < NonOverlappingJob
22+
limits_concurrency to: 2, key: ->(job_result, **) { job_result }, on_conflict: :block, max_blocked: 2
23+
end
24+
25+
class DiscardableWithMaxBlockedJob < NonOverlappingJob
26+
limits_concurrency key: ->(job_result, **) { job_result }, on_conflict: :discard, max_blocked: 5
27+
end
28+
1729
class DiscardableThrottledJob < NonOverlappingJob
1830
limits_concurrency to: 2, key: ->(job_result, **) { job_result }, on_conflict: :discard
1931
end
@@ -137,6 +149,90 @@ class DiscardableNonOverlappingGroupedJob2 < NonOverlappingJob
137149
end
138150
end
139151

152+
test "block jobs up to max_blocked, then discard further conflicts" do
153+
assert_ready do
154+
BoundedBlockedJob.perform_later(@result, name: "ready")
155+
end
156+
157+
assert_blocked do
158+
BoundedBlockedJob.perform_later(@result, name: "blocked")
159+
end
160+
161+
# Cap is hit: the next two jobs should be silently discarded
162+
assert_job_counts do
163+
BoundedBlockedJob.perform_later(@result, name: "discarded-1")
164+
BoundedBlockedJob.perform_later(@result, name: "discarded-2")
165+
end
166+
end
167+
168+
test "max_blocked cap reopens after the blocked job is released" do
169+
ready_active_job = BoundedBlockedJob.perform_later(@result, name: "ready")
170+
BoundedBlockedJob.perform_later(@result, name: "blocked-1")
171+
ready_job, blocked_job = SolidQueue::Job.last(2)
172+
173+
# Cap full: this enqueue is discarded
174+
assert_job_counts do
175+
BoundedBlockedJob.perform_later(@result, name: "discarded")
176+
end
177+
178+
# Discarding the ready job signals the semaphore and promotes the blocked one.
179+
# Net: ready_job destroyed (-1 Job), BlockedExecution -1, ReadyExecution unchanged
180+
# (one destroyed for ready_job, one created from promotion).
181+
assert_job_counts(blocked: -1) do
182+
ready_job.discard
183+
end
184+
assert blocked_job.reload.ready?
185+
186+
# Cap is now empty; a fresh enqueue should block again, not be discarded
187+
assert_blocked do
188+
BoundedBlockedJob.perform_later(@result, name: "newly-blocked")
189+
end
190+
end
191+
192+
test "max_blocked enforces the cap independently per concurrency key" do
193+
another_result = JobResult.create!(queue_name: "default")
194+
195+
BoundedBlockedJob.perform_later(@result, name: "ready-a")
196+
BoundedBlockedJob.perform_later(another_result, name: "ready-b")
197+
198+
assert_job_counts(blocked: 2) do
199+
BoundedBlockedJob.perform_later(@result, name: "blocked-a")
200+
BoundedBlockedJob.perform_later(another_result, name: "blocked-b")
201+
end
202+
203+
# Both keys are at cap; further enqueues for either key get discarded
204+
assert_job_counts do
205+
BoundedBlockedJob.perform_later(@result, name: "discarded-a")
206+
BoundedBlockedJob.perform_later(another_result, name: "discarded-b")
207+
end
208+
end
209+
210+
test "max_blocked respects concurrency_limit > 1" do
211+
# to: 2, max_blocked: 2 => up to 2 ready + up to 2 blocked
212+
assert_job_counts(ready: 2) do
213+
BoundedBlockedThrottledJob.perform_later(@result, name: "ready-1")
214+
BoundedBlockedThrottledJob.perform_later(@result, name: "ready-2")
215+
end
216+
217+
assert_job_counts(blocked: 2) do
218+
BoundedBlockedThrottledJob.perform_later(@result, name: "blocked-1")
219+
BoundedBlockedThrottledJob.perform_later(@result, name: "blocked-2")
220+
end
221+
222+
assert_job_counts do
223+
BoundedBlockedThrottledJob.perform_later(@result, name: "discarded")
224+
end
225+
end
226+
227+
test "max_blocked is a no-op when on_conflict is :discard" do
228+
# max_blocked is only consulted on the :block path; with :discard, every
229+
# conflict is destroyed regardless of how many are allegedly "allowed".
230+
assert_job_counts(ready: 1) do
231+
DiscardableWithMaxBlockedJob.perform_later(@result, name: "ready")
232+
DiscardableWithMaxBlockedJob.perform_later(@result, name: "would-have-blocked-if-cap-applied")
233+
end
234+
end
235+
140236
test "block jobs in the same concurrency group when concurrency limits are reached" do
141237
assert_ready do
142238
active_job = NonOverlappingGroupedJob1.perform_later(@result, name: "A")

0 commit comments

Comments
 (0)