From 370feb765e879eec7e6763a5eb40a71e94a59a5b Mon Sep 17 00:00:00 2001 From: sotashimozono Date: Tue, 22 Sep 2026 08:27:33 +0000 Subject: [PATCH] Observe the sources on every process at run! start, and put each process's token in its .done A point's marker must name the observation of the process that computed it: under pmap that is the worker, whose loaded code need not be the master's (a closure runs the master's body while the named functions it calls resolve on the worker). run! now has the master observe for itself and each worker observe for itself, through DataVault.observe_sources, before any key is dispatched; _run_one_with_retry!, running on that same process, passes that process's token to mark_done!. Tokens are kept per vault identity, so two run!s sharing a process never hand each other's token to a marker. An observation that fails does not stop the run: the event log records observe_failed with the error, and that process's markers read observation=unknown, as they do with observe=false (which also forgets any earlier token rather than reusing it). DataVault floor 0.8.4. Version 0.6.6. Co-Authored-By: Claude Opus 5 (1M context) --- Project.toml | 8 ++-- src/Observe.jl | 93 ++++++++++++++++++++++++++++++++++++ src/Run.jl | 24 ++++++++-- src/SweepRunner.jl | 1 + test/run/test_run_observe.jl | 89 ++++++++++++++++++++++++++++++++++ 5 files changed, 208 insertions(+), 7 deletions(-) create mode 100644 src/Observe.jl create mode 100644 test/run/test_run_observe.jl diff --git a/Project.toml b/Project.toml index 3a9fba8..d50e3be 100644 --- a/Project.toml +++ b/Project.toml @@ -1,6 +1,6 @@ name = "SweepRunner" uuid = "be946ad2-3cb3-4b6e-8f7e-4a5ecc3c255b" -version = "0.6.5" +version = "0.6.6" authors = ["sota shimozono "] [deps] @@ -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" diff --git a/src/Observe.jl b/src/Observe.jl new file mode 100644 index 0000000..baf8f9f --- /dev/null +++ b/src/Observe.jl @@ -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 diff --git a/src/Run.jl b/src/Run.jl index 703f853..ebbe96e 100644 --- a/src/Run.jl +++ b/src/Run.jl @@ -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 @@ -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=`). 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 @@ -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" @@ -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 @@ -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; @@ -903,6 +918,7 @@ function run_loop!( load=nothing, prerequisite=nothing, affinity=nothing, + observe::Bool=true, ) pre = nothing if prerequisite !== nothing @@ -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 diff --git a/src/SweepRunner.jl b/src/SweepRunner.jl index 665a844..4b47399 100644 --- a/src/SweepRunner.jl +++ b/src/SweepRunner.jl @@ -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") diff --git a/test/run/test_run_observe.jl b/test/run/test_run_observe.jl new file mode 100644 index 0000000..5c505e3 --- /dev/null +++ b/test/run/test_run_observe.jl @@ -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