diff --git a/Project.toml b/Project.toml index d50e3be..edfc9d6 100644 --- a/Project.toml +++ b/Project.toml @@ -1,6 +1,6 @@ name = "SweepRunner" uuid = "be946ad2-3cb3-4b6e-8f7e-4a5ecc3c255b" -version = "0.6.6" +version = "0.6.7" authors = ["sota shimozono "] [deps] @@ -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" diff --git a/src/Observe.jl b/src/Observe.jl index baf8f9f..020fefb 100644 --- a/src/Observe.jl +++ b/src/Observe.jl @@ -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 @@ -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 @@ -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) diff --git a/src/Run.jl b/src/Run.jl index ebbe96e..d058ecd 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -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=`). 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 @@ -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 diff --git a/test/run/test_run_observe.jl b/test/run/test_run_observe.jl index 5c505e3..a6b92c7 100644 --- a/test/run/test_run_observe.jl +++ b/test/run/test_run_observe.jl @@ -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 @@ -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"))