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/Project.toml b/Project.toml index a4ee99b..9ed006e 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] @@ -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..3013dfe 100644 --- a/src/EventLog.jl +++ b/src/EventLog.jl @@ -44,12 +44,17 @@ 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 | +| `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 new file mode 100644 index 0000000..94aa9d6 --- /dev/null +++ b/src/Liveness.jl @@ -0,0 +1,155 @@ +# 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. + +""" + owner_token() -> String + +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 = _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 + +`: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 + + # 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 + + 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_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 + # 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 + any(j -> startswith(j, jobid * "_"), live) && return :alive + return :dead +end + +# 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} + return lock(_squeue_lock) do + t, cached = _squeue_cache[] + time() - t < _SQUEUE_TTL && return cached + fresh = try + 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) + return fresh + end +end + +export owner_token, holder_liveness diff --git a/src/Prerequisite.jl b/src/Prerequisite.jl index 86a1f3d..32fa5ca 100644 --- a/src/Prerequisite.jl +++ b/src/Prerequisite.jl @@ -13,6 +13,14 @@ 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 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. """ @@ -54,7 +62,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 @@ -62,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, @@ -96,7 +107,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 +120,25 @@ end _n_undone(p::Prerequisite) = count(k -> !DataVault.is_done(p.vault, k), p.keys) +# `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 + 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 b2134d4..fe388df 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -132,7 +132,19 @@ 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 +# 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!` @@ -201,7 +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`, so a short stage is attributable without re-reading the clock. +`: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: @@ -293,6 +306,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 @@ -301,8 +315,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 @@ -315,7 +333,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; @@ -342,6 +362,32 @@ 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 + # An `isfile` first: `running_owner` opens and reads, and the uncontended case is every key. + DataVault.is_running(vault, key) || return false + # 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 + """ _run_one_with_lock!(work_fn, vault, key, stage, log, opts) -> (DataKey, Symbol) @@ -357,7 +403,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, @@ -371,14 +418,20 @@ 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) - end + # 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) + 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. + _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 +446,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 +456,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 +476,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 +503,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) @@ -470,8 +521,14 @@ 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, 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_outcome(stop) + 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)) @@ -511,7 +568,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) @@ -596,8 +653,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 = _WORKER_DEATH_REDISPATCHES + 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 +677,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( @@ -713,9 +792,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 +834,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,26 +853,55 @@ 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 + 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 + 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 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) @@ -793,7 +910,8 @@ function run_loop!( ran=true, rounds=rounds, done=n_done, - stopped_by=_stop_reason(opts), + busy=n_busy, + stopped_by=stopped, prerequisite=pre, ) end 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/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 5f163bc..df65626 100644 --- a/test/run/test_affinity.jl +++ b/test/run/test_affinity.jl @@ -155,38 +155,87 @@ 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 + # 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 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..95a1896 --- /dev/null +++ b/test/run/test_liveness.jl @@ -0,0 +1,318 @@ +# 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, JSON3 + +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(["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_3") === :alive # array TASK, as squeue names it + @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 "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 + 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 "_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 + _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 + +@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("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 + 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 "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 + # 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_prerequisite.jl b/test/run/test_prerequisite.jl index 8ef5ed4..ead924d 100644 --- a/test/run/test_prerequisite.jl +++ b/test/run/test_prerequisite.jl @@ -206,3 +206,143 @@ end @test r.stopped_by === nothing end end + +@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)) + 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 + + # 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. + 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 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 5fdee43..14a0107 100644 --- a/test/run/test_run_deadline.jl +++ b/test/run/test_run_deadline.jl @@ -120,3 +120,57 @@ 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 + # 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, + 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() + 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 + @test r.done + r.stop + r.err + r.busy == length(keys) + @test r.stopped_by === :deadline + end +end diff --git a/test/run/test_run_loop_busy.jl b/test/run/test_run_loop_busy.jl new file mode 100644 index 0000000..75e8500 --- /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`. +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) + 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