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
6 changes: 3 additions & 3 deletions .github/workflows/FormatCheck.yml
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
name: Format Check

# No `paths:` filter. A filtered required check never runs on a pull request that touches no
# `.jl` file, and a required check that never runs stays pending forever — the pull request
# becomes unmergeable rather than exempt. This repository's `main` requires no status checks
# today, so the filter costs nothing yet; it would cost everything the day one is added.
# `.jl` file, and a required check that never runs stays pending forever: the pull request
# becomes unmergeable rather than exempt. `format / format-check` IS required on this
# repository's `main`, so restoring the filter would cost exactly that.
on: { pull_request: {} }
jobs: { format: { uses: QAtlasHub/.github/.github/workflows/format-check.yml@main } }
9 changes: 7 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.1"
version = "0.6.2"
authors = ["sota shimozono <shimozono-sota631@g.ecc.u-tokyo.ac.jp>"]

[deps]
Expand All @@ -18,6 +18,7 @@ SlurmClusterManager = "c82cd089-7bf7-41d7-976b-6b5d413cbe0a"
TOML = "fa267f1f-6049-4f14-aa54-33bafae1ed76"

[compat]
Aqua = "0.8"
DataVault = "0.7, 0.8"
Dates = "1.11"
Distributed = "1.11"
Expand All @@ -28,14 +29,18 @@ Logging = "1.11"
ParamIO = "0.3, 0.4"
Printf = "1.11"
SHA = "0.7"
Serialization = "1.11"
SlurmClusterManager = "0.1, 1"
TOML = "1"
Test = "1.11"
TestShards = "0.3"
julia = "1.11"

[extras]
Aqua = "4c88cf16-eb10-579e-8560-4a9242c79595"
Serialization = "9e88b42a-f829-5b0c-bbe9-9e923198166b"
Test = "8dfed614-e22c-5e08-85e1-65c5234f0b40"
TestShards = "acceef1d-f5e0-4fe4-a546-818dc56ce7b2"

[targets]
test = ["Test", "TestShards"]
test = ["Aqua", "Serialization", "Test", "TestShards"]
51 changes: 50 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,12 +37,32 @@ 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.
- **Pure work functions** — your physics is a plain
`(DataKey) -> Dict`, IO/locking/logging live in the runtime.
- **Worker affinity** — `run!(...; affinity = k -> ...)` makes a free worker
prefer a key whose group it has already handled, so worker-local memoisation
of a shared setup is hit instead of reloaded. A preference, never a partition:
no worker idles while a key is pending. Measured, 24 keys over 2 groups on 8
workers: **8 group changes without it, 0 with.**
- **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 @@ -67,6 +87,35 @@ 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),
affinity = k -> ParamIO.param(k, "system.L"),
opts = RunOpts(deadline = time() + 25*60))
```

`prerequisite` removes the duplicated *build*; `affinity` removes the repeated *load* of what it
built. The second only matters once the first is in place.

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
59 changes: 45 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 All @@ -35,6 +38,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 All @@ -44,6 +49,7 @@ one `write(io, line)` call to stay within that guarantee.
| `retry` | another attempt will follow |
| `gave_up` | all `max_attempts` attempts exhausted |
| `skip_complete` | full-done early exit (manifest had every key) |
| `worker_lost` | every worker died with keys pending; the key was not attempted |

`: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 All @@ -67,7 +73,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 @@ -92,9 +111,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 @@ -116,10 +136,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
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
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!
Loading
Loading