diff --git a/Project.toml b/Project.toml index abac725..4de0f6c 100644 --- a/Project.toml +++ b/Project.toml @@ -1,7 +1,7 @@ name = "ParallelTestRunner" uuid = "d3525ed8-44d0-4b2c-a655-542cee43accc" authors = ["Valentin Churavy "] -version = "2.7.0" +version = "2.8.0" [deps] Dates = "ade2ca70-3891-5945-98fb-dc099432e06a" diff --git a/docs/src/advanced.md b/docs/src/advanced.md index 2ac2e93..e49acab 100644 --- a/docs/src/advanced.md +++ b/docs/src/advanced.md @@ -174,6 +174,53 @@ duration, longest first) and their results appear in the same overall summary. If the user filters tests via positional arguments (e.g. `julia test/runtests.jl unit`), any serial test names that were filtered out are silently removed from the serial list. +## Failure Handling + +Both options described in this section are opt-in and default to off. + +### Recycling Workers after a Failure + +Workers are reused across tests, so a test that corrupts process-wide state — a wedged GPU driver whose every subsequent allocation fails, a global left in an inconsistent state, a library put in an unusable configuration — can make every later test scheduled on that same worker fail too. + +Setting `recycle_on_failure=true` stops the worker after any test that did not pass, so the next test gets a fresh process: + +```julia +runtests(MyPackage, ARGS; recycle_on_failure=true) +``` + +This complements the existing recycling of workers exceeding `max_worker_rss` and of workers that crashed outright. + +### Retrying Failed Tests + +When several workers compete for a limited resource (usually memory), a failure can mean "lost the race for the resource" rather than "the code is broken". +Such a test typically passes when run on its own. + +The `retries` keyword argument re-runs tests that did not pass, up to `N` times, after the main run has completed: + +```julia +runtests(MyPackage, ARGS; retries=1) +``` + +Retried tests run **sequentially on a single fresh worker**, so a test that failed only because of concurrent resource pressure gets an otherwise-idle system. +If a test fails again, its worker is stopped before the next retry, so one failure cannot contaminate the following one. + +Only the final attempt of each test is recorded in the results, so a test that passes on retry is reported as passing and a persistently broken test is reported as failing. +Retries are visible in the output, so flakiness is surfaced rather than hidden: + +``` +Retrying 1 failed test (1) +fails (8) │ 0.05 │ failed at 2026-08-08T15:10:15.526 +``` + +!!! note + Retries are skipped when the run was interrupted (e.g. `Ctrl+C`) or when `--quickfail` is + in effect, since in both cases the run stopped early on purpose. + +!!! tip + `recycle_on_failure` and `retries` address different halves of the same problem and work + well together: recycling keeps one bad test from cascading onto its worker during the run, + while retries give the tests that did fail a contention-free second chance. + ## Custom Workers For tests that require specific environment variables or Julia flags, you can use the `test_worker` keyword argument to [`runtests`](@ref) to assign tests to custom workers: @@ -303,3 +350,5 @@ function jltest { 1. **Use custom workers sparingly**: Custom workers add overhead. Only use them when tests genuinely require different configurations. 1. **Use `serial` for resource-intensive tests**: If a test allocates significant memory or uses exclusive hardware resources, mark it as serial rather than reducing `--jobs` globally. This keeps the rest of your suite running in parallel. + +1. **Only use `retries` for worker contention-related failures**: Not all intermittent failures are caused by parallel worker resource contention. Ensure you aren't masking real test failures when using this feature. diff --git a/docs/src/index.md b/docs/src/index.md index cb89f41..b7d569c 100644 --- a/docs/src/index.md +++ b/docs/src/index.md @@ -114,6 +114,16 @@ The `serial` keyword argument to [`runtests`](@ref) lets you designate specific for sequential execution, either before or after the parallel batch. See [Serial Tests](@ref) in the advanced usage guide for details. +### Failure Recycling and Retries + +Workers are recycled when they crash or exceed the memory threshold. +Additionally, [`runtests`](@ref) has two keyword arguments to further customize +failure hanlding. Setting `recycle_on_failure=true` recycles a worker after any +failed test, so a test that corrupts process-wide state cannot poison later tests, +and `retries=N` re-runs failed tests sequentially up to `N` times to reduce false +failures caused by resource contention. +See [Failure Handling](@ref) in the advanced usage guide for details. + ### Real-time Progress The test runner provides real-time output showing: diff --git a/src/ParallelTestRunner.jl b/src/ParallelTestRunner.jl index 8baee3e..f30548c 100644 --- a/src/ParallelTestRunner.jl +++ b/src/ParallelTestRunner.jl @@ -167,6 +167,7 @@ struct TestIOContext alloc_align::Int rss_align::Int max_worker_rss::Int + nonpass_color::Ref{Symbol} end function test_IOContext(::Type{<:AbstractTestRecord}, stdout::IO, stderr::IO, lock::ReentrantLock, name_align::Int, verbose::Bool, max_worker_rss::Int) @@ -181,7 +182,7 @@ function test_IOContext(::Type{<:AbstractTestRecord}, stdout::IO, stderr::IO, lo return TestIOContext( stdout, stderr, color, verbose, lock, name_align, elapsed_align, compile_align, gc_align, percent_align, - alloc_align, rss_align, max_worker_rss + alloc_align, rss_align, max_worker_rss, Ref(:red) ) end @@ -268,26 +269,26 @@ function print_test_failed(record::AbstractTestRecord, wrkr, test, ctx::TestIOCo base = parent(record) lock(ctx.lock) try - printstyled(ctx.stderr, test, color = :red) + printstyled(ctx.stderr, test, color = ctx.nonpass_color[]) printstyled( ctx.stderr, lpad("($wrkr)", ctx.name_align - textwidth(test) + 1, " "), " │" - , color = :red + , color = ctx.nonpass_color[] ) time_str = @sprintf("%7.2f", base.time) - printstyled(ctx.stderr, lpad(time_str, ctx.elapsed_align + 1, " "), " │", color = :red) + printstyled(ctx.stderr, lpad(time_str, ctx.elapsed_align + 1, " "), " │", color = ctx.nonpass_color[]) if ctx.verbose init_time_str = @sprintf("%7.2f", base.total_time - base.time) - printstyled(ctx.stderr, lpad(init_time_str, ctx.elapsed_align + 1, " "), " │ ", color = :red) + printstyled(ctx.stderr, lpad(init_time_str, ctx.elapsed_align + 1, " "), " │ ", color = ctx.nonpass_color[]) end failed_str = "failed at $(now())\n" # 11 -> 3 from " │ " 3x and 2 for each " " on either side fail_align = (11 + ctx.gc_align + ctx.percent_align + ctx.alloc_align + ctx.rss_align - textwidth(failed_str)) ÷ 2 + textwidth(failed_str) failed_str = lpad(failed_str, fail_align, " ") - printstyled(ctx.stderr, failed_str, color = :red) + printstyled(ctx.stderr, failed_str, color = ctx.nonpass_color[]) # TODO: print other stats? @@ -300,11 +301,11 @@ end function print_test_crashed(::Type{<:AbstractTestRecord}, wrkr, test, ctx::TestIOContext) lock(ctx.lock) try - printstyled(ctx.stderr, test, color = :red) + printstyled(ctx.stderr, test, color = ctx.nonpass_color[]) printstyled( ctx.stderr, lpad("($wrkr)", ctx.name_align - textwidth(test) + 1, " "), " │", - " "^ctx.elapsed_align, " crashed at $(now())\n", color = :red + " "^ctx.elapsed_align, " crashed at $(now())\n", color = ctx.nonpass_color[] ) flush(ctx.stderr) @@ -872,7 +873,9 @@ end stderr = Base.stderr, max_worker_rss = get_max_worker_rss(), serial = String[], - serial_position::Symbol = :before) + serial_position::Symbol = :before, + recycle_on_failure::Bool = false, + retries::Integer = 0) runtests(mod::Module, ARGS; ...) Run Julia tests in parallel across multiple worker processes. @@ -920,6 +923,10 @@ Several keyword arguments are also supported: testsuite; names that are valid but deselected by command-line filtering are ignored. - `serial_position`: When to run serial tests relative to the parallel batch. Must be `:before` (default) or `:after`. +- `recycle_on_failure`: Whether to recycle a worker after any test that did not pass + (default: `false`). See the Failure Handling section below. +- `retries`: How many times to re-run tests that did not pass after the main run completes + (default: `0`). See the Failure Handling section below. ## Command Line Options @@ -999,6 +1006,15 @@ runtests(MyPackage, ARGS; serial=["big_alloc_test", "huge_matrix"]) Workers are automatically recycled when they exceed memory limits to prevent out-of-memory issues during long test runs. The memory limit is set based on system architecture. + +## Failure Handling + +With `recycle_on_failure = true`, a worker is recycled after any test that did not pass, so +a test that corrupts process-wide state (e.g. wedges a GPU driver) cannot poison subsequent +tests on the same worker. + +With `retries = N` (default 0), tests that did not pass are re-run sequantially up to `N` +times after the main run completes. Only the final attempt of each test is reported. """ function runtests(mod::Module, args::ParsedArgs; testsuite::Dict{String,Expr} = find_tests(pwd()), @@ -1013,6 +1029,8 @@ function runtests(mod::Module, args::ParsedArgs; stdout = Base.stdout, stderr = Base.stderr, max_worker_rss = get_max_worker_rss(), + recycle_on_failure::Bool = false, + retries::Integer = 0, ) # # set-up @@ -1070,6 +1088,8 @@ function runtests(mod::Module, args::ParsedArgs; stdout, stderr, max_worker_rss, + recycle_on_failure, + retries, ) end @@ -1092,6 +1112,8 @@ function _runtests(mod::Module, args::ParsedArgs; stdout = Base.stdout, stderr = Base.stderr, max_worker_rss = get_max_worker_rss(), + recycle_on_failure::Bool = false, + retries::Integer = 0, ) # partition into serial and parallel groups @@ -1283,6 +1305,18 @@ function _runtests(mod::Module, args::ParsedArgs; clear_status() print_test_crashed(RecordType, wrkr, test_name, io_ctx) + + elseif msg_type === :retry + tests_n, retry_n = msg[2], msg[3] + + clear_status() + lock(io_ctx.lock) + try + printstyled(io_ctx.stdout, "Retrying $tests_n failed test$(tests_n > 1 ? "s" : " ") ($retry_n)\n", color=:white) + flush(io_ctx.stdout) + finally + unlock(io_ctx.lock) + end end end @@ -1319,6 +1353,7 @@ function _runtests(mod::Module, args::ParsedArgs; # tests_to_start = Threads.Atomic{Int}(length(tests)) + interrupted = false # After parallel-before-serial: stop extra workers so only one process is alive for # serial tests, but keep one parallel worker so we do not add a third addworker (ID_COUNTER). function drain_pool_leaving_one_worker!(pool, njobs) @@ -1340,131 +1375,144 @@ function _runtests(mod::Module, args::ParsedArgs; put!(pool, nothing) end end - try - phases = test_phases - for i in 1:length(phases) - phase_tests, sem, shared_worker = phases[i] - isempty(phase_tests) && continue - # for serial phases, reserve one pool slot for the shared worker - if !isnothing(shared_worker) - shared_worker[] = take!(worker_pool) - end - next_test = Threads.Atomic{Int}(1) - @sync for _ in eachindex(phase_tests) - push!(worker_tasks, Threads.@spawn begin - local p = nothing - acquired = false - try - Base.acquire(sem) - acquired = true - p = !isnothing(shared_worker) ? shared_worker[] : take!(worker_pool) - Threads.atomic_sub!(tests_to_start, 1) - - done[] && return - - # with multiple threads, tasks reach this point in arbitrary order, - # so pick the next test to run only now, rather than at spawn time, - # to preserve the sorted test order (issue #139) - test = phase_tests[Threads.atomic_add!(next_test, 1)] - - test_t0 = @lock running_tests begin - test_t0 = time() - running_tests[][test] = test_t0 - end + function run_test_phase(phase_tests, sem, shared_worker) + # for serial phases, reserve one pool slot for the shared worker + if !isnothing(shared_worker) + shared_worker[] = take!(worker_pool) + end - # pass in init_worker_code to custom worker function if defined - wrkr = if init_worker_code == :() - test_worker(test) - else - test_worker(test, init_worker_code) - end - if wrkr === nothing - wrkr = p - end - # if a worker failed, spawn a new one - if wrkr === nothing || !Malt.isrunning(wrkr) - wrkr = p = addworker(; init_worker_code, io_ctx.color, - exename, exeflags, env) - end + next_test = Threads.Atomic{Int}(1) + @sync for _ in eachindex(phase_tests) + push!(worker_tasks, Threads.@spawn begin + local p = nothing + acquired = false + try + Base.acquire(sem) + acquired = true + p = !isnothing(shared_worker) ? shared_worker[] : take!(worker_pool) + Threads.atomic_sub!(tests_to_start, 1) + + done[] && return + + # with multiple threads, tasks reach this point in arbitrary order, + # so pick the next test to run only now, rather than at spawn time, + # to preserve the sorted test order (issue #139) + test = phase_tests[Threads.atomic_add!(next_test, 1)] + + test_t0 = @lock running_tests begin + test_t0 = time() + running_tests[][test] = test_t0 + end - # run the test - put!(printer_channel, (:started, test, worker_id(wrkr))) - result = try - Malt.remote_eval_wait(Main, wrkr.w, :(import ParallelTestRunner)) - Malt.remote_call_fetch(invokelatest, wrkr.w, runtest, - RecordType, testsuite[test], test, - init_code, test_t0, custom_args) - catch ex - if isa(ex, InterruptException) - # the worker got interrupted, signal other tasks to stop - stop_work() - return - end + # pass in init_worker_code to custom worker function if defined + wrkr = if init_worker_code == :() + test_worker(test) + else + test_worker(test, init_worker_code) + end + if wrkr === nothing + wrkr = p + end + # if a worker failed, spawn a new one + if wrkr === nothing || !Malt.isrunning(wrkr) + wrkr = p = addworker(; init_worker_code, io_ctx.color, + exename, exeflags, env) + end - ex + # run the test + put!(printer_channel, (:started, test, worker_id(wrkr))) + result = try + Malt.remote_eval_wait(Main, wrkr.w, :(import ParallelTestRunner)) + Malt.remote_call_fetch(invokelatest, wrkr.w, runtest, + RecordType, testsuite[test], test, + init_code, test_t0, custom_args) + catch ex + if isa(ex, InterruptException) + # the worker got interrupted, signal other tasks to stop + stop_work() + return end - test_t1 = time() - output = @lock wrkr.io String(take!(wrkr.io[])) - @lock results push!(results[], (; test, result, output, test_t0, test_t1)) - - # act on the results - if result isa AbstractTestRecord - put!(printer_channel, (:finished, test, worker_id(wrkr), result)) - if anynonpass(result[]) && args.quickfail !== nothing - stop_work() - return - end - - if memory_usage(result) > max_worker_rss - # the worker has reached the max-rss limit, recycle it - # so future tests start with a smaller working set - Malt.stop(wrkr) - end - else - # One of Malt.TerminatedWorkerException, Malt.RemoteException, or ErrorException - @assert result isa Exception - put!(printer_channel, (:crashed, test, worker_id(wrkr))) - if args.quickfail !== nothing - stop_work() - return - end - # the worker encountered some serious failure, recycle it - Malt.stop(wrkr) + ex + end + test_t1 = time() + output = @lock wrkr.io String(take!(wrkr.io[])) + @lock results push!(results[], (; test, result, output, test_t0, test_t1)) + + # act on the results + if result isa AbstractTestRecord + put!(printer_channel, (:finished, test, worker_id(wrkr), result)) + if anynonpass(result[]) && args.quickfail !== nothing + stop_work() + return end - # get rid of the custom worker - if wrkr != p + if memory_usage(result) > max_worker_rss + # the worker has reached the max-rss limit, recycle it + # so future tests start with a smaller working set + Malt.stop(wrkr) + elseif recycle_on_failure && anynonpass(result[]) + # a failing test may have left the worker in a bad state + # (e.g. a wedged GPU driver whose every later allocation + # fails); recycle it so future tests get a fresh process Malt.stop(wrkr) end - - @lock running_tests begin - delete!(running_tests[], test) + else + # One of Malt.TerminatedWorkerException, Malt.RemoteException, or ErrorException + @assert result isa Exception + put!(printer_channel, (:crashed, test, worker_id(wrkr))) + if args.quickfail !== nothing + stop_work() + return end - catch ex - isa(ex, InterruptException) || rethrow() - finally - if acquired - if !isnothing(shared_worker) - shared_worker[] = p - else - # stop the worker if no more tests will need one from the pool - if tests_to_start[] == 0 && p !== nothing && Malt.isrunning(p) - Malt.stop(p) - p = nothing - end - put!(worker_pool, p) + + # the worker encountered some serious failure, recycle it + Malt.stop(wrkr) + end + + # get rid of the custom worker + if wrkr != p + Malt.stop(wrkr) + end + + @lock running_tests begin + delete!(running_tests[], test) + end + catch ex + isa(ex, InterruptException) || rethrow() + finally + if acquired + if !isnothing(shared_worker) + shared_worker[] = p + else + # stop the worker if no more tests will need one from the pool + if tests_to_start[] == 0 && p !== nothing && Malt.isrunning(p) + Malt.stop(p) + p = nothing end - Base.release(sem) + put!(worker_pool, p) end + Base.release(sem) end - end) - end - # return the serial worker to the pool for potential reuse - if !isnothing(shared_worker) - put!(worker_pool, shared_worker[]) - shared_worker[] = nothing - end + end + end) + end + + # return the serial worker to the pool for potential reuse + if !isnothing(shared_worker) + put!(worker_pool, shared_worker[]) + shared_worker[] = nothing + end + end + try + phases = test_phases + retries > 0 && (io_ctx.nonpass_color[] = :yellow) + for i in 1:length(phases) + phase_tests, sem, shared_worker = phases[i] + isempty(phase_tests) && continue + + run_test_phase(phase_tests, sem, shared_worker) + # parallel workers are not stopped while serial tests remain (tests_to_start > 0); # drain before serial-after so only one worker is alive for the serial phase if isnothing(shared_worker) && i < length(phases) @@ -1474,7 +1522,25 @@ function _runtests(mod::Module, args::ParsedArgs; end end end + + # retries + if retries > 0 && !interrupted && args.quickfail === nothing + for i in 1:retries + retries == i && (io_ctx.nonpass_color[] = :red) + retry_tests = [r.test for r in results.value + if r.result isa Exception || anynonpass(r.result[])] + isempty(retry_tests) && break + + put!(printer_channel, (:retry, length(retry_tests), i)) + sem = Base.Semaphore(1) + shared_worker = serial_worker + filter!(r -> r.test ∉ retry_tests, results.value) + + run_test_phase(retry_tests, sem, shared_worker) + end + end catch err + interrupted = true if !(err isa InterruptException) println(io_ctx.stderr, "\nCaught an error, stopping...") end diff --git a/test/runtests.jl b/test/runtests.jl index 1cb7314..ee7062a 100644 --- a/test/runtests.jl +++ b/test/runtests.jl @@ -1463,6 +1463,194 @@ end end end +@testset "recycle_on_failure" begin + # Call `_runtests` throughout, so that we can enforce a run order, and use a single job, + # so that all tests share the same pool slot: a test only gets a new worker if the + # previous one was recycled. + testsuite = Dict( + "fail1" => :( @test false ), + "pass1" => :( @test true ), + "fail2" => :( @test false ), + "pass2" => :( @test true ), + ) + tests = ["fail1", "pass1", "fail2", "pass2"] + + @testset "workers are reused across failures by default" begin + io = IOBuffer() + old_id_counter = ParallelTestRunner.ID_COUNTER[] + @test_throws Test.FallbackTestSetException begin + ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=1"]); + testsuite, + tests, + stdout=io, + stderr=io, + ) + end + str = String(take!(io)) + @test contains(str, "FAILURE") + # A failing test does not recycle its worker, so a single one runs all four tests. + @test ParallelTestRunner.ID_COUNTER[] == old_id_counter + 1 + end + + @testset "worker is recycled after a failed test" begin + io = IOBuffer() + old_id_counter = ParallelTestRunner.ID_COUNTER[] + @test_throws Test.FallbackTestSetException begin + ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=1"]); + testsuite, + tests, + stdout=io, + stderr=io, + recycle_on_failure=true, + ) + end + str = String(take!(io)) + @test contains(str, "FAILURE") + # `fail1` and `fail2` recycle their worker, so `pass1` and `pass2` each need a fresh + # one: 1 initial worker + 2 replacements. + @test ParallelTestRunner.ID_COUNTER[] == old_id_counter + 3 + end +end + +@testset "retries" begin + # A test that fails on its first attempt and passes on any subsequent one, by recording + # attempts in a file: the worker running the retry is a different process, so the marker + # has to live outside of it. + flaky_test(marker, body=:( @test true )) = quote + if isfile($marker) + $body + else + touch($marker) + @test false + end + end + + @testset "failed test passing on retry is reported as passing" begin + mktempdir() do dir + testsuite = Dict( + "flaky" => flaky_test(joinpath(dir, "flaky")), + "passes" => :( @test true ), + ) + io = IOBuffer() + @show_if_error io ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=1"]); + testsuite, + tests=["flaky", "passes"], + stdout=io, + stderr=io, + retries=1, + ) + str = String(take!(io)) + # Only the failed test is retried, and its retried result is the one reported. + @test contains(str, "Retrying 1 failed test") + @test contains(str, "SUCCESS") + # Two results in total: the failed attempt of `flaky` was replaced by the + # retried one, rather than reported next to it. + @test contains(str, r"Overall +\| +2 +2 ") + end + end + + @testset "persistent failure is retried and reported once" begin + testsuite = Dict( + "always_fails" => :( @test false ), + "passes" => :( @test true ), + ) + io = IOBuffer() + @test_throws Test.FallbackTestSetException begin + ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=1"]); + testsuite, + tests=["always_fails", "passes"], + stdout=io, + stderr=io, + retries=2, + ) + end + str = String(take!(io)) + @test contains(str, "FAILURE") + # Both retry rounds run, and each of them fails again. + @test length(collect(eachmatch(r"always_fails.*failed", str))) == 3 + # Despite the three attempts, the test is reported exactly once, as a failure. + @test contains(str, r"always_fails +\| +1 +1 ") + end + + @testset "retried test runs alone" begin + mktempdir() do dir + # On its retry, the flaky test checks it is the only worker left alive. + check_alone = quote + children = _count_child_pids($(getpid())) + if children >= 0 + @test children == 1 + end + end + testsuite = Dict( + "flaky" => flaky_test(joinpath(dir, "flaky"), check_alone), + "pass1" => :( @test true ), + "pass2" => :( @test true ), + "pass3" => :( @test true ), + ) + io = IOBuffer() + @show_if_error io ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=3"]); + testsuite, + tests=["flaky", "pass1", "pass2", "pass3"], + init_code=:(include($(joinpath(@__DIR__, "utils.jl")))), + stdout=io, + stderr=io, + retries=1, + ) + str = String(take!(io)) + @test length(collect(eachmatch(r"failed", str))) == 2 + @test contains(str, "SUCCESS") + end + end + + @testset "no retries by default" begin + mktempdir() do dir + testsuite = Dict("flaky" => flaky_test(joinpath(dir, "flaky"))) + io = IOBuffer() + @test_throws Test.FallbackTestSetException begin + ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--jobs=1"]); + testsuite, + tests=["flaky"], + stdout=io, + stderr=io, + ) + end + str = String(take!(io)) + @test !contains(str, "Retrying") + @test contains(str, "FAILURE") + end + end + + @testset "quickfail skips retries" begin + mktempdir() do dir + testsuite = Dict( + "flaky" => flaky_test(joinpath(dir, "flaky")), + "passes" => :( @test true ), + ) + io = IOBuffer() + @test_throws Test.FallbackTestSetException begin + ParallelTestRunner._runtests( + ParallelTestRunner, parse_args(["--quickfail", "--jobs=1"]); + testsuite, + tests=["flaky", "passes"], + stdout=io, + stderr=io, + retries=1, + ) + end + str = String(take!(io)) + # The run stopped early on purpose, retrying would defeat that. + @test !contains(str, "Retrying") + @test contains(str, "FAILURE") + end + end +end + # This testset should always be the last one, don't add anything after this. # We want to make sure there are no running workers at the end of the tests. @testset "no workers running" begin