From dc33fda8c56cbb5e7016bd3dd1b7f5301fc519ec Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 15 Sep 2026 07:55:48 +0000 Subject: [PATCH] feat: a deadline, a durable per-key claim record, and one BLAS warning Closes #43, #44, #45. ## deadline (#43) `RunOpts(deadline = time() + 25*60)`. The loop stops handing out keys once `time() > deadline`. Both docstrings now state the granularity, because it is the same for both and `stop_flag`'s did not: the flag and the deadline are read BETWEEN keys. A key already in `work_fn` runs to completion, so neither bounds when `run!` returns. What the deadline buys is that it is set in ADVANCE. A flag raised 60 s before the wall clock cannot buy back a key that runs for ten minutes; a deadline can be budgeted as `allocation_end - longest_expected_key - summary_time`. #43's stage died to a SLURM time limit 14 minutes after the flag went up, taking its diagnostics with it. `run!` now reports `stopped_by` as `:flag`, `:deadline` or `nothing`, so a short stage is attributable without re-reading the clock. It is appended at the END of the returned NamedTuple: every reader in the workspace uses named access, but a positional reader would otherwise be silently rebound. The full-done early exit previously returned a SHORTER NamedTuple than the normal path (no `busy`, `gave_up`, `stop`); it now returns the same field set. ## key_acquired (#44) One `:info` event written at acquire, flushed by `log_event`'s open/write/close. I measured #44's premise before building to it, and it is partly wrong. With a work_fn that throws on one axis value (the systematic-defect shape): | | status tree | event log | |---|---|---| | `work_fn` throws | only the successes' `.done` | error x6, gave_up x2 | | `kill -9` mid-key | `.running` SURVIVES | only `stage_start` | So the failures WERE durably recorded, in `events__.jsonl`, and `.running` does survive a SIGKILL. What is genuinely missing is the claim itself: `:key_start` is `:debug` and suppressed by default, so a key acquired and then killed leaves no event naming it. `:key_acquired` closes that at O(keys this master worked), which is the order of `:key_done`, not the O(masters x keys) of `:lock_busy` that the level split exists to control. DataVault has no `.failed` marker at all (only `.done` and `.running`), so #44's "failed: 0" was structurally guaranteed rather than evidence. Left alone here: putting failures in the status tree is a DataVault change. ## BLAS warning (#45) 287 identical lines on a 72-node job, interleaved with the rows of the table above them. Now one `@warn` carrying the count, after the table. The docstring no longer presents the setting as diagnosed. It is a known cause of OpenBLAS segfaults in multi-process Julia, and on this workload ms/step was flat from 1 to 36 threads with no segfault in ~250 keys at blas=16, so the warning reports rather than concludes. The existing test pinned the old text (`r"BLAS threads=4"`). It now asserts what the fix is: three hot workers produce exactly ONE warning naming "3 of 3", and a second testset with no hot worker asserts zero, so `== 1` is not just "the warning is unconditional". Verified locally, targeted files: test_run_deadline.jl (new, 27 assertions) and test_verify_workers_adversarial.jl (rewritten), plus every other file under test/, all green. Co-Authored-By: Claude Opus 5 (1M context) --- Project.toml | 2 +- README.md | 11 +- src/EventLog.jl | 2 + src/InitWorkers.jl | 23 +++- src/Run.jl | 55 +++++++- .../test_verify_workers_adversarial.jl | 40 +++++- test/run/test_run_deadline.jl | 122 ++++++++++++++++++ 7 files changed, 238 insertions(+), 17 deletions(-) create mode 100644 test/run/test_run_deadline.jl diff --git a/Project.toml b/Project.toml index 65922a1..62da787 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] diff --git a/README.md b/README.md index 66c0872..ff8fde8 100644 --- a/README.md +++ b/README.md @@ -37,7 +37,16 @@ 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. diff --git a/src/EventLog.jl b/src/EventLog.jl index 35e9dbf..5d5bab3 100644 --- a/src/EventLog.jl +++ b/src/EventLog.jl @@ -35,6 +35,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`) | 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/Run.jl b/src/Run.jl index e183fea..edb30ca 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,10 @@ 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)`. +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!`. @@ -207,7 +236,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)) @@ -256,6 +295,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 +307,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 +318,7 @@ function run!( stop=n_stop, skipped=length(keys) - length(todo), total=length(keys), + stopped_by=stopped_by, ) end @@ -323,6 +365,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) 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/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