From a60a9bb85292386d9c99e2edcb5754564a6c754e Mon Sep 17 00:00:00 2001 From: Ryan Hoffman Date: Wed, 7 Oct 2026 13:47:50 -0700 Subject: [PATCH] Don't report "No leader was elected" while the leader is still streaming In lazy-load mode the leader streams tests to the queue in batches and only marks the queue as initialized after the last batch, while workers start running tests as soon as the first batch arrives. If the report's wait ended before that, because --max-test-failed was reached or the report timed out, `queue_initialized?` was still false, so `report` aborted with "No leader was elected" and exit status 40 instead of reporting the failures. Treat a streaming leader as present, so the report falls through to BuildStatusReporter: it prints the failures and exits with 44 (too many failed tests) or 43 (timed out). Tests the leader hasn't streamed yet aren't in the queue, so the "tests weren't run" count becomes "at least". --- ruby/lib/minitest/queue/runner.rb | 18 +++- ruby/test/integration/minitest_redis_test.rb | 90 ++++++++++++++++++++ 2 files changed, 106 insertions(+), 2 deletions(-) diff --git a/ruby/lib/minitest/queue/runner.rb b/ruby/lib/minitest/queue/runner.rb index 106c7b5f..de39a04e 100644 --- a/ruby/lib/minitest/queue/runner.rb +++ b/ruby/lib/minitest/queue/runner.rb @@ -304,7 +304,11 @@ def report_command step("Waiting for workers to complete") unless supervisor.wait_for_workers { display_warnings(supervisor.build) } - unless supervisor.queue_initialized? + # A lazy-load leader only marks the queue as initialized once it has + # streamed every test, but workers start on the first batch, so the wait + # can end (e.g. on --max-test-failed) while it is still streaming. + # Check streaming? first: the status only moves from streaming to ready. + unless supervisor.streaming? || supervisor.queue_initialized? abort! "No leader was elected. This typically means no worker was able to start. Were there any errors during application boot?", 40 end @@ -314,7 +318,7 @@ def report_command reporter.write_failure_file(queue_config.failure_file) if queue_config.failure_file reporter.write_flaky_tests_file(queue_config.export_flaky_tests_file) if queue_config.export_flaky_tests_file - abort!("#{supervisor.size} tests weren't run.", exit_code) + abort!(unrun_tests_message(supervisor), exit_code) end end @@ -447,6 +451,16 @@ def print_worker_profiles(supervisor) Minitest::Queue::WorkerProfileReporter.new(supervisor).print_summary end + def unrun_tests_message(supervisor) + # Read the status before the size: tests a streaming leader hasn't pushed + # yet aren't in the queue, so the size is only a lower bound. + if supervisor.streaming? + "At least #{supervisor.size} tests weren't run (the leader hadn't finished streaming tests to the queue)." + else + "#{supervisor.size} tests weren't run." + end + end + def parser @parser ||= OptionParser.new do |opts| opts.banner = "Usage: minitest-queue [options] COMMAND [ARGS]" diff --git a/ruby/test/integration/minitest_redis_test.rb b/ruby/test/integration/minitest_redis_test.rb index 75c509f4..974e4745 100644 --- a/ruby/test/integration/minitest_redis_test.rb +++ b/ruby/test/integration/minitest_redis_test.rb @@ -674,6 +674,78 @@ def test_all_workers_died assert_equal expected.strip, normalize(out.lines[0..2].join.strip) end + def test_max_test_failed_while_leader_is_streaming + build_id = 'max-test-failed-while-streaming' + test_file = File.expand_path('../../fixtures/test/failing_test.rb', __FILE__) + entries = 10.times.map { |i| CI::Queue::QueueEntry.format("FailingTest#test_failing_#{i}", test_file) } + + with_streaming_leader(build_id, entries) do + _, err = capture_subprocess_io do + system( + @exe, 'run', + '--queue', @redis_url, + '--seed', 'foobar', + '--build', build_id, + '--worker', '1', + '--timeout', '5', + '--lazy-load', + '--max-test-failed', '1', + '-Itest', + 'test/failing_test.rb', + chdir: 'test/fixtures/', + ) + end + + refute_predicate $?, :success? + assert_equal 'This worker is exiting early because too many failed tests were encountered.', filter_deprecation_warnings(err).chomp + + out, err = capture_subprocess_io do + system( + @exe, 'report', + '--queue', @redis_url, + '--build', build_id, + '--timeout', '5', + '--max-test-failed', '1', + chdir: 'test/fixtures/', + ) + end + + assert_equal 44, $?.exitstatus + assert_empty filter_deprecation_warnings(err) + output = normalize(out) + refute_includes output, 'No leader was elected' + assert_includes output, 'Ran 1 tests, 1 assertions, 1 failures, 0 errors, 0 skips, 0 requeues in X.XXs (aggregated)' + assert_includes output, 'Encountered too many failed tests. Test run was ended early.' + assert_includes output, 'Error 1 of 1' + assert_equal "At least 9 tests weren't run (the leader hadn't finished streaming tests to the queue).", output.lines.last.strip + end + end + + def test_report_timeout_while_leader_is_streaming + build_id = 'report-timeout-while-streaming' + test_file = File.expand_path('../../fixtures/test/passing_test.rb', __FILE__) + entries = 3.times.map { |i| CI::Queue::QueueEntry.format("PassingTest#test_passing_#{i}", test_file) } + + with_streaming_leader(build_id, entries) do + out, err = capture_subprocess_io do + system( + @exe, 'report', + '--queue', @redis_url, + '--build', build_id, + '--timeout', '1', + chdir: 'test/fixtures/', + ) + end + + assert_equal 43, $?.exitstatus + assert_empty filter_deprecation_warnings(err) + output = normalize(out) + refute_includes output, 'No leader was elected' + assert_includes output, 'Timed out waiting for tests to be executed.' + assert_equal "At least 3 tests weren't run (the leader hadn't finished streaming tests to the queue).", output.lines.last.strip + end + end + def test_circuit_breaker out, err = capture_subprocess_io do system( @@ -1966,5 +2038,23 @@ def test_application_error def normalize_xml(output) normalize_backtrace(freeze_xml_timing(rewrite_paths(output))) end + + # Streams `entries` to the queue like a lazy-load leader, then pauses before + # marking the queue ready, so the block runs while the leader is still streaming. + def with_streaming_leader(build_id, entries) + tests = Enumerator.new do |yielder| + entries.each { |entry| yielder << entry } + Fiber.yield # pause stream_populate here until the fiber is resumed + end + leader = CI::Queue::Redis.new( + @redis_url, + CI::Queue::Configuration.new(build_id: build_id, worker_id: 'leader', timeout: 5), + ) + stream = Fiber.new { leader.stream_populate(tests, batch_size: 1) } + capture_io { stream.resume } + yield + ensure + capture_io { stream.resume } if stream&.alive? + end end end