Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion ruby/Gemfile.lock
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
PATH
remote: .
specs:
ci-queue (0.97.0)
ci-queue (0.98.0)
logger

GEM
Expand Down
35 changes: 35 additions & 0 deletions ruby/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion ruby/lib/ci/queue/version.rb
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

module CI
module Queue
VERSION = '0.97.0'
VERSION = '0.98.0'
DEV_SCRIPTS_ROOT = ::File.expand_path('../../../../../redis', __FILE__)
RELEASE_SCRIPTS_ROOT = ::File.expand_path('../redis', __FILE__)
end
Expand Down
63 changes: 63 additions & 0 deletions ruby/lib/minitest/queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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=)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
19 changes: 19 additions & 0 deletions ruby/lib/minitest/queue/test_data.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down
14 changes: 13 additions & 1 deletion ruby/test/integration/minitest_redis_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down Expand Up @@ -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
Expand Down
122 changes: 122 additions & 0 deletions ruby/test/minitest/queue/parallel_worker_metadata_test.rb
Original file line number Diff line number Diff line change
@@ -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
Loading