Skip to content
Draft
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
4 changes: 2 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.6"
version = "0.6.7"
authors = ["sota shimozono <shimozono-sota631@g.ecc.u-tokyo.ac.jp>"]

[deps]
Expand All @@ -24,7 +24,7 @@ Aqua = "0.8"
# (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.4"
DataVault = "0.8.6"
Dates = "1.11"
Distributed = "1.11"
JLD2 = "0.6"
Expand Down
26 changes: 20 additions & 6 deletions src/Observe.jl
Original file line number Diff line number Diff line change
Expand Up @@ -15,18 +15,21 @@ const _OBSERVATIONS_LOCK = ReentrantLock()
_observation_key(vault::Vault) = (vault.outdir, vault.spec.study.project_name, vault.run)

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

Observe the sources from this process and remember the token for `vault`. On failure the token is
Observe the sources from this process and remember the token for `vault`. `work_fn` is named as
the entry code (`observe_sources(...; code=[work_fn])`): the binding vouches for it or says why
not — a closure or a function defined in the driving script cannot be checked. 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)
function _observe_here!(vault::Vault, role::AbstractString, work_fn=nothing)
token, err = try
DataVault.observe_sources(
vault;
phase="run-start",
process=Dict("role" => String(role), "myid" => myid()),
code=work_fn === nothing ? () : (work_fn,),
),
nothing
catch e
Expand Down Expand Up @@ -57,7 +60,9 @@ 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)
function _observe_processes!(
vault::Vault, multi::Bool, observe::Bool, log, stage; work_fn=nothing
)
targets = if multi
vcat([(myid(), "master")], [(w, "worker") for w in workers()])
else
Expand All @@ -75,9 +80,18 @@ function _observe_processes!(vault::Vault, multi::Bool, observe::Bool, log, stag
end
outcomes = asyncmap(targets) do (pid, role)
return if pid == myid()
_observe_here!(vault, role)
_observe_here!(vault, role, work_fn)
else
remotecall_fetch(SweepRunner._observe_here!, pid, vault, role)
# `work_fn` travels to the worker, as it will for pmap: a named function resolves to
# the worker's own, a closure arrives as the master's code. One the worker cannot
# receive is a failed observation, not a failed run!.
try
remotecall_fetch(SweepRunner._observe_here!, pid, vault, role, work_fn)
catch e
e isa InterruptException && rethrow()
remotecall_fetch(SweepRunner._forget_observation!, pid, vault)
nothing, sprint(showerror, e)
end
end
end
for ((pid, role), (token, err)) in zip(targets, outcomes)
Expand Down
6 changes: 4 additions & 2 deletions src/Run.jl
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,9 @@ With `observe=true` (the default) the master and every worker call `DataVault.ob
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
never claims more than was checked. `work_fn` is named as the entry code: the binding can be
`loaded-matches-disk` only when it is a function of a package loaded from the study's sources, not
a closure or a function defined in the driving script (those are `unverified`, with the reason). 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
Expand Down Expand Up @@ -312,7 +314,7 @@ function run!(
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)
_observe_processes!(vault, multi, observe, log, stage; work_fn)
dispatch = ks -> if !multi
_run_sequential!(work_fn, vault, ks, stage, log, opts)
elseif affinity === nothing
Expand Down
15 changes: 12 additions & 3 deletions test/run/test_run_observe.jl
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,16 @@ end
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-")
@test length(tokens) == 1 && startswith(only(tokens), "obs2-")
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")
# The work function is the entry code; defined in this test, not a package, it cannot be
# vouched for, and the observation says so rather than claiming a match.
@test only(r["code"])["name"] == "_obsrun_work"
@test r["binding"] == "unverified"
@test any(contains("_obsrun_work"), r["binding_reasons"])
end
end

Expand Down Expand Up @@ -77,6 +80,12 @@ end
]
@test all(r -> r["process"]["role"] == "worker", recs)
@test all(r -> r["process"]["myid"] in pids, recs)
# A closure is the master's code on the worker: never a match.
@test all(r -> r["binding"] == "unverified", recs)
@test all(
r -> any(contains("arrived from another process"), r["binding_reasons"]),
recs,
)
roles = [
TOML.parsefile(joinpath(_obsrun_dir(v), "observations", f))["process"]["role"]
for f in readdir(joinpath(_obsrun_dir(v), "observations"))
Expand Down
Loading