Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Project.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
name = "SweepRunner"
uuid = "be946ad2-3cb3-4b6e-8f7e-4a5ecc3c255b"
version = "0.6.1"
version = "0.6.2"
authors = ["sota shimozono <shimozono-sota631@g.ecc.u-tokyo.ac.jp>"]

[deps]
Expand Down
11 changes: 10 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 2 additions & 0 deletions src/EventLog.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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`) |
Expand Down
23 changes: 17 additions & 6 deletions src/InitWorkers.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
"""
Expand Down Expand Up @@ -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(
Expand All @@ -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
Expand Down
55 changes: 51 additions & 4 deletions src/Run.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down Expand Up @@ -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
Expand All @@ -69,6 +84,7 @@ struct RunOpts
heartbeat_interval::Float64
stop_flag::Union{String,Nothing}
log_level::Symbol
deadline::Union{Float64,Nothing}
end

function RunOpts(;
Expand All @@ -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(
Expand All @@ -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!`
Expand Down Expand Up @@ -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!`.
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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;
Expand All @@ -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,
Expand All @@ -277,6 +318,7 @@ function run!(
stop=n_stop,
skipped=length(keys) - length(todo),
total=length(keys),
stopped_by=stopped_by,
)
end

Expand Down Expand Up @@ -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)
Expand Down
40 changes: 35 additions & 5 deletions test/init_workers/test_verify_workers_adversarial.jl
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
# ─────────────────────────────────────────────────────────────────────────────

using SweepRunner, Test
using Logging
using Distributed, LinearAlgebra

@testset "verify_workers! survives workers without SweepRunner" begin
Expand Down Expand Up @@ -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
Expand Down
122 changes: 122 additions & 0 deletions test/run/test_run_deadline.jl
Original file line number Diff line number Diff line change
@@ -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
Loading