From ee85216fc926faeb377694a7b9447dffc1ee0f08 Mon Sep 17 00:00:00 2001 From: Christian Guinard <28689358+christiangnrd@users.noreply.github.com> Date: Fri, 7 Aug 2026 17:16:19 -0300 Subject: [PATCH] Make test phase run into function --- src/ParallelTestRunner.jl | 231 ++++++++++++++++++++------------------ 1 file changed, 119 insertions(+), 112 deletions(-) diff --git a/src/ParallelTestRunner.jl b/src/ParallelTestRunner.jl index 8baee3e..7c1723e 100644 --- a/src/ParallelTestRunner.jl +++ b/src/ParallelTestRunner.jl @@ -1340,131 +1340,138 @@ 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) 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 + 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)