Skip to content

Commit 837a51b

Browse files
GoodForOneFareshopify-riverclaude
committed
Settle futures on any exception; validate WorkQueue/Progress input
WorkQueue workers only rescued StandardError and Interrupt. Any other exception (ScriptError family - including NotImplementedError and LoadError from a failed require - SystemExit, etc.) escaped the worker thread without completing its future. The dead worker also kept counting against max_concurrent, so queued work never started. Callers blocked in Future#value forever, and SpinGroup#wait spun indefinitely because the task's future never completed. WorkQueue changes: * Workers rescue Exception and fail the future, so waiters always observe the error and the worker stays alive for queued work. Interrupt handling is unchanged (fail the future, re-raise). * Every future is settled structurally, from an ensure: Thread#raise and Thread#kill can land outside the worker's rescues (e.g. WorkQueue#interrupt hitting a worker between Queue#pop and the begin). A future abandoned that way is failed with WorkQueue::AbandonedTaskError instead of blocking Future#value forever. * interrupt tolerates Interrupt re-raised by Thread#join: the Interrupt it raises into a worker can land after that worker has already left its rescues (it was terminating), killing it with an unhandled Interrupt that join would otherwise re-raise in the interrupting thread - replacing whatever that thread was propagating. * WorkQueue.new rejects max_concurrent < 1, which produced a queue that accepted work, never started it, and returned from #wait immediately while futures stayed incomplete. SpinGroup changes: * Task#check treats any exception other than Interrupt/SystemExit as a task failure, reported through the normal debrief. Raising out of the render loop skipped work_queue.interrupt, the final render, and the debrief, leaving sibling workers running after #wait returned. * wait interrupts the queue before letting SystemExit propagate, so remaining workers aren't orphaned. * initialize falls back to 1024 for any non-positive max_concurrent instead of raising from WorkQueue's new guard. Progress changes: * tick validates percent/set_percent explicitly (finite numbers only, named in the error message) instead of relying on Array#min raising on an incomparable value, and validates before assigning so a raising tick leaves the bar renderable from its last valid state. * The result is clamped to 0.0..1.0; the lower bound was previously unguarded, so a negative percent rendered a negative suffix. Co-authored-by: River <river@shopify.com> Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 83cbecb commit 837a51b

6 files changed

Lines changed: 180 additions & 17 deletions

File tree

‎lib/cli/ui/progress.rb‎

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -77,15 +77,20 @@ def initialize(title = nil, width: Terminal.width, reporter: nil)
7777
# * +:percent+ - Increment progress by a specific percent amount
7878
# * +:set_percent+ - Set progress to a specific percent
7979
#
80-
# *Note:* The +:percent+ and +:set_percent must be between 0.00 and 1.0
80+
# *Note:* The +:percent+ and +:set_percent+ must be finite numbers. The
81+
# resulting progress is clamped to 0.0..1.0.
8182
#
8283
#: (?percent: Numeric?, ?set_percent: Numeric?) -> void
8384
def tick(percent: nil, set_percent: nil)
8485
raise ArgumentError, 'percent and set_percent cannot both be specified' if percent && set_percent
8586

86-
@percent_done += percent || 0.01
87-
@percent_done = set_percent if set_percent
88-
@percent_done = [@percent_done, 1.0].min # Make sure we can't go above 1.0
87+
validate_finite!(percent, :percent)
88+
validate_finite!(set_percent, :set_percent)
89+
90+
# Validate and clamp before assigning: a raising tick leaves the bar
91+
# renderable from its last valid state.
92+
new_percent = set_percent || @percent_done + (percent || 0.01)
93+
@percent_done = new_percent.clamp(0.0, 1.0)
8994

9095
# Update terminal progress reporter with current percentage
9196
@reporter&.set_progress((@percent_done * 100).floor)
@@ -125,6 +130,16 @@ def to_s
125130

126131
[title, bar].compact.join("\n")
127132
end
133+
134+
private
135+
136+
#: (untyped value, Symbol name) -> void
137+
def validate_finite!(value, name)
138+
return if value.nil?
139+
return if value.is_a?(Numeric) && value.real? && value.to_f.finite?
140+
141+
raise ArgumentError, "#{name} must be a finite number (got #{value.inspect})"
142+
end
128143
end
129144
end
130145
end

‎lib/cli/ui/spinner/spin_group.rb‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ def initialize(auto_debrief: true, interrupt_debrief: false, max_concurrent: 0,
6969
@start = Time.new
7070
@stopped = false
7171
@internal_work_queue = work_queue.nil?
72-
@work_queue = work_queue || WorkQueue.new(max_concurrent.zero? ? 1024 : max_concurrent) #: WorkQueue
72+
@work_queue = work_queue || WorkQueue.new(max_concurrent.positive? ? max_concurrent : 1024) #: WorkQueue
7373
if block_given?
7474
yield self
7575
wait(to: to)
@@ -143,7 +143,12 @@ def check
143143
result = @future.value
144144
@success = true
145145
@success = false if result == TASK_FAILED
146-
rescue => exc
146+
rescue Interrupt, SystemExit
147+
raise
148+
rescue Exception => exc # rubocop:disable Lint/RescueException
149+
# Any other exception (including ScriptError and friends) is a task
150+
# failure, reported through the normal debrief rather than raised
151+
# out of the middle of SpinGroup#wait's render loop.
147152
@exception = exc
148153
@success = false
149154
end
@@ -439,6 +444,10 @@ def wait(to: $stdout)
439444
@work_queue.interrupt
440445
debrief(to: to) if @interrupt_debrief
441446
stopped? ? false : raise
447+
rescue SystemExit
448+
# Leaving the render loop mid-flight would orphan the remaining workers.
449+
@work_queue.interrupt
450+
raise
442451
end
443452

444453
#: (String message) -> void

‎lib/cli/ui/work_queue.rb‎

Lines changed: 29 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,10 @@
44
module CLI
55
module UI
66
class WorkQueue
7+
# Raised into a future whose worker exited without settling it: Thread#raise
8+
# and Thread#kill can land outside the worker's rescues.
9+
class AbandonedTaskError < StandardError; end
10+
711
class Future
812
#: -> void
913
def initialize
@@ -66,6 +70,8 @@ def start
6670

6771
#: (Integer max_concurrent) -> void
6872
def initialize(max_concurrent)
73+
raise ArgumentError, "max_concurrent must be at least 1 (got #{max_concurrent})" if max_concurrent < 1
74+
6975
@max_concurrent = max_concurrent
7076
@queue = Queue.new #: Queue
7177
@mutex = Mutex.new #: Mutex
@@ -105,7 +111,15 @@ def interrupt
105111
end
106112
# Interrupt all worker threads
107113
@workers.each { |worker| worker.raise(Interrupt) if worker.alive? }
108-
@workers.each(&:join)
114+
@workers.each do |worker|
115+
worker.join
116+
rescue Interrupt
117+
# The Interrupt raised above can land after a worker has left its
118+
# rescues (it was already terminating); the worker then dies with
119+
# it and join re-raises it here. That must not replace whatever
120+
# this thread is propagating (e.g. the SystemExit that triggered
121+
# this interrupt).
122+
end
109123
@workers.clear
110124
end
111125
end
@@ -116,21 +130,25 @@ def interrupt
116130
def start_worker
117131
@workers << Thread.new do
118132
loop do
119-
work = @queue.pop
120-
break if work.nil?
121-
122-
future, block = work
133+
future = nil #: Future?
123134

124135
begin
136+
work = @queue.pop
137+
break if work.nil?
138+
139+
future, block = work
125140
future.start
126-
result = block.call
127-
future.complete(result)
141+
future.complete(block.call)
128142
rescue Interrupt => e
129-
future.fail(e)
143+
future&.fail(e)
130144
raise # Always re-raise interrupts to terminate the worker
131-
rescue StandardError => e
132-
future.fail(e)
133-
# Don't re-raise standard errors - allow worker to continue
145+
rescue Exception => e # rubocop:disable Lint/RescueException
146+
future&.fail(e)
147+
# Don't re-raise - the future carries the error to callers; allow worker to continue
148+
ensure
149+
# A worker must never abandon a future: an unsettled future blocks
150+
# Future#value forever. No-op once the future is settled.
151+
future&.fail(AbandonedTaskError.new('worker exited before completing this task'))
134152
end
135153
end
136154
rescue Interrupt

‎test/cli/ui/progress_test.rb‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,39 @@ def test_tick_with_set_percent_and_percent_raises
2828
end
2929
end
3030

31+
def test_tick_with_set_percent_below_zero_is_clamped_to_zero
32+
assert_bar(set_percent: -0.5, expected_filled: 0, expected_unfilled: 10, suffix: ' 0% ')
33+
end
34+
35+
def test_tick_with_percent_change_below_zero_is_clamped_to_zero
36+
assert_bar(percent: -0.5, expected_filled: 0, expected_unfilled: 10, suffix: ' 0% ')
37+
end
38+
39+
def test_tick_with_invalid_value_raises_without_corrupting_state
40+
capture_io do
41+
bar = Progress.new(width: 15)
42+
bar.tick(set_percent: 0.4)
43+
44+
['0.5', Float::NAN, Float::INFINITY, -Float::INFINITY, nil.to_a].each do |bad|
45+
assert_raises(ArgumentError) { bar.tick(set_percent: bad) }
46+
assert_raises(ArgumentError) { bar.tick(percent: bad) }
47+
end
48+
49+
# The bar still renders from its last valid state
50+
assert_match(/ 40% /, bar.to_s)
51+
end
52+
end
53+
54+
def test_tick_error_message_names_the_bad_argument
55+
bar = Progress.new(width: 15)
56+
57+
error = assert_raises(ArgumentError) { bar.tick(set_percent: '0.5') }
58+
assert_match(/set_percent must be a finite number/, error.message)
59+
60+
error = assert_raises(ArgumentError) { bar.tick(percent: Float::NAN) }
61+
assert_match(/percent must be a finite number/, error.message)
62+
end
63+
3164
def test_with_title
3265
assert_bar(title: 'Title', set_percent: 0.1, expected_filled: 1, expected_unfilled: 9, suffix: ' 10% ')
3366
end

‎test/cli/ui/spinner/spin_group_test.rb‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
# frozen_string_literal: true
22

33
require 'test_helper'
4+
require 'timeout'
45

56
module CLI
67
module UI
@@ -36,6 +37,64 @@ def test_spin_group_auto_debrief_false
3637
assert_equal('', err)
3738
end
3839

40+
def test_spin_group_non_standard_error_is_reported_as_a_task_failure
41+
capture_io do
42+
CLI::UI::StdoutRouter.ensure_activated
43+
44+
failures = []
45+
sibling_ran = false
46+
47+
sg = SpinGroup.new
48+
sg.failure_debrief { |title, exception, _out, _err| failures << [title, exception] }
49+
sg.add('raises') { raise NotImplementedError, 'not done yet' }
50+
sg.add('sibling') { sibling_ran = true }
51+
52+
# Before WorkQueue settled futures for non-StandardError exceptions
53+
# this spun forever; the timeout keeps a regression from wedging the
54+
# suite instead of failing it.
55+
refute(Timeout.timeout(10) { sg.wait })
56+
57+
assert(sibling_ran, 'sibling task should still run to completion')
58+
assert_equal(1, failures.size)
59+
title, error = failures.first
60+
assert_equal('raises', title)
61+
assert_instance_of(NotImplementedError, error)
62+
assert_equal('not done yet', error.message)
63+
end
64+
end
65+
66+
def test_spin_group_system_exit_propagates_and_interrupts_remaining_work
67+
capture_io do
68+
CLI::UI::StdoutRouter.ensure_activated
69+
70+
late = Queue.new
71+
sg = SpinGroup.new(auto_debrief: false)
72+
sg.add('exits') { exit(1) }
73+
sg.add('slow') do
74+
sleep(5)
75+
late << :ran
76+
end
77+
78+
error = Timeout.timeout(10) { assert_raises(SystemExit) { sg.wait } }
79+
80+
assert_equal(1, error.status)
81+
assert(late.empty?, 'remaining workers should be interrupted, not orphaned')
82+
end
83+
end
84+
85+
def test_spin_group_non_positive_max_concurrent_uses_the_default
86+
capture_io do
87+
CLI::UI::StdoutRouter.ensure_activated
88+
89+
[0, -1].each do |max_concurrent|
90+
sg = SpinGroup.new(max_concurrent: max_concurrent, auto_debrief: false)
91+
sg.add('s') { true }
92+
93+
assert(Timeout.timeout(10) { sg.wait }, "max_concurrent: #{max_concurrent} should run tasks")
94+
end
95+
end
96+
end
97+
3998
def test_spin_group_success_debrief
4099
capture_io do
41100
CLI::UI::StdoutRouter.ensure_activated

‎test/cli/ui/work_queue_test.rb‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,35 @@ def test_future_error
6565
assert_raises(StandardError, 'Test error') { future.value }
6666
end
6767

68+
def test_future_non_standard_error_completes_future_and_keeps_worker_alive
69+
@work_queue = WorkQueue.new(1)
70+
failed = @work_queue.enqueue { raise ScriptError, 'boom' }
71+
replacement = @work_queue.enqueue { :ran }
72+
73+
@work_queue.wait
74+
75+
assert(failed.completed?, 'failing future should be completed')
76+
error = assert_raises(ScriptError) { failed.value }
77+
assert_equal('boom', error.message)
78+
assert(replacement.completed?, 'subsequent task should still run on the same worker')
79+
assert_equal(:ran, replacement.value)
80+
end
81+
82+
def test_future_is_settled_when_worker_exits_without_finishing
83+
@work_queue = WorkQueue.new(1)
84+
abandoned = @work_queue.enqueue { Thread.exit }
85+
86+
@work_queue.wait
87+
88+
assert(abandoned.completed?, 'a worker must never leave a future unsettled')
89+
assert_raises(WorkQueue::AbandonedTaskError) { abandoned.value }
90+
end
91+
92+
def test_rejects_non_positive_max_concurrent
93+
assert_raises(ArgumentError) { WorkQueue.new(0) }
94+
assert_raises(ArgumentError) { WorkQueue.new(-1) }
95+
end
96+
6897
def test_max_concurrent
6998
max_concurrent = 2
7099
@work_queue = WorkQueue.new(max_concurrent)

0 commit comments

Comments
 (0)