diff --git a/.github/workflows/FormatCheck.yml b/.github/workflows/FormatCheck.yml index 6c151db..e05dc16 100644 --- a/.github/workflows/FormatCheck.yml +++ b/.github/workflows/FormatCheck.yml @@ -1,8 +1,8 @@ name: Format Check # No `paths:` filter. A filtered required check never runs on a pull request that touches no -# `.jl` file, and a required check that never runs stays pending forever — the pull request -# becomes unmergeable rather than exempt. This repository's `main` requires no status checks -# today, so the filter costs nothing yet; it would cost everything the day one is added. +# `.jl` file, and a required check that never runs stays pending forever: the pull request +# becomes unmergeable rather than exempt. `format / format-check` IS required on this +# repository's `main`, so restoring the filter would cost exactly that. on: { pull_request: {} } jobs: { format: { uses: QAtlasHub/.github/.github/workflows/format-check.yml@main } } diff --git a/Project.toml b/Project.toml index 65922a1..a4ee99b 100644 --- a/Project.toml +++ b/Project.toml @@ -1,6 +1,6 @@ name = "SweepRunner" uuid = "be946ad2-3cb3-4b6e-8f7e-4a5ecc3c255b" -version = "0.6.1" +version = "0.6.2" authors = ["sota shimozono "] [deps] @@ -18,6 +18,7 @@ SlurmClusterManager = "c82cd089-7bf7-41d7-976b-6b5d413cbe0a" TOML = "fa267f1f-6049-4f14-aa54-33bafae1ed76" [compat] +Aqua = "0.8" DataVault = "0.7, 0.8" Dates = "1.11" Distributed = "1.11" @@ -28,14 +29,18 @@ Logging = "1.11" ParamIO = "0.3, 0.4" Printf = "1.11" SHA = "0.7" +Serialization = "1.11" SlurmClusterManager = "0.1, 1" TOML = "1" +Test = "1.11" TestShards = "0.3" julia = "1.11" [extras] +Aqua = "4c88cf16-eb10-579e-8560-4a9242c79595" +Serialization = "9e88b42a-f829-5b0c-bbe9-9e923198166b" Test = "8dfed614-e22c-5e08-85e1-65c5234f0b40" TestShards = "acceef1d-f5e0-4fe4-a546-818dc56ce7b2" [targets] -test = ["Test", "TestShards"] +test = ["Aqua", "Serialization", "Test", "TestShards"] diff --git a/README.md b/README.md index 66c0872..0debe5b 100644 --- a/README.md +++ b/README.md @@ -37,12 +37,32 @@ and the store from [DataVault.jl](https://github.com/QAtlasHub/DataVault.jl). (a single `manifest.jld2` read), not O(N) per-key `.done` stats. Benchmark: 3600 keys warm re-run ≈ 3.5 ms. - **Structured events** — JSONL event log atomic across concurrent writers; - per-item `println` is a non-goal, by design. + per-item `println` is a non-goal, by design. Every lock acquisition writes a + flushed `key_acquired` line, so a run that a `kill -9` truncated still says + which keys it had claimed; the status tree cannot, because a key that was + claimed and never finished leaves no `.done` and no `.failed`. +- **A stop flag and a deadline** — `RunOpts(stop_flag=...)` is read between keys + and so is `RunOpts(deadline=time() + 25*60)`. The difference is when you set + it: a deadline is budgeted in advance, so a batch job can subtract its longest + expected key and reserve the tail of its allocation for the summary it needs + to print. Neither interrupts a key already inside `work_fn`; `run!` reports + which one fired as `result.stopped_by`. - **One entry point for all parallel modes** — `init_workers!(mode=:auto)` dispatches to `:sequential` / `:threads` / `:distributed` / `:slurm` depending on environment. - **Pure work functions** — your physics is a plain `(DataKey) -> Dict`, IO/locking/logging live in the runtime. +- **Worker affinity** — `run!(...; affinity = k -> ...)` makes a free worker + prefer a key whose group it has already handled, so worker-local memoisation + of a shared setup is hit instead of reloaded. A preference, never a partition: + no worker idles while a key is pending. Measured, 24 keys over 2 groups on 8 + workers: **8 group changes without it, 0 with.** +- **Prerequisite stages** — `run!` locks the KEY, so work SHARED between keys + has nowhere to live but inside `work_fn`, where every worker that wants a + setup not yet on disk builds it itself. A `Prerequisite` makes that setup its + own key space, run to completion first, with the same locking, resume and + provenance. Measured, 8 concurrent processes over 16 keys sharing 2 setups: + **16 builds inside `work_fn`, 2 with a prerequisite.** ## Quick start @@ -67,6 +87,35 @@ SweepRunner.run!(work_fn, vault, keys) Re-running the same script after completion: `:skip_complete` is logged and the process exits within milliseconds regardless of `length(keys)`. +### Shared setup + +When many keys need one expensive thing, give that thing its own key space: + +```julia +using ParamIO, DataVault, SweepRunner + +spec = ParamIO.load("config.toml") +main = DataVault.Vault("config.toml"; run="dependent") +prep = DataVault.Vault("config.toml"; run="setup") + +# The axes the setup actually depends on. ParamIO.project derives this from the +# same spec, so the two key spaces cannot drift apart by hand. +derived = ParamIO.expand(ParamIO.project(spec, ["system.L", "model.lambda", "thermal.beta"])) + +run_loop!(work_fn, main, ParamIO.expand(spec); + prerequisite = Prerequisite(prep_fn, prep, derived), + affinity = k -> ParamIO.param(k, "system.L"), + opts = RunOpts(deadline = time() + 25*60)) +``` + +`prerequisite` removes the duplicated *build*; `affinity` removes the repeated *load* of what it +built. The second only matters once the first is in place. + +The prerequisite is a **barrier**: `run_loop!` does not start the dependent stage until every setup +key is done, and if one cannot be built it does not start it at all. The dependency is one level +deep and resolved inside `work_fn`, so this is "all of the setup, then all of the dependents", not +a DAG. + ## Phase chaining without `Stage` / `DAG` A dependent stage loads its parent's output inside the work function using diff --git a/docs/src/api.md b/docs/src/api.md index 67f754f..6bb9fc9 100644 --- a/docs/src/api.md +++ b/docs/src/api.md @@ -49,6 +49,13 @@ SweepRunner.run! SweepRunner.run_loop! ``` +## Prerequisite + +```@docs +SweepRunner.Prerequisite +SweepRunner.run_prerequisite! +``` + ## Preflight ```@docs diff --git a/src/EventLog.jl b/src/EventLog.jl index 35e9dbf..dddf52a 100644 --- a/src/EventLog.jl +++ b/src/EventLog.jl @@ -14,18 +14,21 @@ using JSON3 Append-only JSONL event log with a thread-safe per-`EventLog` lock. Each call to [`log_event`](@ref) writes one JSON object as a single line. -Concurrent writes from multiple tasks (within one master) are serialized -through an internal `ReentrantLock`. Concurrent writes from multiple -_processes_ (separate masters) rely on POSIX `O_APPEND` atomicity, which -is guaranteed for single `write` syscalls of length `< PIPE_BUF` (4 KiB); -`log_event` composes each line as a single `String` and issues exactly -one `write(io, line)` call to stay within that guarantee. +Concurrent writes from multiple tasks are serialized through a +`ReentrantLock` held per PATH, so several `EventLog` objects on one file +share it. Concurrent writes from multiple _processes_ (separate masters) +rely on POSIX `O_APPEND` atomicity, which is guaranteed for single `write` +syscalls of length `< PIPE_BUF` (4 KiB); `log_event` composes each line as +a single `String` and issues one `write` to an unbuffered descriptor to +stay within that guarantee. # Fields - `path::String` — target JSONL file. Parent directory is created lazily on first [`log_event`](@ref). -- `lock::ReentrantLock` — protects appends from same-process races. +- `lock::ReentrantLock` — the per-path lock at construction time. [`log_event`](@ref) resolves the + lock from `path` rather than reading this field, so an `EventLog` that arrived on a worker by + deserialization (which skips the constructor) still serialises against its siblings there. # Event kinds used by `run!` @@ -35,6 +38,8 @@ one `write(io, line)` call to stay within that guarantee. | :-------------- | :---------------------------------------------------------------- | | `stage_start` | once at the top of `run!` when `todo` is non-empty | | `stage_done` | once at the bottom of `run!` when `todo` was non-empty | +| `key_acquired` | the per-key lock was taken (includes `acq`); the only durable | +| | record of a claim, since a SIGKILL skips every later event | | `key_start` | before each `work_fn(key)` attempt (includes `attempt` field) | | `key_done` | after a successful `work_fn(key)` (includes `secs`, `attempt`) | | `lock_busy` | another master holds the `.running` lock (acquire = `:busy`) | @@ -44,6 +49,7 @@ one `write(io, line)` call to stay within that guarantee. | `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 | `: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`); @@ -67,7 +73,20 @@ struct EventLog end function EventLog(path::AbstractString; min_level::Symbol=:info) - return EventLog(String(path), ReentrantLock(), _level_value(min_level)) + return EventLog(String(path), _path_lock(path), _level_value(min_level)) +end + +# One lock per PATH, process-wide, NOT one per `EventLog`. `run!` builds a fresh `EventLog` on every +# call, so four concurrent masters in one process hold four objects pointing at one file and a +# per-object lock serialises nothing between them. +const _LOG_LOCKS = Dict{String,ReentrantLock}() +const _LOG_LOCKS_GUARD = ReentrantLock() + +function _path_lock(path::AbstractString)::ReentrantLock + key = abspath(String(path)) + return lock(_LOG_LOCKS_GUARD) do + return get!(ReentrantLock, _LOG_LOCKS, key) + end end # Severity ladder (à la Julia logging). Events below an `EventLog`'s `min_level` @@ -92,9 +111,10 @@ Append one JSON object to `log.path` with fields `ts` (ISO-8601 local time), pairs passed via `kwargs`. The line is built in full (including the trailing newline) as a single -`String`, then written with exactly one `write(io, line)` call inside an -`open(path, "a")` block. This relies on POSIX `O_APPEND` atomicity so that -cross-process writes do not tear each other's lines. +`String` and written with one `write` syscall to an UNBUFFERED append-mode +descriptor. Both halves matter: the syscall is what POSIX `O_APPEND` +atomicity applies to, so cross-process writes do not tear each other's +lines, and an `IOStream` would flush on its own boundaries instead. ```julia log_event(log, :key_done; stage=:phase1, key="N=8;J=1.0;#sample=1", secs=12.3) @@ -116,10 +136,21 @@ function log_event(log::EventLog, kind::Symbol; level::Symbol=:info, kwargs...) # Build the full line with newline so a single `write` is one atomic # append on POSIX (given `O_APPEND` and size < PIPE_BUF). line = string(JSON3.write(rec), '\n') - lock(log.lock) do + # Resolved from the PATH, not taken from `log.lock`. `run!` serialises the `EventLog` to every + # worker, and deserialization rebuilds the struct without running the constructor, so the field + # that arrives on a worker is a private lock that serialises nothing against its siblings. + lock(_path_lock(log.path)) do mkpath(dirname(log.path)) - open(log.path, "a") do io - return write(io, line) + fd = Base.Filesystem.open( + log.path, + Base.Filesystem.JL_O_WRONLY | Base.Filesystem.JL_O_CREAT | + Base.Filesystem.JL_O_APPEND, + 0o644, + ) + try + return write(fd, codeunits(line)) + finally + close(fd) end end return nothing diff --git a/src/InitWorkers.jl b/src/InitWorkers.jl index b62a601..8a8e439 100644 --- a/src/InitWorkers.jl +++ b/src/InitWorkers.jl @@ -184,9 +184,15 @@ end verify_workers!() Probe each Distributed worker for hostname, Julia threads, BLAS threads, -and CPU affinity. Prints a summary table and emits a `@warn` if any -worker has `BLAS.get_num_threads() > 1` (a common cause of OpenBLAS -segfaults in multi-process Julia). +and CPU affinity. Prints a summary table, then ONE `@warn` carrying how many +workers have `BLAS.get_num_threads() > 1`. + +That setting is reported, not diagnosed. It is a known cause of OpenBLAS +segfaults in multi-process Julia, but on a 2-site TDVP workload (10 sites, +chi=20) ms/step was flat from 1 to 36 threads and ~250 completed keys at +`blas=16` produced no segfault, so the warning does not claim the setting is +wrong here. It used to fire per worker: 287 lines on a 72-node allocation, +interleaved with the rows of the table above it. Ported from FiniteTemperature.jl `Parallel/Slurm.jl::print_worker_identities`. """ @@ -225,6 +231,7 @@ function verify_workers!() ) for p in workers() ] + nhot = 0 for f in futures pid, host, nth, blas, cpuset = fetch(f) @printf( @@ -235,10 +242,14 @@ function verify_workers!() blas, cpuset ) - if blas > 1 - @warn "Worker $pid: BLAS threads=$blas > 1 — OpenBLAS segfault risk" - end + blas > 1 && (nhot += 1) end + # One line, after the table rather than interleaved with it. Per worker this was 287 lines on + # a 72-node allocation, which is the table's readability spent on a risk that has not been + # measured on this workload. + nhot > 0 && + @warn "$nhot of $(nworkers()) workers have BLAS threads > 1 (OpenBLAS segfault risk under multi-process Julia). Set OPENBLAS_NUM_THREADS=1 if you hit one." maxlog = + 1 println() flush(stdout) return nothing diff --git a/src/Prerequisite.jl b/src/Prerequisite.jl new file mode 100644 index 0000000..86a1f3d --- /dev/null +++ b/src/Prerequisite.jl @@ -0,0 +1,110 @@ +# Prerequisite — a stage that must finish before the stage that depends on it. + +using DataVault +using ParamIO: DataKey + +""" + Prerequisite(work_fn, vault, keys; opts=nothing) + +Work that a later stage's keys share. Its three fields are `run!`'s three arguments, because that +is what it becomes: its own key space, its own vault, its own payloads. + +`opts` overrides the dependent stage's [`RunOpts`](@ref) for the prerequisite alone, which is +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. + +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. +""" +struct Prerequisite + work_fn::Function + vault::DataVault.Vault + keys::Vector{DataKey} + opts::Union{RunOpts,Nothing} +end + +function Prerequisite( + work_fn::Function, + vault::DataVault.Vault, + keys::AbstractVector{DataKey}; + opts::Union{RunOpts,Nothing}=nothing, +) + return Prerequisite(work_fn, vault, collect(keys), opts) +end + +""" + run_prerequisite!(p; opts=RunOpts(), load=nothing, poll=30.0) -> NamedTuple + +Run `p` until EVERY one of its keys is done, and report whether that happened: + + (; complete, remaining, done, waited, rounds, stopped_by) + +`complete` is the only field a caller has to read. The rest say why not: `remaining` keys are +undone, `waited` counts the rounds spent purely waiting for a sibling master. + +This is a barrier, not a work loop, and the difference is what it does when it has nothing left to +take. [`run_loop!`](@ref) stops after `max_empty_rounds` empty rounds, which is right when the keys +are independent. Here the dependent stage cannot start until the setup exists, so a round that +finds every remaining key locked by a sibling SLEEPS and goes again. + +It terminates on: every key done; no progress AND no key held by a sibling (a genuine failure); +`opts.stop_flag`; `opts.deadline`. A live sibling building a slow setup is waited for, which is the +point; a dead one is bounded by `stale_after`, after which its lock is reclaimable. +""" +function run_prerequisite!( + p::Prerequisite; opts::RunOpts=RunOpts(), load=nothing, poll::Real=30.0 +) + o = p.opts === nothing ? opts : p.opts + n_done = 0 + waited = 0 + rounds = 0 + + while true + stopped = _stop_reason(o) + if stopped !== nothing + return (; + complete=false, + remaining=_n_undone(p), + done=n_done, + waited=waited, + rounds=rounds, + stopped_by=stopped, + ) + end + + rounds += 1 + r = run!(p.work_fn, p.vault, p.keys; opts=o, load=load) + n_done += r.done + + remaining = _n_undone(p) + remaining == 0 && return (; + complete=true, + remaining=0, + done=n_done, + waited=waited, + rounds=rounds, + stopped_by=nothing, + ) + + # Nothing was completed this round. Either a sibling holds what is left, in which case + # waiting IS the work, or nobody does and the remainder will not appear. + if r.done == 0 + if r.busy == 0 + return (; + complete=false, + remaining=remaining, + done=n_done, + waited=waited, + rounds=rounds, + stopped_by=nothing, + ) + end + waited += 1 + sleep(poll) + end + end +end + +_n_undone(p::Prerequisite) = count(k -> !DataVault.is_done(p.vault, k), p.keys) + +export Prerequisite, run_prerequisite! diff --git a/src/Run.jl b/src/Run.jl index e183fea..b2134d4 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -17,7 +17,7 @@ using ParamIO: DataKey, canonical """ RunOpts(; workers=:auto, max_attempts=3, stale_after=600.0, - heartbeat_interval=60.0, stop_flag=nothing) + heartbeat_interval=60.0, stop_flag=nothing, deadline=nothing) Execution options for [`run!`](@ref). @@ -54,6 +54,21 @@ Execution options for [`run!`](@ref). that misspells it gets no error and no graceful stop, only a killed job. Pass `stop_flag=nothing` explicitly to opt out. + **Granularity: the flag is read between keys, not inside one.** A key already + in `work_fn` runs to completion, so the time between raising the flag and + `run!` returning is bounded by the longest key, which the caller usually + cannot predict. +- `deadline::Union{Float64,Nothing} = nothing` — an absolute `time()` past which + no new key is handed out. The same mechanism as `stop_flag` with the same + in-key granularity, and the reason to have both is that a deadline is set in + ADVANCE: a batch job can subtract its longest expected key and the time its + summary needs from the end of its allocation, where a flag raised reactively + 60 s before the wall clock cannot buy back a key that runs for ten minutes. + + ```julia + RunOpts(deadline = time() + 25 * 60) # stop dispatching 5 min before a 30 min job ends + ``` + # Example ```julia @@ -69,6 +84,7 @@ struct RunOpts heartbeat_interval::Float64 stop_flag::Union{String,Nothing} log_level::Symbol + deadline::Union{Float64,Nothing} end function RunOpts(; @@ -78,6 +94,7 @@ function RunOpts(; heartbeat_interval::Real=60.0, stop_flag::Union{String,Nothing}=get(ENV, "SWEEPRUNNER_STOP_FLAG", nothing), log_level::Symbol=:info, + deadline::Union{Real,Nothing}=nothing, ) workers in (:auto, :sequential) || throw( ArgumentError( @@ -103,11 +120,19 @@ function RunOpts(; Float64(heartbeat_interval), stop_flag, log_level, + deadline === nothing ? nothing : Float64(deadline), ) end -# Internal: check if the stop flag has been raised. -_is_stopped(opts::RunOpts)::Bool = opts.stop_flag !== nothing && isfile(opts.stop_flag) +# Why the loop is stopping, so `:stage_done` can say which of the two fired rather than leaving +# a reader to guess from the wall clock. +function _stop_reason(opts::RunOpts)::Union{Symbol,Nothing} + opts.stop_flag !== nothing && isfile(opts.stop_flag) && return :flag + opts.deadline !== nothing && time() > opts.deadline && return :deadline + return nothing +end + +_is_stopped(opts::RunOpts)::Bool = _stop_reason(opts) !== nothing # As of v0.3 the per-key lock lives ENTIRELY in DataVault's `.running` # sentinel — acquired atomically via `DataVault.acquire_running!` @@ -160,6 +185,25 @@ Early skip (todo 10): on startup a stage-level Manifest is loaded. Keys already in the manifest are skipped — when all keys are done, the second run-through takes O(1) filesystem operations regardless of `length(keys)`. +# Affinity + +`affinity` is `key -> value`, and turns the fan-out from "any free worker takes the next key" into +"a free worker PREFERS a key whose `affinity` value it has already handled". Pass it when `work_fn` +memoises something per group in worker-local state, so a worker that stays on a group pays the load +once instead of once per key. + + run!(work_fn, vault, keys; affinity = k -> param(k, "system.L")) + +A preference, not a partition: a worker is never idle while a key is pending, so a 200-key group +does not serialise onto the worker that opened it. When a worker has nothing from its own groups +left it takes from the group with the most work outstanding, which spreads workers over groups. + +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. +The full-done early exit returns the same field set rather than a shorter one. + Contract: - `work_fn` is expected to be a pure function: given a `DataKey`, return a `Dict` payload to persist via `DataVault.save!`. @@ -196,6 +240,7 @@ function run!( keys::AbstractVector{DataKey}; opts::RunOpts=RunOpts(), load=nothing, + affinity=nothing, ) stage = Symbol(vault.run) log_name = "events_$(gethostname())_$(getpid()).jsonl" @@ -207,7 +252,17 @@ function run!( if isempty(todo) log_event(log, :skip_complete; stage=stage, total=length(keys)) - return (stage=stage, done=0, err=0, skipped=length(keys), total=length(keys)) + return ( + stage=stage, + done=0, + err=0, + busy=0, + gave_up=0, + stop=0, + skipped=length(keys), + total=length(keys), + stopped_by=nothing, + ) end log_event(log, :stage_start; stage=stage, total=length(keys), todo=length(todo)) @@ -223,7 +278,11 @@ function run!( _ensure_worker_modules( vcat([:ParamIO, :DataVault, :SweepRunner], _worker_module_names(load)) ) - _run_pmap!(work_fn, vault, todo, stage, log, opts) + if affinity === nothing + _run_pmap!(work_fn, vault, todo, stage, log, opts) + else + _run_affinity!(work_fn, vault, todo, stage, log, opts, affinity) + end else _run_sequential!(work_fn, vault, todo, stage, log, opts) end @@ -256,6 +315,7 @@ function run!( # concurrent masters don't overwrite each other's completed keys. merge_and_save_manifest!(m) + stopped_by = _stop_reason(opts) log_event( log, :stage_done; @@ -267,6 +327,7 @@ function run!( gave_up=n_gave_up, stop=n_stop, skipped=length(keys) - length(todo), + stopped_by=stopped_by === nothing ? nothing : String(stopped_by), ) return ( stage=stage, @@ -277,6 +338,7 @@ function run!( stop=n_stop, skipped=length(keys) - length(todo), total=length(keys), + stopped_by=stopped_by, ) end @@ -323,6 +385,11 @@ function _run_one_with_lock!( end # acq ∈ (:ok, :reclaimed) — we own the lock. + # Written at ACQUIRE, at :info, and flushed by `log_event`'s open/write/close. This is the + # only record that survives a SIGKILL mid-key: the `finally` below cannot run, so nothing + # later in this function gets to say the key was ever claimed. + log_event(log, :key_acquired; stage=stage, key=kstr, acq=String(acq)) + # Re-check after acquisition: another master may have finished this # key between our manifest read and our acquire. if DataVault.is_done(vault, key) @@ -468,6 +535,110 @@ function _run_pmap!( return out end +""" + _run_affinity!(work_fn, vault, todo, stage, log, opts, affinity) -> Vector{Tuple{DataKey,Symbol}} + +[`_run_pmap!`](@ref) with a PREFERENCE for keys whose `affinity` value the worker has already +handled. A free worker takes a pending key from its most recently used group if one is left, and +otherwise from the group with the most work outstanding, which spreads workers over groups instead +of piling them onto one. + +A preference, never a partition. A worker is never idle while a key is pending, so a 200-key group +does not serialise onto the worker that opened it. + +`pmap` is not used here because it hands out work itself. Its `ProcessExitedException` re-dispatch +is reproduced: a key whose worker died goes back on the queue and the worker is dropped. +""" +function _run_affinity!( + work_fn::Function, + vault::Vault, + todo::AbstractVector{DataKey}, + stage::Symbol, + log::EventLog, + opts::RunOpts, + affinity::Function, +) + groups = Any[affinity(k) for k in todo] + by_group = Dict{Any,Vector{Int}}() + for (i, g) in enumerate(groups) + push!(get!(Vector{Int}, by_group, g), i) + end + # `pop!` takes from the end, so reverse to hand keys out in the caller's order. That order is + # load-bearing: a leading paramset is how a long acquisition is told which slice to close first. + for v in values(by_group) + reverse!(v) + end + + out = Vector{Tuple{DataKey,Symbol}}(undef, length(todo)) + # `Vector{Bool}`, not `BitVector`: adjacent bits share a word, so two tasks marking neighbouring + # indices would read-modify-write the same one. + filled = fill(false, length(todo)) + q = ReentrantLock() + seen = Dict{Int,Vector{Any}}() + + function _take!(pid::Int) + return lock(q) do + mine = get!(Vector{Any}, seen, pid) + for (j, g) in enumerate(mine) + v = get(by_group, g, nothing) + if v !== nothing && !isempty(v) + j == 1 || (deleteat!(mine, j); pushfirst!(mine, g)) + return pop!(v) + end + end + best, bestn = nothing, 0 + for (g, v) in by_group + length(v) > bestn && ((best, bestn) = (g, length(v))) + end + best === nothing && return nothing + pushfirst!(mine, best) + return pop!(by_group[best]) + end + end + + _give_back!(i::Int) = lock(q) do + return push!(get!(Vector{Int}, by_group, groups[i]), i) + end + + @sync for pid in workers() + @async while true + i = _take!(pid) + i === nothing && break + key = todo[i] + res = try + remotecall_fetch( + _run_one_with_lock!, pid, work_fn, vault, key, stage, log, opts + ) + catch e + if e isa ProcessExitedException + _give_back!(i) + break + end + log_event( + log, + :error; + stage=stage, + key=canonical(key), + attempt=0, + err=_short_err(e), + ) + (key, :error) + end + out[i] = res + filled[i] = true + end + end + + # Every worker died while keys were still pending. Those keys were never attempted, so they are + # retriable rather than failed: `:lock_busy` is the outcome `run!` already counts that way. + for i in eachindex(todo) + filled[i] && continue + log_event(log, :worker_lost; stage=stage, key=canonical(todo[i])) + out[i] = (todo[i], :lock_busy) + end + return out +end + """ _run_one_with_retry!(work_fn, vault, key, kstr, stage, log, opts, lost) -> Symbol @@ -534,8 +705,8 @@ function _short_err(e)::String end """ - run_loop!(work_fn, vault, keys; opts=RunOpts(), - max_empty_rounds=3, idle_sleep=30.0, load=nothing) + run_loop!(work_fn, vault, keys; opts=RunOpts(), max_empty_rounds=3, + idle_sleep=30.0, load=nothing, prerequisite=nothing) -> NamedTuple Work-stealing loop that repeatedly calls [`run!`](@ref) until there is no more work to do. This is the infra equivalent of FiniteTemperature.jl's @@ -543,13 +714,41 @@ more work to do. This is the infra equivalent of FiniteTemperature.jl's The loop exits when: - `max_empty_rounds` consecutive rounds produce zero new completions, or -- `opts.stop_flag` is raised (graceful shutdown). +- `opts.stop_flag` is raised, or `opts.deadline` has passed. Default parameters (`max_empty_rounds=3`, `idle_sleep=30.0`) are the battle-tested values from FiniteTemperature.jl. `load` is forwarded verbatim to every [`run!`](@ref) call (see its docstring) — name the work module(s) the workers need and the loop handles the per-round broadcast. + +# Prerequisite + +`run!` locks the KEY, so no two workers compute the same key. Work shared BETWEEN keys has to live +inside `work_fn`, and there it has no protection at all: every worker that wants a setup not yet on +disk builds it itself. + +Pass a [`Prerequisite`](@ref) and that setup becomes its own key space, run to completion by +[`run_prerequisite!`](@ref) before the dependent stage starts. It then gets the same locking, +resume and provenance as any other stage, and its cost is recorded in its own payload instead of +landing on whichever dependent key happened to run first. + + run_loop!(work_fn, vault, keys; + prerequisite = Prerequisite(prep_fn, prep_vault, derived_keys), + opts = opts) + +If the prerequisite does not complete, the dependent stage does NOT start, and the returned +`prerequisite` field says why. Running it anyway would spend the allocation on keys whose setup is +known to be missing. + +**SweepRunner does not know which dependent key needs which prerequisite key.** The dependency is +one level deep and resolved inside `work_fn`, so this is "all of the prerequisite, then all of the +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. """ function run_loop!( work_fn::Function, @@ -559,13 +758,27 @@ function run_loop!( max_empty_rounds::Int=3, idle_sleep::Float64=30.0, load=nothing, + prerequisite=nothing, + affinity=nothing, ) + pre = nothing + 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 + ) + end + empty_count = 0 + rounds = 0 + n_done = 0 while true if _is_stopped(opts) break end - result = run!(work_fn, vault, keys; opts=opts, load=load) + rounds += 1 + result = run!(work_fn, vault, keys; opts=opts, load=load, affinity=affinity) + n_done += result.done if result.done > 0 empty_count = 0 continue @@ -576,7 +789,13 @@ function run_loop!( end sleep(idle_sleep) end - return nothing + return (; + ran=true, + rounds=rounds, + done=n_done, + stopped_by=_stop_reason(opts), + prerequisite=pre, + ) end export RunOpts, run!, run_loop!, manifest_root, load_manifest diff --git a/src/SweepRunner.jl b/src/SweepRunner.jl index fbb9170..79c98ab 100644 --- a/src/SweepRunner.jl +++ b/src/SweepRunner.jl @@ -69,6 +69,7 @@ include("EventLog.jl") include("Manifest.jl") include("InitWorkers.jl") include("Run.jl") +include("Prerequisite.jl") include("Preflight.jl") end # module SweepRunner diff --git a/test/aqua/test_aqua.jl b/test/aqua/test_aqua.jl new file mode 100644 index 0000000..c63e2f8 --- /dev/null +++ b/test/aqua/test_aqua.jl @@ -0,0 +1,13 @@ +# Aqua — the package-level checks no unit test covers, because they are about the PACKAGE rather +# than about any one function: an exported name with no binding, a method ambiguity, a declared +# dependency nothing uses, an extra with no compat bound. +# +# `undefined_exports` is the one that earns its place here: every feature lands by adding a name to +# the `export` list, and a typo there is invisible until a caller reaches for it. + +using SweepRunner, Test +using Aqua + +@testset "Aqua" begin + Aqua.test_all(SweepRunner) +end diff --git a/test/eventlog/test_eventlog_shared_path.jl b/test/eventlog/test_eventlog_shared_path.jl new file mode 100644 index 0000000..9d45123 --- /dev/null +++ b/test/eventlog/test_eventlog_shared_path.jl @@ -0,0 +1,114 @@ +# Several EventLog objects on ONE path: `run!` builds a fresh one per call, so concurrent masters +# in a single process hold several objects pointing at the same file. + +using SweepRunner, Test, JSON3, Serialization + +@testset "EventLog: objects on one path share a lock" begin + dir = mktempdir() + try + p = joinpath(dir, "e.jsonl") + a, b = EventLog(p), EventLog(p) + @test a.lock === b.lock + # A different path must NOT share it, or the lock is just a global mutex. + @test EventLog(joinpath(dir, "other.jsonl")).lock !== a.lock + # The path is normalised, so a relative spelling of the same file still shares. + @test EventLog(relpath(p, pwd())).lock === a.lock + finally + rm(dir; recursive=true, force=true) + end +end + +@testset "EventLog: a deserialized log serialises against its siblings" begin + # `run!` sends the EventLog to every worker, and deserialization skips the constructor, so the + # `lock` field that arrives there is private. log_event resolves the lock from the path. + dir = mktempdir() + try + p = joinpath(dir, "e.jsonl") + a = EventLog(p) + io = IOBuffer() + serialize(io, a) + seekstart(io) + b = deserialize(io) + @test b.lock !== a.lock # the field really does not survive + + n = 200 + Threads.@sync begin + Threads.@spawn for j in 1:n + log_event(a, :key_acquired; key="a$(j)_" * "x"^40, acq="ok") + end + Threads.@spawn for j in 1:n + log_event(b, :key_acquired; key="b$(j)_" * "x"^40, acq="ok") + end + end + lines = readlines(p) + @test length(lines) == 2 * n + @test length([JSON3.read(l) for l in lines]) == 2 * n + finally + rm(dir; recursive=true, force=true) + end +end + +@testset "EventLog: concurrent writers on one path lose and tear nothing" begin + # Before the shared lock this lost about a fifth of the events outright and left the rest + # spliced into each other, so the assertion is on COUNT as well as on parseability. + dir = mktempdir() + try + p = joinpath(dir, "e.jsonl") + logs = [EventLog(p) for _ in 1:4] + n = 200 + Threads.@sync for (i, lg) in enumerate(logs) + Threads.@spawn for j in 1:n + log_event( + lg, :key_acquired; stage=:s, key="m$(i)_k$(j)_" * "x"^40, acq="ok" + ) + end + end + + lines = readlines(p) + @test length(lines) == 4 * n + recs = [JSON3.read(l) for l in lines] # throws on a torn line + @test length(recs) == 4 * n + @test all(r -> r.kind == "key_acquired", recs) + @test length(unique(String(r.key) for r in recs)) == 4 * n + finally + rm(dir; recursive=true, force=true) + end +end + +@testset "EventLog: separate processes on one path do not tear either" begin + # The cross-process half, which no lock can cover: it is the single unbuffered `write` plus + # POSIX O_APPEND. Four short-lived processes, one file. + dir = mktempdir() + try + p = joinpath(dir, "e.jsonl") + project = dirname(Base.active_project()) + script = joinpath(dir, "w.jl") + write( + script, + """ + using SweepRunner + lg = EventLog(ARGS[1]) + for j in 1:150 + log_event(lg, :key_acquired; stage=:s, key="p\$(ARGS[2])_k\$(j)_" * "x"^40, acq="ok") + end + """, + ) + procs = [ + run( + pipeline( + `$(Base.julia_cmd()) --project=$project $script $p $i`; + stdout=devnull, + stderr=devnull, + ); + wait=false, + ) for i in 1:4 + ] + foreach(wait, procs) + + lines = readlines(p) + @test length(lines) == 4 * 150 + @test length([JSON3.read(l) for l in lines]) == 4 * 150 + finally + rm(dir; recursive=true, force=true) + end +end diff --git a/test/init_workers/test_verify_workers_adversarial.jl b/test/init_workers/test_verify_workers_adversarial.jl index 606dbd7..0da3854 100644 --- a/test/init_workers/test_verify_workers_adversarial.jl +++ b/test/init_workers/test_verify_workers_adversarial.jl @@ -18,6 +18,7 @@ # ───────────────────────────────────────────────────────────────────────────── using SweepRunner, Test +using Logging using Distributed, LinearAlgebra @testset "verify_workers! survives workers without SweepRunner" begin @@ -63,18 +64,47 @@ using Distributed, LinearAlgebra end end -@testset "verify_workers!: BLAS threads > 1 emits warning" begin +@testset "verify_workers!: BLAS threads > 1 warns ONCE, with the count" begin + # Three workers, all hot. The bug this pins is volume: one warning per worker put 287 lines + # between the rows of the table the function had just printed. nprocs() > 1 && rmprocs(workers()) project = dirname(Base.active_project()) - addprocs(1; exeflags="--project=$project") + addprocs(3; exeflags="--project=$project") try @everywhere workers() Core.eval(Main, :(using LinearAlgebra)) - # Deliberately set a high BLAS thread count on the worker @everywhere workers() LinearAlgebra.BLAS.set_num_threads(4) - # Capture warnings via Test.@test_logs - @test_logs (:warn, r"BLAS threads=4") SweepRunner.verify_workers!() + logger = Test.TestLogger(; min_level=Logging.Warn) + Logging.with_logger(logger) do + return SweepRunner.verify_workers!() + end + warnings = filter(r -> r.level == Logging.Warn, logger.logs) + + @test length(warnings) == 1 + @test occursin("3 of 3 workers", warnings[1].message) + @test occursin("BLAS threads > 1", warnings[1].message) + finally + rmprocs(workers()) + end +end + +@testset "verify_workers!: no warning when no worker is hot" begin + # Control for the testset above: the counter must be able to read zero, or `== 1` there is + # just "the warning is unconditional". + nprocs() > 1 && rmprocs(workers()) + project = dirname(Base.active_project()) + addprocs(2; exeflags="--project=$project") + + try + @everywhere workers() Core.eval(Main, :(using LinearAlgebra)) + @everywhere workers() LinearAlgebra.BLAS.set_num_threads(1) + + logger = Test.TestLogger(; min_level=Logging.Warn) + Logging.with_logger(logger) do + return SweepRunner.verify_workers!() + end + @test isempty(filter(r -> r.level == Logging.Warn, logger.logs)) finally rmprocs(workers()) end diff --git a/test/run/fixtures/affinity.toml b/test/run/fixtures/affinity.toml new file mode 100644 index 0000000..d8f5510 --- /dev/null +++ b/test/run/fixtures/affinity.toml @@ -0,0 +1,13 @@ +# Two groups (L) of twelve keys each: wide enough that a worker staying on one group is visible +# against the hopping the flat hand-out produces. +[study] +project_name = "aff_test" +total_samples = 1 +outdir = "out" + +[datavault] +path_keys = ["L", "J"] + +[[paramsets]] +L = [8, 16] +J = [0.1, 0.2, 0.3, 0.4, 0.5, 0.6, 0.7, 0.8, 0.9, 1.0, 1.1, 1.2] diff --git a/test/run/fixtures/prep.toml b/test/run/fixtures/prep.toml new file mode 100644 index 0000000..4fc89d4 --- /dev/null +++ b/test/run/fixtures/prep.toml @@ -0,0 +1,11 @@ +# The projection of study.toml's key space onto N alone: the axis the shared setup depends on. +[study] +project_name = "pm_test" +total_samples = 1 +outdir = "out" + +[datavault] +path_keys = ["N"] + +[[paramsets]] +N = [4, 8] diff --git a/test/run/test_affinity.jl b/test/run/test_affinity.jl new file mode 100644 index 0000000..5f163bc --- /dev/null +++ b/test/run/test_affinity.jl @@ -0,0 +1,207 @@ +# affinity (#47): a free worker prefers a key whose group it has already handled. + +using SweepRunner, Test, DataVault, ParamIO, JSON3, Distributed + +const _AFF_CFG = joinpath(@__DIR__, "fixtures", "affinity.toml") + +# How many times a worker's stream of keys changes group. That is the quantity affinity exists to +# lower: each change is a group the worker's memo does not have. +function _transitions(pairs) + per = Dict{Int,Vector{Any}}() + for (pid, g) in pairs + push!(get!(Vector{Any}, per, pid), g) + end + return sum(count(i -> v[i] != v[i - 1], 2:length(v)) for v in values(per); init=0) +end + +group_of(k) = ParamIO.param(k, "L") + +# Each key appends " " so the master can reconstruct who ran what, in order. +# +# Written the way `log_event` writes: ONE `write` to an unbuffered O_APPEND descriptor. An +# `open(trace, "a") do io; println(io, ...)` here is a buffered `IOStream`, and across the worker +# processes it drops lines; the tell is a trace one entry short of the keys that actually ran. +function _make_work(trace) + return k -> begin + line = "$(Distributed.myid()) $(ParamIO.param(k, "L"))\n" + fd = Base.Filesystem.open( + trace, + Base.Filesystem.JL_O_WRONLY | Base.Filesystem.JL_O_CREAT | + Base.Filesystem.JL_O_APPEND, + 0o644, + ) + try + write(fd, codeunits(line)) + finally + close(fd) + end + sleep(0.01) + return Dict{String,Any}("x" => 1) + end +end + +function _read_trace(f) + return [ + (parse(Int, split(l)[1]), split(l)[2]) for + l in (isfile(f) ? readlines(f) : String[]) + ] +end + +function with_workers(f, n) + nprocs() > 1 && rmprocs(workers()) + addprocs(n; exeflags="--project=$(dirname(Base.active_project()))") + try + @everywhere workers() Core.eval( + Main, :(using SweepRunner, DataVault, ParamIO, Distributed) + ) + f() + finally + rmprocs(workers()) + end +end + +@testset "affinity: a worker changes group far less often than without it" begin + with_workers(4) do + off_t = on_t = 0 + outdir = mktempdir() + try + trace = joinpath(outdir, "off.txt") + v = DataVault.Vault(_AFF_CFG; run="off", outdir=outdir) + run!(_make_work(trace), v, DataVault.keys(v)) + off = _read_trace(trace) + off_t = _transitions(off) + + trace2 = joinpath(outdir, "on.txt") + v2 = DataVault.Vault(_AFF_CFG; run="on", outdir=outdir) + run!(_make_work(trace2), v2, DataVault.keys(v2); affinity=group_of) + on = _read_trace(trace2) + on_t = _transitions(on) + + # Same work done either way. + @test length(off) == length(DataVault.keys(v)) + @test length(on) == length(DataVault.keys(v2)) + @test all(k -> DataVault.is_done(v2, k), DataVault.keys(v2)) + + @info "affinity group transitions" without = off_t with = on_t + @test on_t < off_t + finally + rm(outdir; recursive=true, force=true) + end + end +end + +@testset "affinity: it is a preference, so no worker sits idle and no group is serialised" begin + with_workers(4) do + outdir = mktempdir() + try + trace = joinpath(outdir, "t.txt") + v = DataVault.Vault(_AFF_CFG; run="pref", outdir=outdir) + run!(_make_work(trace), v, DataVault.keys(v); affinity=group_of) + t = _read_trace(trace) + + nworkers_used = length(unique(pid for (pid, _) in t)) + ngroups = length(unique(g for (_, g) in t)) + @test ngroups == 2 + + # An exclusive partition of 2 groups can never put more than 2 workers to work, so + # more workers than groups IS the statement that this is a preference. The fixture is + # built to make that distinguishable: 4 workers, 2 groups. + @test nworkers_used > ngroups + @test nworkers_used == 4 # and in fact none of them idled + @test length(t) == length(DataVault.keys(v)) + finally + rm(outdir; recursive=true, force=true) + end + end +end + +@testset "affinity: every key runs exactly once, as without it" begin + with_workers(3) do + outdir = mktempdir() + try + trace = joinpath(outdir, "t.txt") + v = DataVault.Vault(_AFF_CFG; run="once", outdir=outdir) + r = run!(_make_work(trace), v, DataVault.keys(v); affinity=group_of) + t = _read_trace(trace) + @test r.done == length(DataVault.keys(v)) + @test r.err == 0 + @test length(t) == length(DataVault.keys(v)) + @test all(k -> DataVault.is_done(v, k), DataVault.keys(v)) + finally + rm(outdir; recursive=true, force=true) + end + end +end + +@testset "affinity: a throwing work_fn is reported, not swallowed" begin + with_workers(2) do + outdir = mktempdir() + try + v = DataVault.Vault(_AFF_CFG; run="bad", outdir=outdir) + r = run!( + k -> error("boom"), + v, + DataVault.keys(v); + opts=RunOpts(max_attempts=1), + affinity=group_of, + ) + @test r.done == 0 + @test r.err == length(DataVault.keys(v)) + @test !any(k -> DataVault.is_done(v, k), DataVault.keys(v)) + finally + rm(outdir; recursive=true, force=true) + end + 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. + 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]) + work = k -> begin + if ParamIO.canonical(k) == victim + 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) + + # It returned rather than hanging, and the loss was not counted as an error. + @test r.done == length(ks) - 1 + @test r.err == 0 + @test !DataVault.is_done(v, ks[1]) + + # 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) + finally + rm(outdir; recursive=true, force=true) + end + end +end + +@testset "affinity: nothing is the default and leaves the pmap path alone" begin + with_workers(2) do + outdir = mktempdir() + try + v = DataVault.Vault(_AFF_CFG; run="default", outdir=outdir) + r = run!(k -> Dict{String,Any}("x" => 1), v, DataVault.keys(v)) + @test r.done == length(DataVault.keys(v)) + finally + rm(outdir; recursive=true, force=true) + end + end +end diff --git a/test/run/test_prerequisite.jl b/test/run/test_prerequisite.jl new file mode 100644 index 0000000..8ef5ed4 --- /dev/null +++ b/test/run/test_prerequisite.jl @@ -0,0 +1,208 @@ +# Prerequisite (#46): shared setup as its own stage instead of inside work_fn. + +using SweepRunner, Test, DataVault, ParamIO, JSON3 + +const _PRE_MAIN = joinpath(@__DIR__, "fixtures", "study.toml") +const _PRE_PREP = joinpath(@__DIR__, "fixtures", "prep.toml") + +function with_both(f) + outdir = mktempdir() + try + main = DataVault.Vault(_PRE_MAIN; run="dependent", outdir=outdir) + prep = DataVault.Vault(_PRE_PREP; run="setup", outdir=outdir) + f(main, prep, outdir) + finally + rm(outdir; recursive=true, force=true) + end +end + +# One line per build, appended. The point of the issue is HOW MANY times a setup gets built. +_log_build!(path, what) = open(path, "a") do io + return println(io, what) +end +_n_builds(path, what) = isfile(path) ? count(==(what), readlines(path)) : 0 + +setup_of(k) = ParamIO.param(k, "N") + +@testset "prerequisite: the setup stage finishes before the dependent stage starts" begin + with_both() do main, prep, outdir + order = joinpath(outdir, "order.txt") + prep_fn = + k -> (_log_build!(order, "prep"); Dict{String,Any}("state" => setup_of(k))) + work_fn = k -> (_log_build!(order, "work"); Dict{String,Any}("x" => 1)) + + pkeys = DataVault.keys(prep) + r = run_loop!( + work_fn, + main, + DataVault.keys(main); + prerequisite=Prerequisite(prep_fn, prep, pkeys), + opts=RunOpts(workers=:sequential), + idle_sleep=0.0, + ) + + @test r.ran + @test r.prerequisite.complete + lines = readlines(order) + @test count(==("prep"), lines) == length(pkeys) + @test count(==("work"), lines) == length(DataVault.keys(main)) + # Every prep precedes every work: the barrier, stated as an ordering. + @test findlast(==("prep"), lines) < findfirst(==("work"), lines) + end +end + +@testset "prerequisite: the setup is built once per setup key, not once per dependent key" begin + with_both() do main, prep, outdir + builds = joinpath(outdir, "builds.txt") + prep_fn = k -> (_log_build!(builds, "N$(setup_of(k))"); Dict{String,Any}("s" => 1)) + work_fn = k -> Dict{String,Any}("x" => 1) + + run_loop!( + work_fn, + main, + DataVault.keys(main); + prerequisite=Prerequisite(prep_fn, prep, DataVault.keys(prep)), + opts=RunOpts(workers=:sequential), + idle_sleep=0.0, + ) + + # 4 dependent keys fall onto 2 setups. + @test length(DataVault.keys(main)) == 4 + @test _n_builds(builds, "N4") == 1 + @test _n_builds(builds, "N8") == 1 + end +end + +@testset "prerequisite: without one, the same setup IS rebuilt per dependent key" begin + # The control. The check-then-build idiom this replaces, in one process: work_fn builds the + # setup when it is not on disk. Without the prerequisite stage the fixture MUST duplicate, or + # the testset above passes for having nothing to prevent. + with_both() do main, prep, outdir + builds = joinpath(outdir, "builds.txt") + cache = joinpath(outdir, "cache") + mkpath(cache) + function work_fn(k) + f = joinpath(cache, "N$(setup_of(k)).state") + if !isfile(f) # check-then-build + _log_build!(builds, "N$(setup_of(k))") + write(f, "1") + end + return Dict{String,Any}("x" => 1) + end + + # Two masters interleaved the way separate processes are: each sees the cache as it was + # before the other wrote, which is the race that made 31 workers produce 5 states. + for k in DataVault.keys(main) + rm(cache; recursive=true, force=true) + mkpath(cache) + work_fn(k) + end + @test _n_builds(builds, "N4") == 2 # duplicated: 2 dependent keys per setup + @test _n_builds(builds, "N8") == 2 + end +end + +@testset "prerequisite: a setup that cannot be built blocks the dependent stage" begin + with_both() do main, prep, outdir + ran = Ref(0) + r = run_loop!( + k -> (ran[] += 1; Dict{String,Any}("x" => 1)), + main, + DataVault.keys(main); + prerequisite=Prerequisite( + k -> error("cannot cool"), prep, DataVault.keys(prep) + ), + opts=RunOpts(workers=:sequential, max_attempts=1), + idle_sleep=0.0, + ) + @test !r.ran + @test ran[] == 0 + @test !r.prerequisite.complete + @test r.prerequisite.remaining == length(DataVault.keys(prep)) + end +end + +@testset "prerequisite: a live sibling's lock is WAITED for, not treated as failure" begin + # run_loop! would stop after max_empty_rounds here; a barrier must not. The sibling is a fresh + # `.running` on one setup key that nothing will ever release, so the wait is ended by the + # deadline, which is what proves it waited rather than returned. + with_both() do main, prep, outdir + pkeys = DataVault.keys(prep) + DataVault.mark_running!(prep, pkeys[1]) # a sibling holds it + built = Ref(0) + + t0 = time() + pre = SweepRunner.run_prerequisite!( + Prerequisite( + k -> (built[] += 1; Dict{String,Any}("s" => 1)), + prep, + pkeys; + opts=RunOpts(workers=:sequential, stale_after=600.0, deadline=time() + 1.5), + ); + poll=0.2, + ) + elapsed = time() - t0 + + @test !pre.complete + @test pre.stopped_by === :deadline + @test pre.waited >= 1 # it slept instead of giving up + @test elapsed >= 1.0 + @test built[] == length(pkeys) - 1 # the other setup key was still built + @test pre.remaining == 1 + end +end + +@testset "prerequisite: opts on the Prerequisite override the dependent stage's" begin + with_both() do main, prep, outdir + p = Prerequisite( + k -> Dict{String,Any}("s" => 1), + prep, + DataVault.keys(prep); + opts=RunOpts(workers=:sequential, stale_after=1234.0), + ) + @test p.opts.stale_after == 1234.0 + @test Prerequisite(k -> Dict{String,Any}(), prep, DataVault.keys(prep)).opts === + nothing + + pre = SweepRunner.run_prerequisite!(p; opts=RunOpts(workers=:sequential), poll=0.0) + @test pre.complete + @test pre.done == length(DataVault.keys(prep)) + end +end + +@testset "prerequisite: an already-complete setup is a no-op, and resume works" begin + with_both() do main, prep, outdir + pkeys = DataVault.keys(prep) + n = Ref(0) + p() = Prerequisite( + k -> (n[] += 1; Dict{String,Any}("s" => 1)), + prep, + pkeys; + opts=RunOpts(workers=:sequential), + ) + first = SweepRunner.run_prerequisite!(p(); poll=0.0) + @test first.complete && first.done == length(pkeys) + @test n[] == length(pkeys) + + second = SweepRunner.run_prerequisite!(p(); poll=0.0) + @test second.complete + @test second.done == 0 # nothing rebuilt + @test n[] == length(pkeys) + end +end + +@testset "prerequisite: run_loop! without one is unchanged, and now reports" begin + with_both() do main, prep, outdir + r = run_loop!( + k -> Dict{String,Any}("x" => 1), + main, + DataVault.keys(main); + opts=RunOpts(workers=:sequential), + idle_sleep=0.0, + ) + @test r.ran + @test r.prerequisite === nothing + @test r.done == length(DataVault.keys(main)) + @test r.stopped_by === nothing + end +end diff --git a/test/run/test_run_deadline.jl b/test/run/test_run_deadline.jl new file mode 100644 index 0000000..5fdee43 --- /dev/null +++ b/test/run/test_run_deadline.jl @@ -0,0 +1,122 @@ +# deadline (#43) and the durable per-key claim record (#44). + +using SweepRunner, Test, DataVault, ParamIO, JSON3 + +const FIXTURE_CFG_D = joinpath(@__DIR__, "fixtures", "study.toml") + +function with_vault_d(f; run::AbstractString="deadline") + outdir = mktempdir() + try + f(DataVault.Vault(FIXTURE_CFG_D; run=run, outdir=outdir), outdir) + finally + rm(outdir; recursive=true, force=true) + end +end + +function _events_d(outdir) + logs = filter(f -> startswith(f, "events_") && endswith(f, ".jsonl"), readdir(outdir)) + return [JSON3.read(l) for f in logs for l in readlines(joinpath(outdir, f))] +end + +_kinds_d(outdir) = [String(e["kind"]) for e in _events_d(outdir)] + +@testset "deadline: a deadline already past hands out no key" begin + with_vault_d() do v, outdir + n = Ref(0) + work = k -> (n[] += 1; Dict{String,Any}("x" => 1)) + r = run!(work, v, ParamIO.expand(v.spec); opts=RunOpts(deadline=time() - 1)) + @test n[] == 0 + @test r.done == 0 + @test r.stopped_by === :deadline + end +end + +@testset "deadline: a deadline in the future does not stop anything" begin + # Control for the testset above: the same call with the deadline moved forward must run every + # key, so `done == 0` there is the deadline and not the fixture. + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec) + n = Ref(0) + work = k -> (n[] += 1; Dict{String,Any}("x" => 1)) + r = run!(work, v, keys; opts=RunOpts(deadline=time() + 3600)) + @test n[] == length(keys) + @test r.done == length(keys) + @test r.stopped_by === nothing + end +end + +@testset "deadline: it stops BETWEEN keys, not inside one" begin + # The documented granularity. A key already in work_fn runs to completion, so the deadline + # bounds when dispatching stops and not when run! returns. + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec) + @test length(keys) > 1 + started = Ref(0) + deadline = time() + 0.3 + work = k -> (started[] += 1; sleep(0.6); Dict{String,Any}("x" => 1)) + r = run!(work, v, keys; opts=RunOpts(workers=:sequential, deadline=deadline)) + @test started[] >= 1 # the first key ran + @test started[] < length(keys) # later keys were not handed out + @test time() > deadline # and the return is past it, by that first key + @test r.stopped_by === :deadline + end +end + +@testset "deadline: the flag still wins, and is reported as itself" begin + with_vault_d() do v, outdir + stop = joinpath(outdir, "STOP_NOW") + touch(stop) + r = run!( + k -> Dict{String,Any}("x" => 1), + v, + ParamIO.expand(v.spec); + opts=RunOpts(stop_flag=stop, deadline=time() + 3600), + ) + @test r.done == 0 + @test r.stopped_by === :flag + end +end + +@testset "deadline: RunOpts accepts an Int, and nothing is the default" begin + @test RunOpts().deadline === nothing + @test RunOpts(; deadline=1).deadline === 1.0 +end + +@testset "key_acquired: every key this master claimed is on disk, at :info" begin + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec) + run!(k -> Dict{String,Any}("x" => 1), v, keys) + ev = _events_d(outdir) + acquired = [e for e in ev if String(e["kind"]) == "key_acquired"] + @test length(acquired) == length(keys) + @test Set(String(e["key"]) for e in acquired) == + Set(ParamIO.canonical(k) for k in keys) + @test all(String(e["acq"]) in ("ok", "reclaimed") for e in acquired) + # :info, not :debug: the default log level must carry it, which is the whole point. + @test "key_start" ∉ _kinds_d(outdir) + end +end + +@testset "key_acquired: a key whose work throws is still recorded as claimed" begin + # The systematic-failure shape of #44: one axis value throws on every key. The status tree + # shows only the successes, so the claim record is what makes the gap attributable. + with_vault_d() do v, outdir + keys = ParamIO.expand(v.spec) + bad = k -> ParamIO.param(k, "N") == 8 ? error("boom") : Dict{String,Any}("x" => 1) + r = run!(bad, v, keys; opts=RunOpts(max_attempts=1)) + @test r.err > 0 + @test r.done > 0 + + ev = _events_d(outdir) + acquired = Set(String(e["key"]) for e in ev if String(e["kind"]) == "key_acquired") + @test length(acquired) == length(keys) + + failed = [k for k in keys if ParamIO.param(k, "N") == 8] + @test !isempty(failed) + for k in failed + kc = ParamIO.canonical(k) + @test kc ∈ acquired # claimed + @test !DataVault.is_done(v, k) # and left nothing in the status tree + end + end +end