diff --git a/Project.toml b/Project.toml index 62da787..080fc5a 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] @@ -28,14 +28,16 @@ Logging = "1.11" ParamIO = "0.3, 0.4" Printf = "1.11" SHA = "0.7" +Serialization = "1.11" SlurmClusterManager = "0.1, 1" TOML = "1" TestShards = "0.3" julia = "1.11" [extras] +Serialization = "9e88b42a-f829-5b0c-bbe9-9e923198166b" Test = "8dfed614-e22c-5e08-85e1-65c5234f0b40" TestShards = "acceef1d-f5e0-4fe4-a546-818dc56ce7b2" [targets] -test = ["Test", "TestShards"] +test = ["Serialization", "Test", "TestShards"] diff --git a/README.md b/README.md index ff8fde8..cad192f 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,12 @@ and the store from [DataVault.jl](https://github.com/QAtlasHub/DataVault.jl). depending on environment. - **Pure work functions** — your physics is a plain `(DataKey) -> Dict`, IO/locking/logging live in the runtime. +- **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 @@ -76,6 +82,31 @@ 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), + opts = RunOpts(deadline = time() + 25*60)) +``` + +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 5d5bab3..9eb0baf 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!` @@ -69,7 +72,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` @@ -94,9 +110,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) @@ -118,10 +135,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/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 edb30ca..4c31ba7 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -581,8 +581,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 @@ -590,13 +590,39 @@ 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. + +Returns `(; ran, rounds, done, stopped_by, prerequisite)`. `ran` is `false` exactly when a +prerequisite blocked the stage. """ function run_loop!( work_fn::Function, @@ -606,13 +632,26 @@ function run_loop!( max_empty_rounds::Int=3, idle_sleep::Float64=30.0, load=nothing, + prerequisite=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 + rounds += 1 result = run!(work_fn, vault, keys; opts=opts, load=load) + n_done += result.done if result.done > 0 empty_count = 0 continue @@ -623,7 +662,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/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/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_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