Skip to content

fix(replication): bound the fresh-offer backlog and send offers without copying - #233

Open
mickvandijke wants to merge 24 commits into
mainfrom
perf/replication-send-path
Open

mickvandijke wants to merge 24 commits into
mainfrom
perf/replication-send-path

Conversation

@mickvandijke

@mickvandijke mickvandijke commented Sep 22, 2026 •

Copy link
Copy Markdown
Member

Summary

Under sustained client uploads on a 60-node testnet, nodes grew to 1–2 GiB within an hour. Heap profiles attributed ~80% of live memory to encoded FreshReplicationOffer messages queued in the fresh-write drainer: every accepted PUT was encoded immediately (~4–5 MiB) and its per-peer send tasks then waited for one of MAX_CONCURRENT_REPLICATION_SENDS (3) permits while pinning that buffer, with nothing bounding the backlog.

Rebased onto main (0.20.0 release baseline, pointers) on 2026-09-29; the saorsa-core and saorsa-transport branches were rebased onto their main too (see Pins).

Commits:

  1. fix(replication): bound encoded fresh offers waiting behind send permits — FreshWriteEvent carries only key + proof; a MAX_PENDING_FRESH_OFFERS (8) permit is taken before the chunk is read back and encoded, and the permit lives with the encoded offer until its last send finishes.
  2. perf(replication): cut copies of fresh offers on the send path — the chunk moves into the offer instead of being copied; ReplicationMessage::encode serializes into an exactly-sized buffer.
  3. perf(replication): share fresh offers with the transport as Bytes — with saorsa-core accepting impl Into<Bytes>, each send attempt hands out a reference-counted handle; an in-flight send costs one frame instead of four copies. Adds ADR-0017.
  4. fix(replication): send PaidNotify before waiting for an offer permit.
  5. fix(replication): two-stage fresh replication with un-gated paid-list evidence — the drainer records PaidForList(self) and sends PaidNotify for every write at arrival rate, then hands it to the offer dispatcher, the only permit-gated stage. Fresh-offer and PaidNotify payload fields are byte strings (serde_bytes, wire-identical). One offer pipeline serves the dispatcher and replicate_fresh.
  6. test(replication): saturation and read-retry coverage for the fresh-write pipeline.

Review follow-ups (2026-09-29, grumbach's review):

  1. fix(replication): read fresh offers back through the verified store path — the dispatcher reads with ChunkStore::get, the read the fetch path serves from, instead of get_raw. A chunk that no longer hashes to its key is quarantined by the read and never offered (every receiver would reject it and charge the sender).
  2. fix(replication): retry a failed read-back without stalling the dispatcher — the retry delay runs on its own tracked task that re-queues the write; the dispatcher no longer sleeps, so a shared read fault (exhausted descriptors) cannot hold every healthy write behind it.
  3. perf(replication): encode fetched chunks and subtree slices as byte strings — FetchResponse::Success::data and SubtreeSliceItem::Present::bao_slice get serde_bytes too (wire-identical). For a 4 MiB fetch response, encode ~1.6 → ~0.07 ms and decode ~2.6 → ~0.05 ms, and the decoded chunk is exactly sized. (Measured with postcard 1.1 in a release build; the size-then-write encode itself costs no second pass there — the sizing loop compiles to a length.)
  4. test(replication): drive the fresh-write pipeline through the PUT handler and hold the offer budget — the e2e harness now wires the fresh-write sender into each node's AntProtocol as a node does, so handler PUTs feed the pipeline. The pipeline tests were rewritten to prove their claims (details under Test evidence). New test-utils accessors: ReplicationEngine::fresh_offers_dispatched, ReplicationEngine::hold_replication_sends.
  5. docs(adr-0017) — verified read-back, retries off the dispatcher, what brings an unadvertised key back, byte-string fields, and "bounded" made precise (see below).

Deep-review follow-ups (2026-09-29; 33 candidates, each verified before fixing):

  1. fix(replication): back off read-back retries so a store-wide fault spares the backlog — with the retry off the dispatcher, a store-wide read fault (EMFILE) lasting over ~2 s spent every queued write's three attempts at once and dropped the whole backlog's offers. Retries now back off over about a minute (1, 2, 4, 8, 16, 32 s; 7 attempts), still on their own tasks; only the first failure and the give-up log at WARN.
  2. test: run the pipeline tests serially — without #[serial] the burst test's receiver refusals leaked into the capacity test's process-wide counters (reproduced 3/3 with two test threads; CI's single thread hid it).
  3. test: drive the capacity test through the fresh-write channel — replicate_fresh now waits for a permit, which paced the capacity loop on fast storage and hid PaidNotify drops the production drainer would cause.
  4. test: time the retry test from the fault and bound the handler PUT — the healthy write is queued the moment the fault is seen (10 ms poll), so a dispatcher that sleeps the delay out cannot pass on a slow runner; the handler PUT helper has the harness's 90 s timeout.
  5. perf(storage): read a chunk into a buffer sized from its file — reads grew from empty through Take<File>, leaving a 4 MiB chunk in an 8 MiB buffer on every GET, audit and read-back.
  6. perf(node): hand chunk GET responses to the transport without copying — the client response path did to_vec() on Bytes.
  7. refactor(replication): tidy the fresh-offer sender — the permit is tied to every handle of the encoded bytes (Bytes::from_owner), the proof moves instead of being copied, and the sender is renamed fresh::send_fresh_offers (it shared a name with the receiver-side dispatch_fresh_offer).
  8. refactor: share one exactly-sized postcard encoder — codec::encode_exact serves replication and the WebRTC path (which stops zero-filling a chunk-sized buffer), keeping the experimental serialized_size call in one place.
  9. refactor(replication) — stale comments that put the permit and read-back in the drainer, the offer channel scoped to start(), no-op drops; docs for the never-closed semaphores.
  10. test(e2e): reuse the helpers — TestHarness::prepopulate_payment_cache_everywhere, ChunkStore::test_chunk_path (test-utils), shared dummy proof / replication waits; the harness skips the engine loudly instead of inventing a closed fresh-write channel.
  11. docs(adr-0017) — the request/response envelope still copies per byte (left for saorsa-core), the permit's lifetime, the new tests.

Nothing is dropped by the permit: both queues are unbounded and FIFO, and every offer keeps the same fan-out, retries and delayed possession check; only the bulk chunk transfer is deferred under load. A read-back that keeps failing for about a minute (7 attempts) is given up, and a store-wide fault longer than that drops the queued writes' offers; other holders and neighbor sync cover them. Chunk buffers are bounded by configuration; the queues ahead of the permit are not a hard bound — under a sustained write rate above the send rate they grow by one key + stripped proof per write (~40 KB single-node, ~130 KB merkle) instead of one encoded chunk. The pending-offer budget bounds the sender, not receiver admission: a one-source burst of small chunks can exceed a receiver's per-source fresh-offer admission cap (12), and neighbor sync fills that gap as for any refused offer. Both are recorded in the ADR.

Pins: saorsa-core b18e3754 (#166), saorsa-transport 5da61737 (#171), both perf/replication-send-path rebased onto the 0.28.0 / 0.37.0 releases; ant-protocol stays on main's pointer pin d11d1010 — with the release baseline the root [patch] reaches ant-protocol's saorsa-core dependency, so no ant-protocol change is needed (ant-protocol #35 closed). The core rebase needed one fix: main's V2-834 failed-send accounting read message_data.len() after the payload had moved into the send; it now records the already-computed wire_len.

Linear issue

Closes V2-1359

Risk tier

  • T0 — docs / tooling / CI / pure UX-output. Repo CI only.
  • T1 — client-only, no network-facing behavior change. CI + prod compat smoke.
  • T2 — node/client logic with behavioral surface, no protocol/format/economics change. Dev testnet + ADR.
  • T3 — protocol / storage format / payments / routing. T2 evidence + adversarial testing.

Compatibility

  • Wire: none — replication message encoding (all serde_bytes fields checked against their u8-sequence layout in both directions) and the saorsa-core frame are byte-identical.
  • Storage: none — the dispatcher reads chunks the PUT handler already stored, through the ordinary verified read.
  • API: in the published ant_node::replication::fresh module, FreshWriteEvent drops its pub data field and the pub async fn replicate_fresh is gone (its crate-internal successor is fresh::send_fresh_offers). No sibling repo uses either, but strictly that is a public-API break — see the semver note below. ReplicationEngine::replicate_fresh keeps its signature but now waits for a pending-offer permit. Test-utils only: AntProtocol::fresh_write_sender, ReplicationEngine::{pending_offer_permits_available, fresh_offers_dispatched, hold_replication_sends}, ChunkStore::test_chunk_path.
  • Tests: the e2e harness now wires each node's PUT handler to its fresh-write pipeline, as src/node.rs does, so any paid PUT in an e2e test fans out. The opt-in first_audit_ab A/B driver therefore compares a different workload if only one side contains this PR.

Semver impact

  • breaking
  • feature
  • fix

(Proposed "fix" for the node's behaviour; the two removed replication::fresh items above would make it "breaking" if the library API is held to semver — reviewer to confirm.)

Test evidence

Full CI, mirrored locally at 4dfbb521 (every ci.yml job step; all green):

  • cargo fmt --all -- --check ✓; cargo check --all-targets --all-features --locked ✓; cargo clippy --all-targets --all-features -- -D warnings ✓ (clean — main fixed the old migration_signal.rs lint); cargo doc --all-features --no-deps with RUSTDOCFLAGS=-Dwarnings ✓; cargo check --lib --no-default-features ✓.
  • cargo test --lib --features test-utils: 1235 passed; cargo test --lib --no-default-features: 1190 passed.
  • cargo test --test e2e --features test-utils -- --test-threads=1: 114 passed, 3 ignored (with the harness now feeding every handler PUT into the fresh-write pipeline).
  • poc_commitment_audit_attacks 19, poc_audit_handler_live 16, poc_bootstrap_stall 3, poc_shutdown_lmdb_drain 1, migration_reclaims_disk 2, migration_crash_safety 5 (3 ignored), migration_shared_volume 5, storage_scale 2, webrtc_direct_devnet 3 + 1 (--ignored) — all passed.
  • Every commit on the branch builds on its own (cargo check --all-targets --all-features).
  • GitHub CI on da51a65c (the PR now targets main, so ci.yml runs): builds on all three OSes, clippy, docs, MSRV, security audit, storage on btrfs/ext4/xfs, WebRTC devnet, and the macOS and Ubuntu test jobs passed.

Fresh-write pipeline e2e tests (rewritten). The previous burst test was flaky — 7 of 10 isolated runs failed on the pre-rebase head: a one-source burst of tiny chunks overran the receivers' per-source admission cap, and when all four receivers refused the same key the write never landed. The tests now assert at the sender, so receiver admission cannot make them flaky:

  • …replicates_a_put_and_skips_missing_chunks — a PUT through the handler replicates; a missing-chunk write queued ahead of it is not offered and its permit comes back.
  • …holds_a_burst_at_the_offer_budget — with the send stage held, a 24-write burst encodes exactly 8 offers and stops; after release all 24 are offered and all 8 permits return.
  • …retries_a_failed_read_off_the_dispatcher — the fault is a directory where the chunk file should be (works as root and on every platform; no silent skip); the failed read releases its permit, a healthy write queued behind it is offered within half the retry delay, and the failed write is offered once the fault clears.
  • …never_offers_a_corrupt_chunk — a chunk corrupted on disk is quarantined and never offered.
  • Repeated runs of the four pipeline tests all passed, plus the full suite. Each test was also checked against a deliberately broken build and fails on it, re-checked after the deep-review refactors: inline retry sleep (now timed from the fault), no permit gating, get_raw read-back, a permit leaked on skip, and a handler that does not emit.

Final-head smoke testnet, 2026-09-23 (b0263b32, the head before the rebase and the review follow-ups). One 60-node DigitalOcean fleet with the layout of the 0922 comparison (12 droplets, 9 regions, 18/60 nodes behind symmetric NAT, dedicated Anvil), 4 native uploaders (50 MiB files, own wallets) + 1 continuous sha256-verified downloader, 30 minutes measured. Stack: ant-node b0263b32, saorsa-core 53f1fc6f, saorsa-transport 1766037e, ant-protocol a66ddcfb, ant-client af8885f (pins only); binaries hash-checked on every droplet. The two right-hand columns are the first 30 minutes of the 0922 run, re-cut so the windows match.

b0263b32 (0923) previous head f55b7df8 (0922, first 30 min) web-support 50167d39 (0922, first 30 min)
Uploads OK / failed 207 / 0 202 / 0 192 / 0
Upload throughput 5.75 MiB/s 5.61 MiB/s 5.33 MiB/s
Upload time mean / p95 / worst 34.0 / 43.3 / 47.9 s 34.8 / 47.6 / 55.3 s 36.7 / 56.5 / 132.0 s
Upload throughput per 10 min 5.6 → 5.9 → 5.8 MiB/s 5.5 → 5.8 → 5.6 5.6 → 5.1 → 5.3
Chunk store p50 / p95 4.9 / 14.1 s 4.3 / 16.0 s 4.8 / 16.7 s
Downloads OK / failed 63 / 0 61 / 0 64 / 0
Download time mean / p95 27.1 / 34.6 s 28.0 / 38.8 s 26.7 / 36.6 s
Node RSS mean / p95 over the window 172 / 247 MiB 171 / 247 MiB 191 / 290 MiB
Largest node RSS at the end 304 MiB 286 MiB 372 MiB
Fleet mean node RSS per 5 min 161 → 170 → 173 → 172 → 175 → 179 MiB
Node CPU mean / p95 27.9 / 47.3 % 26.4 / 48.6 % 25.2 / 44.8 %
Lowest free memory on any node droplet 2,134 MiB 2,170 MiB 1,392 MiB
Uploader client peak / p95 memory 433 / 303 MiB 436 / 312 MiB 671 / 465 MiB
Node restarts 0 0 0
Node WARN / ERROR log lines 2,944 / 0 (30 min) 3,773 / 12 (60 min) 52,613 / 38 (60 min)

WARN breakdown for the 30 minutes (60 nodes): 236 relay-canary backoff, 196 happy-eyeballs connect failures, 91 PoP verification errors on PaidNotify (median-quote check; 148/h on the previous head), 84 possession-probe timeouts (1,877/h on web-support, 286/h on the previous head), 80 dial failures, 80 "NAT traversal poll took Nms", 57 possession-check unreachable, 51 stream resets by peer, 36 channel-send failures, and 35 "Paid notify dropped at admission (global_pool_full)" — the only counter that moved the wrong way against the previous head (21 in 60 min there; 2,795/h on web-support), because commit 5 sends PaidNotify at arrival rate; each drop is re-derived by the receiver's verification cycle. Raw data: ant-testnet/state/comparisons/web-support-final-head-smoke-30m-0923/.

  • Testnet validation of commit 1 (2026-09-21, two 60-node fleets): fleet peak live memory 1068 → 298 MiB and 674 → 262 MiB; the node that had reached 1051 MiB stayed flat at 0.0 MiB/min; the worst bootstrap's queued offers fell from 101 (463 MiB) to 6 (17.7 MiB). Data: ant-testnet/state/comparisons/web-support-memory-diag-0921/.
  • Testnet comparison, 2026-09-22 (60 minutes, two identical 60-node fleets, stack at f55b7df8 vs web-support 50167d39): uploads 375 → 412 (+10%), upload throughput 5.21 → 5.72 MiB/s, upload p95 63.4 → 46.1 s, node RSS mean / p95 / largest 199 / 305 / 428 → 182 / 262 / 364 MiB, WARN lines 52,613 → 3,773, downloads 130 → 120 (−8%, the expected cost of deferring replication of just-uploaded files). Data: ant-testnet/state/comparisons/web-support-send-path-vs-base-60m-0922/.

The testnet runs predate the rebase and both rounds of follow-ups. Those change the read-back (verified, retried with backoff), the encoding of two more fields and a few copies on the read and send paths; none of them loosens the memory bound those runs measured.

New dependency

serde_bytes 0.11 (already in the dependency graph transitively; now direct, for the byte-string encoding of payload fields). The bytes requirement rises from 1 to 1.10.1 for Bytes::from_owner (the lockfile already has 1.12.1).

ADR

https://github.com/WithAutonomi/ant-node/blob/perf/replication-send-path/docs/adr/ADR-0017-bounded-fresh-offers-and-copy-free-sends.md

Mitigation / rollback

Revert the branch and drop the saorsa-core / saorsa-transport [patch] entries to return to the 0.28.0 / 0.37.0 releases. Wire format is unchanged, so a partial rollout or rollback is safe. If the pending-offer cap ever throttles replication too hard, raising MAX_PENDING_FRESH_OFFERS trades memory for burst absorption without touching the wire.

🤖 Generated with Claude Code

@mickvandijke
mickvandijke marked this pull request as ready for review September 22, 2026 15:07
Base automatically changed from web-support to main September 22, 2026 17:29

@grumbach grumbach left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed at b0263b32. The memory result is real and the measurement behind it is thorough: 1068 MiB down to 298 MiB on the diagnostic fleet, and warnings dropping from 52,613 to 2,944 is a good independent signal that the fan-out really was the problem. The two-stage split is the right shape, and keeping the paid-list evidence off the back-pressure path is the right call.

Four things I would want addressed or at least answered.

1. The chunk is read back without being checked, and a bad read-back is charged to the sender

src/replication/mod.rs:2522 reads with ChunkStore::get_raw. That skips two things get does: the content-address check, and the store's own view of whether the chunk is servable. FileStore::exists refuses a key in suspect or known_wrong (src/storage/file_store.rs:1231), but get_raw goes straight to read_file and hands back the bytes either way.

Before this change the offered bytes came from the PUT request, which the handler had already content-checked, so an address mismatch was impossible for an honest sender. Now they come off disk, possibly much later under a backlog. If they no longer hash to the key, every receiver rejects at src/replication/mod.rs:5732 and charges the sender ApplicationFailure(REPLICATION_TRUST_WEIGHT) at src/replication/mod.rs:5795. With CLOSE_GROUP_SIZE of 7 that is up to 6 application failures for one bad local file, against a node that did nothing wrong.

To be fair about the scale: the weight is 1.0, not the audit weight, and the reciprocal half does not currently fire, because the possession check's "did not hold it" penalty goes through penalise_unheld_close_group_chunk and this release suspends it (RELEASE_SUSPEND_CLOSE_GROUP_STORAGE_PENALTY = true). So it is a correctness hole rather than an immediate fleet risk. But FileStore::get's own doc says a chunk proven wrong is one the node "cannot claim, commit or offer", and this is now the one path that offers it.

Cheap version: check compute_address(&data) == event.key in the dispatcher and skip without retry on a mismatch. That is one BLAKE3 pass over a chunk we are about to copy into a message anyway. Thorough version: a storage call that only returns verified, servable bytes, so it does not depend on the verify_on_read setting. Either way it is worth deciding rather than leaving implicit.

2. The read retry sleeps inside the only dispatcher

src/replication/mod.rs:2542 sleeps FRESH_READ_RETRY_DELAY in the dispatcher loop before re-queueing. Dropping the permit first at line 2529 does not help, because there is no second dispatcher to pick it up. Nothing is drained while that sleep runs.

The doc on MAX_FRESH_READ_ATTEMPTS names exhausted descriptors as the expected cause, and that fault is correlated: if one read fails that way the next ones will too. So the realistic case is not a single 1 s pause but N of them back to back, then again on the second retry round, with every healthy write queued behind them. A delayed re-queue on its own task, or a small timer queue, keeps the consumer free.

Related, and worth a line in the ADR either way: after the third failed read the write is dropped, and the store's read path has left the key marked suspect, so the node has stopped advertising a chunk it was just paid for and has sent no replica. What brings that key back?

3. encode() now costs two passes for the two payload fields that did not get serde_bytes

postcard::experimental::serialized_size at src/replication/protocol.rs:53 is a real serialization through the Size flavor, and to_extend at line 61 is a second one. For the three fields this PR annotated that is free: serialize_bytes reaches the flavor as one try_extend, and Size::try_extend is +len. For a plain Vec<u8> it is one try_push per byte, on both passes.

Two fields are still plain: FetchResponse::Success::data at line 946, which is a whole chunk, and SubtreeSliceItem::Present::bao_slice at line 1228, one per sampled block in a round 2 reply. send_replication_response_checked already calls FetchResponse a heavy serve path carrying "~99% of served bytes", so as it stands this PR makes the most expensive encode on the node about twice as expensive.

I measured the single pass on current main, where it is already the dominant cost for that message: a 1 MiB fetch response is 1,048,608 per-byte pushes, and a 3 MiB one lands on 4,194,304 bytes of capacity for a 3,145,890 byte message. That growth slack is the same problem this PR fixes for fresh offers, on a buffer held for the whole upload, three at a time against a limit documented as "about 12 MiB of simultaneous chunk data". The decoded chunk on the receiving side has it too, because serde's cautious size hint caps the initial allocation and the rest is regrown.

The fix is the same one annotation each, wire identical by exactly the argument your byte_string_fields_encode_like_u8_sequences test makes. I have it on a local branch off main with tests that count the copies and check both wire directions across the varint boundaries. Happy to hand it over, or land it separately so it does not grow this PR.

4. The new tests do not yet prove what the PR claims for them

  • The harness creates the fresh-write channel and keeps the sender on TestNode (tests/e2e/testnet.rs:1342) but never calls set_fresh_write_sender on the AntProtocol. All three tests post FreshWriteEvent by hand, so the handler change in src/storage/handler.rs is not covered by them. The fresh_write_sender() accessor added for this is not called anywhere.
  • The missing-chunk test says it covers permit leaks but never reads the permit count. Leaking one still leaves seven, the next write replicates, and the test passes.
  • The burst test writes short strings and never observes the cap being reached. A build that dropped the permit immediately after acquiring it, or did not gate at all, would pass it. The closing "all permits back" check says nothing about whether the cap ever bit.
  • The burst test reuses the replication deadline for the permit-drain wait, so a run that replicates close to the deadline reports a leak that is not one.
  • The retry test does not check that the permit was released during the sleep, which is the claim the retry path rests on. It also returns quietly with no coverage when file modes are not enforced, which will happen on some CI images.

Recording a low-water mark of available permits would fix most of these with one accessor.

What held up

I went looking for a permit leak and did not find one. It is released on a missing chunk, on exhausted retries, on encode failure, on an empty peer set, on cancellation, and when the last per-peer send drops the shared EncodedOffer. Holding the send permit across delivery retries is the right call and the comment explains why. The serde_bytes annotations are wire safe, and the compatibility test is the right way to show it.

One thing to note, not a blocker

Both hand-off queues are unbounded, so "memory must stay bounded under any client write rate" is not quite what this delivers. A FreshOfferEvent still carries its payment proof, and sidecar-stripped single-node proofs run around 40 KB, so a sustained backlog now grows at roughly 40 KB per write instead of 4 MiB. That is a 100x improvement and clearly worth having, but it is a slower leak rather than a bound, and the ADR reads as though it is a bound. Worth softening the wording, or saying what the intended ceiling is.

@mickvandijke
mickvandijke force-pushed the perf/replication-send-path branch from b0263b3 to 35132c1 Compare September 29, 2026 12:34

@dirvine dirvine left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed exact head 4dfbb521128d4530ed53026aa7dfbf5894976889, together with WithAutonomi/saorsa-core#166 at b18e3754 and WithAutonomi/saorsa-transport#171 at 5da61737.

No merge-blocking regression found in the reviewed diff. The pending-offer permit is acquired before disk read/encoding, retained through the shared Bytes owner, and released on the last handle. Paid-write evidence remains off the backpressured dispatcher. Read retries release their permit and requeue off-dispatcher; missing and corrupt chunks do not obstruct healthy writes under the default verified-read configuration. Wire annotations preserve postcard bytes.

Direct local verification (macOS, frozen head):

  • cargo fmt --check passed.
  • Focused codec, permit-owner lifetime, exact file-read allocation, serde_bytes compatibility and exact fetch-decode allocation tests passed.
  • cargo test --locked --features test-utils --test e2e fresh_write_pipeline -- --test-threads=1: 4 passed, including burst saturation/drain, PUT-to-replication/missing chunk, transient read failure/recovery, and corrupt-chunk quarantine.
  • Dependency pins match the reviewed core/transport heads. GitHub build/test/security and latest required metadata checks passed; historical metadata failures remain visible in the rollup.

Non-blocking release notes: the bound is on chunk-sized offers, not the key/payment-proof backlog (both hand-off channels remain unbounded). The documented Rust API change removes FreshWriteEvent.data and the standalone fresh::replicate_fresh; retain that caveat in release classification. Verification still honours verify_on_read; do not describe it as unconditional when disabled.

Publish the dependency chain transport → core → node, replacing temporary pins with released versions before publication. This review does not authorise merging or deploying. Fleet/RSS/throughput measurements in the PR were not independently reproduced; local regression tests are not a production soak.

Review team: six completed review roles (five independent GLM-5.2 responses plus a separate source adjudicator), reconciled against source and lead-executed tests. All six returned APPROVE. Local DS4 integration attempts timed out and were replaced; timeouts were not counted as opinions.

mickvandijke and others added 23 commits October 1, 2026 22:12
Under sustained client writes on a real network, nodes grew to 1-2 GiB
within an hour. Heap profiles on the testnet attribute ~80% of live
memory to encoded FreshReplicationOffer messages queued in the
fresh-write drainer: every accepted PUT was encoded immediately (chunk
plus proof, ~4-5 MiB) and its per-peer send tasks then waited for one
of MAX_CONCURRENT_REPLICATION_SENDS (3) permits while pinning that
buffer. When WAN sends hold permits longer than writes arrive, nothing
bounded the backlog, so the number of encoded offers kept growing.

- FreshWriteEvent no longer carries the chunk bytes; the chunk is on
  disk already and the drainer reads it back when it is ready to send.
- The drainer acquires a pending-offer permit (MAX_PENDING_FRESH_OFFERS)
  before reading and encoding, and the permit lives with the encoded
  buffer until the last per-peer send drops it. A backlog now waits as
  small queued events instead of chunk-sized buffers.
- The direct ReplicationEngine::replicate_fresh entry point takes the
  same permit so tests and callers share the bound.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Two of the copies each queued fresh offer carried were avoidable inside
this crate:

- The chunk read from storage now moves into FreshReplicationOffer
  instead of being copied, and the offer is dropped as soon as it has
  been encoded, so only the encoded bytes stay alive while sends queue.
- ReplicationMessage::encode serializes into a buffer sized from
  postcard's serialized_size. A doubling Vec left chunk-sized messages
  with up to twice their length in capacity, retained by every queued
  offer for as long as it waited for a send permit.

The remaining copies per in-flight send live in saorsa-core (payload
clone per channel attempt, signing re-serialization, wire frame) and
saorsa-transport (stream buffer copy in SendStream::write).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
With saorsa-core accepting `impl Into<Bytes>` on `send_message`, the
encoded fresh offer is now held as `Bytes` and each per-peer send
attempt hands out a reference-counted handle instead of cloning the
multi-MiB buffer. Together with the exactly-sized frame and the
transport's owned-buffer write, an in-flight send now costs one frame
instead of the previous four copies.

Adds ADR-0017 describing the bounded fresh-offer backlog and the
copy-free send path across ant-node, saorsa-core and saorsa-transport.

Pins: saorsa-core 0da3260c, saorsa-transport c2b9f3b0 (both
perf/replication-send-path).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
PaidNotify carries the paid-list evidence the paid close group needs to
repair a key later. Since the pending-offer permit was introduced it was
sent from replicate_fresh, i.e. only after the drainer had waited for a
permit, so a chunk backlog also delayed the evidence. Send it as soon as
a write is dequeued (and from the direct replicate_fresh entry point),
before any permit wait; only the bulk chunk offers are back-pressured.

Nothing is dropped by the permit: the fresh-write queue is unbounded and
FIFO, permits are released whenever a send terminates, and each offer
keeps the same fan-out, retries and delayed possession check. ADR-0017
now says so explicitly and records the measured download-latency cost.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
… evidence

Review follow-up. Sending PaidNotify before the permit wait only helped
the head-of-line event: the drainer was one serial loop, so every write
queued behind a blocked permit still had its PaidNotify and PaidForList
insert delayed by chunk back-pressure.

- The fresh-write drainer now never waits for a permit. For every event,
  at arrival rate, it records PaidForList(self) and sends PaidNotify,
  then forwards the event to a new offer dispatcher, the only stage that
  takes a pending-offer permit.
- The dispatcher reads the chunk back with `get_raw` (it was
  content-checked when stored), retries a failed read up to
  MAX_FRESH_READ_ATTEMPTS times with the permit released in between,
  and skips only a chunk that is no longer stored.
- The offer pipeline is one function (`dispatch_fresh_offer`) shared by
  the dispatcher and the direct `replicate_fresh` entry point; the
  8-argument helper and its clippy allow are gone.
- Chunk-carrying protocol fields are encoded as byte strings
  (`serde_bytes`), which has the same postcard layout as a u8 sequence
  (unit-tested) but sizes and serializes in one memcpy pass; an
  oversized body is now refused before anything is allocated.
- PaidNotify shares one `Bytes` buffer across its recipients.
- A new e2e test drives the PUT pipeline through the real channel with
  a missing-chunk event queued ahead of a real one; the harness keeps
  the fresh-write sender so the drainer stays alive in tests.

Pins: saorsa-core c5aae116, saorsa-transport 4e0fffc5 (both
perf/replication-send-path).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…rite pipeline

Review "patch first" follow-up for the two-stage fresh replication.

Two e2e tests over the real harness:

- `fresh_write_pipeline_drains_a_burst_larger_than_the_offer_budget`
  queues 3 x MAX_PENDING_FRESH_OFFERS writes in one go, asserts every one
  replicates to another node, and asserts all pending-offer permits are
  back afterwards, so a burst neither loses a write to back-pressure nor
  leaks a permit.
- `fresh_write_pipeline_retries_a_transient_read_failure` makes the chunk
  file unreadable on disk, waits until the dispatcher's failed read has
  marked the chunk suspect (observable through `exists`), restores it
  inside FRESH_READ_RETRY_DELAY, and asserts the write replicates and the
  suspect mark clears on the source.

Adds the test-utils accessor `ReplicationEngine::pending_offer_permits_available`
and shares the store/poll helpers with the existing pipeline test.

Pins: saorsa-transport 5da61737 (legacy `send` copies only on the QUIC
path, `send_bytes` regression tests), saorsa-core b18e3754 (pin bump
only).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The offer dispatcher read the chunk back with `get_raw`, which skips the
content-address check every other serve path applies. The bytes were
checked when the PUT stored them, but under a backlog they come off disk
much later, and a chunk that no longer hashes to its key would be offered
to the whole close group: every receiver rejects it and charges the
sender an application failure, and the store's own view that the chunk
is wrong (`known_wrong`) was never consulted.

The dispatcher now reads with `ChunkStore::get`, the read the fetch path
serves from. A chunk that fails verification is quarantined by that
read, so the node stops advertising it and ordinary repair replaces it;
its read-back retry finds it gone and skips it. Verification follows the
store's `verify_on_read` setting, like every serve.

Review follow-up (point 1).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…tcher

The offer dispatcher slept `FRESH_READ_RETRY_DELAY` inline before
re-queueing a write whose chunk could not be read. It is the only
dispatcher, so releasing the permit first did not help: nothing was
drained while it slept. The expected cause of a failed read, exhausted
descriptors, is shared by the reads that follow it, so a fault turned
into one delay per queued write, twice over, with every healthy write
held behind them.

The delay now runs on a tracked task of its own that re-queues the
write when it expires (or does nothing at shutdown), and the dispatcher
moves straight on to the next event.

Review follow-up (point 2).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…trings

`FetchResponse::Success::data` (a whole chunk) and
`SubtreeSliceItem::Present::bao_slice` were plain `Vec<u8>` fields, which
serde walks one byte at a time. The size-then-write `encode` compiles the
sizing pass down to a length in optimised builds, so it did not make these
more expensive, but the per-byte write pass is still the bulk of encoding
a fetch response, and decoding one grows the buffer from serde's cautious
size hint, leaving a 3 MiB chunk in a 4 MiB allocation.

Both are now byte strings (`serde_bytes`), like the fresh-offer fields:
postcard lays them out exactly like a `u8` sequence, so the wire is
unchanged, and each direction is one copy. For a 4 MiB fetch response,
encoding drops from ~1.6 ms to ~0.07 ms and decoding from ~2.6 ms to
~0.05 ms (postcard 1.1, release build).

The wire-equivalence test now checks every annotated field against its
pre-annotation layout in both directions and on either side of the
varint length boundaries, and a new test pins exact-capacity decoding of
a fetched chunk.

Review follow-up (point 3).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…dler and hold the offer budget

The pipeline tests did not prove what they claimed (review point 4):

- The harness never set the fresh-write sender on the `AntProtocol`, so
  the handler's emit path was not exercised; the tests posted events by
  hand. `create_node` now wires the sender into the protocol as a node
  does and parks the receiver for the engine `start_node` creates, so
  every PUT through a test node's handler feeds its pipeline.
- The burst test wrote tiny chunks and never observed the cap binding,
  and it was flaky: a one-source burst of small chunks overruns the
  receivers' per-source fresh-offer admission cap (12), and when all four
  receivers refused the same key the write never landed (7 failures in 10
  isolated runs on the previous head). It now holds the send stage
  (`ReplicationEngine::hold_replication_sends`), shows the dispatcher
  encodes exactly `MAX_PENDING_FRESH_OFFERS` offers and stops, then that
  every write is offered once sends resume and every permit comes back.
  Offers are counted at the sender (`fresh_offers_dispatched`), so
  receiver admission cannot make the test flaky.
- The missing-chunk test now PUTs through the handler, and checks the
  skipped write was not offered and its permit came back.
- The retry test injects the fault with a directory where the chunk file
  should be, which the store refuses to read for any user on any
  platform, so it no longer skips silently as root. It checks the failed
  read released its permit, and that a healthy write queued behind it is
  offered within half the retry delay, which a dispatcher sleeping out
  the delay cannot do.
- A new test corrupts a stored chunk on disk and checks it is quarantined
  and never offered.
- Each wait has its own deadline.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…hat is bounded

- The read-back is the verified serve path, and why.
- Read-back retries wait off the dispatcher; what happens after the last
  one, and what brings an unadvertised key back.
- The fetch-response and subtree-slice payloads are byte strings too.
- Soften "bounded": chunk buffers are bounded, the queues ahead of them
  grow by one stripped proof per write under sustained overload.
- The pending-offer budget bounds the sender, not receiver admission.
- Validation lists the new tests.

Review follow-up.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The four pipeline tests were the only harness tests in the file without
`#[serial]`, so serial_test let them overlap the serial ones. The burst
test drives 24 offers per receiver past the per-source admission cap,
and those refusals land in the process-wide counters that
`normal_upload_never_reaches_fresh_offer_capacity` asserts stay at zero:
run together with two test threads, the capacity test failed 3 of 3
times. CI runs e2e with one thread, which hid it.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…osed

`replicate_fresh` said its permit wait "only fails at shutdown" and
`hold_replication_sends` said it returns `None` only when shutting down.
Neither semaphore is ever closed, so neither can fail; say so.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ares the backlog

Read-back retries were a fixed three attempts one second apart. A read
fault is usually store-wide (exhausted descriptors), and then every
queued write fails together; since each failure frees its permit for the
next write at once, the whole backlog spent its attempts within about
two seconds and lost its fresh offers and possession checks, with every
key left marked suspect. origin/main was immune (the bytes travelled in
the event), and the earlier inline sleep paced it at about one attempt a
second at the cost of stalling healthy writes.

Retries now back off: seven attempts, the pause doubling from one second
(1, 2, 4, 8, 16, 32 s), still on their own tasks so a single unreadable
chunk holds nothing behind it. A store-wide fault shorter than about a
minute now costs no offers. Only the first failure and the final give-up
log at WARN, so a fault across a large backlog does not flood the log.

ADR-0017 describes the window and the whole-backlog case, and lists the
readers that actually clear a suspect mark.

Deep-review follow-up (K1).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…annel

The capacity driver looped over `ReplicationEngine::replicate_fresh`,
which since the two-stage split waits for a pending-offer permit after
announcing each write. On fast storage that paced the loop: 33-39 of the
48 calls blocked, spreading PaidNotify over ~0.8 s where the production
drainer sends it at arrival rate (~0.2 s). With receiver PaidNotify
handling slowed to 30 ms the production path dropped 77-79 notifies at
the per-peer cap while this test still reported zero, so it no longer
guarded the pressure it says it applies.

It now stores each chunk and queues a `FreshWriteEvent` on the channel
the harness wires into the PUT handler, the path production takes. It
still sees no refusals with the current caps.

Deep-review follow-up (K3).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…andler PUT

- The retry test's half-delay window opened only after a handler PUT
  returned, and the fault was detected on a 200 ms poll, so on a slow
  runner a dispatcher that slept the delay out inline could still pass.
  The healthy chunk is now stored up front and queued straight after the
  fault is seen, the fault is detected on a 10 ms poll, and the window is
  counted from the fault. The restore now lands well before the first
  retry, too.
- The corrupt-chunk test no longer claims to observe the retry; it
  checks nothing is offered by the time the retry is due.
- `wait_until_every` takes the poll interval, so the retry test does not
  hand-roll its own loop.
- The handler PUT helper is bounded by the harness's chunk-operation
  timeout (now public) instead of hanging CI if the handler stalls.

Deep-review follow-up (K6, K8, K25, K27).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
`read_bounded` read into `Vec::new()` through `Take<File>`, which gets
no size hint and grows by doubling: a full 4 MiB chunk ended in an
8 MiB allocation, a 3 MiB one in 4 MiB. Every GET serve, audit read and
now every fresh-offer read-back went through it, and a fetch response
holds that buffer while it is sent.

The buffer is now sized from the file's length, capped at the chunk
ceiling so a planted oversized file still cannot make it allocate more
than a chunk up front. A test pins exact capacity for a full chunk,
verified and raw.

Deep-review follow-up (K15).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The client response path copied each encoded response with `to_vec()`
before `send_message`, although it is already `Bytes` and
`send_message` now takes `impl Into<Bytes>`. A GET response carries up
to a whole chunk, so that was a full extra copy per GET, held for the
send. The length used for traffic accounting is read before the hand-off.

Deep-review follow-up (K12).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Tie the pending-offer permit to every handle of the encoded offer:
  `EncodedOffer` now owns the buffer and the permit and is wrapped with
  `Bytes::from_owner`, so the per-peer sends and the transport share one
  reference count and the permit is released with the last handle
  anywhere, not the last send task's `Arc`. With today's saorsa-core,
  which frames into its own buffer, nothing changes; a transport that
  queued the `Bytes` itself would otherwise have loosened the bound.
  `bytes` needs 1.10.1 for this (from_owner without its `to_vec` leak).
- Move the payment proof into the offer instead of copying it (up to
  512 KiB); the dispatcher owns it and dropped it right after.
- Rename the sender-side `fresh::dispatch_fresh_offer` to
  `fresh::send_fresh_offers`, which no longer shares a name with the
  receiver-side admission function `dispatch_fresh_offer`.
- Drop the explicit `drop(offer_msg)`: no await follows it any more, so
  it freed nothing earlier and its comment claimed otherwise.
- `send_paid_notify` is private again and `FreshOfferContext` loses an
  unused `Clone`; both were left from intermediate commits.

Deep-review follow-up (K16, K20, K21, K22, K32).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
`ReplicationMessage::encode` and `web_rtc::encode_binary_response` each
sized a message, refused it over a limit and then encoded it into
exactly that space. Both now call `codec::encode_exact`, which keeps the
one use of `postcard::experimental::serialized_size` in a single place,
and the WebRTC path stops zero-filling a buffer of up to a chunk before
encoding into it.

In `encode`, the family-ceiling rationale had been left below the
encode after the check moved above it; it sits with the check again.

Deep-review follow-up (K10, K19, K30).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ents; scope the offer channel to start()

- Comments in config.rs, handler.rs and mod.rs still said the
  fresh-write drainer takes the pending-offer permit, reads the chunk
  back and schedules possession checks, and that the direct entry point
  schedules them itself. The offer dispatcher and
  `fresh::send_fresh_offers` do all of that; the drainer never waits by
  design. The `refuse_fresh_offer` doc named the deleted
  `fresh::replicate_fresh`. `replicate_fresh`'s doc now says how it
  differs from the PUT path (it waits for a permit and offers the
  caller's bytes).
- The drainer-to-dispatcher channel was kept as two engine fields used
  only by the launchers `start()` calls; it is now created in
  `start_fresh_replication` and handed to both stages.
- The dispatcher's explicit `drop(pending_offer)` did nothing the end of
  the iteration does not; a comment says when the permit goes instead.

Deep-review follow-up (K18, K19, K20, K23).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- `TestHarness::prepopulate_payment_cache_everywhere` replaces the
  all-nodes `cache_insert` loop copied through replication.rs and the
  capacity test.
- The inline dummy-proof literals use `dummy_payment_proof()`, and
  `test_fresh_replication_propagates_to_close_group` uses
  `store_paid_chunk` and `wait_until_replicated` instead of its own
  copies of them.
- `ChunkStore::test_chunk_path` (test-utils) exposes the store's own
  path for a chunk, so the tests that damage or remove a chunk file no
  longer re-derive the on-disk layout.
- The fresh-replication source index is documented as a regular node,
  not the first one (the minimal harness has two bootstraps).
- `start_node` no longer invents a closed channel when a node has no
  fresh-write receiver, which would silently end the drainer; it skips
  the engine with a warning, like a node without an identity, so tests
  that need replication fail on the missing engine.

Deep-review follow-up (K24, K26, K28, K29, K31).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…, and the new tests

- "A fetched chunk is encoded and decoded in one copy each" overclaimed:
  that holds for the replication message, but a fetch answered over
  request/response is wrapped in saorsa-core's envelope, which still
  serializes and decodes the payload per byte. Say so, and that the
  wire-identical fix belongs in saorsa-core.
- The pending-offer permit is released with the last handle to the
  encoded bytes (`Bytes::from_owner`), the proof moves into the offer,
  chunk reads are exactly sized, the exact encoder is shared, and client
  GET responses are not copied before the send.
- Validation lists the new unit tests, the fault-timed retry test and
  the capacity driver's switch to the production channel.

Deep-review follow-up (K11).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The rebase onto main moved the saorsa-core and saorsa-transport revs this
branch patches; the lock now names the rebased commits. ant-protocol
stays on the published 3.1.0 that main took as its release baseline.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@jacderida
jacderida force-pushed the perf/replication-send-path branch from 4dfbb52 to 289b89d Compare October 1, 2026 21:17
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.

4 participants