Skip to content

feat: deadline, claim record, prerequisite stage, worker affinity, plus the EventLog corruption they uncovered - #51

Merged
sotashimozono merged 8 commits into
mainfrom
ci/aqua
Sep 15, 2026
Merged

sotashimozono merged 8 commits into
mainfrom
ci/aqua

Conversation

@sotashimozono

@sotashimozono sotashimozono commented Sep 15, 2026 •

Copy link
Copy Markdown
Member

Closes #43. Closes #44. Closes #45. Closes #46. Closes #47.

The stack that was #48 + #49 + #50 + #51, collapsed into one pull request so it ships as one version step (0.6.1 -> 0.6.2) rather than four. The closed PRs keep their bodies and the per-change reasoning stays in the commits.

deadline (#43)

RunOpts(deadline = time() + 25 * 60)

Both docstrings now state the granularity, because stop_flag's did not: the flag and the deadline are read between keys. A key already inside work_fn runs to completion, so neither bounds when run! returns.

What a deadline buys is that it is set in advance. A flag raised 60 s before the wall clock cannot buy back a ten-minute key; a deadline can be budgeted as allocation_end - longest_expected_key - summary_time. In #43 the stage died to a SLURM limit 14 minutes after the flag went up, taking its diagnostics with it.

run! reports stopped_by as :flag, :deadline or nothing, appended at the end of the returned NamedTuple so a positional reader is not silently rebound. The full-done early exit used to return a shorter NamedTuple; it now returns the same field set.

key_acquired (#44)

I measured #44's premise first, and it is partly wrong:

termination status tree event log
work_fn throws only the successes' .done error x6, gave_up x2
kill -9 mid-key .running survives only stage_start

The failures were durably recorded, and .running does survive a SIGKILL. What is missing is the claim: :key_start is :debug and suppressed, so a key acquired then killed leaves no event naming it. :key_acquired closes that at O(keys this master worked).

DataVault has no .failed marker at all, so #44's failed: 0 was structurally guaranteed rather than evidence.

One BLAS warning (#45)

287 identical lines on a 72-node job, between the rows of the table. Now one @warn with the count. The docstring no longer presents the setting as diagnosed: ms/step was flat from 1 to 36 threads with no segfault in ~250 keys at blas=16. The existing test pinned the old text and now asserts the fix: three hot workers produce exactly one warning naming "3 of 3", and a second testset with no hot worker asserts zero.

Prerequisite stage (#46)

run_loop!(work_fn, vault, keys; prerequisite = Prerequisite(prep_fn, prep_vault, derived_keys))

Measured, 8 concurrent processes over 16 keys sharing 2 setups, gated to start together:

setup builds
check-then-build inside work_fn 16
prerequisite stage 2

The stronger argument is correctness: work_fn's wall clock included the setup for one key per fibre and excluded it for the other ~199, so a cost claim built on it depended on scheduling luck.

It is a barrier, not a work loop: a round finding every remaining key held by a sibling sleeps and goes again. If the prerequisite does not complete, the dependent stage does not start. Not a DAG — SweepRunner does not know which dependent key needs which prerequisite key.

Worker affinity (#47)

run!(work_fn, vault, keys; affinity = k -> param(k, "system.L"))

24 keys over 2 groups, six repetitions, identical every time: 4 group changes on 4 workers without, 0 with; 8 against 0 on 8 workers.

A preference, never a partition. The test states it unflakeably: more workers used than there are groups — an exclusive partition of 2 groups could never put 4 workers to work. pmap's ProcessExitedException re-dispatch is reproduced by hand, with a test that kills a worker outright via ccall(:_exit, ...).

The EventLog corruption this uncovered

:key_acquired doubled the event rate and turned a rare failure into a frequent one. It is not a flake and it was not new:

4 EventLog objects on one path, 300 events each, 4 threads
  before: 947 of 1200 lines, 82 malformed
  after:  1200 of 1200,      0 malformed

test_run_keylock.jl, 10 runs each
  base branch: 2/10 tore a line  |  with key_acquired: 3/10  |  with this fix: 0/10

run! builds a fresh EventLog per call, so concurrent masters in one process hold several objects on one file and several different locks. Events were being lost, not merely spliced. Three fixes, all needed: the lock is held per path; the write goes to an unbuffered append descriptor (open(path, "a") is an IOStream that flushes on its own boundaries, so write(io, line) was never the single syscall the O_APPEND guarantee is about); and the lock is resolved at call time, because deserialization skips the constructor and workers otherwise get a private lock.

The same trap then bit the affinity test's own trace file, which is how I found the instrumentation was lying rather than the code.

Aqua

Found one thing, the same in all three infra packages: [extras] without [compat]. Now 11/11. FormatCheck's comment, which said main required no status checks, has outlived its condition and is corrected.

Verification

New: test_run_deadline.jl (27), test_prerequisite.jl (33), test_affinity.jl (21), test_eventlog_shared_path.jl (9), plus a rewritten test_verify_workers_adversarial.jl. All 21 files under test/ green; the affinity file run 4x standalone for stability after its instrumentation fix.

Additive: no export removed or renamed, new RunOpts field has a default. Patch bump, 0.6.1 -> 0.6.2.

🤖 Generated with Claude Code

sotashimozono and others added 6 commits September 15, 2026 07:55
Closes #43, #44, #45.

## deadline (#43)

`RunOpts(deadline = time() + 25*60)`. The loop stops handing out keys once `time() > deadline`.

Both docstrings now state the granularity, because it is the same for both and `stop_flag`'s did
not: the flag and the deadline are read BETWEEN keys. A key already in `work_fn` runs to
completion, so neither bounds when `run!` returns.

What the deadline buys is that it is set in ADVANCE. A flag raised 60 s before the wall clock
cannot buy back a key that runs for ten minutes; a deadline can be budgeted as
`allocation_end - longest_expected_key - summary_time`. #43's stage died to a SLURM time limit
14 minutes after the flag went up, taking its diagnostics with it.

`run!` now reports `stopped_by` as `:flag`, `:deadline` or `nothing`, so a short stage is
attributable without re-reading the clock. It is appended at the END of the returned NamedTuple:
every reader in the workspace uses named access, but a positional reader would otherwise be
silently rebound. The full-done early exit previously returned a SHORTER NamedTuple than the
normal path (no `busy`, `gave_up`, `stop`); it now returns the same field set.

## key_acquired (#44)

One `:info` event written at acquire, flushed by `log_event`'s open/write/close.

I measured #44's premise before building to it, and it is partly wrong. With a work_fn that throws
on one axis value (the systematic-defect shape):

| | status tree | event log |
|---|---|---|
| `work_fn` throws | only the successes' `.done` | error x6, gave_up x2 |
| `kill -9` mid-key | `.running` SURVIVES | only `stage_start` |

So the failures WERE durably recorded, in `events_<host>_<pid>.jsonl`, and `.running` does survive
a SIGKILL. What is genuinely missing is the claim itself: `:key_start` is `:debug` and suppressed
by default, so a key acquired and then killed leaves no event naming it.

`:key_acquired` closes that at O(keys this master worked), which is the order of `:key_done`, not
the O(masters x keys) of `:lock_busy` that the level split exists to control.

DataVault has no `.failed` marker at all (only `.done` and `.running`), so #44's "failed: 0" was
structurally guaranteed rather than evidence. Left alone here: putting failures in the status tree
is a DataVault change.

## BLAS warning (#45)

287 identical lines on a 72-node job, interleaved with the rows of the table above them. Now one
`@warn` carrying the count, after the table.

The docstring no longer presents the setting as diagnosed. It is a known cause of OpenBLAS
segfaults in multi-process Julia, and on this workload ms/step was flat from 1 to 36 threads with
no segfault in ~250 keys at blas=16, so the warning reports rather than concludes.

The existing test pinned the old text (`r"BLAS threads=4"`). It now asserts what the fix is: three
hot workers produce exactly ONE warning naming "3 of 3", and a second testset with no hot worker
asserts zero, so `== 1` is not just "the warning is unconditional".

Verified locally, targeted files: test_run_deadline.jl (new, 27 assertions) and
test_verify_workers_adversarial.jl (rewritten), plus every other file under test/, all green.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Closes #46.

## Prerequisite

`run!` locks the KEY, so work shared BETWEEN keys has nowhere to live but inside `work_fn`, where
it has no protection at all: every worker that wants a setup not yet on disk builds it itself.

    run_loop!(work_fn, vault, keys;
              prerequisite = Prerequisite(prep_fn, prep_vault, derived_keys))

The setup becomes its own key space and gets the same locking, resume and provenance as any other
stage, so its cost is recorded in its own payload instead of landing on whichever dependent key
happened to run first. That is the stronger half of #46: `work_fn`'s wall-clock included the setup
for one key in each fibre and excluded it for the other ~199, so a cost claim built on it depended
on scheduling luck.

Measured here, 8 concurrent processes over 16 keys sharing 2 setups (8 keys per fibre), gated so
they start together:

| | setup builds |
|---|---|
| check-then-build inside work_fn | 16 |
| prerequisite stage | 2 |

`run_prerequisite!` is a BARRIER, not a work loop, and that is the whole difference from
`run_loop!`. A round that finds every remaining key held by a sibling SLEEPS and goes again rather
than counting an empty round; the dependent stage cannot start until the setup exists. It
terminates on: all done; no progress AND nothing held by a sibling; the stop flag; the deadline.

If the prerequisite does not complete, the dependent stage does not start at all. Running it anyway
would spend the allocation on keys whose setup is known to be missing, which is what `Preflight`
already exists to prevent one layer up.

SweepRunner does NOT know which dependent key needs which prerequisite key: the dependency is
resolved inside `work_fn`. This is "all of the prerequisite, then all of the dependents", not a
DAG, and the docstring says so.

`run_loop!` now returns `(; ran, rounds, done, stopped_by, prerequisite)` instead of `nothing`.

## EventLog corruption, found on the way

The new `:key_acquired` event from #44 roughly doubled the event rate and turned a rare failure in
`test_run_keylock.jl` into a frequent one. It is not a flake and it is not mine:

    4 EventLog objects on one path, 300 events each, 4 threads
      before: 947 of 1200 lines, 82 malformed
      after:  1200 of 1200, 0 malformed

    test_run_keylock.jl, 10 runs
      base branch:        2/10 tore a line
      with key_acquired:  3/10
      with this fix:      0/10

`run!` builds a fresh `EventLog` on every call, so four concurrent masters in one process hold four
objects pointing at one file and four DIFFERENT `ReentrantLock`s. The documented "serialized
through an internal ReentrantLock" was therefore false in exactly the configuration the repo's own
tests use. Two fixes, both needed:

- the lock is now held per PATH, process-wide, so several `EventLog`s on one file share it;
- the write goes to an UNBUFFERED append descriptor. `open(path, "a")` gives an `IOStream`, which
  flushes on its own boundaries, so `write(io, line)` was never the single syscall the O_APPEND
  guarantee is about. The docstring claimed it was.

The regression test covers both halves: threads within one process, and four separate processes
appending to one file.

Verified locally: all 20 test files under test/, green.

Stacked on #48. Merge that first.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The per-path lock landed in the previous commit but did not reach the workers, which is where the
concurrency actually is. `run!` serialises the `EventLog` to every worker, and deserialization
rebuilds the struct WITHOUT running the constructor, so the `lock` field that arrives there is a
fresh private lock that serialises nothing against its siblings on that worker.

Measured: `deserialize(serialize(log)).lock === log.lock` is false, and `=== _path_lock(path)` is
also false.

`log_event` now looks the lock up by path on every call, so it is correct however the object got
there. The test asserts the field really does not survive serialization, then writes 200 events
from each of the original and the deserialized log concurrently and requires all 400 lines back.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The EventLog deserialization test used Serialization without it being declared, which works in a
dev environment that happens to have it and fails in the isolated test env CI builds:
`ArgumentError: Package Serialization not found in current path` on shard s1.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
`pmap` hands `todo` out flat. When `work_fn` memoises something per group in worker-local state,
which key a worker draws is luck, and a worker that hops throws the memo away.

    run!(work_fn, vault, keys; affinity = k -> param(k, "system.L"))

Measured, 24 keys over 2 groups on 4 workers, six repetitions, identical every time:

| | group changes |
|---|---|
| without affinity | 4 |
| with affinity | 0 |

and on 8 workers, 8 against 0.

A PREFERENCE, never a partition, which is the part the issue is explicit about: an exclusive
partition would serialise a 200-key fibre onto one worker, which is worse than the hopping it
fixes. A free worker takes from its most recently used group if anything is left there, and
otherwise from the group with the most work outstanding, which spreads workers over groups rather
than piling them onto one. No worker is ever idle while a key is pending.

This is secondary to #46 and only makes sense after it: a prerequisite stage removes the duplicated
BUILD, affinity removes the repeated LOAD of what it built.

`pmap` is not used here because it hands out work itself. Its `ProcessExitedException`
re-dispatch is reproduced by hand: a key whose worker died goes back on the queue and the worker is
dropped. If every worker dies with keys still pending, those keys were never attempted, so they are
reported `:lock_busy` (which `run!` already counts as retriable) rather than `:error`, and a
`:worker_lost` event names each one.

`filled` is a `Vector{Bool}` rather than a `BitVector`: adjacent bits share a word, so two tasks
marking neighbouring indices would read-modify-write the same one. `@async` is sticky today so it
would not actually race, but the dispatcher should not depend on that.

Tests, 21 assertions:

- the transition count, with a control run of the same fixture WITHOUT affinity;
- more workers used than there are groups, which is the statement that this is not a partition (an
  exclusive partition of 2 groups could never put 4 workers to work);
- every key runs exactly once;
- a throwing work_fn is still reported;
- a worker killed outright with `_exit` mid-key: `run!` returns, the loss is not counted as an
  error, and a later pass completes the key once the abandoned lock goes stale.

The affinity test's own trace file had to be written the way `log_event` now writes: one `write` to
an unbuffered O_APPEND descriptor. An `open(f, "a")` there is a buffered `IOStream` and drops lines
across the worker processes, which is the same defect this branch fixed in `EventLog` and showed up
as a trace one entry short of the keys that actually ran.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The package-level checks no unit test covers, because they are about the PACKAGE rather than about
any one function. `undefined_exports` is the one that earns its place: every feature in this
repository lands by adding a name to the `export` list, and a typo there is invisible until a caller
reaches for it.

Aqua found one thing, the same one in all three infra packages: `[extras]` entries with no
`[compat]` bound. `Test` here, and the `Serialization` I added last week had one only by accident.
Adding `Aqua = "0.8"` and `Test = "1.11"` makes `test_all` green at 11/11.

The other seven checks were already clean: method ambiguities, unbound type parameters, undefined
exports, project-extras agreement, stale dependencies, piracy, persistent tasks.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@github-actions github-actions Bot added the chore label Sep 15, 2026
@github-actions

Copy link
Copy Markdown
Contributor

📚 Docs preview: https://qatlashub.github.io/SweepRunner.jl/previews/PR51/

(updates on each push to this PR)

@codecov

codecov Bot commented Sep 15, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 97.22222% with 3 lines in your changes missing coverage. Please review.

Files with missing lines Patch % Lines
src/Run.jl 95.83% 3 Missing ⚠️

📢 Thoughts on this report? Let us know!

sotashimozono and others added 2 commits September 15, 2026 12:16
The comment ended "This repository's \`main\` requires no status checks today, so the filter costs
nothing yet; it would cost everything the day one is added." That day has come and gone: `main`
now requires `test / All shards passed` AND `format / format-check`, so the clause that said the
filter was harmless is the one sentence a reader would act on wrongly.

The mechanism it explains is unchanged and stays. Only the conditional half is rewritten to state
what is true now.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The branches were stacked, each one patch above its own base, which is what `version-check`
requires of a pull request. Merging them in sequence would have published four patch versions to
General for one body of work, one of them for adding Aqua.

Collapsed to a single pull request against main, so the stack ships as one step and AutoRegister
fires once: it detects a release by comparing the version line at HEAD~1 against HEAD, so only the
one bump that lands is registered.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@sotashimozono
sotashimozono changed the base branch from feat/worker-affinity to main September 15, 2026 12:25
@sotashimozono sotashimozono changed the title test: Aqua feat: deadline, claim record, prerequisite stage, worker affinity, plus the EventLog corruption they uncovered Sep 15, 2026
@github-actions github-actions Bot added the enhancement New feature or request label Sep 15, 2026
@sotashimozono
sotashimozono merged commit 945440e into main Sep 15, 2026
19 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment