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
7 changes: 6 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
docs/build/
docs/src/assets/favicon.ico
docs/src/assets/logo.png
Manifest.toml
Manifest.toml

# Coverage counters from `julia --code-coverage`. Written next to the source they measure,
# so a `git add -A` after a local coverage run sweeps them in.
*.cov
*.info
7 changes: 5 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 All @@ -19,7 +19,10 @@ TOML = "fa267f1f-6049-4f14-aa54-33bafae1ed76"

[compat]
Aqua = "0.8"
DataVault = "0.7, 0.8"
# 0.8.1 is the floor, not a preference: 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"
Dates = "1.11"
Distributed = "1.11"
JLD2 = "0.6"
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!
```

## Liveness

```@docs
SweepRunner.owner_token
SweepRunner.holder_liveness
```

## Prerequisite

```@docs
Expand Down
7 changes: 6 additions & 1 deletion src/EventLog.jl
Original file line number Diff line number Diff line change
Expand Up @@ -44,12 +44,17 @@ stay within that guarantee.
| `key_done` | after a successful `work_fn(key)` (includes `secs`, `attempt`) |
| `lock_busy` | another master holds the `.running` lock (acquire = `:busy`) |
| `lock_lost` | our lock was reclaimed mid-work; result discarded (no double-run) |
| `lock_reaped` | a lock whose holder was shown dead was cleared without waiting |
| `reap_failed` | reaping threw; the key falls back to the `stale_after` timeout |
| `lock_reclaimed`| (reserved, not currently emitted) |
| `error` | `work_fn` threw on this attempt |
| `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 |
| `worker_died` | the worker exited on this key every time it was dispatched, up to |
| | 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 |

`: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
155 changes: 155 additions & 0 deletions src/Liveness.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
# Liveness — is the master that holds this lock still running?
#
# `stale_after` answers that by waiting. It has to, because a heartbeat can only ever say
# "recently alive"; the absence of one is indistinguishable from a slow filesystem until enough
# time has passed. Where the question CAN be asked outright, asking is worth several minutes of a
# batch allocation.
#
# Nothing here is required for correctness: every path degrades to `:unknown`, and `stale_after`
# then decides exactly as before. `:slurm` is one of four worker modes, so a mechanism that only
# worked there would cover a quarter of the runs.

"""
owner_token() -> String

The identity of ONE acquisition, stamped into `.running` by [`run!`](@ref) so a sibling can ask
about it later. `host:pid:nonce` from DataVault, with `:slurm<jobid>` appended inside a Slurm
allocation.

Call it per acquisition and do not cache it. The nonce is what tells a master's current hold from
a hold it lost and retook, and a token broadcast once per process cannot make that distinction.

The Slurm field is what makes the question answerable ACROSS hosts: `/proc` only works for a
holder on this machine, and the master of the job that was killed is usually somewhere else.
"""
function owner_token()::String
base = DataVault.new_owner_token()
job = _slurm_queue_id()
return isempty(job) ? base : string(base, ":slurm", job)
end

# The id as `squeue` PRINTS it, which is not always `SLURM_JOB_ID`. In an array task that variable
# holds a distinct raw id per task while the queue lists `<array id>_<task index>`, so stamping it
# makes every task but the first unfindable, and unfindable reads as `:dead`.
function _slurm_queue_id()::String
arr = get(ENV, "SLURM_ARRAY_JOB_ID", "")
idx = get(ENV, "SLURM_ARRAY_TASK_ID", "")
(isempty(arr) || isempty(idx)) && return get(ENV, "SLURM_JOB_ID", "")
return string(arr, "_", idx)
end

"""
holder_liveness(owner) -> Symbol

`:alive`, `:dead`, or `:unknown` for the holder named by an [`owner_token`](@ref).

`:dead` is only ever returned on positive evidence that the process is gone. Everything else is
`:unknown`, including every error path: a wrong `:dead` would hand a live master's key to someone
else, which is the one outcome the lock exists to prevent.

Two sources, in order:

1. a Slurm job id, when this process is ITSELF inside a Slurm allocation. The job is absent from
the queue, so it has finished, been cancelled, or hit its wall clock.
2. the pid, when the holder is on THIS host and `/proc` exists. A recycled pid reads as `:alive`,
which is the safe direction.
"""
function holder_liveness(owner::AbstractString)::Symbol
parts = split(String(owner), ':')
length(parts) >= 3 || return :unknown

# The Slurm field is the FOURTH part and nowhere else, because that is where `owner_token`
# writes it. Searching every part for a `slurm` prefix matched a HOSTNAME beginning with
# `slurm`, read the rest of the hostname as a job id, found it absent from the queue, and
# returned `:dead` for a master that was alive.
if length(parts) >= 4 && startswith(parts[4], "slurm")
s = _slurm_liveness(chopprefix(parts[4], "slurm"))
s === :unknown || return s
end

parts[1] == gethostname() || return :unknown
pid = tryparse(Int, parts[2])
pid === nothing && return :unknown
return _pid_liveness(pid)
end

# `/proc/<pid>` is the whole check on Linux. Elsewhere there is no equally cheap answer that does
# not risk a false `:dead`, so there is no answer.
function _pid_liveness(pid::Int)::Symbol
Sys.islinux() || return :unknown
return isdir("/proc/$(pid)") ? :alive : :dead
end

# The live-job set, cached: `squeue` is a scheduler RPC and this is consulted per contended key.
const _SQUEUE_TTL = 15.0
const _SQUEUE_TIMEOUT = 10.0
const _squeue_cache = Ref{Tuple{Float64,Union{Set{String},Nothing}}}((-Inf, nothing))
const _squeue_lock = ReentrantLock()

function _slurm_liveness(jobid::AbstractString)::Symbol
isempty(jobid) && return :unknown
# A job id Slurm would never issue cannot be looked up, and "not in the queue" must not be the
# answer for it: that is a `:dead` conjured out of a corrupt `.running` file. `_` is legal
# because an array task's queue id is `<array id>_<task index>`.
all(c -> isdigit(c) || c == '_', jobid) || return :unknown
# Only from INSIDE an allocation: a `squeue` on PATH may be a wrapper pointed at an unrelated
# cluster, where every foreign job id is absent and would therefore read as `:dead`.
haskey(ENV, "SLURM_JOB_ID") || return :unknown
live = _live_slurm_jobs()
live === nothing && return :unknown
# The queue has to be able to see THIS process before its silence about anyone else means
# anything. A `squeue` pointed at a different cluster answers successfully and lists none of
# our ids, and every holder would then read as `:dead`.
_slurm_queue_id() in live || return :unknown
# An array task is `12345_7` in the queue while `SLURM_JOB_ID` is `12345`, so the whole job is
# alive if any of its tasks is.
jobid in live && return :alive
any(j -> startswith(j, jobid * "_"), live) && return :alive
return :dead
end

# Bounded, because this runs on the hot path of every contended key and `squeue` can be a wrapper
# around a network call. An unbounded one that stalls holds `_squeue_lock` and takes the allocation
# with it, and neither `stop_flag` nor `deadline` is read inside a key.
function _squeue_output(timeout::Real)::Union{String,Nothing}
Sys.which("squeue") === nothing && return nothing
tmp = tempname()
proc = try
run(pipeline(`squeue -h -o %i`; stdout=tmp, stderr=devnull); wait=false)
catch e
e isa InterruptException && rethrow()
rm(tmp; force=true)
return nothing
end
try
if timedwait(() -> !process_running(proc), Float64(timeout); pollint=0.05) !== :ok
kill(proc, Base.SIGKILL)
return nothing
end
success(proc) || return nothing
return read(tmp, String)
finally
rm(tmp; force=true)
end
end

# `nothing` means "no opinion": squeue is absent, timed out, or failed. Listing the queue and testing
# membership is deliberate. `squeue -j <id>` exits non-zero BOTH for a finished job and for a
# broken scheduler, and those must not collapse into the same answer.
function _live_slurm_jobs()::Union{Set{String},Nothing}
return lock(_squeue_lock) do
t, cached = _squeue_cache[]
time() - t < _SQUEUE_TTL && return cached
fresh = try
out = _squeue_output(_SQUEUE_TIMEOUT)
out === nothing ? nothing : Set(String.(split(out; keepempty=false)))
catch e
e isa InterruptException && rethrow()
nothing
end
_squeue_cache[] = (time(), fresh)
return fresh
end
end

export owner_token, holder_liveness
42 changes: 38 additions & 4 deletions src/Prerequisite.jl
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,14 @@ is what it becomes: its own key space, its own vault, its own payloads.
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.

`stop_flag` and `deadline` are not stage knobs: they bound the JOB, and a stage may not loosen a
bound the job set. `deadline` therefore takes the TIGHTER of the two, and `stop_flag` takes the
caller's whenever it has one. A prerequisite cannot redirect or outlive either.

That asymmetry is deliberate for `stop_flag`: [`RunOpts`](@ref) resolves its default from
`ENV["SWEEPRUNNER_STOP_FLAG"]`, so an `opts` here that never mentions `stop_flag` can still carry
one, and `=== nothing` does not mean "the author left it unset" for that field.

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.
"""
Expand Down Expand Up @@ -54,17 +62,20 @@ point; a dead one is bounded by `stale_after`, after which its lock is reclaimab
function run_prerequisite!(
p::Prerequisite; opts::RunOpts=RunOpts(), load=nothing, poll::Real=30.0
)
o = p.opts === nothing ? opts : p.opts
o = _merged_opts(p, opts)
n_done = 0
waited = 0
rounds = 0

while true
stopped = _stop_reason(o)
if stopped !== nothing
# A barrier whose keys are all done is satisfied whether or not a stop is pending;
# `complete=false` beside `remaining=0` would gate the dependent stage on nothing.
undone = _n_undone(p)
return (;
complete=false,
remaining=_n_undone(p),
complete=undone == 0,
remaining=undone,
done=n_done,
waited=waited,
rounds=rounds,
Expand Down Expand Up @@ -96,7 +107,9 @@ function run_prerequisite!(
done=n_done,
waited=waited,
rounds=rounds,
stopped_by=nothing,
# A round that handed out nothing may have been cut short rather than empty,
# and those are different failures: one retries, the other will not.
stopped_by=r.stopped_by,
)
end
waited += 1
Expand All @@ -107,4 +120,25 @@ end

_n_undone(p::Prerequisite) = count(k -> !DataVault.is_done(p.vault, k), p.keys)

# `nothing` is "no bound", so it loses to any real one.
function _tighter(a::Union{Float64,Nothing}, b::Union{Float64,Nothing})
return a === nothing ? b : (b === nothing ? a : min(a, b))
end

# `p.opts` replaces the stage's knobs wholesale, but not the two that bound the job. Taking the
# tighter deadline rather than the explicit one matters in both directions: reverting it to
# `nothing` lets the barrier outlive the allocation, and so does honouring a more generous ceiling
# the Prerequisite was built with months earlier.
function _merged_opts(p::Prerequisite, opts::RunOpts)::RunOpts
p.opts === nothing && return opts
o = p.opts
job = (
stop_flag=opts.stop_flag === nothing ? o.stop_flag : opts.stop_flag,
deadline=_tighter(o.deadline, opts.deadline),
)
# Everything else forwarded BY NAME, so a field added to `RunOpts` later cannot be silently
# reset to its constructor default here.
return RunOpts(; (f => get(job, f, getfield(o, f)) for f in fieldnames(RunOpts))...)
end

export Prerequisite, run_prerequisite!
Loading
Loading