Skip to content

feat: worker affinity - #50

Closed
sotashimozono wants to merge 1 commit into
feat/prerequisite-stagefrom
feat/worker-affinity
Closed

sotashimozono wants to merge 1 commit into
feat/prerequisite-stagefrom
feat/worker-affinity

Conversation

@sotashimozono

Copy link
Copy Markdown
Member

Closes #47.

Stacked on #49, which is stacked on #48. Merge order: SweepRunner #48, #49, then this.

What it does

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, six repetitions, identical every time:

workers group changes without with
4 4 0
8 8 0

A preference, never a partition

The issue is explicit that this is the point: 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.

The test states this as a claim that cannot flake: more workers used than there are groups. An exclusive partition of 2 groups could never put 4 workers to work.

Secondary to #46, as the issue says

A prerequisite stage removes the duplicated build; affinity removes the repeated load of what it built. The second only matters once the first is in place.

Fault tolerance

pmap is not used here because it hands out work itself, so its ProcessExitedException re-dispatch had to be reproduced: 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.

There is a test for this: one key kills its worker outright with ccall(:_exit, ...), so remotecall_fetch sees a real ProcessExitedException rather than a caught error. run! returns, the loss is not counted as an error, and a later pass completes the key once the abandoned lock goes stale.

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.

One thing worth flagging

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. That is the same defect #49 fixed in EventLog, and it showed up here as a trace one entry short of the keys that actually ran (23 == 24). Worth knowing for anyone instrumenting a multi-process test in this repo.

Verification

21 assertions in test_affinity.jl, plus all other files under test/. Local stability, repeated whole-file runs: 5/5 clean after the trace fix, and clean again on the final code. CI runs it per shard.

Patch bump, 0.6.3 -> 0.6.4.

🤖 Generated with Claude Code

@github-actions github-actions Bot added the enhancement New feature or request label Sep 15, 2026
@github-actions

Copy link
Copy Markdown
Contributor

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

(updates on each push to this PR)

`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>
@codecov

codecov Bot commented Sep 15, 2026

Copy link
Copy Markdown

Codecov Report

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

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

📢 Thoughts on this report? Let us know!

@sotashimozono

Copy link
Copy Markdown
Member Author

Folded into #51.

The stack was one patch bump per pull request, each a single step above its own base, which is what version-check asks of a pull request. Merging them in sequence would have published a patch version to General for each one, including one whose content is test: Aqua. Collapsing the stack onto main ships it as one step instead.

Every commit from this branch is in #51 unchanged, so nothing here is lost; this body stays readable as the per-change account. The branch is kept.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant