fix(replication): bound the fresh-offer backlog and send offers without copying - #233
mickvandijke wants to merge 24 commits into
Conversation
grumbach
left a comment
There was a problem hiding this comment.
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 callsset_fresh_write_senderon theAntProtocol. All three tests postFreshWriteEventby hand, so the handler change insrc/storage/handler.rsis not covered by them. Thefresh_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.
b0263b3 to
35132c1
Compare
dirvine
left a comment
There was a problem hiding this comment.
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 --checkpassed.- 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.
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>
4dfbb52 to
289b89d
Compare
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
FreshReplicationOffermessages 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 ofMAX_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 theirmaintoo (see Pins).Commits:
FreshWriteEventcarries only key + proof; aMAX_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.ReplicationMessage::encodeserializes into an exactly-sized buffer.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.PaidForList(self)and sendsPaidNotifyfor every write at arrival rate, then hands it to the offer dispatcher, the only permit-gated stage. Fresh-offer andPaidNotifypayload fields are byte strings (serde_bytes, wire-identical). One offer pipeline serves the dispatcher andreplicate_fresh.Review follow-ups (2026-09-29, grumbach's review):
ChunkStore::get, the read the fetch path serves from, instead ofget_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).FetchResponse::Success::dataandSubtreeSliceItem::Present::bao_slicegetserde_bytestoo (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-writeencodeitself costs no second pass there — the sizing loop compiles to a length.)AntProtocolas 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.Deep-review follow-ups (2026-09-29; 33 candidates, each verified before fixing):
#[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).replicate_freshnow waits for a permit, which paced the capacity loop on fast storage and hid PaidNotify drops the production drainer would cause.Take<File>, leaving a 4 MiB chunk in an 8 MiB buffer on every GET, audit and read-back.to_vec()onBytes.Bytes::from_owner), the proof moves instead of being copied, and the sender is renamedfresh::send_fresh_offers(it shared a name with the receiver-sidedispatch_fresh_offer).codec::encode_exactserves replication and the WebRTC path (which stops zero-filling a chunk-sized buffer), keeping the experimentalserialized_sizecall in one place.start(), no-op drops; docs for the never-closed semaphores.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.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-transport5da61737(#171), bothperf/replication-send-pathrebased onto the 0.28.0 / 0.37.0 releases; ant-protocol stays on main's pointer pind11d1010— 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 readmessage_data.len()after the payload had moved into the send; it now records the already-computedwire_len.Linear issue
Closes V2-1359
Risk tier
Compatibility
serde_bytesfields checked against theiru8-sequence layout in both directions) and the saorsa-core frame are byte-identical.ant_node::replication::freshmodule,FreshWriteEventdrops itspub datafield and thepub async fn replicate_freshis gone (its crate-internal successor isfresh::send_fresh_offers). No sibling repo uses either, but strictly that is a public-API break — see the semver note below.ReplicationEngine::replicate_freshkeeps 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.src/node.rsdoes, so any paid PUT in an e2e test fans out. The opt-infirst_audit_abA/B driver therefore compares a different workload if only one side contains this PR.Semver impact
(Proposed "fix" for the node's behaviour; the two removed
replication::freshitems above would make it "breaking" if the library API is held to semver — reviewer to confirm.)Test evidence
Full CI, mirrored locally at
4dfbb521(everyci.ymljob 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 oldmigration_signal.rslint);cargo doc --all-features --no-depswithRUSTDOCFLAGS=-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_attacks19,poc_audit_handler_live16,poc_bootstrap_stall3,poc_shutdown_lmdb_drain1,migration_reclaims_disk2,migration_crash_safety5 (3 ignored),migration_shared_volume5,storage_scale2,webrtc_direct_devnet3 + 1 (--ignored) — all passed.cargo check --all-targets --all-features).da51a65c(the PR now targetsmain, soci.ymlruns): 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.get_rawread-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-nodeb0263b32, saorsa-core53f1fc6f, saorsa-transport1766037e, ant-protocola66ddcfb, ant-clientaf8885f(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)f55b7df8(0922, first 30 min)web-support50167d39(0922, first 30 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 onweb-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/.ant-testnet/state/comparisons/web-support-memory-diag-0921/.f55b7df8vsweb-support50167d39): 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). Thebytesrequirement rises from1to1.10.1forBytes::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, raisingMAX_PENDING_FRESH_OFFERStrades memory for burst absorption without touching the wire.🤖 Generated with Claude Code