Skip to content
Closed
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
13 changes: 13 additions & 0 deletions ruby/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,19 @@ minitest-queue --queue redis://example.com run -Itest test/**/*_test.rb

Additionally you can configure the requeue settings (see main README) with `--max-requeues` and `--requeue-tolerance`.

To recover artifacts lost with a retried distributed worker, replay every test
reserved by that worker before resuming the shared queue:

```bash
minitest-queue --queue redis://example.com \
--retry-mode worker-history \
--recovery-manifest log/ci-queue-recovery.json \
run -Itest test/**/*_test.rb
```

This mode requires the retry to retain its worker ID and queue build ID. It
fails when the worker's reservation history or suite chunk metadata is missing.


If you'd like to centralize the error reporting you can do so with:

Expand Down
1 change: 1 addition & 0 deletions ruby/lib/ci/queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
require 'ci/queue/common'
require 'ci/queue/build_record'
require 'ci/queue/static'
require 'ci/queue/worker_history_recovery'
require 'ci/queue/file'
require 'ci/queue/grind'
require 'ci/queue/bisect'
Expand Down
10 changes: 9 additions & 1 deletion ruby/lib/ci/queue/configuration.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ class Configuration
attr_accessor :timing_redis_url
attr_accessor :write_duration_averages
attr_accessor :heartbeat_grace_period, :heartbeat_interval
attr_accessor :retry_mode, :recovery_manifest, :retry_count
attr_reader :circuit_breakers
attr_writer :seed, :build_id
attr_writer :queue_init_timeout, :report_timeout, :inactive_workers_timeout
Expand All @@ -30,6 +31,7 @@ def from_env(env)
redis_ttl: env['CI_QUEUE_REDIS_TTL']&.to_i || 8 * 60 * 60,
known_flaky_tests: load_known_flaky_tests(env['CI_QUEUE_KNOWN_FLAKY_TESTS']),
branch: env['BUILDKITE_BRANCH'],
retry_count: env['BUILDKITE_RETRY_COUNT']&.to_i || 0,
)
end

Expand Down Expand Up @@ -66,7 +68,10 @@ def initialize(
branch: nil,
timing_redis_url: nil,
heartbeat_grace_period: 30,
heartbeat_interval: 10
heartbeat_interval: 10,
retry_mode: :failures,
recovery_manifest: nil,
retry_count: 0
)
@build_id = build_id
@circuit_breakers = [CircuitBreaker::Disabled]
Expand Down Expand Up @@ -105,6 +110,9 @@ def initialize(
@write_duration_averages = false
@heartbeat_grace_period = heartbeat_grace_period
@heartbeat_interval = heartbeat_interval
@retry_mode = retry_mode
@recovery_manifest = recovery_manifest
@retry_count = retry_count
end

def queue_init_timeout
Expand Down
1 change: 1 addition & 0 deletions ruby/lib/ci/queue/redis.rb
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ module Queue
module Redis
Error = Class.new(StandardError)
LostMaster = Class.new(Error)
WorkerHistoryError = Class.new(Error)

class << self

Expand Down
46 changes: 46 additions & 0 deletions ruby/lib/ci/queue/redis/worker.rb
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,15 @@ class << self
self.max_sleep_time = 2

class Worker < Base
class WorkerHistory
attr_reader :history_items, :test_ids

def initialize(history_items:, test_ids:)
@history_items = history_items
@test_ids = test_ids.freeze
end
end

DEFAULT_SLEEP_SECONDS = 0.5
attr_reader :total

Expand Down Expand Up @@ -154,6 +163,25 @@ def retry_queue
Retry.new(log, config, redis: redis)
end

def worker_history
reservations = redis.lrange(key('worker', worker_id, 'queue'), 0, -1)
if reservations.empty?
raise WorkerHistoryError, "Reservation history is missing for worker #{worker_id}"
end

seen = {}
test_ids = reservations.reverse_each.each_with_object([]) do |reservation_id, ids|
expand_reservation(reservation_id).each do |test_id|
next if seen[test_id]

seen[test_id] = true
ids << test_id
end
end

WorkerHistory.new(history_items: reservations.size, test_ids: test_ids)
end

def supervisor
Supervisor.new(redis_url, config)
end
Expand Down Expand Up @@ -525,6 +553,24 @@ def chunk_id?(id)
id.include?(':chunk_')
end

def expand_reservation(id)
return [id] unless chunk_id?(id)

chunk_json = redis.get(key('chunk', id))
unless chunk_json
raise WorkerHistoryError, "Chunk metadata is missing for #{id}"
end

test_ids = CI::Queue::TestChunk.from_json(id, chunk_json).test_ids
if test_ids.empty?
raise WorkerHistoryError, "Chunk metadata contains no tests for #{id}"
end

test_ids
rescue JSON::ParserError => error
raise WorkerHistoryError, "Chunk metadata is invalid for #{id}: #{error.message}"
end

def resolve_executable(id)
# Detect chunk by ID pattern
if chunk_id?(id)
Expand Down
173 changes: 173 additions & 0 deletions ruby/lib/ci/queue/worker_history_recovery.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
# frozen_string_literal: true

module CI
module Queue
class WorkerHistoryRecovery
attr_reader :config, :history_items, :replayed_tests

def initialize(shared_queue, history)
@shared_queue = shared_queue
@config = shared_queue.config
@history_items = history.history_items
@replay_ids = history.test_ids.dup
@replayed_tests = history.test_ids.size
@replay_completed = false
@resumed_shared_queue = false
@replay_failures = 0
@shutdown_required = false
@phase = :replay
end

def distributed?
true
end

def populate(tests, random: Random.new)
@index = tests.map { |test| [test.id, test] }.to_h
shared_queue.populate(tests, random: random)
self
end

def populated?
defined?(@index) && shared_queue.populated?
end

def poll(&block)
while replaying? && replay_allowed? && (id = @replay_ids.shift)
block.call(index.fetch(id))
end

return unless @replay_ids.empty?
return unless replay_allowed?

@replay_completed = true
@resumed_shared_queue = true
@phase = :shared
shared_queue.poll(&block)
end
Comment thread
cursor[bot] marked this conversation as resolved.

def replay_completed?
@replay_completed
end

def resumed_shared_queue?
@resumed_shared_queue
end

def acknowledge(test)
return shared_queue.acknowledge(test) unless replaying?

true
end

def requeue(test, **options)
return false if replaying?

if options.empty?
shared_queue.requeue(test)
else
shared_queue.requeue(test, **options)
end
end

def increment_test_failed
if replaying?
@replay_failures += 1
else
shared_queue.increment_test_failed
end
end

def test_failed
replaying? ? @replay_failures : shared_queue.test_failed
end

def max_test_failed?
return false if config.max_test_failed.nil?

test_failed >= config.max_test_failed
end

def exhausted?
@replay_ids.empty? && shared_queue.exhausted?
end

def size
@replay_ids.size + shared_queue.size
end

def total
shared_queue.total
end

def progress
shared_queue.progress
end

def to_a
@replay_ids.map { |id| index.fetch(id) } + shared_queue.to_a
end

def build
shared_queue.build
end

def supervisor
shared_queue.supervisor
end

def retrying?
true
end

def retry_queue
self
end

def expired?
shared_queue.expired?
end

def created_at=(timestamp)
shared_queue.created_at = timestamp
end

def release!
shared_queue.release!
end

def shutdown!
@shutdown_required = true
shared_queue.shutdown!
end

def flaky?(test)
shared_queue.flaky?(test)
end

def report_failure!
shared_queue.report_failure!
end

def report_success!
shared_queue.report_success!
end

def rescue_connection_errors(handler = ->(_error) { nil }, &block)
shared_queue.rescue_connection_errors(handler, &block)
end

private

attr_reader :index, :shared_queue

def replaying?
@phase == :replay
end

def replay_allowed?
!@shutdown_required && config.circuit_breakers.none?(&:open?) && !max_test_failed?
end
end
end
end
1 change: 1 addition & 0 deletions ruby/lib/minitest/queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
require 'minitest/queue/grind_reporter'
require 'minitest/queue/test_time_recorder'
require 'minitest/queue/test_time_reporter'
require 'minitest/queue/recovery_reporter'

module Minitest
class Requeue < Skip
Expand Down
47 changes: 47 additions & 0 deletions ruby/lib/minitest/queue/recovery_reporter.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
# frozen_string_literal: true

require 'fileutils'
require 'json'
require 'minitest/reporters'

module Minitest
module Queue
class RecoveryReporter < Minitest::Reporters::BaseReporter
def initialize(path:, queue:, config:)
super({})
@path = path
@queue = queue
@config = config
end

def report
super
return unless queue.exhausted?
return if recovery? && !queue.replay_completed?

FileUtils.mkdir_p(File.dirname(path))
File.write(path, JSON.pretty_generate(manifest))
end

private

attr_reader :config, :path, :queue

def recovery?
queue.is_a?(CI::Queue::WorkerHistoryRecovery)
end

def manifest
{
schema_version: 1,
worker_id: config.worker_id.to_s,
retry_count: config.retry_count,
history_items: recovery? ? queue.history_items : 0,
replayed_tests: recovery? ? queue.replayed_tests : 0,
resumed_shared_queue: recovery? && queue.resumed_shared_queue?,
replay_completed: !recovery? || queue.replay_completed?
}
end
end
end
end
Loading