From c30d5ed8ac97b2a6cb23c40b9b2bea186df8442d Mon Sep 17 00:00:00 2001 From: Aiden Storey Date: Thu, 6 Aug 2026 15:19:25 -0400 Subject: [PATCH 1/2] Stamp per-worker-process execution metadata on test results Adds three fields to each recorded result, stamped in the process that ran the test (worker-side, before any DRb send in embedding environments): - parallel_worker_pid: Process.pid at result creation - parallel_worker_test_index: 0-based per-process monotonic counter, incremented per execution (requeued runs get their own index), fork-safe - parallel_worker_id: injected by the embedding environment via Minitest::Queue.parallel_worker_id= or CI_QUEUE_PARALLEL_WORKER_ID; nil when not applicable Stamping is first-writer-wins: embedders that run tests in forked workers and transport results to a central reporting process (e.g. Rails parallel testing over DRb, where handle_test_result runs server-side) must call Minitest::Queue.stamp_parallel_worker_metadata in the worker before sending; pre-stamped results pass through reporting untouched. Otherwise the reporting-side stamp would carry the server's pid and an arrival-order index interleaved across workers. The fields are carried nil-safely through TestData#to_h into log/test_data.json (TestDataReporter unchanged), so per-worker-process execution order is reconstructable downstream: PARTITION BY parallel_worker_id, parallel_worker_pid ORDER BY parallel_worker_test_index Assisted-By: devx/77f6d45d-84ba-4c9e-983d-5ec948229a08 --- ruby/README.md | 35 +++++ ruby/lib/minitest/queue.rb | 63 +++++++++ ruby/lib/minitest/queue/test_data.rb | 19 +++ ruby/test/integration/minitest_redis_test.rb | 14 +- .../queue/parallel_worker_metadata_test.rb | 122 ++++++++++++++++++ ruby/test/minitest/queue/test_data_test.rb | 49 +++++++ 6 files changed, 301 insertions(+), 1 deletion(-) create mode 100644 ruby/test/minitest/queue/parallel_worker_metadata_test.rb diff --git a/ruby/README.md b/ruby/README.md index 52c04683..8a7dc833 100644 --- a/ruby/README.md +++ b/ruby/README.md @@ -158,6 +158,41 @@ rspec-queue --queue redis://example.com --timeout 600 --report Because of how `ci-queue` executes the examples, `before(:all)` and `after(:all)` hooks are not supported. `rspec-queue` will explicitly reject them. +### Parallel worker metadata + +Each recorded result is stamped, in the process that ran the test, with: + +- `parallel_worker_pid`: the pid of the worker process at result creation. +- `parallel_worker_test_index`: a 0-based, per-process monotonic counter incremented as each test runs (requeued executions get their own index). +- `parallel_worker_id`: an identifier injected by the embedding environment (e.g. a Rails parallel-testing worker number); `nil` when not applicable. + +All three fields are included (nil-safe) in `log/test_data.json` emitted by the test data reporter, making per-worker-process execution order reconstructable downstream (`PARTITION BY parallel_worker_id, parallel_worker_pid ORDER BY parallel_worker_test_index`). + +The worker id can be provided either programmatically or via the environment: + +```ruby +Minitest::Queue.parallel_worker_id = 3 +``` + +Stamping is first-writer-wins. When ci-queue runs tests in-process (the normal +`minitest-queue` flow), results are stamped automatically as they are recorded. +Embedding environments that run tests in forked workers and transport results to +another process (e.g. over DRb to a central reporting server) must stamp in the +worker **before** sending: + +```ruby +# in the forked worker, after running the test and before the DRb send +Minitest::Queue.stamp_parallel_worker_metadata(result) +``` + +Otherwise the automatic stamp during reporting would capture the reporting +process's pid and an arrival-order index interleaved across workers. Pre-stamped +results pass through reporting untouched. + +| Variable | Description | +|---|---| +| `CI_QUEUE_PARALLEL_WORKER_ID=N` | Sets `parallel_worker_id` for results produced by this process. The setter takes precedence. | + ### Worker circuit breakers Both runners support two independent, disabled-by-default circuit breakers: diff --git a/ruby/lib/minitest/queue.rb b/ruby/lib/minitest/queue.rb index 4ae1a40c..6cabf571 100644 --- a/ruby/lib/minitest/queue.rb +++ b/ruby/lib/minitest/queue.rb @@ -111,6 +111,15 @@ module ResultMetadata attr_accessor :queue_id, :queue_entry end + # Per-worker-process execution metadata, stamped on each result in the + # process that ran the test (before any DRb send in embedding + # environments), so that per-process execution order is reconstructable + # downstream: PARTITION BY parallel_worker_id, parallel_worker_pid + # ORDER BY parallel_worker_test_index. + module ParallelWorkerMetadata + attr_accessor :parallel_worker_id, :parallel_worker_test_index, :parallel_worker_pid + end + module Queue extend ::CI::Queue::OutputHelpers attr_writer :run_command_formatter, :project_root @@ -155,6 +164,45 @@ def self.relative_path(path, root: project_root) end class << self + # Identifies which parallel worker this process is (e.g. Rails' + # parallel testing worker number). Injected by the embedding + # environment via this setter or the CI_QUEUE_PARALLEL_WORKER_ID + # environment variable; nil when not applicable. + attr_writer :parallel_worker_id + + def parallel_worker_id + @parallel_worker_id || parallel_worker_id_from_env + end + + # Stamps per-process execution metadata on a result. Must be called + # in the process that ran the test so pid and per-process test index + # are captured worker-side. + # + # First writer wins: results that were already stamped are left + # untouched. Embedding environments that run tests in forked workers + # and transport results to another process (e.g. over DRb) MUST call + # this in the worker before sending, otherwise the stamp applied + # during reporting would carry the reporting process's pid and an + # arrival-order index interleaved across workers. + def stamp_parallel_worker_metadata(result) + return unless result.respond_to?(:parallel_worker_pid=) + return if result.parallel_worker_pid # already stamped worker-side + + pid = Process.pid + if @parallel_worker_metadata_pid != pid + # Restart the per-process counter in forked workers so that + # (parallel_worker_id, parallel_worker_pid, parallel_worker_test_index) + # always reflects execution order within a single process. + @parallel_worker_metadata_pid = pid + @parallel_worker_next_test_index = 0 + end + + result.parallel_worker_id = parallel_worker_id + result.parallel_worker_pid = pid + result.parallel_worker_test_index = @parallel_worker_next_test_index + @parallel_worker_next_test_index += 1 + end + def queue Minitest.queue end @@ -179,6 +227,8 @@ def run(reporter, *) end def handle_test_result(reporter, example, result) + stamp_parallel_worker_metadata(result) + if result.respond_to?(:queue_id=) result.queue_id = example.id result.queue_entry = example.queue_entry if result.respond_to?(:queue_entry=) @@ -213,6 +263,17 @@ def handle_test_result(reporter, example, result) private + def parallel_worker_id_from_env + value = ENV['CI_QUEUE_PARALLEL_WORKER_ID'] + return nil if value.nil? || value.empty? + + begin + Integer(value) + rescue ArgumentError + value + end + end + def report_load_stats(queue) return unless CI::Queue.debug? return unless queue.respond_to?(:file_loader) @@ -608,11 +669,13 @@ def loaded_tests Minitest::Result.prepend(Minitest::Flakiness) Minitest::Result.prepend(Minitest::WithTimestamps) Minitest::Result.prepend(Minitest::ResultMetadata) + Minitest::Result.prepend(Minitest::ParallelWorkerMetadata) else Minitest::Test.prepend(Minitest::Requeueing) Minitest::Test.prepend(Minitest::Flakiness) Minitest::Test.prepend(Minitest::WithTimestamps) Minitest::Test.prepend(Minitest::ResultMetadata) + Minitest::Test.prepend(Minitest::ParallelWorkerMetadata) module MinitestBackwardCompatibility def source_location diff --git a/ruby/lib/minitest/queue/test_data.rb b/ruby/lib/minitest/queue/test_data.rb index 0043d220..d4a2d0f2 100644 --- a/ruby/lib/minitest/queue/test_data.rb +++ b/ruby/lib/minitest/queue/test_data.rb @@ -73,6 +73,22 @@ def test_file_line_number @test.source_location.last end + # The parallel_worker_* fields are stamped by + # Minitest::Queue.stamp_parallel_worker_metadata in the worker process + # that ran the test. They are nil for embedders that don't produce + # them (results lacking the accessors, or no worker id configured). + def parallel_worker_id + @test.parallel_worker_id if @test.respond_to?(:parallel_worker_id) + end + + def parallel_worker_test_index + @test.parallel_worker_test_index if @test.respond_to?(:parallel_worker_test_index) + end + + def parallel_worker_pid + @test.parallel_worker_pid if @test.respond_to?(:parallel_worker_pid) + end + # Error class only considers failures wheras the other error fields also consider skips def error_class return nil unless @test.failure @@ -119,6 +135,9 @@ def to_h test_finish_timestamp: test_finish_timestamp, test_file_path: test_file_path, test_file_line_number: test_file_line_number, + parallel_worker_id: parallel_worker_id, + parallel_worker_test_index: parallel_worker_test_index, + parallel_worker_pid: parallel_worker_pid, error_class: error_class, error_message: error_message, error_file_path: error_file_path, diff --git a/ruby/test/integration/minitest_redis_test.rb b/ruby/test/integration/minitest_redis_test.rb index b3d450e4..7660f54e 100644 --- a/ruby/test/integration/minitest_redis_test.rb +++ b/ruby/test/integration/minitest_redis_test.rb @@ -1477,7 +1477,10 @@ def test_down_redis def test_test_data_reporter out, err = capture_subprocess_io do system( - {'CI_QUEUE_FLAKY_TESTS' => 'test/ci_queue_flaky_tests_list.txt'}, + { + 'CI_QUEUE_FLAKY_TESTS' => 'test/ci_queue_flaky_tests_list.txt', + 'CI_QUEUE_PARALLEL_WORKER_ID' => '7', + }, @exe, 'run', '--queue', @redis_url, '--seed', 'foobar', @@ -1544,6 +1547,15 @@ def test_test_data_reporter assert_equal 'ATest#test_flaky_passes', failures[4][:test_id] assert_equal 'success', failures[4][:test_result] + + # Parallel worker metadata is stamped per execution, in the worker + # process, so per-process execution order is reconstructable. + assert_equal [7], failures.map { |f| f[:parallel_worker_id] }.uniq + assert_equal 1, failures.map { |f| f[:parallel_worker_pid] }.uniq.size + assert_kind_of Integer, failures.first[:parallel_worker_pid] + assert_equal (0...failures.size).to_a, failures.map { |f| f[:parallel_worker_test_index] }.sort + # The requeued execution of ATest#test_bar ran before its final one. + assert failures[0][:parallel_worker_test_index] < failures[1][:parallel_worker_test_index] end def test_test_data_time_reporter diff --git a/ruby/test/minitest/queue/parallel_worker_metadata_test.rb b/ruby/test/minitest/queue/parallel_worker_metadata_test.rb new file mode 100644 index 00000000..49568583 --- /dev/null +++ b/ruby/test/minitest/queue/parallel_worker_metadata_test.rb @@ -0,0 +1,122 @@ +# frozen_string_literal: true + +require 'test_helper' + +module Minitest::Queue + class ParallelWorkerMetadataTest < Minitest::Test + include ReporterTestHelper + + ENV_KEY = 'CI_QUEUE_PARALLEL_WORKER_ID' + + def setup + @original_env = ENV.delete(ENV_KEY) + @original_worker_id = Minitest::Queue.instance_variable_get(:@parallel_worker_id) + Minitest::Queue.parallel_worker_id = nil + end + + def teardown + @original_env ? ENV[ENV_KEY] = @original_env : ENV.delete(ENV_KEY) + Minitest::Queue.parallel_worker_id = @original_worker_id + end + + def test_stamps_pid_and_monotonic_per_process_index + first = result('test_foo') + second = result('test_bar') + + Minitest::Queue.stamp_parallel_worker_metadata(first) + Minitest::Queue.stamp_parallel_worker_metadata(second) + + assert_equal Process.pid, first.parallel_worker_pid + assert_equal Process.pid, second.parallel_worker_pid + assert_kind_of Integer, first.parallel_worker_test_index + assert_equal first.parallel_worker_test_index + 1, second.parallel_worker_test_index + end + + def test_worker_id_is_nil_by_default + test = result('test_foo') + Minitest::Queue.stamp_parallel_worker_metadata(test) + + assert_nil test.parallel_worker_id + end + + def test_worker_id_from_setter + Minitest::Queue.parallel_worker_id = 3 + + test = result('test_foo') + Minitest::Queue.stamp_parallel_worker_metadata(test) + + assert_equal 3, test.parallel_worker_id + end + + def test_worker_id_from_env + ENV[ENV_KEY] = '7' + + test = result('test_foo') + Minitest::Queue.stamp_parallel_worker_metadata(test) + + assert_equal 7, test.parallel_worker_id + end + + def test_setter_takes_precedence_over_env + ENV[ENV_KEY] = '7' + Minitest::Queue.parallel_worker_id = 3 + + assert_equal 3, Minitest::Queue.parallel_worker_id + end + + def test_non_numeric_env_worker_id_is_passed_through + ENV[ENV_KEY] = 'worker-a' + + assert_equal 'worker-a', Minitest::Queue.parallel_worker_id + end + + def test_empty_env_worker_id_is_nil + ENV[ENV_KEY] = '' + + assert_nil Minitest::Queue.parallel_worker_id + end + + def test_index_restarts_when_pid_changes + # Prime the counter in this process. + Minitest::Queue.stamp_parallel_worker_metadata(result('test_foo')) + + # Simulate a fork: pretend the counter belongs to another process. + Minitest::Queue.instance_variable_set(:@parallel_worker_metadata_pid, Process.pid - 1) + + test = result('test_bar') + Minitest::Queue.stamp_parallel_worker_metadata(test) + + assert_equal 0, test.parallel_worker_test_index + assert_equal Process.pid, test.parallel_worker_pid + end + + def test_ignores_results_without_accessors + plain = Object.new + Minitest::Queue.stamp_parallel_worker_metadata(plain) # does not raise + end + + def test_does_not_overwrite_worker_side_stamps + # Simulates an embedding environment (e.g. Rails parallel testing over + # DRb) where the worker stamped the result before sending it to the + # process running the reporters: the reporting-side stamp in + # handle_test_result must not clobber it. + stamped = result('test_foo') + stamped.parallel_worker_id = 9 + stamped.parallel_worker_test_index = 4 + stamped.parallel_worker_pid = 4242 + + before = result('test_before') + Minitest::Queue.stamp_parallel_worker_metadata(before) + Minitest::Queue.stamp_parallel_worker_metadata(stamped) + after = result('test_after') + Minitest::Queue.stamp_parallel_worker_metadata(after) + + assert_equal 9, stamped.parallel_worker_id + assert_equal 4, stamped.parallel_worker_test_index + assert_equal 4242, stamped.parallel_worker_pid + + # The local per-process counter must not advance for pre-stamped results. + assert_equal before.parallel_worker_test_index + 1, after.parallel_worker_test_index + end + end +end diff --git a/ruby/test/minitest/queue/test_data_test.rb b/ruby/test/minitest/queue/test_data_test.rb index 1ded1ea0..a0821d9b 100644 --- a/ruby/test/minitest/queue/test_data_test.rb +++ b/ruby/test/minitest/queue/test_data_test.rb @@ -41,5 +41,54 @@ def test_error_location_uses_nested_exception_backtrace assert_equal 'test/nested_error_test.rb', data[:error_file_path].to_s assert_equal 42, data[:error_file_number] end + + def test_parallel_worker_metadata_defaults_to_nil + test = result('test_foo') + + data = TestData.new( + test: test, + index: 0, + namespace: 'namespace', + base_path: Minitest::Queue.project_root, + ).to_h + + assert data.key?(:parallel_worker_id) + assert data.key?(:parallel_worker_test_index) + assert data.key?(:parallel_worker_pid) + assert_nil data[:parallel_worker_id] + assert_nil data[:parallel_worker_test_index] + assert_nil data[:parallel_worker_pid] + end + + def test_parallel_worker_metadata_carried_when_stamped + test = result('test_foo') + test.parallel_worker_id = 3 + test.parallel_worker_test_index = 42 + test.parallel_worker_pid = 12345 + + data = TestData.new( + test: test, + index: 0, + namespace: 'namespace', + base_path: Minitest::Queue.project_root, + ).to_h + + assert_equal 3, data[:parallel_worker_id] + assert_equal 42, data[:parallel_worker_test_index] + assert_equal 12345, data[:parallel_worker_pid] + end + + def test_parallel_worker_metadata_nil_safe_without_accessors + data = TestData.new( + test: Object.new, + index: 0, + namespace: 'namespace', + base_path: Minitest::Queue.project_root, + ) + + assert_nil data.parallel_worker_id + assert_nil data.parallel_worker_test_index + assert_nil data.parallel_worker_pid + end end end From c6d40a8bea4bbc1100754111aea69ef34263cad0 Mon Sep 17 00:00:00 2001 From: Aiden Storey Date: Tue, 11 Aug 2026 11:56:32 -0400 Subject: [PATCH 2/2] Address review feedback - Guard stamp state with a mutex: results recorded from multiple threads (e.g. a DRb server dispatching each call on its own thread) get unique, gap-free per-process indexes. Documented that per-worker order reconstruction is only meaningful with forked workers; under thread-based parallelization all threads share one partition. - Memoize the CI_QUEUE_PARALLEL_WORKER_ID env lookup per process (keyed on pid, so forked workers re-read it) instead of re-parsing on every test. - README: moved the section under Minitest, fixed stamp-timing wording, documented the relationship to ci-queue's own --worker/worker_id, and where build/job identity is expected to come from. - Tests: stamp now returns the result or nil so skip paths are assertable; added thread-safety and env-memoization tests; reset all stamp state ivars in setup/teardown; style fix in the integration assertions. Assisted-By: devx/d62b5d38-ac27-42b3-8627-64db4c652781 --- ruby/README.md | 70 ++++++++++++------- ruby/lib/minitest/queue.rb | 44 ++++++++---- ruby/test/integration/minitest_redis_test.rb | 2 +- .../queue/parallel_worker_metadata_test.rb | 50 +++++++++++-- 4 files changed, 120 insertions(+), 46 deletions(-) diff --git a/ruby/README.md b/ruby/README.md index 8a7dc833..6c101bbb 100644 --- a/ruby/README.md +++ b/ruby/README.md @@ -138,32 +138,12 @@ The runner also comes with a tool to investigate leaky tests: minitest-queue --queue path/to/test_order.log --failing-test 'SomeTest#test_something' bisect -Itest test/**/*_test.rb ``` -### RSpec [DEPRECATED] - -The rspec-queue runner is deprecated. The minitest-queue runner continues to be supported and is actively being improved. At Shopify, we strongly recommend that new projects set up their test suite using Minitest rather than RSpec. - -Assuming you use one of the supported CI providers, the command can be as simple as: - -```bash -rspec-queue --queue redis://example.com -``` +#### Parallel worker metadata -If you'd like to centralize the error reporting you can do so with: +Each result is stamped, as it is recorded in the process that ran the test, with: -```bash -rspec-queue --queue redis://example.com --timeout 600 --report -``` - -#### Limitations - -Because of how `ci-queue` executes the examples, `before(:all)` and `after(:all)` hooks are not supported. `rspec-queue` will explicitly reject them. - -### Parallel worker metadata - -Each recorded result is stamped, in the process that ran the test, with: - -- `parallel_worker_pid`: the pid of the worker process at result creation. -- `parallel_worker_test_index`: a 0-based, per-process monotonic counter incremented as each test runs (requeued executions get their own index). +- `parallel_worker_pid`: the pid of the process that ran the test. +- `parallel_worker_test_index`: a 0-based monotonic counter of results recorded by that process (requeued executions get their own index). Restarts at 0 in each forked process. - `parallel_worker_id`: an identifier injected by the embedding environment (e.g. a Rails parallel-testing worker number); `nil` when not applicable. All three fields are included (nil-safe) in `log/test_data.json` emitted by the test data reporter, making per-worker-process execution order reconstructable downstream (`PARTITION BY parallel_worker_id, parallel_worker_pid ORDER BY parallel_worker_test_index`). @@ -174,6 +154,10 @@ The worker id can be provided either programmatically or via the environment: Minitest::Queue.parallel_worker_id = 3 ``` +| Variable | Description | +|---|---| +| `CI_QUEUE_PARALLEL_WORKER_ID=N` | Sets `parallel_worker_id` for results produced by this process. Read once per process; the setter takes precedence. | + Stamping is first-writer-wins. When ci-queue runs tests in-process (the normal `minitest-queue` flow), results are stamped automatically as they are recorded. Embedding environments that run tests in forked workers and transport results to @@ -189,9 +173,41 @@ Otherwise the automatic stamp during reporting would capture the reporting process's pid and an arrival-order index interleaved across workers. Pre-stamped results pass through reporting untouched. -| Variable | Description | -|---|---| -| `CI_QUEUE_PARALLEL_WORKER_ID=N` | Sets `parallel_worker_id` for results produced by this process. The setter takes precedence. | +Notes: + +- `parallel_worker_id` identifies a forked test process *inside* one queue worker. + It is unrelated to ci-queue's own `--worker` / `CI::Queue::Configuration#worker_id`, + which identifies the whole queue worker (typically one CI job). +- The payload intentionally carries no build/job identity — like every other field + in `log/test_data.json`, scoping to a build/job (e.g. a `job_id` column making + `(job_id, parallel_worker_id, parallel_worker_pid)` unique across a build) is + expected to be attached by whatever pipeline ingests the file. +- Stamping is mutex-guarded, so indexes are unique and gap-free even when results + are recorded from multiple threads. However, per-worker order reconstruction is + only meaningful with process-based (forked) workers: under thread-based + parallelization all threads share one `(parallel_worker_id, parallel_worker_pid)` + partition and the index reflects record order across threads. +- This applies to minitest-queue only; rspec-queue does not emit these fields. + +### RSpec [DEPRECATED] + +The rspec-queue runner is deprecated. The minitest-queue runner continues to be supported and is actively being improved. At Shopify, we strongly recommend that new projects set up their test suite using Minitest rather than RSpec. + +Assuming you use one of the supported CI providers, the command can be as simple as: + +```bash +rspec-queue --queue redis://example.com +``` + +If you'd like to centralize the error reporting you can do so with: + +```bash +rspec-queue --queue redis://example.com --timeout 600 --report +``` + +#### Limitations + +Because of how `ci-queue` executes the examples, `before(:all)` and `after(:all)` hooks are not supported. `rspec-queue` will explicitly reject them. ### Worker circuit breakers diff --git a/ruby/lib/minitest/queue.rb b/ruby/lib/minitest/queue.rb index 6cabf571..65803555 100644 --- a/ruby/lib/minitest/queue.rb +++ b/ruby/lib/minitest/queue.rb @@ -124,6 +124,9 @@ module Queue extend ::CI::Queue::OutputHelpers attr_writer :run_command_formatter, :project_root + PARALLEL_WORKER_METADATA_MUTEX = Mutex.new + private_constant :PARALLEL_WORKER_METADATA_MUTEX + def run_command_formatter @run_command_formatter ||= if defined?(Rails) && defined?(Rails::TestUnitRailtie) RAILS_RUN_COMMAND_FORMATTER @@ -171,7 +174,14 @@ class << self attr_writer :parallel_worker_id def parallel_worker_id - @parallel_worker_id || parallel_worker_id_from_env + return @parallel_worker_id if @parallel_worker_id + + # Memoized per process: re-read after a fork, but not on every test. + if @parallel_worker_id_env_pid != Process.pid + @parallel_worker_id_env_pid = Process.pid + @parallel_worker_id_from_env = parallel_worker_id_from_env + end + @parallel_worker_id_from_env end # Stamps per-process execution metadata on a result. Must be called @@ -184,23 +194,31 @@ def parallel_worker_id # this in the worker before sending, otherwise the stamp applied # during reporting would carry the reporting process's pid and an # arrival-order index interleaved across workers. + # + # The counter is mutex-guarded: results recorded from multiple threads + # (e.g. a DRb server dispatching each call on its own thread) get + # unique, gap-free indexes reflecting record order in this process. + # Returns the result when it was stamped, nil when it was skipped. def stamp_parallel_worker_metadata(result) return unless result.respond_to?(:parallel_worker_pid=) return if result.parallel_worker_pid # already stamped worker-side - pid = Process.pid - if @parallel_worker_metadata_pid != pid - # Restart the per-process counter in forked workers so that - # (parallel_worker_id, parallel_worker_pid, parallel_worker_test_index) - # always reflects execution order within a single process. - @parallel_worker_metadata_pid = pid - @parallel_worker_next_test_index = 0 - end + PARALLEL_WORKER_METADATA_MUTEX.synchronize do + pid = Process.pid + if @parallel_worker_metadata_pid != pid + # Restart the per-process counter in forked workers so that + # (parallel_worker_id, parallel_worker_pid, parallel_worker_test_index) + # always reflects execution order within a single process. + @parallel_worker_metadata_pid = pid + @parallel_worker_next_test_index = 0 + end - result.parallel_worker_id = parallel_worker_id - result.parallel_worker_pid = pid - result.parallel_worker_test_index = @parallel_worker_next_test_index - @parallel_worker_next_test_index += 1 + result.parallel_worker_id = parallel_worker_id + result.parallel_worker_pid = pid + result.parallel_worker_test_index = @parallel_worker_next_test_index + @parallel_worker_next_test_index += 1 + end + result end def queue diff --git a/ruby/test/integration/minitest_redis_test.rb b/ruby/test/integration/minitest_redis_test.rb index 7660f54e..75c509f4 100644 --- a/ruby/test/integration/minitest_redis_test.rb +++ b/ruby/test/integration/minitest_redis_test.rb @@ -1553,7 +1553,7 @@ def test_test_data_reporter assert_equal [7], failures.map { |f| f[:parallel_worker_id] }.uniq assert_equal 1, failures.map { |f| f[:parallel_worker_pid] }.uniq.size assert_kind_of Integer, failures.first[:parallel_worker_pid] - assert_equal (0...failures.size).to_a, failures.map { |f| f[:parallel_worker_test_index] }.sort + assert_equal((0...failures.size).to_a, failures.map { |f| f[:parallel_worker_test_index] }.sort) # The requeued execution of ATest#test_bar ran before its final one. assert failures[0][:parallel_worker_test_index] < failures[1][:parallel_worker_test_index] end diff --git a/ruby/test/minitest/queue/parallel_worker_metadata_test.rb b/ruby/test/minitest/queue/parallel_worker_metadata_test.rb index 49568583..d4c8a27f 100644 --- a/ruby/test/minitest/queue/parallel_worker_metadata_test.rb +++ b/ruby/test/minitest/queue/parallel_worker_metadata_test.rb @@ -8,15 +8,23 @@ class ParallelWorkerMetadataTest < Minitest::Test ENV_KEY = 'CI_QUEUE_PARALLEL_WORKER_ID' + STATE_IVARS = %i[ + @parallel_worker_id + @parallel_worker_id_env_pid + @parallel_worker_id_from_env + @parallel_worker_metadata_pid + @parallel_worker_next_test_index + ].freeze + def setup @original_env = ENV.delete(ENV_KEY) - @original_worker_id = Minitest::Queue.instance_variable_get(:@parallel_worker_id) - Minitest::Queue.parallel_worker_id = nil + @original_state = STATE_IVARS.to_h { |ivar| [ivar, Minitest::Queue.instance_variable_get(ivar)] } + reset_stamp_state end def teardown @original_env ? ENV[ENV_KEY] = @original_env : ENV.delete(ENV_KEY) - Minitest::Queue.parallel_worker_id = @original_worker_id + @original_state.each { |ivar, value| Minitest::Queue.instance_variable_set(ivar, value) } end def test_stamps_pid_and_monotonic_per_process_index @@ -57,6 +65,18 @@ def test_worker_id_from_env assert_equal 7, test.parallel_worker_id end + def test_env_worker_id_is_read_once_per_process + ENV[ENV_KEY] = '7' + assert_equal 7, Minitest::Queue.parallel_worker_id + + ENV[ENV_KEY] = '8' + assert_equal 7, Minitest::Queue.parallel_worker_id + + # A forked process re-reads the environment. + Minitest::Queue.instance_variable_set(:@parallel_worker_id_env_pid, Process.pid - 1) + assert_equal 8, Minitest::Queue.parallel_worker_id + end + def test_setter_takes_precedence_over_env ENV[ENV_KEY] = '7' Minitest::Queue.parallel_worker_id = 3 @@ -92,7 +112,21 @@ def test_index_restarts_when_pid_changes def test_ignores_results_without_accessors plain = Object.new - Minitest::Queue.stamp_parallel_worker_metadata(plain) # does not raise + assert_nil Minitest::Queue.stamp_parallel_worker_metadata(plain) + end + + def test_stamping_is_thread_safe + results = Array.new(100) { |i| result("test_#{i}") } + + results.each_slice(20).map do |slice| + Thread.new do + slice.each { |test| Minitest::Queue.stamp_parallel_worker_metadata(test) } + end + end.each(&:join) + + indexes = results.map(&:parallel_worker_test_index).sort + assert_equal((0...results.size).to_a, indexes) + assert_equal [Process.pid], results.map(&:parallel_worker_pid).uniq end def test_does_not_overwrite_worker_side_stamps @@ -107,7 +141,7 @@ def test_does_not_overwrite_worker_side_stamps before = result('test_before') Minitest::Queue.stamp_parallel_worker_metadata(before) - Minitest::Queue.stamp_parallel_worker_metadata(stamped) + assert_nil Minitest::Queue.stamp_parallel_worker_metadata(stamped) after = result('test_after') Minitest::Queue.stamp_parallel_worker_metadata(after) @@ -118,5 +152,11 @@ def test_does_not_overwrite_worker_side_stamps # The local per-process counter must not advance for pre-stamped results. assert_equal before.parallel_worker_test_index + 1, after.parallel_worker_test_index end + + private + + def reset_stamp_state + STATE_IVARS.each { |ivar| Minitest::Queue.instance_variable_set(ivar, nil) } + end end end