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
6 changes: 4 additions & 2 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.2"
version = "0.6.3"
authors = ["sota shimozono <shimozono-sota631@g.ecc.u-tokyo.ac.jp>"]

[deps]
Expand Down Expand Up @@ -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"]
31 changes: 31 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Expand Down
7 changes: 7 additions & 0 deletions docs/src/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,13 @@ SweepRunner.run!
SweepRunner.run_loop!
```

## Prerequisite

```@docs
SweepRunner.Prerequisite
SweepRunner.run_prerequisite!
```

## Preflight

```@docs
Expand Down
56 changes: 42 additions & 14 deletions src/EventLog.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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!`

Expand Down Expand Up @@ -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`
Expand All @@ -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)
Expand All @@ -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
Expand Down
110 changes: 110 additions & 0 deletions src/Prerequisite.jl
Original file line number Diff line number Diff line change
@@ -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!
53 changes: 49 additions & 4 deletions src/Run.jl
Original file line number Diff line number Diff line change
Expand Up @@ -581,22 +581,48 @@ 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
`_work_loop` driver.

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,
Expand All @@ -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
Expand All @@ -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
1 change: 1 addition & 0 deletions src/SweepRunner.jl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading
Loading