From aa4d8604d8a6336a5438711e26d32d50e77e7226 Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 15 Sep 2026 12:55:48 +0000 Subject: [PATCH 1/7] fix: run_loop! abandoned keys a killed sibling still held, and said nothing A job that follows one the wall clock killed cannot complete the campaign, and does not report that it failed to. `acquire_running!` already reclaims a lock whose heartbeat is older than `stale_after`, so the stack CAN recover from a master killed mid-key: something just has to keep trying until the lock expires. `run_loop!` stops before that. Its exit test is `max_empty_rounds` rounds with `done == 0`, and a round where every remaining key came back `:lock_busy` counts as empty. With the defaults that is 90 s against a `stale_after` of 600 s. Measured, the tail of a campaign: one key left, held by a `.running` whose owner was killed. | | before | after | |---|---|---| | run_loop! returned after | 5.3 s | 62.1 s | | stale_after | 60 s | 60 s | | work_fn ran on the key | no | YES | | key done at exit | NO | yes | It gave up 55 s before the lock could be reclaimed, left the key undone, and returned a value with no field that could have said so. A busy round no longer counts toward `max_empty_rounds` until `stale_after` has been waited out. That is exactly the quantity that separates the two cases the caller cannot otherwise tell apart: past it the holder is gone and the lock is reclaimable, before it a live sibling is doing the work. `busy` is now in the return, so "everything is finished" is distinguishable from "someone else still has work out". The wait is bounded by `stale_after + 2 * idle_sleep`, so a sibling that keeps refreshing forever is not waited on forever. The loop is not idle during it: every round still calls `run!`, so any key that frees up is taken immediately, and `stop_flag` / `deadline` are checked each round. The cost is holding an allocation up to 11 min instead of 90 s when a LIVE sibling holds the remainder. That is the right side to err on: exiting early and having the sibling then die leaves those keys with nobody to finish them. Three testsets, one per direction. The dead-holder case completes the key and takes longer than `stale_after`. The control, with nothing held, still exits well under `stale_after`, so this is not "always wait". The live-sibling case never steals the key, returns within the budget, and reports `busy > 0`. Verified locally: all 23 files under test/, green. Co-Authored-By: Claude Opus 5 (1M context) --- Project.toml | 2 +- src/Run.jl | 41 ++++++++++-- test/run/test_run_loop_busy.jl | 112 +++++++++++++++++++++++++++++++++ 3 files changed, 150 insertions(+), 5 deletions(-) create mode 100644 test/run/test_run_loop_busy.jl diff --git a/Project.toml b/Project.toml index a4ee99b..675d71c 100644 --- a/Project.toml +++ b/Project.toml @@ -1,6 +1,6 @@ name = "SweepRunner" uuid = "be946ad2-3cb3-4b6e-8f7e-4a5ecc3c255b" -version = "0.6.2" +version = "0.6.3" authors = ["sota shimozono "] [deps] diff --git a/src/Run.jl b/src/Run.jl index b2134d4..c282f71 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -713,9 +713,17 @@ more work to do. This is the infra equivalent of FiniteTemperature.jl's `_work_loop` driver. The loop exits when: -- `max_empty_rounds` consecutive rounds produce zero new completions, or +- `max_empty_rounds` consecutive rounds produce zero new completions AND leave nothing held by a + sibling, or - `opts.stop_flag` is raised, or `opts.deadline` has passed. +A round that completes nothing but finds keys `:lock_busy` does NOT count toward +`max_empty_rounds` until `opts.stale_after` has been waited out. Those keys are either being +worked on by a live sibling, or held by one the wall clock killed, and `stale_after` is what +separates the two: past it, `acquire_running!` reclaims the lock on the next attempt. Returning +before then leaves the campaign short and reports nothing, because `max_empty_rounds * +idle_sleep` (90 s by default) is an order of magnitude under `stale_after` (600 s). + Default parameters (`max_empty_rounds=3`, `idle_sleep=30.0`) are the battle-tested values from FiniteTemperature.jl. @@ -747,8 +755,9 @@ dependents", not a DAG. `affinity` is forwarded verbatim to every [`run!`](@ref) call. -Returns `(; ran, rounds, done, stopped_by, prerequisite)`. `ran` is `false` exactly when a -prerequisite blocked the stage. +Returns `(; ran, rounds, done, busy, stopped_by, prerequisite)`. `busy` is how many keys the last +round found held by a sibling, so a caller can tell "everything is done" from "someone else still +has work out". `ran` is `false` exactly when a prerequisite blocked the stage. """ function run_loop!( work_fn::Function, @@ -765,13 +774,24 @@ function run_loop!( if prerequisite !== nothing pre = run_prerequisite!(prerequisite; opts=opts, load=load, poll=idle_sleep) pre.complete || return (; - ran=false, rounds=0, done=0, stopped_by=pre.stopped_by, prerequisite=pre + ran=false, + rounds=0, + done=0, + busy=0, + stopped_by=pre.stopped_by, + prerequisite=pre, ) end empty_count = 0 rounds = 0 n_done = 0 + n_busy = 0 + busy_waited = 0.0 + # A lock is reclaimable once its heartbeat is `stale_after` old, so waiting that long is what + # separates "a sibling is working on it" from "the holder is gone". The margin covers the round + # that has to follow the expiry to act on it. + busy_budget = opts.stale_after + 2 * idle_sleep while true if _is_stopped(opts) break @@ -779,8 +799,20 @@ function run_loop!( rounds += 1 result = run!(work_fn, vault, keys; opts=opts, load=load, affinity=affinity) n_done += result.done + n_busy = result.busy if result.done > 0 empty_count = 0 + busy_waited = 0.0 + continue + end + # A round that completed nothing but found keys held by a SIBLING is not an empty round: + # either that sibling finishes them, or it is dead and `acquire_running!` reclaims them + # once its heartbeat passes `stale_after`. Counting it as empty is what made a follow-on + # job return after `max_empty_rounds * idle_sleep` while the locks stayed held for + # `stale_after`, leaving the campaign short and saying nothing. + if result.busy > 0 && busy_waited < busy_budget + busy_waited += idle_sleep + sleep(idle_sleep) continue end empty_count += 1 @@ -793,6 +825,7 @@ function run_loop!( ran=true, rounds=rounds, done=n_done, + busy=n_busy, stopped_by=_stop_reason(opts), prerequisite=pre, ) diff --git a/test/run/test_run_loop_busy.jl b/test/run/test_run_loop_busy.jl new file mode 100644 index 0000000..13f3f0c --- /dev/null +++ b/test/run/test_run_loop_busy.jl @@ -0,0 +1,112 @@ +# run_loop! and keys held by a sibling (#..): a round that completes nothing but finds `:lock_busy` +# is not an empty round until `stale_after` has been waited out. +# +# The bug this pins: `max_empty_rounds * idle_sleep` is 90 s by default and `stale_after` is 600 s, +# so a job following one the wall clock killed returned 8.5 minutes before the dead job's locks +# became reclaimable, left the campaign short, and reported nothing. + +using SweepRunner, Test, DataVault, ParamIO + +const _BUSY_CFG = joinpath(@__DIR__, "fixtures", "study.toml") + +# Fast constants so the wait is seconds. `heartbeat_interval` must stay under `stale_after`. +_opts(; kw...) = RunOpts(; + workers=:sequential, stale_after=3.0, heartbeat_interval=1.0, kw... +) +const _BUDGET = 3.0 + 2 * 0.5 # opts.stale_after + 2 * idle_sleep + +function with_tail(f) + outdir = mktempdir() + try + v = DataVault.Vault(_BUSY_CFG; run="tail", outdir=outdir) + ks = DataVault.keys(v) + for k in ks[2:end] # everything done but the first + DataVault.save!(v, k, Dict("x" => 1)) + DataVault.mark_done!(v, k) + end + f(v, ks) + finally + rm(outdir; recursive=true, force=true) + end +end + +@testset "run_loop!: a dead holder's key is waited out and completed" begin + with_tail() do v, ks + DataVault.mark_running!(v, ks[1]) # killed mid-key: the marker outlived the process + n = Ref(0) + t0 = time() + r = run_loop!( + k -> (n[] += 1; Dict{String,Any}("x" => 1)), + v, + ks; + opts=_opts(), + max_empty_rounds=2, + idle_sleep=0.5, + ) + elapsed = time() - t0 + + @test n[] == 1 # it ran the key + @test DataVault.is_done(v, ks[1]) # and the campaign is complete + @test r.done == 1 + @test elapsed >= 3.0 # it waited stale_after out rather than giving up + end +end + +@testset "run_loop!: with nothing held, it still exits promptly" begin + # The control. Without this the fix would read as "always wait stale_after", which would make + # every finished campaign pay 600 s to notice it is finished. + with_tail() do v, ks + DataVault.save!(v, ks[1], Dict("x" => 1)) + DataVault.mark_done!(v, ks[1]) # everything done, nothing running + + t0 = time() + r = run_loop!( + k -> Dict{String,Any}("x" => 1), + v, + ks; + opts=_opts(), + max_empty_rounds=2, + idle_sleep=0.5, + ) + elapsed = time() - t0 + + @test r.done == 0 + @test r.busy == 0 + @test elapsed < 3.0 # well under stale_after + end +end + +@testset "run_loop!: a LIVE sibling is not waited on forever" begin + # The other side of the bound. A holder that keeps refreshing is alive and doing the work, so + # this master has nothing to contribute and must return rather than spin. + with_tail() do v, ks + DataVault.mark_running!(v, ks[1]) + alive = Threads.Atomic{Bool}(true) + beater = Threads.@spawn while alive[] + DataVault.touch_running!(v, ks[1]) # the sibling's heartbeat, never going stale + sleep(0.3) + end + try + n = Ref(0) + t0 = time() + r = run_loop!( + k -> (n[] += 1; Dict{String,Any}("x" => 1)), + v, + ks; + opts=_opts(), + max_empty_rounds=2, + idle_sleep=0.5, + ) + elapsed = time() - t0 + + @test n[] == 0 # never stole the live sibling's key + @test !DataVault.is_done(v, ks[1]) + @test r.busy > 0 # and SAYS it left work held, which it could not before + @test elapsed >= _BUDGET # waited the budget + @test elapsed < _BUDGET + 10 # then returned + finally + alive[] = false + wait(beater) + end + end +end From 8fae397456ecf5429e6caecb52bd423ba536471e Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 15 Sep 2026 13:19:59 +0000 Subject: [PATCH 2/7] feat: ask whether the holder is gone, instead of waiting to find out Builds on the previous commit, which made `run_loop!` wait `stale_after` out rather than abandon keys a killed sibling held. That is correct on its own and stays the fallback. This removes the WAIT, where the answer can be had. holder_liveness(owner) -> :alive | :dead | :unknown Measured: a key whose holder is a dead pid on this host now completes in under 30 s with `stale_after = 600.0`. Before, nothing could distinguish it from a live holder for ten minutes. `:dead` is returned only on positive evidence. Every error path, every unparseable token, every holder that cannot be asked about is `:unknown`, and `stale_after` then decides exactly as before. A false `:dead` hands a live master's key to someone else, which is the one outcome the lock exists to prevent. Two sources: - `/proc/` when the holder is on THIS host. A recycled pid reads as `:alive`, the safe side. - `squeue` when the holder carries a Slurm job id AND this process is itself inside an allocation. The Slurm guard is not defensive programming, it is a measurement. On the development box behind this package, `/home/souta/.local/bin/squeue` is a wrapper that answers about a REMOTE cluster's queue, and it lists live jobs there while this process has no SLURM_* variables at all. Any job id from anywhere else is absent from that queue and would read as `:dead`. Requiring that we are ourselves a Slurm job means the queue being consulted is demonstrably the one that would have run the holder. The queue is read whole and tested for membership rather than asked with `squeue -j `, because `-j` exits non-zero BOTH for a finished job and for a broken scheduler, and those must not collapse into one answer. The result is cached for 15 s: this is consulted per contended key and `squeue` is a scheduler RPC. SweepRunner now stamps an owner into `.running` and uses the owner-aware `refresh_running!` / `clear_running!`. It had been calling the two-argument forms, so DataVault 0.8.1's owner-stamped locks were unused here. The heartbeat's loss detection is no longer best-effort: the previous existence-based refresh returned `true` against the reclaimer's own file. The comment in `Run.jl` saying a race-free guarantee needed owner-stamped locks in DataVault is resolved rather than moved. DataVault's compat floor moves to 0.8.1 for it. ## A defect this surfaced `_run_affinity!` gave a key back to the queue on `ProcessExitedException` with NO bound, where the `pmap` path it replaced uses `retry_delays=ExponentialBackOff(; n=2)`. Unbounded, a key that kills whoever takes it is handed to worker after worker forever. It was invisible while a dead holder's lock sat until `stale_after`, because only one worker died per round; reclaiming immediately turned it into an immediate cascade. Bounded to the same two give-backs, after which the key is logged `:gave_up` and reported `:error`. Measured on a key that kills every worker: 3 deaths (one dispatch plus two give-backs), `err = 1`, and the other 23 keys still complete. The worker-death test modelled a POISON key, which is not what the feature is for. It now models PREEMPTION: a fuse file makes the death happen exactly once, and the assertion is that the next worker completes the key inside the same run. The poison case is its own testset, asserting the bound. Verified locally: all 24 files under test/, green. Co-Authored-By: Claude Opus 5 (1M context) --- Project.toml | 5 +- docs/src/api.md | 7 ++ src/EventLog.jl | 1 + src/Liveness.jl | 113 ++++++++++++++++++++++++++++++++ src/Run.jl | 70 +++++++++++++++----- src/SweepRunner.jl | 1 + test/run/test_affinity.jl | 79 ++++++++++++++++++----- test/run/test_liveness.jl | 114 +++++++++++++++++++++++++++++++++ test/run/test_run_loop_busy.jl | 6 +- 9 files changed, 361 insertions(+), 35 deletions(-) create mode 100644 src/Liveness.jl create mode 100644 test/run/test_liveness.jl diff --git a/Project.toml b/Project.toml index 675d71c..9ed006e 100644 --- a/Project.toml +++ b/Project.toml @@ -19,7 +19,10 @@ TOML = "fa267f1f-6049-4f14-aa54-33bafae1ed76" [compat] Aqua = "0.8" -DataVault = "0.7, 0.8" +# 0.8.1 is the floor, not a preference: the owner-stamped `.running` API +# (`new_owner_token`, `running_owner`, and the three-argument `refresh_running!` / +# `clear_running!`) arrived there, and the liveness reaper is built on it. +DataVault = "0.8.1" Dates = "1.11" Distributed = "1.11" JLD2 = "0.6" diff --git a/docs/src/api.md b/docs/src/api.md index 6bb9fc9..881989a 100644 --- a/docs/src/api.md +++ b/docs/src/api.md @@ -49,6 +49,13 @@ SweepRunner.run! SweepRunner.run_loop! ``` +## Liveness + +```@docs +SweepRunner.owner_token +SweepRunner.holder_liveness +``` + ## Prerequisite ```@docs diff --git a/src/EventLog.jl b/src/EventLog.jl index dddf52a..228fe9b 100644 --- a/src/EventLog.jl +++ b/src/EventLog.jl @@ -44,6 +44,7 @@ stay within that guarantee. | `key_done` | after a successful `work_fn(key)` (includes `secs`, `attempt`) | | `lock_busy` | another master holds the `.running` lock (acquire = `:busy`) | | `lock_lost` | our lock was reclaimed mid-work; result discarded (no double-run) | +| `lock_reaped` | a lock whose holder was shown dead was cleared without waiting | | `lock_reclaimed`| (reserved, not currently emitted) | | `error` | `work_fn` threw on this attempt | | `retry` | another attempt will follow | diff --git a/src/Liveness.jl b/src/Liveness.jl new file mode 100644 index 0000000..21c147b --- /dev/null +++ b/src/Liveness.jl @@ -0,0 +1,113 @@ +# Liveness — is the master that holds this lock still running? +# +# `stale_after` answers that by waiting. It has to, because a heartbeat can only ever say +# "recently alive"; the absence of one is indistinguishable from a slow filesystem until enough +# time has passed. Where the question CAN be asked outright, asking is worth several minutes of a +# batch allocation. +# +# Nothing here is required for correctness: every path degrades to `:unknown`, and `stale_after` +# then decides exactly as before. `:slurm` is one of four worker modes, so a mechanism that only +# worked there would cover a quarter of the runs. + +using Dates + +""" + owner_token() -> String + +This master's identity, stamped into `.running` by [`run!`](@ref) so a sibling can ask about it +later. `host:pid:nonce` from DataVault, with `:slurm` appended inside a Slurm allocation. + +The Slurm field is what makes the question answerable ACROSS hosts: `/proc` only works for a +holder on this machine, and the master of the job that was killed is usually somewhere else. +""" +function owner_token()::String + base = DataVault.new_owner_token() + job = get(ENV, "SLURM_JOB_ID", "") + return isempty(job) ? base : string(base, ":slurm", job) +end + +""" + holder_liveness(owner) -> Symbol + +`:alive`, `:dead`, or `:unknown` for the holder named by an [`owner_token`](@ref). + +`:dead` is only ever returned on positive evidence that the process is gone. Everything else is +`:unknown`, including every error path: a wrong `:dead` would hand a live master's key to someone +else, which is the one outcome the lock exists to prevent. + +Two sources, in order: + +1. a Slurm job id, when this process is ITSELF inside a Slurm allocation. The job is absent from + the queue, so it has finished, been cancelled, or hit its wall clock. +2. the pid, when the holder is on THIS host and `/proc` exists. A recycled pid reads as `:alive`, + which is the safe direction. +""" +function holder_liveness(owner::AbstractString)::Symbol + parts = split(String(owner), ':') + length(parts) >= 3 || return :unknown + + job = findfirst(p -> startswith(p, "slurm"), parts) + if job !== nothing + s = _slurm_liveness(parts[job][6:end]) + s === :unknown || return s + end + + parts[1] == gethostname() || return :unknown + pid = tryparse(Int, parts[2]) + pid === nothing && return :unknown + return _pid_liveness(pid) +end + +# `/proc/` is the whole check on Linux. Elsewhere there is no equally cheap answer that does +# not risk a false `:dead`, so there is no answer. +function _pid_liveness(pid::Int)::Symbol + Sys.islinux() || return :unknown + return isdir("/proc/$(pid)") ? :alive : :dead +end + +# The live-job set, cached: `squeue` is a scheduler RPC and this is consulted per contended key. +const _SQUEUE_TTL = 15.0 +const _squeue_cache = Ref{Tuple{Float64,Union{Set{String},Nothing}}}((-Inf, nothing)) +const _squeue_lock = ReentrantLock() + +function _slurm_liveness(jobid::AbstractString)::Symbol + isempty(jobid) && return :unknown + # Ask Slurm only from INSIDE a Slurm allocation, so the queue being consulted is demonstrably + # the one that would have run the holder. `squeue` existing proves nothing: measured on the + # development box behind this package, `/home/…/.local/bin/squeue` is a wrapper that answers + # about a REMOTE cluster's queue, and every job id from anywhere else reads as absent there. + # Absent would then mean `:dead`, and a false `:dead` hands a live master's key to someone + # else, which is the one outcome the lock exists to prevent. + haskey(ENV, "SLURM_JOB_ID") || return :unknown + live = _live_slurm_jobs() + live === nothing && return :unknown + # An array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so the whole job is + # alive if any of its tasks is. + jobid in live && return :alive + any(j -> startswith(j, jobid * "_"), live) && return :alive + return :dead +end + +# `nothing` means "no opinion": squeue is absent, or it failed. Listing the queue and testing +# membership is deliberate. `squeue -j ` exits non-zero BOTH for a finished job and for a +# broken scheduler, and those must not collapse into the same answer. +function _live_slurm_jobs()::Union{Set{String},Nothing} + return lock(_squeue_lock) do + t, cached = _squeue_cache[] + time() - t < _SQUEUE_TTL && return cached + fresh = try + if Sys.which("squeue") === nothing + nothing + else + out = read(pipeline(`squeue -h -o %i`; stderr=devnull), String) + Set(String.(split(out; keepempty=false))) + end + catch + nothing + end + _squeue_cache[] = (time(), fresh) + return fresh + end +end + +export owner_token, holder_liveness diff --git a/src/Run.jl b/src/Run.jl index c282f71..b600b0e 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -342,6 +342,21 @@ function run!( ) end +# Clear a `.running` whose holder is provably gone, so the key is retriable NOW rather than in +# `stale_after`. Returns whether anything was cleared. +# +# Only `:dead` acts. `:unknown` is the common answer (a holder on another host with no Slurm id) +# and leaves the timeout to decide, exactly as before. +function _reap_if_dead!(vault::Vault, key::DataKey, stage::Symbol, log::EventLog)::Bool + DataVault.is_running(vault, key) || return false + owner = DataVault.running_owner(vault, key) + owner === nothing && return false # unstamped: cannot be attributed, so cannot be judged + holder_liveness(owner) === :dead || return false + cleared = DataVault.clear_running!(vault, key, owner) + cleared && log_event(log, :lock_reaped; stage=stage, key=canonical(key), owner=owner) + return cleared +end + """ _run_one_with_lock!(work_fn, vault, key, stage, log, opts) -> (DataKey, Symbol) @@ -375,10 +390,15 @@ function _run_one_with_lock!( return (key, :stop) end + # A lock whose holder can be SHOWN to be gone does not have to wait out `stale_after`. The + # clear is owner-checked, so it is a no-op if the holder changed since the question was asked. + _reap_if_dead!(vault, key, stage, log) + # DataVault owns the lock file. `acquire_running!` is atomic on # NFS via POSIX `link()`: concurrent masters see at most one # `:ok` / `:reclaimed`; the losers see `:busy`. - acq = DataVault.acquire_running!(vault, key; stale_after=opts.stale_after) + tok = owner_token() + acq = DataVault.acquire_running!(vault, key, tok; stale_after=opts.stale_after) if acq === :busy log_event(log, :lock_busy; level=:debug, stage=stage, key=kstr) return (key, :lock_busy) @@ -393,7 +413,7 @@ function _run_one_with_lock!( # Re-check after acquisition: another master may have finished this # key between our manifest read and our acquire. if DataVault.is_done(vault, key) - DataVault.clear_running!(vault, key) + DataVault.clear_running!(vault, key, tok) return (key, :already_done) end @@ -403,14 +423,12 @@ function _run_one_with_lock!( # the stop signal thread-safe; a short sleep tick keeps finally # cleanup responsive (the earlier fixed 60-s sleep would block the # whole shutdown until the next heartbeat tick). - # `lost[]` is raised by the heartbeat task if it observes that a sibling - # reclaimed our lock (we stalled past `stale_after`). `_run_one_with_retry!` - # checks it before `save!`, and the `finally` below skips `clear_running!` - # when lost, so we neither commit on top of nor delete the lock now owned by - # the reclaiming master. Detection is BEST-EFFORT: `refresh_running!` is - # existence-based, so a reclaim is reliably caught only via stale-reclaim or - # the brief window the lock file is absent. A fully race-free guarantee needs - # owner-stamped locks in DataVault (see PR notes). + # `lost[]` is raised by the heartbeat task if it observes that a sibling reclaimed our lock (we + # stalled past `stale_after`). `_run_one_with_retry!` checks it before `save!`, and the + # `finally` below skips `clear_running!` when lost, so we neither commit on top of nor delete + # the lock now owned by the reclaiming master. The refresh is OWNER-CHECKED (DataVault 0.8.1), + # so a reclaim that has already happened is seen on the next beat; the previous + # existence-based form returned `true` against the reclaimer's own file. hb_stop = Threads.Atomic{Bool}(false) lost = Threads.Atomic{Bool}(false) hb_task = Threads.@spawn begin @@ -425,7 +443,7 @@ function _run_one_with_lock!( # lock, e.g. NFS hiccup) both mean "treat as lost": stop # heartbeating and signal it, rather than silently dying. alive = try - DataVault.refresh_running!(vault, key) + DataVault.refresh_running!(vault, key, tok) catch false end @@ -452,7 +470,7 @@ function _run_one_with_lock!( # `clear_running!` is owner-blind (`isfile && rm`), so clearing it would # delete THEIR lock and re-open double-execution. On `:ok`, `mark_done!` # already removed our `.running`; `clear_running!` is otherwise idempotent. - lost[] || DataVault.clear_running!(vault, key) + lost[] || DataVault.clear_running!(vault, key, tok) end return (key, outcome) @@ -596,8 +614,17 @@ function _run_affinity!( end end - _give_back!(i::Int) = lock(q) do - return push!(get!(Vector{Int}, by_group, groups[i]), i) + # `pmap` bounds its own `ProcessExitedException` re-dispatch with + # `retry_delays=ExponentialBackOff(; n=2)`; this dispatcher has to bound it too. Unbounded, a + # key that reliably kills whoever takes it is handed to worker after worker forever, and the + # faster a dead holder's lock is reclaimed the faster that cascade runs. + const_giveback_limit = 2 + givebacks = zeros(Int, length(todo)) + _give_back!(i::Int)::Bool = lock(q) do + givebacks[i] += 1 + givebacks[i] > const_giveback_limit && return false + push!(get!(Vector{Int}, by_group, groups[i]), i) + return true end @sync for pid in workers() @@ -611,7 +638,20 @@ function _run_affinity!( ) catch e if e isa ProcessExitedException - _give_back!(i) + # Requeued, or out of attempts: a key that has taken down `const_giveback_limit` + # workers is reported rather than handed to the next one. + if !_give_back!(i) + log_event( + log, + :gave_up; + stage=stage, + key=canonical(key), + attempts=const_giveback_limit + 1, + err="worker exited on this key every time it was dispatched", + ) + out[i] = (key, :error) + filled[i] = true + end break end log_event( diff --git a/src/SweepRunner.jl b/src/SweepRunner.jl index 79c98ab..06fbc2b 100644 --- a/src/SweepRunner.jl +++ b/src/SweepRunner.jl @@ -68,6 +68,7 @@ include("AtomicIO.jl") include("EventLog.jl") include("Manifest.jl") include("InitWorkers.jl") +include("Liveness.jl") include("Run.jl") include("Prerequisite.jl") include("Preflight.jl") diff --git a/test/run/test_affinity.jl b/test/run/test_affinity.jl index 5f163bc..969b16e 100644 --- a/test/run/test_affinity.jl +++ b/test/run/test_affinity.jl @@ -155,38 +155,85 @@ end end @testset "affinity: a worker that dies hands its key back instead of losing it" begin - # The fault tolerance `pmap` gives through `retry_check`, which this dispatcher has to - # reproduce by hand. One key kills its worker outright with `_exit`, so remotecall_fetch sees - # ProcessExitedException rather than a caught error. + # The fault tolerance `pmap` gives through `retry_check`, which this dispatcher reproduces by + # hand. The model is PREEMPTION: the wall clock takes one worker out mid-key, and the key must + # survive that. A marker file makes the death happen exactly once, so the next worker to take + # the key completes it, which is what a preempted key does. with_workers(3) do outdir = mktempdir() try v = DataVault.Vault(_AFF_CFG; run="died", outdir=outdir) ks = DataVault.keys(v) victim = ParamIO.canonical(ks[1]) + fuse = joinpath(outdir, "fuse") work = k -> begin - if ParamIO.canonical(k) == victim + if ParamIO.canonical(k) == victim && !isfile(fuse) + touch(fuse) ccall(:_exit, Cvoid, (Cint,), 1) end sleep(0.01) return Dict{String,Any}("x" => 1) end - # stale_after short enough that the dead worker's abandoned .running is reclaimable on - # the next pass; without that the key stays :busy for its whole duration. - o = RunOpts(; stale_after=2.0, heartbeat_interval=1.0) - r = run!(work, v, ks; opts=o, affinity=group_of) + r = run!( + work, + v, + ks; + opts=RunOpts(; stale_after=2.0, heartbeat_interval=1.0), + affinity=group_of, + ) - # It returned rather than hanging, and the loss was not counted as an error. - @test r.done == length(ks) - 1 + # It returned rather than hanging, and the preempted key was re-dispatched and + # finished inside the SAME run: the dead worker's lock carries its owner, and that + # owner is a pid on this host that no longer exists. + @test isfile(fuse) # the death really happened + @test DataVault.is_done(v, ks[1]) + @test r.done == length(ks) @test r.err == 0 - @test !DataVault.is_done(v, ks[1]) + finally + rm(outdir; recursive=true, force=true) + end + end +end + +@testset "affinity: a key that kills its worker is re-dispatched a BOUNDED number of times" begin + # `pmap` bounds its own ProcessExitedException re-dispatch (`ExponentialBackOff(; n=2)`); this + # dispatcher has to as well. Unbounded, a key that kills whoever takes it is handed to worker + # after worker with no limit, and reclaiming a dead holder's lock quickly, which this branch + # adds, only makes that spin faster. + # + # Four workers so the limit binds before the pool is exhausted: one initial dispatch plus two + # give-backs is three deaths, leaving a live worker to observe the give-up. + with_workers(4) do + outdir = mktempdir() + try + v = DataVault.Vault(_AFF_CFG; run="poison", outdir=outdir) + ks = DataVault.keys(v) + poison = ParamIO.canonical(ks[1]) + deaths = joinpath(outdir, "deaths") + mkpath(deaths) + work = k -> begin + if ParamIO.canonical(k) == poison + touch(joinpath(deaths, "d$(Distributed.myid())")) + ccall(:_exit, Cvoid, (Cint,), 1) + end + return Dict{String,Any}("x" => 1) + end + + r = run!( + work, + v, + ks; + opts=RunOpts(; stale_after=2.0, heartbeat_interval=1.0), + affinity=group_of, + ) - # And the key is retriable: a later pass, once the abandoned lock is stale, gets it. - sleep(2.2) - r2 = run!(k -> Dict{String,Any}("x" => 1), v, ks; opts=o, affinity=group_of) - @test r2.done == 1 - @test all(k -> DataVault.is_done(v, k), ks) + ndeaths = length(readdir(deaths)) + @info "poison key" ndeaths err = r.err done = r.done + @test ndeaths >= 1 # the fixture really does kill workers + @test ndeaths <= 3 # 1 dispatch + 2 give-backs, and no more + @test !DataVault.is_done(v, ks[1]) # it cannot complete, and did not + @test r.done == length(ks) - 1 # every OTHER key still finished finally rm(outdir; recursive=true, force=true) end diff --git a/test/run/test_liveness.jl b/test/run/test_liveness.jl new file mode 100644 index 0000000..216cd5c --- /dev/null +++ b/test/run/test_liveness.jl @@ -0,0 +1,114 @@ +# holder_liveness: asking whether the holder is gone, instead of waiting to find out. +# +# `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only +# where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a +# live master's key to someone else, which is what the lock exists to prevent. + +using SweepRunner, Test, DataVault, ParamIO + +const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") +const _HOST = gethostname() +const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running + +@testset "holder_liveness: :dead only on positive evidence" begin + @test holder_liveness(owner_token()) === :alive # this very process + @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host + + # Everything uncertain is :unknown, so stale_after decides exactly as before. + @test holder_liveness("otherhost:123:abcd") === :unknown # another machine + @test holder_liveness("garbage") === :unknown + @test holder_liveness("") === :unknown + @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown +end + +@testset "holder_liveness: Slurm is only asked from inside an allocation" begin + # `squeue` existing proves nothing about WHICH queue it answers for. On the machine this was + # written on it is a wrapper around a remote cluster, where every foreign job id is absent and + # would read as :dead. + tok = "otherhost:1:abcd:slurm999999" + withenv("SLURM_JOB_ID" => nothing) do + @test holder_liveness(tok) === :unknown + end +end + +@testset "owner_token: carries the Slurm job so another HOST can ask" begin + withenv("SLURM_JOB_ID" => "4242") do + t = owner_token() + @test occursin(":slurm4242", t) + @test startswith(t, gethostname() * ":") + end + withenv("SLURM_JOB_ID" => nothing) do + @test !occursin("slurm", owner_token()) + end +end + +function with_one_left(f) + outdir = mktempdir() + try + v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) + ks = DataVault.keys(v) + for k in ks[2:end] + DataVault.save!(v, k, Dict("x" => 1)) + DataVault.mark_done!(v, k) + end + f(v, ks) + finally + rm(outdir; recursive=true, force=true) + end +end + +@testset "a dead holder's lock is reaped at once, not after stale_after" begin + with_one_left() do v, ks + @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + @test DataVault.running_owner(v, ks[1]) == _DEAD + + t0 = time() + r = run_loop!( + k -> Dict{String,Any}("x" => 1), + v, + ks; + opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), + max_empty_rounds=2, + idle_sleep=0.5, + ) + elapsed = time() - t0 + + @test DataVault.is_done(v, ks[1]) # completed + @test r.done == 1 + @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after + end +end + +@testset "a LIVE holder's lock is not reaped" begin + # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a + # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in + # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. + with_one_left() do v, ks + log = EventLog(joinpath(v.outdir, "e.jsonl")) + mine = owner_token() # this process: provably alive + @test DataVault.acquire_running!(v, ks[1], mine) === :ok + + @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false + @test DataVault.is_running(v, ks[1]) + @test DataVault.running_owner(v, ks[1]) == mine # untouched + + # And the same reaper DOES clear it once the owner is one that cannot be running, so the + # `false` above is the liveness answer and not an inert function. + DataVault.clear_running!(v, ks[1], mine) + @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true + @test !DataVault.is_running(v, ks[1]) + end +end + +@testset "an unstamped lock is left to the timeout" begin + # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. + with_one_left() do v, ks + DataVault.mark_running!(v, ks[1]) + @test DataVault.running_owner(v, ks[1]) === nothing + @test SweepRunner._reap_if_dead!( + v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) + ) == false + @test DataVault.is_running(v, ks[1]) + end +end diff --git a/test/run/test_run_loop_busy.jl b/test/run/test_run_loop_busy.jl index 13f3f0c..75e8500 100644 --- a/test/run/test_run_loop_busy.jl +++ b/test/run/test_run_loop_busy.jl @@ -10,9 +10,9 @@ using SweepRunner, Test, DataVault, ParamIO const _BUSY_CFG = joinpath(@__DIR__, "fixtures", "study.toml") # Fast constants so the wait is seconds. `heartbeat_interval` must stay under `stale_after`. -_opts(; kw...) = RunOpts(; - workers=:sequential, stale_after=3.0, heartbeat_interval=1.0, kw... -) +function _opts(; kw...) + return RunOpts(; workers=:sequential, stale_after=3.0, heartbeat_interval=1.0, kw...) +end const _BUDGET = 3.0 + 2 * 0.5 # opts.stale_after + 2 * idle_sleep function with_tail(f) From 9f486543835a9c0007df0679f041fd7377d753d5 Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 15 Sep 2026 13:32:24 +0000 Subject: [PATCH 3/7] test: cover the Slurm queue decision, which the guard makes unreachable outside an allocation codecov flagged the patch at 77.14% against a target of 87.01%. The gap was the whole Slurm path: `holder_liveness` returns early unless the process is itself inside an allocation, so nothing below that guard ran. The membership rule is the part with real content and it was pure assertion until now: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is alive if ANY of its tasks is. Only the squeue CALL is stubbed, by priming the cache; the decision under test is the package's own. A stubbed empty queue plus no allocation also pins that the guard is what produces `:unknown`, not the absence of jobs. One test forces the cache miss so the real call runs. Its assertion is the safety contract rather than a value: this executes with no scheduler (CI), with a `squeue` that answers about a remote cluster (the development box), and inside a real allocation, and in all three it has to come back with an answer rather than an exception. Each stub restores `_squeue_cache` in a `finally`, so no stubbed queue leaks into another file on the same shard. Co-Authored-By: Claude Opus 5 (1M context) --- test/run/test_liveness.jl | 61 +++++++- test/run/test_liveness.jl.3860359.cov | 114 +++++++++++++++ test/run/test_liveness.jl.3861352.cov | 141 ++++++++++++++++++ test/run/test_liveness.jl.3861939.cov | 161 +++++++++++++++++++++ test/run/test_run_loop_busy.jl.3860359.cov | 112 ++++++++++++++ 5 files changed, 582 insertions(+), 7 deletions(-) create mode 100644 test/run/test_liveness.jl.3860359.cov create mode 100644 test/run/test_liveness.jl.3861352.cov create mode 100644 test/run/test_liveness.jl.3861939.cov create mode 100644 test/run/test_run_loop_busy.jl.3860359.cov diff --git a/test/run/test_liveness.jl b/test/run/test_liveness.jl index 216cd5c..75b1449 100644 --- a/test/run/test_liveness.jl +++ b/test/run/test_liveness.jl @@ -21,13 +21,60 @@ const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown end -@testset "holder_liveness: Slurm is only asked from inside an allocation" begin - # `squeue` existing proves nothing about WHICH queue it answers for. On the machine this was - # written on it is a wrapper around a remote cluster, where every foreign job id is absent and - # would read as :dead. - tok = "otherhost:1:abcd:slurm999999" - withenv("SLURM_JOB_ID" => nothing) do - @test holder_liveness(tok) === :unknown +@testset "holder_liveness: the queue decision, with the fetch stubbed" begin + # The membership rule cannot be reached without a scheduler, and it is the part with real + # content: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is + # alive if ANY of its tasks is. Only the squeue CALL is stubbed; the decision under test is the + # package's own. + saved = SweepRunner._squeue_cache[] + try + SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) + withenv("SLURM_JOB_ID" => "1") do + @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued + @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued + @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed + end + + # No opinion from squeue is never :dead, however long ago it was asked. + SweepRunner._squeue_cache[] = (time(), nothing) + withenv("SLURM_JOB_ID" => "1") do + @test holder_liveness("h:1:ab:slurm999") === :unknown + end + finally + SweepRunner._squeue_cache[] = saved # never leak a stub into a sibling test file + end +end + +@testset "holder_liveness: outside an allocation the queue is not consulted at all" begin + saved = SweepRunner._squeue_cache[] + try + SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent + withenv("SLURM_JOB_ID" => nothing) do + # Would be :dead if the guard were not there, since the job is absent from this queue. + @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown + end + finally + SweepRunner._squeue_cache[] = saved + end +end + +@testset "the squeue fetch cannot throw and cannot invent a :dead" begin + # Forces the cache miss so the real call runs. The assertion is the safety contract rather + # than a value: this executes on a machine with no scheduler (CI), on one whose `squeue` is a + # wrapper around a remote cluster (the development box), and inside a real allocation, and in + # every one of those it must come back with an answer rather than an exception. + saved = SweepRunner._squeue_cache[] + try + SweepRunner._squeue_cache[] = (-Inf, nothing) + live = SweepRunner._live_slurm_jobs() + @test live === nothing || live isa Set{String} + + SweepRunner._squeue_cache[] = (-Inf, nothing) + withenv("SLURM_JOB_ID" => "1") do + @test holder_liveness("h:1:ab:slurm999999") in (:alive, :dead, :unknown) + end + finally + SweepRunner._squeue_cache[] = saved end end diff --git a/test/run/test_liveness.jl.3860359.cov b/test/run/test_liveness.jl.3860359.cov new file mode 100644 index 0000000..6c7ccc5 --- /dev/null +++ b/test/run/test_liveness.jl.3860359.cov @@ -0,0 +1,114 @@ + - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. + - # + - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only + - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a + - # live master's key to someone else, which is what the lock exists to prevent. + - + - using SweepRunner, Test, DataVault, ParamIO + - + - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") + - const _HOST = gethostname() + - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running + - + - @testset "holder_liveness: :dead only on positive evidence" begin + - @test holder_liveness(owner_token()) === :alive # this very process + - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host + - + - # Everything uncertain is :unknown, so stale_after decides exactly as before. + - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine + - @test holder_liveness("garbage") === :unknown + - @test holder_liveness("") === :unknown + - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown + - end + - + - @testset "holder_liveness: Slurm is only asked from inside an allocation" begin + - # `squeue` existing proves nothing about WHICH queue it answers for. On the machine this was + - # written on it is a wrapper around a remote cluster, where every foreign job id is absent and + - # would read as :dead. + - tok = "otherhost:1:abcd:slurm999999" + - withenv("SLURM_JOB_ID" => nothing) do + 1 @test holder_liveness(tok) === :unknown + - end + - end + - + - @testset "owner_token: carries the Slurm job so another HOST can ask" begin + - withenv("SLURM_JOB_ID" => "4242") do + 2 t = owner_token() + 1 @test occursin(":slurm4242", t) + 1 @test startswith(t, gethostname() * ":") + - end + - withenv("SLURM_JOB_ID" => nothing) do + 1 @test !occursin("slurm", owner_token()) + - end + - end + - + 3 function with_one_left(f) + 3 outdir = mktempdir() + 3 try + 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) + 3 ks = DataVault.keys(v) + 3 for k in ks[2:end] + 9 DataVault.save!(v, k, Dict("x" => 1)) + 9 DataVault.mark_done!(v, k) + 9 end + 3 f(v, ks) + - finally + 3 rm(outdir; recursive=true, force=true) + - end + - end + - + - @testset "a dead holder's lock is reaped at once, not after stale_after" begin + - with_one_left() do v, ks + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test DataVault.running_owner(v, ks[1]) == _DEAD + - + 1 t0 = time() + 1 r = run_loop!( + 1 k -> Dict{String,Any}("x" => 1), + - v, + - ks; + - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 2 @test DataVault.is_done(v, ks[1]) # completed + 1 @test r.done == 1 + 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after + - end + - end + - + - @testset "a LIVE holder's lock is not reaped" begin + - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a + - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in + - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. + - with_one_left() do v, ks + 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) + 1 mine = owner_token() # this process: provably alive + 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok + - + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false + 2 @test DataVault.is_running(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched + - + - # And the same reaper DOES clear it once the owner is one that cannot be running, so the + - # `false` above is the liveness answer and not an inert function. + 1 DataVault.clear_running!(v, ks[1], mine) + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true + 2 @test !DataVault.is_running(v, ks[1]) + - end + - end + - + - @testset "an unstamped lock is left to the timeout" begin + - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. + - with_one_left() do v, ks + 1 DataVault.mark_running!(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) === nothing + 1 @test SweepRunner._reap_if_dead!( + - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) + - ) == false + 2 @test DataVault.is_running(v, ks[1]) + - end + - end diff --git a/test/run/test_liveness.jl.3861352.cov b/test/run/test_liveness.jl.3861352.cov new file mode 100644 index 0000000..6669771 --- /dev/null +++ b/test/run/test_liveness.jl.3861352.cov @@ -0,0 +1,141 @@ + - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. + - # + - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only + - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a + - # live master's key to someone else, which is what the lock exists to prevent. + - + - using SweepRunner, Test, DataVault, ParamIO + - + - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") + - const _HOST = gethostname() + - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running + - + - @testset "holder_liveness: :dead only on positive evidence" begin + - @test holder_liveness(owner_token()) === :alive # this very process + - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host + - + - # Everything uncertain is :unknown, so stale_after decides exactly as before. + - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine + - @test holder_liveness("garbage") === :unknown + - @test holder_liveness("") === :unknown + - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown + - end + - + - @testset "holder_liveness: the queue decision, with the fetch stubbed" begin + - # The membership rule cannot be reached without a scheduler, and it is the part with real + - # content: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is + - # alive if ANY of its tasks is. Only the squeue CALL is stubbed; the decision under test is the + - # package's own. + - saved = SweepRunner._squeue_cache[] + - try + - SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) + - withenv("SLURM_JOB_ID" => "1") do + 1 @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued + 1 @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued + 1 @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed + - end + - + - # No opinion from squeue is never :dead, however long ago it was asked. + - SweepRunner._squeue_cache[] = (time(), nothing) + - withenv("SLURM_JOB_ID" => "1") do + 1 @test holder_liveness("h:1:ab:slurm999") === :unknown + - end + - finally + - SweepRunner._squeue_cache[] = saved # never leak a stub into a sibling test file + - end + - end + - + - @testset "holder_liveness: outside an allocation the queue is not consulted at all" begin + - saved = SweepRunner._squeue_cache[] + - try + - SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent + - withenv("SLURM_JOB_ID" => nothing) do + - # Would be :dead if the guard were not there, since the job is absent from this queue. + 1 @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown + - end + - finally + - SweepRunner._squeue_cache[] = saved + - end + - end + - + - @testset "owner_token: carries the Slurm job so another HOST can ask" begin + - withenv("SLURM_JOB_ID" => "4242") do + 2 t = owner_token() + 1 @test occursin(":slurm4242", t) + 1 @test startswith(t, gethostname() * ":") + - end + - withenv("SLURM_JOB_ID" => nothing) do + 1 @test !occursin("slurm", owner_token()) + - end + - end + - + 3 function with_one_left(f) + 3 outdir = mktempdir() + 3 try + 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) + 3 ks = DataVault.keys(v) + 3 for k in ks[2:end] + 9 DataVault.save!(v, k, Dict("x" => 1)) + 9 DataVault.mark_done!(v, k) + 9 end + 3 f(v, ks) + - finally + 3 rm(outdir; recursive=true, force=true) + - end + - end + - + - @testset "a dead holder's lock is reaped at once, not after stale_after" begin + - with_one_left() do v, ks + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test DataVault.running_owner(v, ks[1]) == _DEAD + - + 1 t0 = time() + 1 r = run_loop!( + 1 k -> Dict{String,Any}("x" => 1), + - v, + - ks; + - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 2 @test DataVault.is_done(v, ks[1]) # completed + 1 @test r.done == 1 + 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after + - end + - end + - + - @testset "a LIVE holder's lock is not reaped" begin + - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a + - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in + - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. + - with_one_left() do v, ks + 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) + 1 mine = owner_token() # this process: provably alive + 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok + - + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false + 2 @test DataVault.is_running(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched + - + - # And the same reaper DOES clear it once the owner is one that cannot be running, so the + - # `false` above is the liveness answer and not an inert function. + 1 DataVault.clear_running!(v, ks[1], mine) + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true + 2 @test !DataVault.is_running(v, ks[1]) + - end + - end + - + - @testset "an unstamped lock is left to the timeout" begin + - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. + - with_one_left() do v, ks + 1 DataVault.mark_running!(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) === nothing + 1 @test SweepRunner._reap_if_dead!( + - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) + - ) == false + 2 @test DataVault.is_running(v, ks[1]) + - end + - end diff --git a/test/run/test_liveness.jl.3861939.cov b/test/run/test_liveness.jl.3861939.cov new file mode 100644 index 0000000..29ff023 --- /dev/null +++ b/test/run/test_liveness.jl.3861939.cov @@ -0,0 +1,161 @@ + - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. + - # + - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only + - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a + - # live master's key to someone else, which is what the lock exists to prevent. + - + - using SweepRunner, Test, DataVault, ParamIO + - + - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") + - const _HOST = gethostname() + - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running + - + - @testset "holder_liveness: :dead only on positive evidence" begin + - @test holder_liveness(owner_token()) === :alive # this very process + - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host + - + - # Everything uncertain is :unknown, so stale_after decides exactly as before. + - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine + - @test holder_liveness("garbage") === :unknown + - @test holder_liveness("") === :unknown + - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown + - end + - + - @testset "holder_liveness: the queue decision, with the fetch stubbed" begin + - # The membership rule cannot be reached without a scheduler, and it is the part with real + - # content: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is + - # alive if ANY of its tasks is. Only the squeue CALL is stubbed; the decision under test is the + - # package's own. + - saved = SweepRunner._squeue_cache[] + - try + - SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) + - withenv("SLURM_JOB_ID" => "1") do + 1 @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued + 1 @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued + 1 @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed + - end + - + - # No opinion from squeue is never :dead, however long ago it was asked. + - SweepRunner._squeue_cache[] = (time(), nothing) + - withenv("SLURM_JOB_ID" => "1") do + 1 @test holder_liveness("h:1:ab:slurm999") === :unknown + - end + - finally + - SweepRunner._squeue_cache[] = saved # never leak a stub into a sibling test file + - end + - end + - + - @testset "holder_liveness: outside an allocation the queue is not consulted at all" begin + - saved = SweepRunner._squeue_cache[] + - try + - SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent + - withenv("SLURM_JOB_ID" => nothing) do + - # Would be :dead if the guard were not there, since the job is absent from this queue. + 1 @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown + - end + - finally + - SweepRunner._squeue_cache[] = saved + - end + - end + - + - @testset "the squeue fetch cannot throw and cannot invent a :dead" begin + - # Forces the cache miss so the real call runs. The assertion is the safety contract rather + - # than a value: this executes on a machine with no scheduler (CI), on one whose `squeue` is a + - # wrapper around a remote cluster (the development box), and inside a real allocation, and in + - # every one of those it must come back with an answer rather than an exception. + - saved = SweepRunner._squeue_cache[] + - try + - SweepRunner._squeue_cache[] = (-Inf, nothing) + - live = SweepRunner._live_slurm_jobs() + - @test live === nothing || live isa Set{String} + - + - SweepRunner._squeue_cache[] = (-Inf, nothing) + - withenv("SLURM_JOB_ID" => "1") do + 1 @test holder_liveness("h:1:ab:slurm999999") in (:alive, :dead, :unknown) + - end + - finally + - SweepRunner._squeue_cache[] = saved + - end + - end + - + - @testset "owner_token: carries the Slurm job so another HOST can ask" begin + - withenv("SLURM_JOB_ID" => "4242") do + 2 t = owner_token() + 1 @test occursin(":slurm4242", t) + 1 @test startswith(t, gethostname() * ":") + - end + - withenv("SLURM_JOB_ID" => nothing) do + 1 @test !occursin("slurm", owner_token()) + - end + - end + - + 3 function with_one_left(f) + 3 outdir = mktempdir() + 3 try + 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) + 3 ks = DataVault.keys(v) + 3 for k in ks[2:end] + 9 DataVault.save!(v, k, Dict("x" => 1)) + 9 DataVault.mark_done!(v, k) + 9 end + 3 f(v, ks) + - finally + 3 rm(outdir; recursive=true, force=true) + - end + - end + - + - @testset "a dead holder's lock is reaped at once, not after stale_after" begin + - with_one_left() do v, ks + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test DataVault.running_owner(v, ks[1]) == _DEAD + - + 1 t0 = time() + 1 r = run_loop!( + 1 k -> Dict{String,Any}("x" => 1), + - v, + - ks; + - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 2 @test DataVault.is_done(v, ks[1]) # completed + 1 @test r.done == 1 + 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after + - end + - end + - + - @testset "a LIVE holder's lock is not reaped" begin + - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a + - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in + - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. + - with_one_left() do v, ks + 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) + 1 mine = owner_token() # this process: provably alive + 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok + - + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false + 2 @test DataVault.is_running(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched + - + - # And the same reaper DOES clear it once the owner is one that cannot be running, so the + - # `false` above is the liveness answer and not an inert function. + 1 DataVault.clear_running!(v, ks[1], mine) + 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true + 2 @test !DataVault.is_running(v, ks[1]) + - end + - end + - + - @testset "an unstamped lock is left to the timeout" begin + - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. + - with_one_left() do v, ks + 1 DataVault.mark_running!(v, ks[1]) + 1 @test DataVault.running_owner(v, ks[1]) === nothing + 1 @test SweepRunner._reap_if_dead!( + - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) + - ) == false + 2 @test DataVault.is_running(v, ks[1]) + - end + - end diff --git a/test/run/test_run_loop_busy.jl.3860359.cov b/test/run/test_run_loop_busy.jl.3860359.cov new file mode 100644 index 0000000..6f56c62 --- /dev/null +++ b/test/run/test_run_loop_busy.jl.3860359.cov @@ -0,0 +1,112 @@ + - # run_loop! and keys held by a sibling (#..): a round that completes nothing but finds `:lock_busy` + - # is not an empty round until `stale_after` has been waited out. + - # + - # The bug this pins: `max_empty_rounds * idle_sleep` is 90 s by default and `stale_after` is 600 s, + - # so a job following one the wall clock killed returned 8.5 minutes before the dead job's locks + - # became reclaimable, left the campaign short, and reported nothing. + - + - using SweepRunner, Test, DataVault, ParamIO + - + - const _BUSY_CFG = joinpath(@__DIR__, "fixtures", "study.toml") + - + - # Fast constants so the wait is seconds. `heartbeat_interval` must stay under `stale_after`. + 3 function _opts(; kw...) + 3 return RunOpts(; workers=:sequential, stale_after=3.0, heartbeat_interval=1.0, kw...) + - end + - const _BUDGET = 3.0 + 2 * 0.5 # opts.stale_after + 2 * idle_sleep + - + 3 function with_tail(f) + 3 outdir = mktempdir() + 3 try + 3 v = DataVault.Vault(_BUSY_CFG; run="tail", outdir=outdir) + 3 ks = DataVault.keys(v) + 3 for k in ks[2:end] # everything done but the first + 9 DataVault.save!(v, k, Dict("x" => 1)) + 9 DataVault.mark_done!(v, k) + 9 end + 3 f(v, ks) + - finally + 3 rm(outdir; recursive=true, force=true) + - end + - end + - + - @testset "run_loop!: a dead holder's key is waited out and completed" begin + - with_tail() do v, ks + 1 DataVault.mark_running!(v, ks[1]) # killed mid-key: the marker outlived the process + 1 n = Ref(0) + 1 t0 = time() + 1 r = run_loop!( + 1 k -> (n[] += 1; Dict{String,Any}("x" => 1)), + - v, + - ks; + - opts=_opts(), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 1 @test n[] == 1 # it ran the key + 2 @test DataVault.is_done(v, ks[1]) # and the campaign is complete + 1 @test r.done == 1 + 1 @test elapsed >= 3.0 # it waited stale_after out rather than giving up + - end + - end + - + - @testset "run_loop!: with nothing held, it still exits promptly" begin + - # The control. Without this the fix would read as "always wait stale_after", which would make + - # every finished campaign pay 600 s to notice it is finished. + - with_tail() do v, ks + 1 DataVault.save!(v, ks[1], Dict("x" => 1)) + 1 DataVault.mark_done!(v, ks[1]) # everything done, nothing running + - + 1 t0 = time() + 1 r = run_loop!( + - k -> Dict{String,Any}("x" => 1), + - v, + - ks; + - opts=_opts(), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 1 @test r.done == 0 + 1 @test r.busy == 0 + 1 @test elapsed < 3.0 # well under stale_after + - end + - end + - + - @testset "run_loop!: a LIVE sibling is not waited on forever" begin + - # The other side of the bound. A holder that keeps refreshing is alive and doing the work, so + - # this master has nothing to contribute and must return rather than spin. + - with_tail() do v, ks + 1 DataVault.mark_running!(v, ks[1]) + 1 alive = Threads.Atomic{Bool}(true) + 2 beater = Threads.@spawn while alive[] + 15 DataVault.touch_running!(v, ks[1]) # the sibling's heartbeat, never going stale + 15 sleep(0.3) + 15 end + 1 try + 1 n = Ref(0) + 1 t0 = time() + 1 r = run_loop!( + - k -> (n[] += 1; Dict{String,Any}("x" => 1)), + - v, + - ks; + - opts=_opts(), + - max_empty_rounds=2, + - idle_sleep=0.5, + - ) + 1 elapsed = time() - t0 + - + 1 @test n[] == 0 # never stole the live sibling's key + 2 @test !DataVault.is_done(v, ks[1]) + 1 @test r.busy > 0 # and SAYS it left work held, which it could not before + 1 @test elapsed >= _BUDGET # waited the budget + 1 @test elapsed < _BUDGET + 10 # then returned + - finally + 1 alive[] = false + 1 wait(beater) + - end + - end + - end From ead7060189b34f6762193911037edb42280f82ce Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 15 Sep 2026 16:53:03 +0000 Subject: [PATCH 4/7] fix(review): three ways holder_liveness answered :dead for a live master A multi-agent review of this PR found three Critical defects, all of them in the direction the PR body claimed could not happen: "`:dead` is returned only on positive evidence". Each would hand a live master's key to another worker, which is the one outcome the lock exists to prevent. Each now has a regression test. 1. A HOSTNAME beginning with `slurm`. `findfirst(p -> startswith(p, "slurm"), parts)` scanned every field of the owner token, so `slurm-node-01:12345:ab01cd23` had `-node-01` read out of its hostname as a job id, found absent from the queue, and answered `:dead` without ever reaching the pid check. Clusters name nodes `slurm*` routinely. The Slurm field is written by `owner_token` at the FOURTH position and nowhere else, so that is the only position now read. 2. Slurm ARRAY tasks. Slurm gives each task its own raw `SLURM_JOB_ID` (36, 37, 38 ...) while `squeue -o %i` lists them as `_` (36_0, 36_1 ...). Stamping the raw id made every task but the first unfindable, and unfindable read as `:dead`. `owner_token` now stamps the id the queue PRINTS. The old test could not have caught this: it hand-wrote the token with the ARRAY id, sharing the implementation's own assumption about what `SLURM_JOB_ID` holds. 3. A `squeue` pointed at a DIFFERENT cluster. It answers successfully and lists none of our ids, so every holder was absent and therefore `:dead`. The queue must now be shown to see THIS process before its silence about anyone else counts as evidence. Two more, Important: 4. `_reap_if_dead!` had no error handling, so an unlink that failed escaped `run!` and took every other key in the round with it. Measured: with one lock in a read-only directory, `run!` threw and 0 of 3 healthy keys were attempted; it now returns with all 3 done. Reaping is an optimisation over `stale_after`, so nothing in it may be fatal. 5. `squeue` had no timeout, on the hot path of every contended key, holding the cache lock, where the binary can be an SSH wrapper. Neither `stop_flag` nor `deadline` is read inside a key, so a stall took the allocation silently. Bounded to 10 s, then killed, then `:unknown`. Also from the review: - `:gave_up` had acquired a second, unrelated source. `_run_affinity!`'s bound is independent of `opts.max_attempts`, so a key could carry `:gave_up` under `max_attempts = 1`, contradicting three docstrings and miscounting for anything aggregating the log by `kind`. Split out as `:worker_died`. - The re-dispatch bound was written twice through unrelated mechanisms (`ExponentialBackOff(; n=2)` and a hand counter). One `_WORKER_DEATH_REDISPATCHES` now feeds both; expressing one bound twice is how the two dispatchers drift apart, which is how the unbounded give-back got in. - `owner_token`'s docstring did not say to call it fresh per acquisition. It is exported and in api.md, and a caller caching it per process would defeat the nonce. - Comments narrating how a trap was found, rather than the trap, are cut. One named a contributor's machine. - `InterruptException` is no longer swallowed by the squeue catch. - Four `.cov` coverage artifacts (528 lines) were committed by a `git add -A`; removed, and `*.cov`/`*.info` added to .gitignore. Test gaps the review named, now closed: - The real fetch-and-parse had NO value-level test: every testset seeded the cache, so the parser never ran. It is now exercised against the shapes `squeue -o %i` actually prints (plain job, running array task, pending array range, throttled array) through a fake `squeue` on PATH. - That also removes a live network call from the suite. The `squeue` on the machine this was written on is `exec ssh -o BatchMode=yes -- squeue`, and CI runs on that machine, so the previous test SSHed to a remote cluster on every run. - `:lock_reaped` was emitted but never asserted. - Slurm evidence outranking the pid on the same host was never exercised: every Slurm test used a foreign hostname, so the two branches were only ever tested apart. - `test_eventlog.jl`'s enumerated-kinds contract was missing four kinds. - The deaths bound asserted a range where the count is structurally exact. Verified by EXIT CODE, not by grepping output: 24 of 24 files under test/ pass. Three times today a grep over test output reported green where the run was red or red where it was noise. Co-Authored-By: Claude Opus 5 (1M context) --- .gitignore | 7 +- src/EventLog.jl | 6 +- src/Liveness.jl | 86 ++++++--- src/Run.jl | 32 +++- test/eventlog/test_eventlog.jl | 8 +- test/run/test_affinity.jl | 6 +- test/run/test_liveness.jl | 193 +++++++++++++++++++-- test/run/test_liveness.jl.3860359.cov | 114 ------------ test/run/test_liveness.jl.3861352.cov | 141 --------------- test/run/test_liveness.jl.3861939.cov | 161 ----------------- test/run/test_run_loop_busy.jl.3860359.cov | 112 ------------ 11 files changed, 285 insertions(+), 581 deletions(-) delete mode 100644 test/run/test_liveness.jl.3860359.cov delete mode 100644 test/run/test_liveness.jl.3861352.cov delete mode 100644 test/run/test_liveness.jl.3861939.cov delete mode 100644 test/run/test_run_loop_busy.jl.3860359.cov diff --git a/.gitignore b/.gitignore index 3053905..64f03ea 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,9 @@ docs/build/ docs/src/assets/favicon.ico docs/src/assets/logo.png -Manifest.toml \ No newline at end of file +Manifest.toml + +# Coverage counters from `julia --code-coverage`. Written next to the source they measure, +# so a `git add -A` after a local coverage run sweeps them in. +*.cov +*.info diff --git a/src/EventLog.jl b/src/EventLog.jl index 228fe9b..3013dfe 100644 --- a/src/EventLog.jl +++ b/src/EventLog.jl @@ -45,12 +45,16 @@ stay within that guarantee. | `lock_busy` | another master holds the `.running` lock (acquire = `:busy`) | | `lock_lost` | our lock was reclaimed mid-work; result discarded (no double-run) | | `lock_reaped` | a lock whose holder was shown dead was cleared without waiting | +| `reap_failed` | reaping threw; the key falls back to the `stale_after` timeout | | `lock_reclaimed`| (reserved, not currently emitted) | | `error` | `work_fn` threw on this attempt | | `retry` | another attempt will follow | | `gave_up` | all `max_attempts` attempts exhausted | | `skip_complete` | full-done early exit (manifest had every key) | -| `worker_lost` | every worker died with keys pending; the key was not attempted | +| `worker_died` | the worker exited on this key every time it was dispatched, up to | +| | the re-dispatch bound (includes `deaths`) | +| `worker_lost` | every worker died with keys still queued; this key was left for a | +| | later run rather than completed or failed | `:key_start` and `:lock_busy` are emitted at `:debug` level and are suppressed unless the `EventLog` is created with `min_level=:debug` (see `RunOpts.log_level`); diff --git a/src/Liveness.jl b/src/Liveness.jl index 21c147b..94aa9d6 100644 --- a/src/Liveness.jl +++ b/src/Liveness.jl @@ -9,23 +9,35 @@ # then decides exactly as before. `:slurm` is one of four worker modes, so a mechanism that only # worked there would cover a quarter of the runs. -using Dates - """ owner_token() -> String -This master's identity, stamped into `.running` by [`run!`](@ref) so a sibling can ask about it -later. `host:pid:nonce` from DataVault, with `:slurm` appended inside a Slurm allocation. +The identity of ONE acquisition, stamped into `.running` by [`run!`](@ref) so a sibling can ask +about it later. `host:pid:nonce` from DataVault, with `:slurm` appended inside a Slurm +allocation. + +Call it per acquisition and do not cache it. The nonce is what tells a master's current hold from +a hold it lost and retook, and a token broadcast once per process cannot make that distinction. The Slurm field is what makes the question answerable ACROSS hosts: `/proc` only works for a holder on this machine, and the master of the job that was killed is usually somewhere else. """ function owner_token()::String base = DataVault.new_owner_token() - job = get(ENV, "SLURM_JOB_ID", "") + job = _slurm_queue_id() return isempty(job) ? base : string(base, ":slurm", job) end +# The id as `squeue` PRINTS it, which is not always `SLURM_JOB_ID`. In an array task that variable +# holds a distinct raw id per task while the queue lists `_`, so stamping it +# makes every task but the first unfindable, and unfindable reads as `:dead`. +function _slurm_queue_id()::String + arr = get(ENV, "SLURM_ARRAY_JOB_ID", "") + idx = get(ENV, "SLURM_ARRAY_TASK_ID", "") + (isempty(arr) || isempty(idx)) && return get(ENV, "SLURM_JOB_ID", "") + return string(arr, "_", idx) +end + """ holder_liveness(owner) -> Symbol @@ -46,9 +58,12 @@ function holder_liveness(owner::AbstractString)::Symbol parts = split(String(owner), ':') length(parts) >= 3 || return :unknown - job = findfirst(p -> startswith(p, "slurm"), parts) - if job !== nothing - s = _slurm_liveness(parts[job][6:end]) + # The Slurm field is the FOURTH part and nowhere else, because that is where `owner_token` + # writes it. Searching every part for a `slurm` prefix matched a HOSTNAME beginning with + # `slurm`, read the rest of the hostname as a job id, found it absent from the queue, and + # returned `:dead` for a master that was alive. + if length(parts) >= 4 && startswith(parts[4], "slurm") + s = _slurm_liveness(chopprefix(parts[4], "slurm")) s === :unknown || return s end @@ -67,20 +82,25 @@ end # The live-job set, cached: `squeue` is a scheduler RPC and this is consulted per contended key. const _SQUEUE_TTL = 15.0 +const _SQUEUE_TIMEOUT = 10.0 const _squeue_cache = Ref{Tuple{Float64,Union{Set{String},Nothing}}}((-Inf, nothing)) const _squeue_lock = ReentrantLock() function _slurm_liveness(jobid::AbstractString)::Symbol isempty(jobid) && return :unknown - # Ask Slurm only from INSIDE a Slurm allocation, so the queue being consulted is demonstrably - # the one that would have run the holder. `squeue` existing proves nothing: measured on the - # development box behind this package, `/home/…/.local/bin/squeue` is a wrapper that answers - # about a REMOTE cluster's queue, and every job id from anywhere else reads as absent there. - # Absent would then mean `:dead`, and a false `:dead` hands a live master's key to someone - # else, which is the one outcome the lock exists to prevent. + # A job id Slurm would never issue cannot be looked up, and "not in the queue" must not be the + # answer for it: that is a `:dead` conjured out of a corrupt `.running` file. `_` is legal + # because an array task's queue id is `_`. + all(c -> isdigit(c) || c == '_', jobid) || return :unknown + # Only from INSIDE an allocation: a `squeue` on PATH may be a wrapper pointed at an unrelated + # cluster, where every foreign job id is absent and would therefore read as `:dead`. haskey(ENV, "SLURM_JOB_ID") || return :unknown live = _live_slurm_jobs() live === nothing && return :unknown + # The queue has to be able to see THIS process before its silence about anyone else means + # anything. A `squeue` pointed at a different cluster answers successfully and lists none of + # our ids, and every holder would then read as `:dead`. + _slurm_queue_id() in live || return :unknown # An array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so the whole job is # alive if any of its tasks is. jobid in live && return :alive @@ -88,7 +108,32 @@ function _slurm_liveness(jobid::AbstractString)::Symbol return :dead end -# `nothing` means "no opinion": squeue is absent, or it failed. Listing the queue and testing +# Bounded, because this runs on the hot path of every contended key and `squeue` can be a wrapper +# around a network call. An unbounded one that stalls holds `_squeue_lock` and takes the allocation +# with it, and neither `stop_flag` nor `deadline` is read inside a key. +function _squeue_output(timeout::Real)::Union{String,Nothing} + Sys.which("squeue") === nothing && return nothing + tmp = tempname() + proc = try + run(pipeline(`squeue -h -o %i`; stdout=tmp, stderr=devnull); wait=false) + catch e + e isa InterruptException && rethrow() + rm(tmp; force=true) + return nothing + end + try + if timedwait(() -> !process_running(proc), Float64(timeout); pollint=0.05) !== :ok + kill(proc, Base.SIGKILL) + return nothing + end + success(proc) || return nothing + return read(tmp, String) + finally + rm(tmp; force=true) + end +end + +# `nothing` means "no opinion": squeue is absent, timed out, or failed. Listing the queue and testing # membership is deliberate. `squeue -j ` exits non-zero BOTH for a finished job and for a # broken scheduler, and those must not collapse into the same answer. function _live_slurm_jobs()::Union{Set{String},Nothing} @@ -96,13 +141,10 @@ function _live_slurm_jobs()::Union{Set{String},Nothing} t, cached = _squeue_cache[] time() - t < _SQUEUE_TTL && return cached fresh = try - if Sys.which("squeue") === nothing - nothing - else - out = read(pipeline(`squeue -h -o %i`; stderr=devnull), String) - Set(String.(split(out; keepempty=false))) - end - catch + out = _squeue_output(_SQUEUE_TIMEOUT) + out === nothing ? nothing : Set(String.(split(out; keepempty=false))) + catch e + e isa InterruptException && rethrow() nothing end _squeue_cache[] = (time(), fresh) diff --git a/src/Run.jl b/src/Run.jl index b600b0e..0bb1abd 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -134,6 +134,11 @@ end _is_stopped(opts::RunOpts)::Bool = _stop_reason(opts) !== nothing +# How many times a key whose worker DIED is handed to another one. Both dispatchers bound this, +# `_run_pmap!` through `pmap`'s `retry_delays` and `_run_affinity!` by counting, and they have to +# agree: the same bound expressed twice through unrelated mechanisms is how they drift apart. +const _WORKER_DEATH_REDISPATCHES = 2 + # As of v0.3 the per-key lock lives ENTIRELY in DataVault's `.running` # sentinel — acquired atomically via `DataVault.acquire_running!` # (implemented with POSIX `link()`). There is no longer a separate @@ -348,13 +353,24 @@ end # Only `:dead` acts. `:unknown` is the common answer (a holder on another host with no Slurm id) # and leaves the timeout to decide, exactly as before. function _reap_if_dead!(vault::Vault, key::DataKey, stage::Symbol, log::EventLog)::Bool + # An `isfile` first: `running_owner` opens and reads, and the uncontended case is every key. DataVault.is_running(vault, key) || return false - owner = DataVault.running_owner(vault, key) - owner === nothing && return false # unstamped: cannot be attributed, so cannot be judged - holder_liveness(owner) === :dead || return false - cleared = DataVault.clear_running!(vault, key, owner) - cleared && log_event(log, :lock_reaped; stage=stage, key=canonical(key), owner=owner) - return cleared + # Reaping is an OPTIMISATION over `stale_after`, so nothing in it may be fatal. Without this, + # an unlink that fails (a read-only status directory, an NFS hiccup) escapes `run!` and takes + # every other key in the round with it, none of which was attempted. + try + owner = DataVault.running_owner(vault, key) + owner === nothing && return false # unstamped: cannot be attributed, so cannot be judged + holder_liveness(owner) === :dead || return false + cleared = DataVault.clear_running!(vault, key, owner) + cleared && + log_event(log, :lock_reaped; stage=stage, key=canonical(key), owner=owner) + return cleared + catch e + e isa InterruptException && rethrow() + log_event(log, :reap_failed; stage=stage, key=canonical(key), err=_short_err(e)) + return false + end end """ @@ -529,7 +545,7 @@ function _run_pmap!( pool, todo; on_error=identity, - retry_delays=ExponentialBackOff(; n=2), + retry_delays=ExponentialBackOff(; n=_WORKER_DEATH_REDISPATCHES), retry_check=(s, e) -> (s, e isa ProcessExitedException), ) do key return _run_one_with_lock!(work_fn, vault, key, stage, log, opts) @@ -618,7 +634,7 @@ function _run_affinity!( # `retry_delays=ExponentialBackOff(; n=2)`; this dispatcher has to bound it too. Unbounded, a # key that reliably kills whoever takes it is handed to worker after worker forever, and the # faster a dead holder's lock is reclaimed the faster that cascade runs. - const_giveback_limit = 2 + const_giveback_limit = _WORKER_DEATH_REDISPATCHES givebacks = zeros(Int, length(todo)) _give_back!(i::Int)::Bool = lock(q) do givebacks[i] += 1 diff --git a/test/eventlog/test_eventlog.jl b/test/eventlog/test_eventlog.jl index 187cc81..020cd32 100644 --- a/test/eventlog/test_eventlog.jl +++ b/test/eventlog/test_eventlog.jl @@ -28,11 +28,17 @@ end :gave_up, :retry, :skip_complete, + :key_acquired, + :lock_lost, + :lock_reaped, + :reap_failed, + :worker_died, + :worker_lost, ) log_event(log, k; info="test") end lines = readlines(log.path) - @test length(lines) == 10 + @test length(lines) == 16 for line in lines rec = JSON3.read(line) @test haskey(rec, :kind) diff --git a/test/run/test_affinity.jl b/test/run/test_affinity.jl index 969b16e..df65626 100644 --- a/test/run/test_affinity.jl +++ b/test/run/test_affinity.jl @@ -230,8 +230,10 @@ end ndeaths = length(readdir(deaths)) @info "poison key" ndeaths err = r.err done = r.done - @test ndeaths >= 1 # the fixture really does kill workers - @test ndeaths <= 3 # 1 dispatch + 2 give-backs, and no more + # Exactly 3, not a range: each worker's dispatch loop breaks after its first + # ProcessExitedException, so one initial dispatch plus two give-backs is the only + # reachable count. A range would still pass if the bound silently shrank to 1. + @test ndeaths == 1 + SweepRunner._WORKER_DEATH_REDISPATCHES @test !DataVault.is_done(v, ks[1]) # it cannot complete, and did not @test r.done == length(ks) - 1 # every OTHER key still finished finally diff --git a/test/run/test_liveness.jl b/test/run/test_liveness.jl index 75b1449..95a1896 100644 --- a/test/run/test_liveness.jl +++ b/test/run/test_liveness.jl @@ -4,7 +4,7 @@ # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a # live master's key to someone else, which is what the lock exists to prevent. -using SweepRunner, Test, DataVault, ParamIO +using SweepRunner, Test, DataVault, ParamIO, JSON3 const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") const _HOST = gethostname() @@ -28,10 +28,10 @@ end # package's own. saved = SweepRunner._squeue_cache[] try - SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) + SweepRunner._squeue_cache[] = (time(), Set(["1", "12345", "777_3", "777_4"])) withenv("SLURM_JOB_ID" => "1") do @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued - @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued + @test holder_liveness("h:1:ab:slurm777_3") === :alive # array TASK, as squeue names it @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed end @@ -45,33 +45,157 @@ end end end -@testset "holder_liveness: outside an allocation the queue is not consulted at all" begin +@testset "owner_token: an array task stamps the id squeue PRINTS, not SLURM_JOB_ID" begin + # Found in review, and the previous test could not have caught it: it hand-wrote the token with + # the ARRAY id, sharing the implementation's assumption about what `SLURM_JOB_ID` holds. + # + # Slurm gives every task of an array its own raw `SLURM_JOB_ID` (36, 37, 38 ...) while the + # queue lists them as `_` (36_0, 36_1 ...). Stamping + # the raw id makes every task but the first unfindable, and unfindable reads as `:dead`. + withenv( + "SLURM_JOB_ID" => "38", "SLURM_ARRAY_JOB_ID" => "36", "SLURM_ARRAY_TASK_ID" => "2" + ) do + @test occursin(":slurm36_2", owner_token()) + @test !occursin(":slurm38", owner_token()) + end + # A plain job has no array variables and keeps using SLURM_JOB_ID. + withenv( + "SLURM_JOB_ID" => "4242", + "SLURM_ARRAY_JOB_ID" => nothing, + "SLURM_ARRAY_TASK_ID" => nothing, + ) do + @test occursin(":slurm4242", owner_token()) + end +end + +@testset "holder_liveness: a queue that cannot see US is not evidence about anyone" begin + # A `squeue` pointed at another cluster answers successfully and lists none of our ids. Every + # holder would then be absent, and absent would mean `:dead`. saved = SweepRunner._squeue_cache[] try - SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent - withenv("SLURM_JOB_ID" => nothing) do - # Would be :dead if the guard were not there, since the job is absent from this queue. - @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown + withenv("SLURM_JOB_ID" => "1", "SLURM_ARRAY_JOB_ID" => nothing) do + SweepRunner._squeue_cache[] = (time(), Set(["9001", "9002"])) # our "1" is absent + @test holder_liveness("h:1:ab:slurm9999") === :unknown + # The control: the identical call becomes :dead once the queue does list us. + SweepRunner._squeue_cache[] = (time(), Set(["1", "9001"])) + @test holder_liveness("h:1:ab:slurm9999") === :dead + end + finally + SweepRunner._squeue_cache[] = saved + end +end + +# A `squeue` on PATH that prints what we tell it to. The real one on the machine this was written +# on is `exec ssh -o BatchMode=yes -- squeue`, so calling it from a test would make a live +# network call on every CI run, and CI runs on that same machine. +function with_fake_squeue(f, output::AbstractString; exitcode::Int=0, sleep_for::Real=0) + dir = mktempdir() + try + bin = joinpath(dir, "squeue") + write( + bin, + """ + #!/bin/sh + [ $(sleep_for) = 0 ] || sleep $(sleep_for) + printf '%s' '$(output)' + exit $(exitcode) + """, + ) + chmod(bin, 0o755) + withenv("PATH" => dir * ":" * get(ENV, "PATH", "")) do + return f() + end + finally + rm(dir; recursive=true, force=true) + end +end + +_uncached() = (SweepRunner._squeue_cache[] = (-Inf, nothing)) + +@testset "_live_slurm_jobs: parses the shapes squeue -o %i actually prints" begin + # The only value-level test of the real fetch-and-parse. Every other testset seeds the cache, + # so the parser itself never runs in them: a regression in the `-o` format or the split would + # go unnoticed. + saved = SweepRunner._squeue_cache[] + try + _uncached() + out = with_fake_squeue("623186\n623186_1\n624087_[0-1]\n595411_[0-3%4]\n") do + SweepRunner._live_slurm_jobs() end + @test out == Set(["623186", "623186_1", "624087_[0-1]", "595411_[0-3%4]"]) + + _uncached() + @test with_fake_squeue("") do + SweepRunner._live_slurm_jobs() + end == Set(String[]) # an empty queue is a real answer, not an error finally SweepRunner._squeue_cache[] = saved end end -@testset "the squeue fetch cannot throw and cannot invent a :dead" begin - # Forces the cache miss so the real call runs. The assertion is the safety contract rather - # than a value: this executes on a machine with no scheduler (CI), on one whose `squeue` is a - # wrapper around a remote cluster (the development box), and inside a real allocation, and in - # every one of those it must come back with an answer rather than an exception. +@testset "_live_slurm_jobs: a failing or absent squeue is no opinion, never an empty queue" begin + # `nothing` and `Set()` must not be confused: an empty queue says every holder is dead, while + # no opinion falls back to stale_after. saved = SweepRunner._squeue_cache[] try - SweepRunner._squeue_cache[] = (-Inf, nothing) - live = SweepRunner._live_slurm_jobs() - @test live === nothing || live isa Set{String} + _uncached() + @test with_fake_squeue("garbage"; exitcode=1) do + SweepRunner._live_slurm_jobs() + end === nothing + + _uncached() + @test withenv("PATH" => "") do # squeue not on PATH at all + SweepRunner._live_slurm_jobs() + end === nothing + finally + SweepRunner._squeue_cache[] = saved + end +end + +@testset "_squeue_output: a hanging squeue is killed, not waited on" begin + # It runs on the hot path of every contended key, under the cache lock, and the real binary can + # be a wrapper around a network call. Neither stop_flag nor deadline is read inside a key. + t0 = time() + out = with_fake_squeue("never"; sleep_for=30) do + SweepRunner._squeue_output(1.0) + end + elapsed = time() - t0 + @test out === nothing # timed out: no opinion + @test elapsed < 15.0 # and did not sit there for the full 30 s +end - SweepRunner._squeue_cache[] = (-Inf, nothing) +@testset "holder_liveness: a hostname beginning with slurm is not a job id" begin + # Found in review. The Slurm field is written by `owner_token` as the FOURTH part; scanning + # every part for a `slurm` prefix matched the HOSTNAME, read `-node-01` out of it as a job id, + # found it absent from the queue and answered `:dead` for a master that was alive. Clusters + # name nodes `slurm*` often enough that this is a real configuration, and a false `:dead` is + # the one answer that costs a double execution. + saved = SweepRunner._squeue_cache[] + try + SweepRunner._squeue_cache[] = (time(), Set(["1", "777"])) withenv("SLURM_JOB_ID" => "1") do - @test holder_liveness("h:1:ab:slurm999999") in (:alive, :dead, :unknown) + @test holder_liveness("slurm-node-01:12345:ab01cd23") === :unknown + @test holder_liveness("slurm:1:ab") === :unknown + # A job id Slurm could never issue is not looked up either: "absent from the queue" + # must not be the answer for a corrupt `.running`. + @test holder_liveness("h:1:ab:slurm-node-01") === :unknown + @test holder_liveness("h:1:ab:slurm" * "\u3042") === :unknown + # ...while a well-formed one on the same code path still resolves both ways. + @test holder_liveness("h:1:ab:slurm777") === :alive + @test holder_liveness("h:1:ab:slurm999") === :dead + end + finally + SweepRunner._squeue_cache[] = saved + end +end + +@testset "holder_liveness: outside an allocation the queue is not consulted at all" begin + saved = SweepRunner._squeue_cache[] + try + SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent + withenv("SLURM_JOB_ID" => nothing) do + # Would be :dead if the guard were not there, since the job is absent from this queue. + @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown end finally SweepRunner._squeue_cache[] = saved @@ -126,6 +250,39 @@ end end end +@testset "reaping emits :lock_reaped, naming the owner it cleared" begin + # `:gave_up` is asserted from the JSONL elsewhere; this kind had no such check, so a reap that + # silently stopped logging would look identical to one that worked. + with_one_left() do v, ks + log = EventLog(joinpath(v.outdir, "reap.jsonl")) + @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok + @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true + + recs = [JSON3.read(l) for l in readlines(log.path)] + reaped = [r for r in recs if String(r["kind"]) == "lock_reaped"] + @test length(reaped) == 1 + @test String(reaped[1]["owner"]) == _DEAD + @test String(reaped[1]["key"]) == ParamIO.canonical(ks[1]) + end +end + +@testset "holder_liveness: Slurm evidence outranks the pid on the same host" begin + # Every other Slurm test uses a foreign hostname, so the two branches are only ever exercised + # apart. Here they disagree: the pid is this live process, the job is absent from the queue. + # Slurm must win, because it can see across the whole allocation and /proc cannot. + saved = SweepRunner._squeue_cache[] + try + withenv("SLURM_JOB_ID" => "1", "SLURM_ARRAY_JOB_ID" => nothing) do + SweepRunner._squeue_cache[] = (time(), Set(["1"])) + tok = "$(_HOST):$(getpid()):abcd:slurm999999" + @test holder_liveness("$(_HOST):$(getpid()):abcd") === :alive # pid alone: alive + @test holder_liveness(tok) === :dead # with Slurm: dead + end + finally + SweepRunner._squeue_cache[] = saved + end +end + @testset "a LIVE holder's lock is not reaped" begin # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in diff --git a/test/run/test_liveness.jl.3860359.cov b/test/run/test_liveness.jl.3860359.cov deleted file mode 100644 index 6c7ccc5..0000000 --- a/test/run/test_liveness.jl.3860359.cov +++ /dev/null @@ -1,114 +0,0 @@ - - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. - - # - - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only - - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a - - # live master's key to someone else, which is what the lock exists to prevent. - - - - using SweepRunner, Test, DataVault, ParamIO - - - - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") - - const _HOST = gethostname() - - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running - - - - @testset "holder_liveness: :dead only on positive evidence" begin - - @test holder_liveness(owner_token()) === :alive # this very process - - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host - - - - # Everything uncertain is :unknown, so stale_after decides exactly as before. - - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine - - @test holder_liveness("garbage") === :unknown - - @test holder_liveness("") === :unknown - - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown - - end - - - - @testset "holder_liveness: Slurm is only asked from inside an allocation" begin - - # `squeue` existing proves nothing about WHICH queue it answers for. On the machine this was - - # written on it is a wrapper around a remote cluster, where every foreign job id is absent and - - # would read as :dead. - - tok = "otherhost:1:abcd:slurm999999" - - withenv("SLURM_JOB_ID" => nothing) do - 1 @test holder_liveness(tok) === :unknown - - end - - end - - - - @testset "owner_token: carries the Slurm job so another HOST can ask" begin - - withenv("SLURM_JOB_ID" => "4242") do - 2 t = owner_token() - 1 @test occursin(":slurm4242", t) - 1 @test startswith(t, gethostname() * ":") - - end - - withenv("SLURM_JOB_ID" => nothing) do - 1 @test !occursin("slurm", owner_token()) - - end - - end - - - 3 function with_one_left(f) - 3 outdir = mktempdir() - 3 try - 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) - 3 ks = DataVault.keys(v) - 3 for k in ks[2:end] - 9 DataVault.save!(v, k, Dict("x" => 1)) - 9 DataVault.mark_done!(v, k) - 9 end - 3 f(v, ks) - - finally - 3 rm(outdir; recursive=true, force=true) - - end - - end - - - - @testset "a dead holder's lock is reaped at once, not after stale_after" begin - - with_one_left() do v, ks - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test DataVault.running_owner(v, ks[1]) == _DEAD - - - 1 t0 = time() - 1 r = run_loop!( - 1 k -> Dict{String,Any}("x" => 1), - - v, - - ks; - - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 2 @test DataVault.is_done(v, ks[1]) # completed - 1 @test r.done == 1 - 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after - - end - - end - - - - @testset "a LIVE holder's lock is not reaped" begin - - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a - - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in - - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. - - with_one_left() do v, ks - 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) - 1 mine = owner_token() # this process: provably alive - 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok - - - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false - 2 @test DataVault.is_running(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched - - - - # And the same reaper DOES clear it once the owner is one that cannot be running, so the - - # `false` above is the liveness answer and not an inert function. - 1 DataVault.clear_running!(v, ks[1], mine) - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true - 2 @test !DataVault.is_running(v, ks[1]) - - end - - end - - - - @testset "an unstamped lock is left to the timeout" begin - - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. - - with_one_left() do v, ks - 1 DataVault.mark_running!(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) === nothing - 1 @test SweepRunner._reap_if_dead!( - - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) - - ) == false - 2 @test DataVault.is_running(v, ks[1]) - - end - - end diff --git a/test/run/test_liveness.jl.3861352.cov b/test/run/test_liveness.jl.3861352.cov deleted file mode 100644 index 6669771..0000000 --- a/test/run/test_liveness.jl.3861352.cov +++ /dev/null @@ -1,141 +0,0 @@ - - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. - - # - - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only - - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a - - # live master's key to someone else, which is what the lock exists to prevent. - - - - using SweepRunner, Test, DataVault, ParamIO - - - - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") - - const _HOST = gethostname() - - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running - - - - @testset "holder_liveness: :dead only on positive evidence" begin - - @test holder_liveness(owner_token()) === :alive # this very process - - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host - - - - # Everything uncertain is :unknown, so stale_after decides exactly as before. - - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine - - @test holder_liveness("garbage") === :unknown - - @test holder_liveness("") === :unknown - - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown - - end - - - - @testset "holder_liveness: the queue decision, with the fetch stubbed" begin - - # The membership rule cannot be reached without a scheduler, and it is the part with real - - # content: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is - - # alive if ANY of its tasks is. Only the squeue CALL is stubbed; the decision under test is the - - # package's own. - - saved = SweepRunner._squeue_cache[] - - try - - SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) - - withenv("SLURM_JOB_ID" => "1") do - 1 @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued - 1 @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued - 1 @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed - - end - - - - # No opinion from squeue is never :dead, however long ago it was asked. - - SweepRunner._squeue_cache[] = (time(), nothing) - - withenv("SLURM_JOB_ID" => "1") do - 1 @test holder_liveness("h:1:ab:slurm999") === :unknown - - end - - finally - - SweepRunner._squeue_cache[] = saved # never leak a stub into a sibling test file - - end - - end - - - - @testset "holder_liveness: outside an allocation the queue is not consulted at all" begin - - saved = SweepRunner._squeue_cache[] - - try - - SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent - - withenv("SLURM_JOB_ID" => nothing) do - - # Would be :dead if the guard were not there, since the job is absent from this queue. - 1 @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown - - end - - finally - - SweepRunner._squeue_cache[] = saved - - end - - end - - - - @testset "owner_token: carries the Slurm job so another HOST can ask" begin - - withenv("SLURM_JOB_ID" => "4242") do - 2 t = owner_token() - 1 @test occursin(":slurm4242", t) - 1 @test startswith(t, gethostname() * ":") - - end - - withenv("SLURM_JOB_ID" => nothing) do - 1 @test !occursin("slurm", owner_token()) - - end - - end - - - 3 function with_one_left(f) - 3 outdir = mktempdir() - 3 try - 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) - 3 ks = DataVault.keys(v) - 3 for k in ks[2:end] - 9 DataVault.save!(v, k, Dict("x" => 1)) - 9 DataVault.mark_done!(v, k) - 9 end - 3 f(v, ks) - - finally - 3 rm(outdir; recursive=true, force=true) - - end - - end - - - - @testset "a dead holder's lock is reaped at once, not after stale_after" begin - - with_one_left() do v, ks - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test DataVault.running_owner(v, ks[1]) == _DEAD - - - 1 t0 = time() - 1 r = run_loop!( - 1 k -> Dict{String,Any}("x" => 1), - - v, - - ks; - - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 2 @test DataVault.is_done(v, ks[1]) # completed - 1 @test r.done == 1 - 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after - - end - - end - - - - @testset "a LIVE holder's lock is not reaped" begin - - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a - - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in - - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. - - with_one_left() do v, ks - 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) - 1 mine = owner_token() # this process: provably alive - 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok - - - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false - 2 @test DataVault.is_running(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched - - - - # And the same reaper DOES clear it once the owner is one that cannot be running, so the - - # `false` above is the liveness answer and not an inert function. - 1 DataVault.clear_running!(v, ks[1], mine) - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true - 2 @test !DataVault.is_running(v, ks[1]) - - end - - end - - - - @testset "an unstamped lock is left to the timeout" begin - - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. - - with_one_left() do v, ks - 1 DataVault.mark_running!(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) === nothing - 1 @test SweepRunner._reap_if_dead!( - - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) - - ) == false - 2 @test DataVault.is_running(v, ks[1]) - - end - - end diff --git a/test/run/test_liveness.jl.3861939.cov b/test/run/test_liveness.jl.3861939.cov deleted file mode 100644 index 29ff023..0000000 --- a/test/run/test_liveness.jl.3861939.cov +++ /dev/null @@ -1,161 +0,0 @@ - - # holder_liveness: asking whether the holder is gone, instead of waiting to find out. - - # - - # `stale_after` is the fallback and stays correct on its own. This only removes the WAIT, and only - - # where the answer can be had. Every uncertain path must read `:unknown`: a false `:dead` hands a - - # live master's key to someone else, which is what the lock exists to prevent. - - - - using SweepRunner, Test, DataVault, ParamIO - - - - const _LIV_CFG = joinpath(@__DIR__, "fixtures", "study.toml") - - const _HOST = gethostname() - - const _DEAD = "$(_HOST):999999:deadbeef" # a pid that cannot be running - - - - @testset "holder_liveness: :dead only on positive evidence" begin - - @test holder_liveness(owner_token()) === :alive # this very process - - @test holder_liveness(_DEAD) === :dead # /proc says no such pid, on this host - - - - # Everything uncertain is :unknown, so stale_after decides exactly as before. - - @test holder_liveness("otherhost:123:abcd") === :unknown # another machine - - @test holder_liveness("garbage") === :unknown - - @test holder_liveness("") === :unknown - - @test holder_liveness("$(_HOST):notanumber:abcd") === :unknown - - end - - - - @testset "holder_liveness: the queue decision, with the fetch stubbed" begin - - # The membership rule cannot be reached without a scheduler, and it is the part with real - - # content: an array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so a job is - - # alive if ANY of its tasks is. Only the squeue CALL is stubbed; the decision under test is the - - # package's own. - - saved = SweepRunner._squeue_cache[] - - try - - SweepRunner._squeue_cache[] = (time(), Set(["12345", "777_3", "777_4"])) - - withenv("SLURM_JOB_ID" => "1") do - 1 @test holder_liveness("h:1:ab:slurm12345") === :alive # plain job, queued - 1 @test holder_liveness("h:1:ab:slurm777") === :alive # array job, a task queued - 1 @test holder_liveness("h:1:ab:slurm999") === :dead # absent: finished or killed - - end - - - - # No opinion from squeue is never :dead, however long ago it was asked. - - SweepRunner._squeue_cache[] = (time(), nothing) - - withenv("SLURM_JOB_ID" => "1") do - 1 @test holder_liveness("h:1:ab:slurm999") === :unknown - - end - - finally - - SweepRunner._squeue_cache[] = saved # never leak a stub into a sibling test file - - end - - end - - - - @testset "holder_liveness: outside an allocation the queue is not consulted at all" begin - - saved = SweepRunner._squeue_cache[] - - try - - SweepRunner._squeue_cache[] = (time(), Set(String[])) # an empty queue: everything absent - - withenv("SLURM_JOB_ID" => nothing) do - - # Would be :dead if the guard were not there, since the job is absent from this queue. - 1 @test holder_liveness("otherhost:1:ab:slurm12345") === :unknown - - end - - finally - - SweepRunner._squeue_cache[] = saved - - end - - end - - - - @testset "the squeue fetch cannot throw and cannot invent a :dead" begin - - # Forces the cache miss so the real call runs. The assertion is the safety contract rather - - # than a value: this executes on a machine with no scheduler (CI), on one whose `squeue` is a - - # wrapper around a remote cluster (the development box), and inside a real allocation, and in - - # every one of those it must come back with an answer rather than an exception. - - saved = SweepRunner._squeue_cache[] - - try - - SweepRunner._squeue_cache[] = (-Inf, nothing) - - live = SweepRunner._live_slurm_jobs() - - @test live === nothing || live isa Set{String} - - - - SweepRunner._squeue_cache[] = (-Inf, nothing) - - withenv("SLURM_JOB_ID" => "1") do - 1 @test holder_liveness("h:1:ab:slurm999999") in (:alive, :dead, :unknown) - - end - - finally - - SweepRunner._squeue_cache[] = saved - - end - - end - - - - @testset "owner_token: carries the Slurm job so another HOST can ask" begin - - withenv("SLURM_JOB_ID" => "4242") do - 2 t = owner_token() - 1 @test occursin(":slurm4242", t) - 1 @test startswith(t, gethostname() * ":") - - end - - withenv("SLURM_JOB_ID" => nothing) do - 1 @test !occursin("slurm", owner_token()) - - end - - end - - - 3 function with_one_left(f) - 3 outdir = mktempdir() - 3 try - 3 v = DataVault.Vault(_LIV_CFG; run="liv", outdir=outdir) - 3 ks = DataVault.keys(v) - 3 for k in ks[2:end] - 9 DataVault.save!(v, k, Dict("x" => 1)) - 9 DataVault.mark_done!(v, k) - 9 end - 3 f(v, ks) - - finally - 3 rm(outdir; recursive=true, force=true) - - end - - end - - - - @testset "a dead holder's lock is reaped at once, not after stale_after" begin - - with_one_left() do v, ks - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test DataVault.running_owner(v, ks[1]) == _DEAD - - - 1 t0 = time() - 1 r = run_loop!( - 1 k -> Dict{String,Any}("x" => 1), - - v, - - ks; - - opts=RunOpts(; workers=:sequential, stale_after=600.0, heartbeat_interval=30.0), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 2 @test DataVault.is_done(v, ks[1]) # completed - 1 @test r.done == 1 - 1 @test elapsed < 30.0 # and did NOT wait out the 600 s stale_after - - end - - end - - - - @testset "a LIVE holder's lock is not reaped" begin - - # The control, on the reaper DIRECTLY. Going through run_loop! here would prove nothing: with a - - # short stale_after the timeout reclaims the lock legitimately (nothing is heartbeating it in - - # this test), and with a long one the loop just waits. Neither exercises the reaper's decision. - - with_one_left() do v, ks - 1 log = EventLog(joinpath(v.outdir, "e.jsonl")) - 1 mine = owner_token() # this process: provably alive - 1 @test DataVault.acquire_running!(v, ks[1], mine) === :ok - - - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == false - 2 @test DataVault.is_running(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) == mine # untouched - - - - # And the same reaper DOES clear it once the owner is one that cannot be running, so the - - # `false` above is the liveness answer and not an inert function. - 1 DataVault.clear_running!(v, ks[1], mine) - 1 @test DataVault.acquire_running!(v, ks[1], _DEAD) === :ok - 1 @test SweepRunner._reap_if_dead!(v, ks[1], :liv, log) == true - 2 @test !DataVault.is_running(v, ks[1]) - - end - - end - - - - @testset "an unstamped lock is left to the timeout" begin - - # `mark_running!` writes no owner, and a lock that cannot be attributed cannot be judged. - - with_one_left() do v, ks - 1 DataVault.mark_running!(v, ks[1]) - 1 @test DataVault.running_owner(v, ks[1]) === nothing - 1 @test SweepRunner._reap_if_dead!( - - v, ks[1], :liv, EventLog(joinpath(v.outdir, "e.jsonl")) - - ) == false - 2 @test DataVault.is_running(v, ks[1]) - - end - - end diff --git a/test/run/test_run_loop_busy.jl.3860359.cov b/test/run/test_run_loop_busy.jl.3860359.cov deleted file mode 100644 index 6f56c62..0000000 --- a/test/run/test_run_loop_busy.jl.3860359.cov +++ /dev/null @@ -1,112 +0,0 @@ - - # run_loop! and keys held by a sibling (#..): a round that completes nothing but finds `:lock_busy` - - # is not an empty round until `stale_after` has been waited out. - - # - - # The bug this pins: `max_empty_rounds * idle_sleep` is 90 s by default and `stale_after` is 600 s, - - # so a job following one the wall clock killed returned 8.5 minutes before the dead job's locks - - # became reclaimable, left the campaign short, and reported nothing. - - - - using SweepRunner, Test, DataVault, ParamIO - - - - const _BUSY_CFG = joinpath(@__DIR__, "fixtures", "study.toml") - - - - # Fast constants so the wait is seconds. `heartbeat_interval` must stay under `stale_after`. - 3 function _opts(; kw...) - 3 return RunOpts(; workers=:sequential, stale_after=3.0, heartbeat_interval=1.0, kw...) - - end - - const _BUDGET = 3.0 + 2 * 0.5 # opts.stale_after + 2 * idle_sleep - - - 3 function with_tail(f) - 3 outdir = mktempdir() - 3 try - 3 v = DataVault.Vault(_BUSY_CFG; run="tail", outdir=outdir) - 3 ks = DataVault.keys(v) - 3 for k in ks[2:end] # everything done but the first - 9 DataVault.save!(v, k, Dict("x" => 1)) - 9 DataVault.mark_done!(v, k) - 9 end - 3 f(v, ks) - - finally - 3 rm(outdir; recursive=true, force=true) - - end - - end - - - - @testset "run_loop!: a dead holder's key is waited out and completed" begin - - with_tail() do v, ks - 1 DataVault.mark_running!(v, ks[1]) # killed mid-key: the marker outlived the process - 1 n = Ref(0) - 1 t0 = time() - 1 r = run_loop!( - 1 k -> (n[] += 1; Dict{String,Any}("x" => 1)), - - v, - - ks; - - opts=_opts(), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 1 @test n[] == 1 # it ran the key - 2 @test DataVault.is_done(v, ks[1]) # and the campaign is complete - 1 @test r.done == 1 - 1 @test elapsed >= 3.0 # it waited stale_after out rather than giving up - - end - - end - - - - @testset "run_loop!: with nothing held, it still exits promptly" begin - - # The control. Without this the fix would read as "always wait stale_after", which would make - - # every finished campaign pay 600 s to notice it is finished. - - with_tail() do v, ks - 1 DataVault.save!(v, ks[1], Dict("x" => 1)) - 1 DataVault.mark_done!(v, ks[1]) # everything done, nothing running - - - 1 t0 = time() - 1 r = run_loop!( - - k -> Dict{String,Any}("x" => 1), - - v, - - ks; - - opts=_opts(), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 1 @test r.done == 0 - 1 @test r.busy == 0 - 1 @test elapsed < 3.0 # well under stale_after - - end - - end - - - - @testset "run_loop!: a LIVE sibling is not waited on forever" begin - - # The other side of the bound. A holder that keeps refreshing is alive and doing the work, so - - # this master has nothing to contribute and must return rather than spin. - - with_tail() do v, ks - 1 DataVault.mark_running!(v, ks[1]) - 1 alive = Threads.Atomic{Bool}(true) - 2 beater = Threads.@spawn while alive[] - 15 DataVault.touch_running!(v, ks[1]) # the sibling's heartbeat, never going stale - 15 sleep(0.3) - 15 end - 1 try - 1 n = Ref(0) - 1 t0 = time() - 1 r = run_loop!( - - k -> (n[] += 1; Dict{String,Any}("x" => 1)), - - v, - - ks; - - opts=_opts(), - - max_empty_rounds=2, - - idle_sleep=0.5, - - ) - 1 elapsed = time() - t0 - - - 1 @test n[] == 0 # never stole the live sibling's key - 2 @test !DataVault.is_done(v, ks[1]) - 1 @test r.busy > 0 # and SAYS it left work held, which it could not before - 1 @test elapsed >= _BUDGET # waited the budget - 1 @test elapsed < _BUDGET + 10 # then returned - - finally - 1 alive[] = false - 1 wait(beater) - - end - - end - - end From 3e5fc169ad728934a988075187f8e0175114d3fd Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Wed, 16 Sep 2026 07:01:59 +0000 Subject: [PATCH 5/7] fix(review): a prerequisite could outlive its allocation, and stopped_by named the wrong reason MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two findings from the review of the published 0.6.2, both about a control answering for a condition it was not measuring. `Prerequisite.opts` replaced the dependent stage's `RunOpts` WHOLESALE. `deadline` defaults to `nothing`, so a prerequisite built for the documented reason — raising `stale_after`, because the shared setup is the slow half — ran with no deadline at all. The barrier then waits for a sibling past the end of the allocation that is running it, which is the exact failure the deadline exists to prevent. `stop_flag` and `deadline` are now inherited from the caller when the prerequisite leaves them unset, and still take precedence when it sets them. `stopped_by` was re-read from the clock at return time, in `run!` and again in `run_loop!`. That answers "is a stop condition true NOW", not "why did this stop", and the two differ in both directions: a loop that sleeps `idle_sleep` between rounds crosses the deadline while doing so, so giving up on `max_empty_rounds` was reported as `:deadline` — a retryable answer for a key that cannot be produced; and a flag file removed in the meantime turned a real flag stop into `nothing`. The reason now travels back with the outcome that carried it (`:stop_flag` / `:stop_deadline`) and is recorded where the stop happened. Fixing that exposed a third: the sequential dispatcher `break`s on a stop and emitted NO outcome for the keys it dropped, while `pmap` hands every remaining key back as stopped. The same stop reported a different `stop` count depending on which dispatcher ran. Both now attribute them, and `done + stop + err + busy == length(keys)` holds on either path. Three tests, each mutation-checked: reverting the fix individually turns the intended testset red, and the unmutated control stays green. Co-Authored-By: Claude Opus 5 (1M context) --- src/Prerequisite.jl | 28 +++++++++++++++++-- src/Run.jl | 51 ++++++++++++++++++++++++++-------- test/run/test_prerequisite.jl | 24 ++++++++++++++++ test/run/test_run_deadline.jl | 52 +++++++++++++++++++++++++++++++++++ 4 files changed, 141 insertions(+), 14 deletions(-) diff --git a/src/Prerequisite.jl b/src/Prerequisite.jl index 86a1f3d..4569d87 100644 --- a/src/Prerequisite.jl +++ b/src/Prerequisite.jl @@ -13,6 +13,11 @@ is what it becomes: its own key space, its own vault, its own payloads. usually about `stale_after`: the shared setup is typically the slow half, and a lock reclaimed mid-build is the thing this exists to prevent. +`stop_flag` and `deadline` are NOT overridden by omission. They say when the job must stop rather +than how this stage runs, so leaving either unset here inherits the caller's; setting one takes +precedence as any other field does. Without that, raising `stale_after` alone silently dropped the +caller's deadline and let the barrier outlive the allocation running it. + Build `keys` by projecting the dependent key space onto the axes the setup actually depends on (`ParamIO.project`), so the two spaces cannot drift apart by hand. """ @@ -54,7 +59,7 @@ point; a dead one is bounded by `stale_after`, after which its lock is reclaimab function run_prerequisite!( p::Prerequisite; opts::RunOpts=RunOpts(), load=nothing, poll::Real=30.0 ) - o = p.opts === nothing ? opts : p.opts + o = _merged_opts(p, opts) n_done = 0 waited = 0 rounds = 0 @@ -96,7 +101,9 @@ function run_prerequisite!( done=n_done, waited=waited, rounds=rounds, - stopped_by=nothing, + # A round that handed out nothing may have been cut short rather than empty, + # and those are different failures: one retries, the other will not. + stopped_by=r.stopped_by, ) end waited += 1 @@ -107,4 +114,21 @@ end _n_undone(p::Prerequisite) = count(k -> !DataVault.is_done(p.vault, k), p.keys) +# `p.opts` replaces the stage's knobs wholesale. `stop_flag` and `deadline` are not stage knobs: +# they bound the JOB. An unset one therefore inherits the caller's rather than reverting to the +# `RunOpts` default, which for `deadline` is `nothing` — no bound at all. +function _merged_opts(p::Prerequisite, opts::RunOpts)::RunOpts + p.opts === nothing && return opts + o = p.opts + return RunOpts(; + workers=o.workers, + max_attempts=o.max_attempts, + stale_after=o.stale_after, + heartbeat_interval=o.heartbeat_interval, + stop_flag=o.stop_flag === nothing ? opts.stop_flag : o.stop_flag, + log_level=o.log_level, + deadline=o.deadline === nothing ? opts.deadline : o.deadline, + ) +end + export Prerequisite, run_prerequisite! diff --git a/src/Run.jl b/src/Run.jl index 0bb1abd..13282ff 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -206,7 +206,9 @@ left it takes from the group with the most work outstanding, which spreads worke Only affects the `pmap` path; the sequential path already visits keys in order. Returns `(; stage, done, err, busy, gave_up, stop, skipped, total, stopped_by)`. `stopped_by` is -`:flag`, `:deadline`, or `nothing`, so a short stage is attributable without re-reading the clock. +`:flag`, `:deadline`, or `nothing`, recorded when a key was actually held back rather than read +off the clock at return, so a stage that finished everything reports `nothing` even if the deadline +passed while its last key ran. The full-done early exit returns the same field set rather than a shorter one. Contract: @@ -298,6 +300,7 @@ function run!( n_busy = 0 n_gave_up = 0 n_stop = 0 + stop_seen = nothing for (key, outcome) in outcomes if outcome === :lock_busy n_busy += 1 @@ -306,8 +309,12 @@ function run!( elseif outcome === :ok add_complete!(m, key) n_done += 1 - elseif outcome === :stop + elseif outcome === :stop_flag n_stop += 1 + stop_seen = :flag # outranks :deadline, as `_stop_reason` does + elseif outcome === :stop_deadline + n_stop += 1 + stop_seen === nothing && (stop_seen = :deadline) elseif outcome === :gave_up n_gave_up += 1 n_err += 1 @@ -320,7 +327,9 @@ function run!( # concurrent masters don't overwrite each other's completed keys. merge_and_save_manifest!(m) - stopped_by = _stop_reason(opts) + # From what the round actually did, so a stage that finished every key is not attributed to a + # deadline that passed while the last one ran. + stopped_by = stop_seen log_event( log, :stage_done; @@ -388,7 +397,8 @@ Outcome symbols: - `:ok` — `work_fn` succeeded and `mark_done!` was called. - `:error` — single-attempt failure (`opts.max_attempts == 1`). - `:gave_up` — all `opts.max_attempts` attempts failed. -- `:stop` — stop flag detected before work started. +- `:stop_flag` / `:stop_deadline` + — a stop condition held before work started, carrying which one. """ function _run_one_with_lock!( work_fn::Function, @@ -402,8 +412,12 @@ function _run_one_with_lock!( # Early exit if stop flag has been raised (checked by both sequential # and pmap paths, so each worker can bail independently). - if _is_stopped(opts) - return (key, :stop) + # The REASON travels back with the outcome. Asking again after the round answers "is a stop + # condition true now", which is a different question: a flag file can be removed and a deadline + # can pass in between, so the answer names something that did not stop this. + stop = _stop_reason(opts) + if stop !== nothing + return (key, stop === :flag ? :stop_flag : :stop_deadline) end # A lock whose holder can be SHOWN to be gone does not have to wait out `stale_after`. The @@ -504,8 +518,15 @@ function _run_sequential!( opts::RunOpts, ) results = Vector{Tuple{DataKey,Symbol}}() - for key in todo - if _is_stopped(opts) + for (i, key) in enumerate(todo) + # The keys a stop drops are ATTRIBUTED, not silently absent. `pmap` hands every remaining + # key to `_run_one_with_lock!`, which returns a stop outcome for each, so leaving them out + # here made the same stop report a different `stop` count depending on the dispatcher — and + # made `stopped_by` unattributable on this path. + stop = _stop_reason(opts) + if stop !== nothing + sym = stop === :flag ? :stop_flag : :stop_deadline + append!(results, ((k, sym) for k in @view todo[i:end])) break end push!(results, _run_one_with_lock!(work_fn, vault, key, stage, log, opts)) @@ -848,10 +869,13 @@ function run_loop!( # separates "a sibling is working on it" from "the holder is gone". The margin covers the round # that has to follow the expiry to act on it. busy_budget = opts.stale_after + 2 * idle_sleep + stopped = nothing while true - if _is_stopped(opts) - break - end + # Captured at the exit rather than re-read at return. A loop that exhausts + # `max_empty_rounds` sleeps `idle_sleep` between rounds and can cross the deadline while + # doing so, and a flag file removed in the meantime turns a real flag stop into `nothing`. + stopped = _stop_reason(opts) + stopped === nothing || break rounds += 1 result = run!(work_fn, vault, keys; opts=opts, load=load, affinity=affinity) n_done += result.done @@ -873,6 +897,9 @@ function run_loop!( end empty_count += 1 if empty_count >= max_empty_rounds + # The round itself may have been cut short rather than empty, and if so that is why + # there was nothing to do. Its own recorded reason, not a fresh clock read. + stopped = result.stopped_by break end sleep(idle_sleep) @@ -882,7 +909,7 @@ function run_loop!( rounds=rounds, done=n_done, busy=n_busy, - stopped_by=_stop_reason(opts), + stopped_by=stopped, prerequisite=pre, ) end diff --git a/test/run/test_prerequisite.jl b/test/run/test_prerequisite.jl index 8ef5ed4..acc7cbc 100644 --- a/test/run/test_prerequisite.jl +++ b/test/run/test_prerequisite.jl @@ -206,3 +206,27 @@ end @test r.stopped_by === nothing end end + +@testset "prerequisite: its own opts do not drop the caller's deadline" begin + # `p.opts` replaced the dependent stage's `RunOpts` wholesale, and `deadline` defaults to + # `nothing`. So a prerequisite built only to raise `stale_after` — the documented reason to + # pass `opts` at all — ran with NO deadline and could outlive the allocation running it. + with_both() do main, prep, outdir + built = Ref(0) + work = k -> (built[] += 1; Dict{String,Any}("s" => 1)) + keys = ParamIO.expand(prep.spec) + + p = Prerequisite(work, prep, keys; opts=RunOpts(stale_after=3600.0)) + r = run_prerequisite!(p; opts=RunOpts(deadline=time() - 1), poll=0.01) + @test r.complete == false + @test r.stopped_by === :deadline # inherited from the caller + @test built[] == 0 # and no setup was handed out + @test p.opts.stale_after == 3600.0 # while the field it WAS given still applies + + # Control: one set on the prerequisite itself still takes precedence over the caller's. + p2 = Prerequisite(work, prep, keys; opts=RunOpts(deadline=time() + 3600)) + r2 = run_prerequisite!(p2; opts=RunOpts(deadline=time() - 1), poll=0.01) + @test r2.complete == true + @test built[] == length(keys) + end +end diff --git a/test/run/test_run_deadline.jl b/test/run/test_run_deadline.jl index 5fdee43..b178645 100644 --- a/test/run/test_run_deadline.jl +++ b/test/run/test_run_deadline.jl @@ -120,3 +120,55 @@ end end end end + +@testset "deadline: an exhausted loop is not attributed to a deadline that passed meanwhile" begin + # `stopped_by` was re-read from the clock when `run_loop!` returned, not recorded when + # something was actually held back. A round that ran past the deadline and then gave up on + # `max_empty_rounds` therefore reported `:deadline` — a retryable answer — for a key that + # cannot be produced and will fail again on the next allocation. + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec)[1:1] # one key, so the stop never holds one back + deadline = time() + 0.2 + work = k -> (sleep(0.5); error("this key cannot be produced")) + r = run_loop!( + work, + v, + keys; + opts=RunOpts(deadline=deadline, max_attempts=1, workers=:sequential), + max_empty_rounds=1, + idle_sleep=0.01, + ) + @test r.done == 0 + @test time() > deadline # the deadline HAS passed by the time it returns + @test r.stopped_by === nothing # but nothing was ever held back by it + + # Control: when the deadline really does hold keys back, it is still reported. + r2 = run_loop!( + k -> Dict{String,Any}("x" => 1), + v, + ParamIO.expand(v.spec); + opts=RunOpts(deadline=time() - 1), + max_empty_rounds=1, + idle_sleep=0.01, + ) + @test r2.done == 0 + @test r2.stopped_by === :deadline + end +end + +@testset "deadline: both dispatchers attribute the keys a stop dropped" begin + # The sequential path used to `break` and emit no outcome for the remaining keys, while `pmap` + # handed every one of them back as stopped. The same stop reported a different `stop` count + # depending on which dispatcher ran, and on this path left `stopped_by` unattributable. + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec) + @test length(keys) > 1 + deadline = time() + 0.3 + work = k -> (sleep(0.5); Dict{String,Any}("x" => 1)) + r = run!(work, v, keys; opts=RunOpts(workers=:sequential, deadline=deadline)) + @test r.done == 1 # the first key ran to completion + @test r.stop == length(keys) - 1 # every later key is accounted for, not absent + @test r.done + r.stop + r.err + r.busy == length(keys) + @test r.stopped_by === :deadline + end +end From 6b17914046f5bcc1211a753a513f23b61f211d9e Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Wed, 16 Sep 2026 07:05:11 +0000 Subject: [PATCH 6/7] test: run the project(...) recipe the Prerequisite docstring recommends MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `Prerequisite`'s docstring says to build `keys` by projecting the dependent key space with `ParamIO.project`, "so the two spaces cannot drift apart by hand". Every test in the file passed `DataVault.keys(prep)` instead — the hand-written projection that IS prep.toml — so the recommended recipe had never been executed, and the fixture it would replace was the thing standing in for it. The recipe reproduces the fixture exactly (`Set(projected) == Set(DataVault.keys(prep))`) and runs end to end. A control grows the sweep's N axis and shows the projection follows while the hand-written file does not, so the equality is not one a projection that ignored the spec could also satisfy. Co-Authored-By: Claude Opus 5 (1M context) --- test/run/test_prerequisite.jl | 36 +++++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/test/run/test_prerequisite.jl b/test/run/test_prerequisite.jl index acc7cbc..ef89c44 100644 --- a/test/run/test_prerequisite.jl +++ b/test/run/test_prerequisite.jl @@ -230,3 +230,39 @@ end @test built[] == length(keys) end end + +@testset "prerequisite: the documented project(...) recipe reproduces the setup key space" begin + # `Prerequisite`'s docstring says to build `keys` by projecting the dependent key space onto + # the axes the setup depends on, "so the two spaces cannot drift apart by hand". Every other + # test in this file passes `DataVault.keys(prep)` — the hand-written projection that IS + # prep.toml — so the recipe the docs recommend was never once executed. + with_both() do main, prep, outdir + projected = ParamIO.expand(ParamIO.project(main.spec, ["N"]; total_samples=1)) + @test Set(projected) == Set(DataVault.keys(prep)) # the hand-written file is the oracle + + builds = joinpath(outdir, "builds.txt") + prep_fn = k -> (_log_build!(builds, "N$(setup_of(k))"); Dict{String,Any}("s" => 1)) + r = run_loop!( + k -> Dict{String,Any}("x" => 1), + main, + DataVault.keys(main); + prerequisite=Prerequisite(prep_fn, prep, projected), + opts=RunOpts(workers=:sequential), + idle_sleep=0.0, + ) + @test r.ran + @test r.prerequisite.complete + @test _n_builds(builds, "N4") == 1 # once per distinct N ... + @test _n_builds(builds, "N8") == 1 + @test length(DataVault.keys(main)) == 4 # ... not once per dependent key + + # Control: the projection tracks a sweep that grows and the hand-written file does not, + # which is the drift the recipe exists to remove. Without it the equality above would also + # hold for a projection that ignored the spec entirely. + grown = joinpath(outdir, "grown.toml") + write(grown, replace(read(_PRE_MAIN, String), "N = [4, 8]" => "N = [4, 8, 12]")) + gproj = ParamIO.expand(ParamIO.project(ParamIO.load(grown), ["N"]; total_samples=1)) + @test length(gproj) == 3 + @test Set(gproj) != Set(DataVault.keys(prep)) + end +end From 3d89e4be54ad1b22b734dfaa0b225d383834a37c Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Wed, 16 Sep 2026 08:07:59 +0000 Subject: [PATCH 7/7] fix(review): a prerequisite could still outlive its allocation, from the other side MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A review of the previous commit found that its own fix reopened the failure it was named for. `_merged_opts` inherited `stop_flag`/`deadline` when the Prerequisite left them unset, but let an explicitly set one win outright. A deadline is not a stage preference: it names when Slurm kills this process tree. A Prerequisite built once with a generous ceiling therefore discarded a caller's tighter, freshly computed one, and `run_prerequisite!`'s busy-wait is unbounded except by that deadline. The previous commit's own "control" asserted this as correct, with the caller's deadline ALREADY EXPIRED and the barrier running to completion anyway. It now takes the tighter of the two. `stop_flag` had a second problem: `=== nothing` is not "unset" for that field, because `RunOpts` resolves its default from `ENV["SWEEPRUNNER_STOP_FLAG"]` — which the shipped `submit_slurm.sh` exports for the whole job. So a Prerequisite built for `stale_after` alone carries a flag it never asked for, and that ambient value outranked the flag the campaign was actually configured with: an operator raising the one they know about would never be seen by the barrier. The caller's now governs whenever it has one. Invisible in CI, which does not set that variable. `run_prerequisite!` hardcoded `complete=false` the moment a stop was observed, computing `remaining` on the same line. A setup finished by an earlier run reported `complete=false, remaining=0`, and `run_loop!` gates the dependent stage on exactly that field: a fully provisioned barrier refused to let the work start. Also: `_merged_opts` rebuilds `RunOpts` by name over `fieldnames`, so a field added later cannot be silently reset to its constructor default; the `:flag`/`:deadline` to `:stop_flag`/`:stop_deadline` map is written once in `_stop_outcome` and THROWS on an unmapped reason, rather than relabelling it as a deadline at each of two hand-written ternaries; `_is_stopped` is removed, having become dead; and the two new wall-clock margins are an order of magnitude wider than the work they bound, since CI here is self-hosted and shares the box. Four mutations, four reds against the intended testset, control green. Co-Authored-By: Claude Opus 5 (1M context) --- src/Prerequisite.jl | 44 +++++++++------ src/Run.jl | 34 ++++++------ test/run/test_prerequisite.jl | 94 +++++++++++++++++++++++++++++--- test/run/test_run_adversarial.jl | 6 +- test/run/test_run_deadline.jl | 12 ++-- 5 files changed, 141 insertions(+), 49 deletions(-) diff --git a/src/Prerequisite.jl b/src/Prerequisite.jl index 4569d87..32fa5ca 100644 --- a/src/Prerequisite.jl +++ b/src/Prerequisite.jl @@ -13,10 +13,13 @@ is what it becomes: its own key space, its own vault, its own payloads. usually about `stale_after`: the shared setup is typically the slow half, and a lock reclaimed mid-build is the thing this exists to prevent. -`stop_flag` and `deadline` are NOT overridden by omission. They say when the job must stop rather -than how this stage runs, so leaving either unset here inherits the caller's; setting one takes -precedence as any other field does. Without that, raising `stale_after` alone silently dropped the -caller's deadline and let the barrier outlive the allocation running it. +`stop_flag` and `deadline` are not stage knobs: they bound the JOB, and a stage may not loosen a +bound the job set. `deadline` therefore takes the TIGHTER of the two, and `stop_flag` takes the +caller's whenever it has one. A prerequisite cannot redirect or outlive either. + +That asymmetry is deliberate for `stop_flag`: [`RunOpts`](@ref) resolves its default from +`ENV["SWEEPRUNNER_STOP_FLAG"]`, so an `opts` here that never mentions `stop_flag` can still carry +one, and `=== nothing` does not mean "the author left it unset" for that field. Build `keys` by projecting the dependent key space onto the axes the setup actually depends on (`ParamIO.project`), so the two spaces cannot drift apart by hand. @@ -67,9 +70,12 @@ function run_prerequisite!( while true stopped = _stop_reason(o) if stopped !== nothing + # A barrier whose keys are all done is satisfied whether or not a stop is pending; + # `complete=false` beside `remaining=0` would gate the dependent stage on nothing. + undone = _n_undone(p) return (; - complete=false, - remaining=_n_undone(p), + complete=undone == 0, + remaining=undone, done=n_done, waited=waited, rounds=rounds, @@ -114,21 +120,25 @@ end _n_undone(p::Prerequisite) = count(k -> !DataVault.is_done(p.vault, k), p.keys) -# `p.opts` replaces the stage's knobs wholesale. `stop_flag` and `deadline` are not stage knobs: -# they bound the JOB. An unset one therefore inherits the caller's rather than reverting to the -# `RunOpts` default, which for `deadline` is `nothing` — no bound at all. +# `nothing` is "no bound", so it loses to any real one. +function _tighter(a::Union{Float64,Nothing}, b::Union{Float64,Nothing}) + return a === nothing ? b : (b === nothing ? a : min(a, b)) +end + +# `p.opts` replaces the stage's knobs wholesale, but not the two that bound the job. Taking the +# tighter deadline rather than the explicit one matters in both directions: reverting it to +# `nothing` lets the barrier outlive the allocation, and so does honouring a more generous ceiling +# the Prerequisite was built with months earlier. function _merged_opts(p::Prerequisite, opts::RunOpts)::RunOpts p.opts === nothing && return opts o = p.opts - return RunOpts(; - workers=o.workers, - max_attempts=o.max_attempts, - stale_after=o.stale_after, - heartbeat_interval=o.heartbeat_interval, - stop_flag=o.stop_flag === nothing ? opts.stop_flag : o.stop_flag, - log_level=o.log_level, - deadline=o.deadline === nothing ? opts.deadline : o.deadline, + job = ( + stop_flag=opts.stop_flag === nothing ? o.stop_flag : opts.stop_flag, + deadline=_tighter(o.deadline, opts.deadline), ) + # Everything else forwarded BY NAME, so a field added to `RunOpts` later cannot be silently + # reset to its constructor default here. + return RunOpts(; (f => get(job, f, getfield(o, f)) for f in fieldnames(RunOpts))...) end export Prerequisite, run_prerequisite! diff --git a/src/Run.jl b/src/Run.jl index 13282ff..fe388df 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -132,7 +132,14 @@ function _stop_reason(opts::RunOpts)::Union{Symbol,Nothing} return nothing end -_is_stopped(opts::RunOpts)::Bool = _stop_reason(opts) !== nothing +# The per-key outcome vocabulary names the same two reasons `_stop_reason` does, for a +# `(key, outcome)` tuple that sits alongside `:ok` / `:error`. Written once, and loudly: a third +# reason added above must fail here rather than be silently relabelled as a deadline. +function _stop_outcome(reason::Symbol)::Symbol + reason === :flag && return :stop_flag + reason === :deadline && return :stop_deadline + return throw(ArgumentError("no per-key outcome for stop reason $(repr(reason))")) +end # How many times a key whose worker DIED is handed to another one. Both dispatchers bound this, # `_run_pmap!` through `pmap`'s `retry_delays` and `_run_affinity!` by counting, and they have to @@ -206,9 +213,8 @@ left it takes from the group with the most work outstanding, which spreads worke Only affects the `pmap` path; the sequential path already visits keys in order. Returns `(; stage, done, err, busy, gave_up, stop, skipped, total, stopped_by)`. `stopped_by` is -`:flag`, `:deadline`, or `nothing`, recorded when a key was actually held back rather than read -off the clock at return, so a stage that finished everything reports `nothing` even if the deadline -passed while its last key ran. +`:flag`, `:deadline`, or `nothing`: a stage that finished every key reports `nothing` even if the +deadline passed while its last key ran, since no key was ever held back by it. The full-done early exit returns the same field set rather than a shorter one. Contract: @@ -398,7 +404,7 @@ Outcome symbols: - `:error` — single-attempt failure (`opts.max_attempts == 1`). - `:gave_up` — all `opts.max_attempts` attempts failed. - `:stop_flag` / `:stop_deadline` - — a stop condition held before work started, carrying which one. + a stop condition held before work started, carrying which one. """ function _run_one_with_lock!( work_fn::Function, @@ -412,13 +418,10 @@ function _run_one_with_lock!( # Early exit if stop flag has been raised (checked by both sequential # and pmap paths, so each worker can bail independently). - # The REASON travels back with the outcome. Asking again after the round answers "is a stop - # condition true now", which is a different question: a flag file can be removed and a deadline - # can pass in between, so the answer names something that did not stop this. + # The reason travels back WITH the outcome: a flag file can be removed and a deadline can pass + # before the outcome is read, so re-deriving it later can name something that did not stop this. stop = _stop_reason(opts) - if stop !== nothing - return (key, stop === :flag ? :stop_flag : :stop_deadline) - end + stop === nothing || return (key, _stop_outcome(stop)) # A lock whose holder can be SHOWN to be gone does not have to wait out `stale_after`. The # clear is owner-checked, so it is a no-op if the holder changed since the question was asked. @@ -519,13 +522,12 @@ function _run_sequential!( ) results = Vector{Tuple{DataKey,Symbol}}() for (i, key) in enumerate(todo) - # The keys a stop drops are ATTRIBUTED, not silently absent. `pmap` hands every remaining - # key to `_run_one_with_lock!`, which returns a stop outcome for each, so leaving them out - # here made the same stop report a different `stop` count depending on the dispatcher — and - # made `stopped_by` unattributable on this path. + # The keys a stop drops are ATTRIBUTED, not silently absent, matching `_run_pmap!`, which + # hands every key to `_run_one_with_lock!` regardless. The result vector has one entry per + # `todo` key on either path. stop = _stop_reason(opts) if stop !== nothing - sym = stop === :flag ? :stop_flag : :stop_deadline + sym = _stop_outcome(stop) append!(results, ((k, sym) for k in @view todo[i:end])) break end diff --git a/test/run/test_prerequisite.jl b/test/run/test_prerequisite.jl index ef89c44..ead924d 100644 --- a/test/run/test_prerequisite.jl +++ b/test/run/test_prerequisite.jl @@ -207,10 +207,9 @@ end end end -@testset "prerequisite: its own opts do not drop the caller's deadline" begin - # `p.opts` replaced the dependent stage's `RunOpts` wholesale, and `deadline` defaults to - # `nothing`. So a prerequisite built only to raise `stale_after` — the documented reason to - # pass `opts` at all — ran with NO deadline and could outlive the allocation running it. +@testset "prerequisite: a stage may not loosen a bound the job set" begin + # A Prerequisite's own `opts` must not let the barrier outlive the allocation running it, + # whether by dropping the caller's deadline or by naming a more generous one of its own. with_both() do main, prep, outdir built = Ref(0) work = k -> (built[] += 1; Dict{String,Any}("s" => 1)) @@ -223,19 +222,100 @@ end @test built[] == 0 # and no setup was handed out @test p.opts.stale_after == 3600.0 # while the field it WAS given still applies - # Control: one set on the prerequisite itself still takes precedence over the caller's. + # The TIGHTER one governs, not the explicit one: a Prerequisite built with a generous + # ceiling must not outlive a caller whose allocation is nearly over. p2 = Prerequisite(work, prep, keys; opts=RunOpts(deadline=time() + 3600)) r2 = run_prerequisite!(p2; opts=RunOpts(deadline=time() - 1), poll=0.01) + @test r2.complete == false + @test r2.stopped_by === :deadline + @test built[] == 0 + + # Control the other way round: with no bound from the caller, the prerequisite's applies + # and the work runs, so `false` above is the tighter bound and not a blanket refusal. + r3 = run_prerequisite!(p2; opts=RunOpts(), poll=0.01) + @test r3.complete == true + @test built[] == length(keys) + end +end + +@testset "prerequisite: the caller's stop flag is the one that governs" begin + # `RunOpts` resolves `stop_flag` from ENV["SWEEPRUNNER_STOP_FLAG"], so a Prerequisite built for + # `stale_after` alone can still carry a flag it never asked for. Treating "non-nothing" as + # "explicitly chosen" then let that ambient value outrank the flag the campaign was configured + # with, and an operator raising the one they know about would never be seen by the barrier. + with_both() do main, prep, outdir + built = Ref(0) + work = k -> (built[] += 1; Dict{String,Any}("s" => 1)) + keys = ParamIO.expand(prep.spec) + + ambient = joinpath(outdir, "AMBIENT_FLAG") + callers = joinpath(outdir, "CALLERS_FLAG") + touch(callers) # only the caller's flag is raised + + p = withenv("SWEEPRUNNER_STOP_FLAG" => ambient) do + Prerequisite(work, prep, keys; opts=RunOpts(stale_after=3600.0)) + end + @test p.opts.stop_flag == ambient # carried without ever being asked for + + r = run_prerequisite!(p; opts=RunOpts(stop_flag=callers), poll=0.01) + @test r.complete == false + @test r.stopped_by === :flag # the caller's flag was seen + @test built[] == 0 + + # Control: with the caller's flag cleared the barrier runs, so `:flag` above is that file + # and not a refusal for some other reason. + rm(callers; force=true) + r2 = run_prerequisite!(p; opts=RunOpts(stop_flag=callers), poll=0.01) @test r2.complete == true @test built[] == length(keys) end end +@testset "prerequisite: an already-satisfied barrier is complete, stop pending or not" begin + # `complete` was a hardcoded `false` the moment a stop was seen, computed on the same line as + # `remaining`. A setup built by an earlier run therefore reported `complete=false, remaining=0`, + # and `run_loop!` gates the dependent stage on exactly that field. + with_both() do main, prep, outdir + built = Ref(0) + work = k -> (built[] += 1; Dict{String,Any}("s" => 1)) + keys = ParamIO.expand(prep.spec) + + p = Prerequisite(work, prep, keys) + @test run_prerequisite!(p; opts=RunOpts(), poll=0.01).complete == true + @test built[] == length(keys) + + r = run_prerequisite!(p; opts=RunOpts(deadline=time() - 1), poll=0.01) + @test r.remaining == 0 + @test r.complete == true # and the two agree + @test built[] == length(keys) # nothing was rebuilt + end +end + +@testset "prerequisite: a round cut short is not reported as a genuine dead end" begin + # The no-progress exit propagates the round's own reason. Without it, a deadline that cut the + # round short read as "nobody holds these and they will not appear", which is the answer that + # tells a resubmitter not to bother. + with_both() do main, prep, outdir + keys = ParamIO.expand(prep.spec) + @test length(keys) > 1 + work = k -> (sleep(3.0); error("cannot build this setup")) + p = Prerequisite(work, prep, keys) + r = run_prerequisite!( + p; + opts=RunOpts(deadline=time() + 2.0, max_attempts=1, workers=:sequential), + poll=0.01, + ) + @test r.complete == false + @test r.remaining > 0 + @test r.stopped_by === :deadline # cut short, not a dead end + end +end + @testset "prerequisite: the documented project(...) recipe reproduces the setup key space" begin # `Prerequisite`'s docstring says to build `keys` by projecting the dependent key space onto # the axes the setup depends on, "so the two spaces cannot drift apart by hand". Every other - # test in this file passes `DataVault.keys(prep)` — the hand-written projection that IS - # prep.toml — so the recipe the docs recommend was never once executed. + # test in this file passes `DataVault.keys(prep)`, the hand-written projection that IS + # prep.toml, so the recipe the docs recommend was never once executed. with_both() do main, prep, outdir projected = ParamIO.expand(ParamIO.project(main.spec, ["N"]; total_samples=1)) @test Set(projected) == Set(DataVault.keys(prep)) # the hand-written file is the oracle diff --git a/test/run/test_run_adversarial.jl b/test/run/test_run_adversarial.jl index 99f4a02..acff541 100644 --- a/test/run/test_run_adversarial.jl +++ b/test/run/test_run_adversarial.jl @@ -55,10 +55,8 @@ allk_a(v) = ParamIO.expand(v.spec) # At least one key ran before stop kicked in, at least one was skipped @test result.done >= 1 @test result.done < length(keys) # stop really cut us short - # Remaining keys are either counted as :stop (if we reached their - # dispatch) or just dropped from the todo list (sequential loop - # break). Either way, fewer done+stop than length means some were - # silently un-processed. + # Every remaining key is attributed on both dispatcher paths, so `done + stop` accounts + # for the whole list. @test counter[] == result.done # work_fn only called for done keys return isfile(stop) && rm(stop; force=true) end diff --git a/test/run/test_run_deadline.jl b/test/run/test_run_deadline.jl index b178645..14a0107 100644 --- a/test/run/test_run_deadline.jl +++ b/test/run/test_run_deadline.jl @@ -124,12 +124,14 @@ end @testset "deadline: an exhausted loop is not attributed to a deadline that passed meanwhile" begin # `stopped_by` was re-read from the clock when `run_loop!` returned, not recorded when # something was actually held back. A round that ran past the deadline and then gave up on - # `max_empty_rounds` therefore reported `:deadline` — a retryable answer — for a key that + # `max_empty_rounds` therefore reported `:deadline`, a retryable answer, for a key that # cannot be produced and will fail again on the next allocation. with_vault_d() do v, outdir keys = ParamIO.expand(v.spec)[1:1] # one key, so the stop never holds one back - deadline = time() + 0.2 - work = k -> (sleep(0.5); error("this key cannot be produced")) + # Margins an order of magnitude wider than the work they bound: CI here is self-hosted and + # shares the box, and the only failure mode is a false RED on good code. + deadline = time() + 2.0 + work = k -> (sleep(3.0); error("this key cannot be produced")) r = run_loop!( work, v, @@ -163,8 +165,8 @@ end with_vault_d() do v, outdir keys = ParamIO.expand(v.spec) @test length(keys) > 1 - deadline = time() + 0.3 - work = k -> (sleep(0.5); Dict{String,Any}("x" => 1)) + deadline = time() + 2.0 + work = k -> (sleep(3.0); Dict{String,Any}("x" => 1)) r = run!(work, v, keys; opts=RunOpts(workers=:sequential, deadline=deadline)) @test r.done == 1 # the first key ran to completion @test r.stop == length(keys) - 1 # every later key is accounted for, not absent