Skip to content
Merged
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
11 changes: 6 additions & 5 deletions 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.3"
version = "0.6.4"
authors = ["sota shimozono <shimozono-sota631@g.ecc.u-tokyo.ac.jp>"]

[deps]
Expand All @@ -19,17 +19,18 @@ TOML = "fa267f1f-6049-4f14-aa54-33bafae1ed76"

[compat]
Aqua = "0.8"
# 0.8.1 is the floor, not a preference: the owner-stamped `.running` API
# 0.8.2 is the floor, not a preference: `artifact!` / `ArtifactBusy` (what the deferral and
# `artifact_affinity` are built for) arrived there. The owner-stamped `.running` API
# (`new_owner_token`, `running_owner`, and the three-argument `refresh_running!` /
# `clear_running!`) arrived there, and the liveness reaper is built on it.
DataVault = "0.8.1"
# `clear_running!`) that the liveness reaper is built on arrived in 0.8.1.
DataVault = "0.8.2"
Dates = "1.11"
Distributed = "1.11"
JLD2 = "0.6"
JSON3 = "1"
LinearAlgebra = "1.11"
Logging = "1.11"
ParamIO = "0.3, 0.4"
ParamIO = "0.4.11"
Printf = "1.11"
SHA = "0.7"
Serialization = "1.11"
Expand Down
29 changes: 29 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,35 @@ key is done, and if one cannot be built it does not start it at all. The depende
deep and resolved inside `work_fn`, so this is "all of the setup, then all of the dependents", not
a DAG.

### Shared setup without a barrier — artifacts

When the setup depends on a subset of the axes, declaring it in the config
(`[artifacts.<name>] depends_on = [...]`, ParamIO ≥ 0.4.11) lets `work_fn` build it on demand
and every other key reuse it, with no second vault and no barrier:

```julia
work_fn = k -> begin
gs = DataVault.artifact!(vault, :ground_state, k; wait=false) do akey # a miss builds
prepare(param(akey, "system.L"))
end
Dict{String,Any}("x" => respond(gs, k))
end
run!(work_fn, vault, keys; affinity = artifact_affinity(vault, :ground_state))
```

- `artifact_affinity` keeps keys that share an artifact on one worker and starts distinct
artifacts on distinct workers.
- With `wait=false`, a key whose artifact another worker or job is still building throws
`DataVault.ArtifactBusy`; `run!` logs `:artifact_busy`, **defers the key without spending an
attempt**, and re-dispatches it once the pass drains (`:deferred_round`; after a pass that
finished nothing it first waits `RunOpts(defer_poll=30.0)`). A key still deferred when the run
stops is counted with `busy`. With `wait=true` (the default) the worker simply blocks until the
artifact exists — the better choice when there are no more keys than workers.

The artifact lives outside the run (`{outdir}/artifacts/...`), so the next job and the next run
under the same `outdir` reuse it too. `Prerequisite` remains for setups that must be complete
before anything else starts.

## Phase chaining without `Stage` / `DAG`

A dependent stage loads its parent's output inside the work function using
Expand Down
6 changes: 6 additions & 0 deletions docs/src/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,12 @@ SweepRunner.run!
SweepRunner.run_loop!
```

## Artifacts

```@docs
SweepRunner.artifact_affinity
```

## Liveness

```@docs
Expand Down
31 changes: 31 additions & 0 deletions src/Artifacts.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# Artifacts — scheduling around `DataVault.artifact!`.
#
# An artifact is built inside `work_fn` by whichever worker first needs it, and reused by the
# rest (DataVault). What the runtime adds is WHERE keys go: keys that share an artifact are
# better on one worker (it loads once, and no second worker waits on the build), and keys that
# need different ones are better spread (the builds run in parallel).

using ParamIO: artifact_identity

"""
artifact_affinity(vault, name) -> Function

An `affinity` for [`run!`](@ref) that groups keys by the artifact `name` they need — its
`ParamIO.artifact_identity` — so a worker keeps drawing keys that reuse the artifact it just
built or loaded, and distinct artifacts start on distinct workers.

```julia
SweepRunner.run!(work_fn, vault, keys; affinity = artifact_affinity(vault, :ground_state))
```

Pair with `DataVault.artifact!(...; wait=false)` inside `work_fn` when keys outnumber workers:
a key whose artifact is mid-build then throws `DataVault.ArtifactBusy`, `run!` defers it without
spending an attempt, and the worker takes another key instead of blocking.
"""
function artifact_affinity(vault::Vault, name)
spec = vault.spec
n = string(name)
return k -> artifact_identity(spec, n, k)
end

export artifact_affinity
3 changes: 3 additions & 0 deletions src/EventLog.jl
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,9 @@ stay within that guarantee.
| | the re-dispatch bound (includes `deaths`) |
| `worker_lost` | every worker died with keys still queued; this key was left for a |
| | later run rather than completed or failed |
| `artifact_busy` | `work_fn` threw `DataVault.ArtifactBusy`; the key is deferred, no |
| | attempt spent (includes `artifact`) |
| `deferred_round`| `run!` re-dispatches its deferred keys (includes `round`, `keys`) |

`: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`);
Expand Down
63 changes: 55 additions & 8 deletions src/Run.jl
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,10 @@ Execution options for [`run!`](@ref).
RunOpts(deadline = time() + 25 * 60) # stop dispatching 5 min before a 30 min job ends
```

- `defer_poll::Float64 = 30.0` — seconds [`run!`](@ref) waits before re-dispatching keys whose
`work_fn` threw `DataVault.ArtifactBusy` (an artifact being built by another worker or job),
when the previous pass made no progress. A deferred key costs no attempt.

# Example

```julia
Expand All @@ -85,6 +89,7 @@ struct RunOpts
stop_flag::Union{String,Nothing}
log_level::Symbol
deadline::Union{Float64,Nothing}
defer_poll::Float64
end

function RunOpts(;
Expand All @@ -95,6 +100,7 @@ function RunOpts(;
stop_flag::Union{String,Nothing}=get(ENV, "SWEEPRUNNER_STOP_FLAG", nothing),
log_level::Symbol=:info,
deadline::Union{Real,Nothing}=nothing,
defer_poll::Real=30.0,
)
workers in (:auto, :sequential) || throw(
ArgumentError(
Expand All @@ -121,6 +127,7 @@ function RunOpts(;
stop_flag,
log_level,
deadline === nothing ? nothing : Float64(deadline),
Float64(defer_poll),
)
end

Expand Down Expand Up @@ -282,7 +289,8 @@ function run!(

# Dispatch strategy: pmap when Distributed workers are present (unless the
# caller forced `workers=:sequential`), otherwise the sequential loop.
outcomes = if opts.workers !== :sequential && nprocs() > 1
multi = opts.workers !== :sequential && nprocs() > 1
if multi
# Ensure the seam packages (+ the user's work module(s) via `load=`) are loaded in `Main`
# on every worker before fan-out. `init_workers!` spawns workers with `--project` but loads
# no packages, so the first pmap task would otherwise die with a cryptic
Expand All @@ -291,14 +299,15 @@ function run!(
_ensure_worker_modules(
vcat([:ParamIO, :DataVault, :SweepRunner], _worker_module_names(load))
)
if affinity === nothing
_run_pmap!(work_fn, vault, todo, stage, log, opts)
else
_run_affinity!(work_fn, vault, todo, stage, log, opts, affinity)
end
end
dispatch = ks -> if !multi
_run_sequential!(work_fn, vault, ks, stage, log, opts)
elseif affinity === nothing
_run_pmap!(work_fn, vault, ks, stage, log, opts)
else
_run_sequential!(work_fn, vault, todo, stage, log, opts)
_run_affinity!(work_fn, vault, ks, stage, log, opts, affinity)
end
outcomes = _redispatch_deferred(dispatch(todo), dispatch, log, stage, opts)

# Aggregate outcomes into counters + manifest updates.
n_done = 0
Expand All @@ -308,7 +317,7 @@ function run!(
n_stop = 0
stop_seen = nothing
for (key, outcome) in outcomes
if outcome === :lock_busy
if outcome === :lock_busy || outcome === :deferred
n_busy += 1
elseif outcome === :already_done
add_complete!(m, key)
Expand Down Expand Up @@ -718,6 +727,38 @@ function _run_affinity!(
return out
end

"""
_redispatch_deferred(outcomes, dispatch, log, stage, opts) -> outcomes

Re-run the keys whose `work_fn` threw `DataVault.ArtifactBusy` until none is left or the run is
stopped. The first re-dispatch follows a pass that finished something, so the artifact it was
waiting on has usually been built by then and it goes at once; after a pass that finished
nothing, it waits `opts.defer_poll` seconds first — the builder is then another job. Returns the
outcomes in the caller's key order. A key still deferred when the run stops is left as
`:deferred`, which `run!` counts with `busy`: it was never attempted.
"""
function _redispatch_deferred(
outcomes, dispatch, log::EventLog, stage::Symbol, opts::RunOpts
)
final = Dict{DataKey,Symbol}(k => o for (k, o) in outcomes)
progressed = any(o -> o === :ok || o === :already_done, last.(outcomes))
round = 0
while true
deferred = DataKey[k for (k, _) in outcomes if final[k] === :deferred]
isempty(deferred) && break
_stop_reason(opts) === nothing || break
progressed || sleep(opts.defer_poll)
round += 1
log_event(log, :deferred_round; stage=stage, round=round, keys=length(deferred))
res = dispatch(deferred)
progressed = any(r -> last(r) === :ok || last(r) === :already_done, res)
for (k, o) in res
final[k] = o
end
end
return [(k, final[k]) for (k, _) in outcomes]
end

"""
_run_one_with_retry!(work_fn, vault, key, kstr, stage, log, opts, lost) -> Symbol

Expand Down Expand Up @@ -761,6 +802,12 @@ function _run_one_with_retry!(
)
return :ok
catch e
# Not a failure: the artifact this key needs is being built elsewhere. Hand the key
# back without spending an attempt; `run!` re-dispatches it once the pass drains.
if e isa DataVault.ArtifactBusy
log_event(log, :artifact_busy; stage=stage, key=kstr, artifact=e.name)
return :deferred
end
last_err = _short_err(e)
log_event(log, :error; stage=stage, key=kstr, attempt=attempt, err=last_err)
if attempt < opts.max_attempts
Expand Down
1 change: 1 addition & 0 deletions src/SweepRunner.jl
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@ include("Manifest.jl")
include("InitWorkers.jl")
include("Liveness.jl")
include("Run.jl")
include("Artifacts.jl") # artifact_affinity; ArtifactBusy deferral lives in Run.jl
include("Prerequisite.jl")
include("Preflight.jl")

Expand Down
134 changes: 134 additions & 0 deletions test/run/test_artifacts.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
# Artifacts: a key whose artifact is mid-build is deferred, not failed; and keys group by artifact.

using SweepRunner, Test, DataVault, ParamIO, JSON3

function _ar_config(dir)
path = joinpath(dir, "config.toml")
write(
path,
"""
[study]
project_name = "ar_study"
[datavault]
path_keys = ["run.U", "run.omega1"]
[[paramsets]]
[paramsets.run]
U = [0.0, 0.2]
omega1 = [1.0, 2.0, 3.0]
[artifacts.gs]
depends_on = ["U"]
""",
)
return path
end

function with_ar(f)
dir = mktempdir()
try
f(DataVault.Vault(_ar_config(dir); outdir=joinpath(dir, "out"), run="a"), dir)
finally
rm(dir; recursive=true, force=true)
end
end

function _events(dir)
return [
JSON3.read(l) for
f in readdir(joinpath(dir, "out"); join=true) if endswith(f, ".jsonl") for
l in eachline(f) if !isempty(l)
]
end

# Hold the build lock of the artifact `k` needs, as another job mid-build would.
function _hold!(v, k)
lock = joinpath(DataVault.artifact_dir(v, :gs, k), ".building")
tok = new_owner_token()
@assert DataVault._acquire_lock_at!(lock, tok) === :ok
return lock, tok
end

@testset "artifacts: a key whose artifact is mid-build is deferred, then done, at no attempt" begin
with_ar() do v, dir
keys = DataVault.keys(v)
held = first(filter(k -> param(k, "run.U") == 0.2, keys))
lock, tok = _hold!(v, held)
# The "other job" finishes its build once all three U = 0.2 keys have been turned away —
# counted, not timed, so a slow first compile cannot release it before they arrive.
busy = Ref(0)
work_fn =
k -> begin
gs = try
DataVault.artifact!(a -> param(a, "run.U"), v, :gs, k; wait=false)
catch e
e isa DataVault.ArtifactBusy &&
(busy[] += 1) == 3 &&
DataVault._clear_lock_at!(lock, tok)
rethrow()
end
Dict{String,Any}("x" => gs + param(k, "run.omega1"))
end
# max_attempts = 1: a deferral that spent an attempt would end these keys as :error.
r = run!(
work_fn,
v,
keys;
opts=RunOpts(
workers=:sequential, max_attempts=1, defer_poll=0.5, stop_flag=nothing
),
)
@test r.done == length(keys)
@test r.err == 0
@test r.busy == 0
@test all(k -> DataVault.is_done(v, k), keys)
kinds = [String(e.kind) for e in _events(dir)]
@test count(==("artifact_busy"), kinds) >= 3 # the three U = 0.2 keys, at least once
@test "deferred_round" in kinds
end
end

@testset "artifacts: a key still deferred when the run stops is busy, not failed" begin
with_ar() do v, dir
keys = DataVault.keys(v)
held = first(filter(k -> param(k, "run.U") == 0.2, keys))
_hold!(v, held) # never released
flag = joinpath(dir, "STOP")
# Raise the stop flag at the first deferral, as a job reaching its deadline would.
work_fn =
k -> Dict{String,Any}(
"x" => try
DataVault.artifact!(a -> 1.0, v, :gs, k; wait=false)
catch e
e isa DataVault.ArtifactBusy && touch(flag)
rethrow()
end
)
r = run!(
work_fn,
v,
keys;
opts=RunOpts(
workers=:sequential, max_attempts=1, defer_poll=0.3, stop_flag=flag
),
)
@test r.done == 3 # U = 0.0
# U = 0.2: the first is deferred, and counted with `busy`; the sequential path does not
# dispatch the rest after the flag, and (as before this change) counts them nowhere.
@test r.busy >= 1
@test r.err == 0
@test !any(k -> param(k, "run.U") == 0.2 && DataVault.is_done(v, k), keys)
@test r.stopped_by === :flag
end
end

@testset "artifacts: artifact_affinity groups keys by the artifact they need" begin
with_ar() do v, dir
aff = artifact_affinity(v, :gs)
keys = DataVault.keys(v)
groups = Dict{Any,Vector{Float64}}()
for k in keys
push!(get!(groups, aff(k), Float64[]), param(k, "run.U"))
end
@test length(groups) == 2
@test all(us -> length(unique(us)) == 1 && length(us) == 3, values(groups))
end
end
Loading