Skip to content

Parallel kernels, one accounting layer, one reverse index - #10

Merged
alexy merged 8 commits into
mainfrom
work/rayon-kernels
Sep 20, 2026
Merged

alexy merged 8 commits into
mainfrom
work/rayon-kernels

Conversation

@alexy

@alexy alexy commented Sep 20, 2026

Copy link
Copy Markdown
Member

Step 3 of the order in codex-to-codex.md (08:15Z): the rayon layer, the two
defects the review reproduced, and the merge of the two parallel layers into one.

What is here

The accounting layer. WorkMeter gives one worker a batched share of the
work budget, so N workers do not retry a compare-exchange against one cache line
once per visited entry. Grants are reclaimable: a meter that cannot be admitted
returns its own idle remainder and takes back every other live meter's, so
whether work fits a budget never depends on how many meters exist.
reserve_for_workers admits per-worker scratch once, before a region, keeping
the memory mutex off the per-element path. ExecutionContext::with_concurrency
is how an embedder says how many threads a query may add; an execution that says
nothing runs exactly the code it ran before.

Parallel kernels: degree, pagerank (weighted and unweighted),
bfs, multiSourceBfs, wcc, projectionStats. Left sequential on purpose:
dfs and topologicalSort, whose emitted order is observable; scc and
dijkstra, which are the oracles their parallel replacements will be tested
against.

One parallel module. The union of both branches': the catalog's
balanced_ranges, map_ranges, ordered_blocks, for_chunks and sort_total,
with this branch's workers, pool and meters. width() survives only until the
catalog kernels move onto workers().

One reverse index. The projection had a reachability-only reverse index and a
transpose with weights; now it has the transpose alone, and edge slots are
omitted from it because no caller reads them — about a gigabyte of arcs on a
directed graph of com-Orkut's size that SCC and PageRank would have paid to
ignore.

The rules the kernels keep

  • Threads come from the execution, never from the machine: nothing reads the CPU
    count, because these kernels also run inside a server with its own pool.
  • Partitions follow the work for disjoint writes, and are fixed for float
    reductions, because regrouping the terms of a float sum changes its low bits.
  • Results are combined in index order, so an answer is bit-for-bit the same at
    one worker and at sixteen.
  • The sequential path stays, below a measured floor and whenever the parallel
    feature is off, and it is the differential oracle in tests.

Tests

grust-algorithms/tests/parallel.rs compares every parallel kernel with its
sequential path on a 120,000-node graph and with itself at 2 and 16 workers;
grust-procedures/tests/contracts/parallel.rs holds the budget contracts at
sixteen threads. Three of those tests arrived from the review as failing
reproductions, cherry-picked unchanged in 70b7003: an idle meter's grant
refusing work that fits, the same under skew at sixteen threads, and PageRank's
scores moving with the worker count. Each fails without its fix on this box.

Numbers

Warm, sixteen workers, against the sequential kernels, on quegee (16 vCPU over 8
physical cores). examples/scaling prints what each kernel produced beside every
timing, and those values agree at every worker count.

graph kernel algorithm alone, 1 worker at 16 workers
roadNet-CA pagerank 4.0x 29.3x
roadNet-CA pagerank, weighted 2.5x 19.0x
roadNet-CA wcc 3.0x 14.5x
roadNet-CA bfs 1.6x 2.6x
com-Orkut pagerank 2.9x 22.0x
com-Orkut wcc 3.8x 28.0x
com-Orkut bfs 1.4x 6.7x

Three of these are better algorithms as well as threaded ones, which is why the
first column exists: reporting only the total would credit threads with work the
algorithm did. Correctness cost what it cost — the fixed reduction chunk took
PageRank on roadNet-CA from 34.4x to 29.3x — and codex-to-codex.md has the
before-and-after table.

Not in this PR

Moving the catalog kernels onto WorkMeter and workers() and retiring Meter
and width() is step 5, and the catalog agent's. Parallelising the projection
build and the transpose is step 6.

alexy and others added 7 commits September 20, 2026 06:12
charge_work admits through a compare-exchange on one counter, which is right
single-threaded and is the shape that degrades when N workers retry against
the same cache line once per visited entry. WorkMeter admits a block at a
time and spends it locally; when a block does not fit it admits exactly what
was asked for, so a budget still fails at the unit that exceeds it, and
unspent units return on drop so final usage counts work performed.
Cancellation is still observed per charge; the deadline samples per block,
the same cadence per worker as before.

reserve_for_workers admits per-worker scratch once, before a parallel
region, keeping the memory mutex off the per-element path.

Concurrency is set on the context rather than added to ExecutionLimits: the
limits struct is built by literal in two dozen places across the workspace,
and a setter that refuses to run once the context is shared gives the same
per-execution semantics with a default of one thread.

Parallel counterparts of the budget contracts cover a budget that exactly
fits, one a unit short, cancellation reaching every worker and an expired
deadline, at sixteen threads.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
… components

The pool is bounded by the execution, never by the machine: an execution
that asks for threads gets them from a pool sized to that request, and one
that asks for nothing runs exactly the code it ran before. Nothing reads the
CPU count, because these kernels also run inside a server that owns its own
runtime and pool.

Results do not move with the thread count. Reductions combine per-chunk
partials in chunk order rather than completion order, so sums are bit-for-bit
equal at one worker and sixteen. Degree gives every node its own output slot.
PageRank pulls into each target over the reverse topology instead of pushing
out of each source, which is what removes the write conflicts; the weighted
projection keeps the push kernel, which needs the source's weight total per
arc. Breadth-first search claims each node with one compare-exchange per
level, so a node enters the next frontier once and its distance is the level
that claimed it. Components link the larger root to the smaller, so the
surviving root is the component's minimum row: the same canonical label the
sequential kernel produces.

Work accounting goes through per-worker meters, and degree's charges are
unchanged unit for unit, so work-unit comparisons against earlier runs still
mean something.

tests/parallel.rs is the differential oracle on a 120k-node graph: the same
answers as the sequential paths, exactly for integers and identifiers and
within a stated tolerance where the parallel path sums floats in a different
order, and identical to itself at two and sixteen threads.

examples/scaling.rs times the kernels on a real edge list, separating the
sequential baseline from the parallel implementation at one worker, so a
change of algorithm cannot be reported as a speedup from threads.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
…statistics

The sequential thresholds are now measurements rather than taste. On quegee,
on roadNet-CA prefixes: PageRank and components are already faster in
parallel at 20,000 edges, so the shared threshold sits just under that size.
Breadth-first search is different, and needed two changes. Its own threshold
is higher, because claiming a node with a compare-exchange costs more than
testing a distance when nothing contends: it ran at 0.78x of sequential on a
20,000-edge graph and only reached 1.37x at 200,000. And it now decides level
by level on the frontier it is about to expand, so a road network whose levels
hold hundreds of nodes expands them in place while com-Orkut, whose frontiers
reach millions, still spreads them.

Projection statistics counts self-loops in parallel, the partials added in
chunk order.

Warm numbers at sixteen workers against the sequential kernels, twenty
PageRank iterations, second call so a one-off reverse-topology build is not
counted: roadNet-CA PageRank 34.4x, components 16.5x, breadth-first 2.7x,
degree 2.5x; com-Orkut PageRank 22.5x, components 37.3x, breadth-first 7.8x.
The results printed beside each timing are identical at every worker count,
so none of it comes from doing less work.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
…d7fb95

Two for WorkMeter refusing work that fits a budget because another meter holds
an unspent block; one for PageRank scores changing bits with the worker count
because chunk length depends on it. See codex-to-codex.md 2026-09-20T07:30Z.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1. A work meter could refuse work that fits. Exact admission was measured
against a counter that still included other meters' unspent grants, and those
came back only when their meter dropped, so whether a query succeeded depended
on how many meters existed. Grants are now reclaimable: a meter that cannot be
admitted returns its own idle remainder, takes back every other live meter's,
and retries. The registry is touched when a meter is created, dropped or
refused, never on a charge.

2. PageRank's scores moved with the worker count. Folding per-chunk sums in
chunk order fixes their order but not their grouping, and the chunk length came
from the worker count, so the dangling mass and through it every score changed
in the low bits. Reductions now use a fixed chunk of 4,096, matching the catalog
kernels; passes that write each item's own slot keep worker-derived chunks,
where grouping cannot affect the answer.

The cost is real and worth stating: on roadNet-CA at sixteen workers PageRank
falls from 34.4x to 28.8x against the sequential kernel, and components from
16.5x to 14.4x. Creating one meter per worker task rather than per chunk
recovered part of it; a fixed reduction chunk is simply more chunks.

Also from the review: the pool cache is capped at eight distinct worker counts
so an embedder passing user-supplied values cannot accumulate idle threads, the
doc now says pools are shared per count per process, the discarded admission
error is explained, and the changelog records that PageRank's low digits differ
from the sequential kernel's once concurrency is requested.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
Components and projection statistics do not sum floats: unions are applied to
shared atomics, and a self-loop count is an integer sum, which is the same
however it is grouped. Both were taking the fixed reduction chunk anyway, which
cost components half its speed on com-Orkut, 37x down to 20x. They follow the
worker count again. Only PageRank's dangling mass and residual, which are float
reductions, need a grouping independent of the worker count.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
The two parallel modules become their union. The catalog's range and block
helpers stay as they are: balanced_ranges beats chunking by count on any skewed
degree distribution, which is every social graph. The worker count, the pool and
the work meters come from this branch, because they are bounded by what the
execution asked for rather than by the machine. width() stays for now, with a
doc note, since seventeen catalog kernels still call it; step 5 retires it.

The projection has one reverse index instead of two. Reachability had a narrow
one and the catalog added a transpose with weights and edge slots; a single
build that carries weights costs less than two builds over the same arcs.
Strongly connected components and PageRank both read incoming() now. Edge slots
became optional and a transpose leaves them out, because no caller of a
transpose reads them: eight bytes an arc that SCC and PageRank would have paid
to ignore, about a gigabyte on a directed graph of com-Orkut's size.

Weighted PageRank follows from that. The in-arc index carries each arc's weight,
so the weighted projection pulls like the unweighted one instead of falling back
to the sequential push, which is now only the oracle. Measured on roadNet-CA
with synthetic weights, twenty iterations, warm: 19.0x at sixteen workers where
there was no parallel path at all before.

From the review: a stress test provokes the reclaim retry window directly, two
hundred rounds of sixteen meters starting together on a budget equal to their
work, with no refusal observed; and the doc comment now says what a caller sees
if all three retries lose, and how to leave headroom against it.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
alexy added a commit that referenced this pull request Sep 20, 2026
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
many_meters_racing_for_the_last_of_a_budget_that_exactly_fits failed in 55 of
80 runs on a ten-core host, and the skewed sixteen-thread contract in 7 of 80.
Meters that ran out together each moved their balance onto their own stack,
found the others' balances empty, reclaimed nothing, and returned
BudgetExceeded on the first attempt; the three retries never ran.

A block is invisible between reaching the shared counter and reaching its
meter's balance. Admitting a block now holds the grant registry shared across
that moment, and a meter about to be refused holds it exclusively, so every
unspent unit is in some balance when it reclaims and the refusal is the one a
single thread would meet. 0 failures in 360 runs, debug and release. Shared
holders never wait for each other; on this host the charge loop and PageRank,
components and breadth-first search at ten workers are within run-to-run noise
of the version without exclusion.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
alexy added a commit that referenced this pull request Sep 20, 2026
Documentation only.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@alexy
alexy merged commit 39a8ca3 into main Sep 20, 2026
1 check passed
alexy added a commit that referenced this pull request Sep 20, 2026
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Sak2FQcdgs5pkL2ruKrSUZ
alexy added a commit that referenced this pull request Sep 20, 2026
Documentation only.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant