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

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

[compat]
Aqua = "0.8"
# 0.8.3 is the floor, not a preference: `save!` returning the digest it wrote, and
# `mark_done!(…; result)` that puts it in `.done`, arrived there. `artifact!` / `ArtifactBusy`
# 0.8.4 is the floor, not a preference: `observe_sources` and `mark_done!(…; observation)` arrived
# there, and `save!` returning the digest it wrote (with `mark_done!(…; result)`) in 0.8.3. `artifact!` / `ArtifactBusy`
# (what the deferral and `artifact_affinity` are built for) came in 0.8.2. The owner-stamped
# `.running` API (`new_owner_token`, `running_owner`, and the three-argument `refresh_running!` /
# `clear_running!`) that the liveness reaper is built on arrived in 0.8.1.
DataVault = "0.8.3"
DataVault = "0.8.4"
Dates = "1.11"
Distributed = "1.11"
JLD2 = "0.6"
Expand Down
93 changes: 93 additions & 0 deletions src/Observe.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
# Observe.jl — one source observation per process per run!, handed to every `.done` it writes.
#
# `DataVault.observe_sources` records what the source looked like and how far THIS process's loaded
# code was checked against it. A point's marker must carry the token of the process that computed
# it: under `pmap` that is the worker, whose loaded code need not be the master's. So each process
# observes for itself at `run!` start, keeps the token here, and `_run_one_with_retry!` — which runs
# on that same process — reads it back.
#
# Keyed by vault identity, not held in one slot, so two `run!`s sharing a process never hand each
# other's token to a marker.

const _OBSERVATIONS = Dict{Tuple{String,String,String},String}()
const _OBSERVATIONS_LOCK = ReentrantLock()

_observation_key(vault::Vault) = (vault.outdir, vault.spec.study.project_name, vault.run)

"""
_observe_here!(vault, role) -> (token, err)

Observe the sources from this process and remember the token for `vault`. On failure the token is
`nothing` and any earlier token for `vault` is forgotten, so a marker written afterwards says
`observation=unknown` rather than naming an observation of some earlier state.
"""
function _observe_here!(vault::Vault, role::AbstractString)
token, err = try
DataVault.observe_sources(
vault;
phase="run-start",
process=Dict("role" => String(role), "myid" => myid()),
),
nothing
catch e
nothing, sprint(showerror, e)
end
lock(_OBSERVATIONS_LOCK) do
if token === nothing
delete!(_OBSERVATIONS, _observation_key(vault))
else
_OBSERVATIONS[_observation_key(vault)] = token
end
end
return token, err
end

"Forget this process's token for `vault` (a `run!` with `observe=false`)."
function _forget_observation!(vault::Vault)
lock(() -> delete!(_OBSERVATIONS, _observation_key(vault)), _OBSERVATIONS_LOCK)
return nothing
end

"This process's token for `vault`, or `nothing`."
function _observation_token(vault::Vault)
return lock(
() -> get(_OBSERVATIONS, _observation_key(vault), nothing), _OBSERVATIONS_LOCK
)
end

# Observe on the master, and on every worker when the run fans out. Each outcome is an event: an
# observation that failed leaves its process's markers at `observation=unknown`, and says why.
function _observe_processes!(vault::Vault, multi::Bool, observe::Bool, log, stage)
targets = if multi
vcat([(myid(), "master")], [(w, "worker") for w in workers()])
else
[(myid(), "master")]
end
if !observe
for (pid, _) in targets
if pid == myid()
_forget_observation!(vault)
else
remotecall_fetch(SweepRunner._forget_observation!, pid, vault)
end
end
return nothing
end
outcomes = asyncmap(targets) do (pid, role)
return if pid == myid()
_observe_here!(vault, role)
else
remotecall_fetch(SweepRunner._observe_here!, pid, vault, role)
end
end
for ((pid, role), (token, err)) in zip(targets, outcomes)
if token === nothing
log_event(
log, :observe_failed; level=:warn, stage=stage, pid=pid, role=role, err=err
)
else
log_event(log, :observed; stage=stage, pid=pid, role=role, token=token)
end
end
return nothing
end
24 changes: 21 additions & 3 deletions src/Run.jl
Original file line number Diff line number Diff line change
Expand Up @@ -185,7 +185,7 @@ derives `(root, stage)` from a `DataVault.Vault`:
load_manifest(vault::Vault) = load_manifest(manifest_root(vault), Symbol(vault.run))

"""
run!(work_fn, vault, keys; opts=RunOpts(), load=nothing) -> NamedTuple
run!(work_fn, vault, keys; opts=RunOpts(), load=nothing, observe=true) -> NamedTuple

Run `work_fn(key) -> Dict` for every `key` in `keys`, persisting through
`vault`. Writes a structured JSONL event log at
Expand All @@ -204,6 +204,15 @@ 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)`.

# Source observations

With `observe=true` (the default) the master and every worker call `DataVault.observe_sources`
before any key is dispatched, and each `.done` a process writes carries that process's token
(`observation=<token>`). The observation records what the source looked like at `run!` start and
its **binding** — how far the code that process had loaded was checked against it — so a marker
never claims more than was checked. An observation that fails does not stop the run: the event log
says why, and that process's markers read `observation=unknown`, as they do with `observe=false`.

# Affinity

`affinity` is `key -> value`, and turns the fan-out from "any free worker takes the next key" into
Expand Down Expand Up @@ -261,6 +270,7 @@ function run!(
opts::RunOpts=RunOpts(),
load=nothing,
affinity=nothing,
observe::Bool=true,
)
stage = Symbol(vault.run)
log_name = "events_$(gethostname())_$(getpid()).jsonl"
Expand Down Expand Up @@ -300,6 +310,9 @@ function run!(
vcat([:ParamIO, :DataVault, :SweepRunner], _worker_module_names(load))
)
end
# Every process that will write markers observes its sources now, so each `.done` names the
# observation of the process that computed it (see Observe.jl).
_observe_processes!(vault, multi, observe, log, stage)
dispatch = ks -> if !multi
_run_sequential!(work_fn, vault, ks, stage, log, opts)
elseif affinity === nothing
Expand Down Expand Up @@ -798,7 +811,9 @@ function _run_one_with_retry!(
# The digest save! took before its rename goes into the marker, so `.done` names the
# bytes this attempt wrote rather than whatever the file holds when someone looks.
saved = DataVault.save!(vault, key, payload)
DataVault.mark_done!(vault, key; result=saved)
DataVault.mark_done!(
vault, key; result=saved, observation=_observation_token(vault)
)
log_event(
log,
:key_done;
Expand Down Expand Up @@ -903,6 +918,7 @@ function run_loop!(
load=nothing,
prerequisite=nothing,
affinity=nothing,
observe::Bool=true,
)
pre = nothing
if prerequisite !== nothing
Expand Down Expand Up @@ -934,7 +950,9 @@ function run_loop!(
stopped = _stop_reason(opts)
stopped === nothing || break
rounds += 1
result = run!(work_fn, vault, keys; opts=opts, load=load, affinity=affinity)
result = run!(
work_fn, vault, keys; opts=opts, load=load, affinity=affinity, observe=observe
)
n_done += result.done
n_busy = result.busy
if result.done > 0
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("Observe.jl") # one source observation per process per run!, for each .done
include("Artifacts.jl") # artifact_affinity; ArtifactBusy deferral lives in Run.jl
include("Prerequisite.jl")
include("Preflight.jl")
Expand Down
89 changes: 89 additions & 0 deletions test/run/test_run_observe.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
# Each `.done` names the source observation of the process that computed its key.

using SweepRunner, Test, DataVault, ParamIO, Distributed, TOML

const _OBSRUN_CFG = joinpath(@__DIR__, "fixtures", "study.toml")

function _obsrun_done(v, key)
pairs = (split(l, '='; limit=2) for l in eachline(DataVault._done_file(v, key)))
return Dict(String(p[1]) => String(p[2]) for p in pairs)
end
_obsrun_dir(v) = joinpath(v.outdir, ".datavault", v.spec.study.project_name)
function _obsrun_record(v, token)
return TOML.parsefile(joinpath(_obsrun_dir(v), "observations", "$token.toml"))
end
_obsrun_work(key) = Dict{String,Any}("N" => key.params["N"])

function with_obsrun_vault(f, run)
outdir = mktempdir()
try
v = DataVault.Vault(_OBSRUN_CFG; run=run, outdir=outdir)
f(v, ParamIO.expand(v.spec))
finally
rm(outdir; recursive=true, force=true)
end
end

@testset "run!: every .done names the master's observation when nothing fans out" begin
with_obsrun_vault("observe") do v, keys
res = run!(_obsrun_work, v, keys; opts=RunOpts(; workers=:sequential))
@test res.done == length(keys)
tokens = unique(_obsrun_done(v, k)["observation"] for k in keys)
@test length(tokens) == 1 && startswith(only(tokens), "obs1-")
r = _obsrun_record(v, only(tokens))
@test r["phase"] == "run-start"
@test r["process"]["role"] == "master" && r["process"]["pid"] == getpid()
@test isfile(joinpath(_obsrun_dir(v), "sources", r["source"], "COMPLETE"))
@test r["binding"] in
("loaded-matches-disk", "loaded-differs-from-disk", "unverified")
end
end

@testset "run!: observe=false writes unknown, never an earlier token" begin
with_obsrun_vault("no-observe") do v, keys
SweepRunner._observe_here!(v, "master") # a token left from before
run!(_obsrun_work, v, keys; opts=RunOpts(; workers=:sequential), observe=false)
@test all(_obsrun_done(v, k)["observation"] == "unknown" for k in keys)
end
end

@testset "run!: an observation that fails does not stop the run, and says why" begin
with_obsrun_vault("observe-fails") do v, keys
mkpath(_obsrun_dir(v))
write(
joinpath(_obsrun_dir(v), "sources"), "a file where the snapshot store should be"
)
res = run!(_obsrun_work, v, keys; opts=RunOpts(; workers=:sequential))
@test res.done == length(keys)
@test all(_obsrun_done(v, k)["observation"] == "unknown" for k in keys)
logs = filter(f -> startswith(f, "events_"), readdir(v.outdir))
lines = vcat((readlines(joinpath(v.outdir, f)) for f in logs)...)
@test any(l -> occursin("observe_failed", l), lines)
end
end

@testset "run! distributed: each marker carries the token of the worker that computed it" begin
proj = dirname(Base.active_project())
pids = addprocs(2; exeflags="--project=$proj")
try
with_obsrun_vault("observe-workers") do v, keys
# A closure, not a named function: a function defined only on the master does not
# exist on the workers.
res = run!(key -> Dict{String,Any}("N" => key.params["N"]), v, keys)
@test res.done == length(keys)
recs = [
_obsrun_record(v, t) for
t in unique(_obsrun_done(v, k)["observation"] for k in keys)
]
@test all(r -> r["process"]["role"] == "worker", recs)
@test all(r -> r["process"]["myid"] in pids, recs)
roles = [
TOML.parsefile(joinpath(_obsrun_dir(v), "observations", f))["process"]["role"]
for f in readdir(joinpath(_obsrun_dir(v), "observations"))
]
@test count(==("master"), roles) == 1 && count(==("worker"), roles) == 2
end
finally
rmprocs(pids)
end
end
Loading