Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
ff0b1f6
fix(replication): bound encoded fresh offers waiting behind send permits
mickvandijke Sep 21, 2026
1e2e233
perf(replication): cut copies of fresh offers on the send path
mickvandijke Sep 22, 2026
408eea0
perf(replication): share fresh offers with the transport as Bytes
mickvandijke Sep 22, 2026
1663d70
fix(replication): send PaidNotify before waiting for an offer permit
mickvandijke Sep 22, 2026
0b3e43e
fix(replication): two-stage fresh replication with un-gated paid-list…
mickvandijke Sep 22, 2026
bfaef64
test(replication): saturation and read-retry coverage for the fresh-w…
mickvandijke Sep 23, 2026
dcb9c58
fix(replication): read fresh offers back through the verified store path
mickvandijke Sep 29, 2026
e5657fb
fix(replication): retry a failed read-back without stalling the dispa…
mickvandijke Sep 29, 2026
551c6c8
perf(replication): encode fetched chunks and subtree slices as byte s…
mickvandijke Sep 29, 2026
060ddd2
test(replication): drive the fresh-write pipeline through the PUT han…
mickvandijke Sep 29, 2026
7ad7bef
docs(adr-0017): verified read-back, retries off the dispatcher, and w…
mickvandijke Sep 29, 2026
afe5cbc
test(replication): run the fresh-write pipeline tests serially
mickvandijke Sep 29, 2026
d5f400e
docs(replication): the pending-offer and send semaphores are never cl…
mickvandijke Sep 29, 2026
338eb3e
fix(replication): back off read-back retries so a store-wide fault sp…
mickvandijke Sep 29, 2026
52bfed0
test(replication): drive the capacity test through the fresh-write ch…
mickvandijke Sep 29, 2026
b0584c6
test(replication): time the retry test from the fault and bound the h…
mickvandijke Sep 29, 2026
5d50836
perf(storage): read a chunk into a buffer sized from its file
mickvandijke Sep 29, 2026
b4e0d76
perf(node): hand chunk GET responses to the transport without copying
mickvandijke Sep 29, 2026
c480908
refactor(replication): tidy the fresh-offer sender
mickvandijke Sep 29, 2026
21eb07f
refactor: share one exactly-sized postcard encoder
mickvandijke Sep 29, 2026
5208c9e
refactor(replication): name the right fresh-replication stage in comm…
mickvandijke Sep 29, 2026
b779d29
test(e2e): reuse the payment, proof, replication and chunk-path helpers
mickvandijke Sep 29, 2026
5149d11
docs(adr-0017): what the envelope still copies, the permit's lifetime…
mickvandijke Sep 29, 2026
289b89d
chore: refresh Cargo.lock for the rebased send-path pins
jacderida Oct 1, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 3 additions & 4 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

13 changes: 12 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,8 @@ rmp-serde = "1"
hex = "0.4"

# Utilities
bytes = "1"
# 1.10.1: `Bytes::from_owner` (1.9) without its `to_vec` leak (fixed in 1.10.1).
bytes = "1.10.1"
chrono = { version = "0.4", features = ["serde"] }
tempfile = "3"
rand = "0.8"
Expand All @@ -107,6 +108,9 @@ page_size = "0.6"

# Protocol serialization
postcard = { version = "1.1.3", features = ["use-std"] }
# Byte-string encoding for chunk payloads in replication messages (same postcard
# wire layout as a u8 sequence, but serialized and sized in one memcpy pass).
serde_bytes = "0.11"
bao = "0.13.1"

# Shared portable browser profile. The native listener is enabled separately
Expand Down Expand Up @@ -228,6 +232,13 @@ webrtc-direct = [
"dep:self_encryption",
]

[patch.crates-io]
# Copy-free sends (ADR-0017): saorsa-core's `send_message` takes
# `impl Into<Bytes>`, and saorsa-transport writes owned buffers to QUIC
# streams without copying. Drop both once releases include them.
saorsa-core = { git = "https://github.com/WithAutonomi/saorsa-core", rev = "7f815141ff8d5e07f443c25c99607586fa7099cf" }
saorsa-transport = { git = "https://github.com/WithAutonomi/saorsa-transport", rev = "6a772bd3cf806e381cf5706cfe34492bede3141e" }

[profile.release]
lto = true
codegen-units = 1
Expand Down
218 changes: 218 additions & 0 deletions docs/adr/ADR-0017-bounded-fresh-offers-and-copy-free-sends.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,218 @@
# ADR-0017: Bounded fresh-replication offers and copy-free message sends

- **Status:** Proposed
- **Date:** 2026-09-22
- **Decision owners:** <pending>
- **Reviewers:** <pending>
- **Supersedes:** none
- **Superseded by:** none
- **Related:** [ADR-0003](ADR-0003-full-node-detection-and-eviction.md)
(best-effort fresh delivery and possession checks),
[ADR-0005](ADR-0005-replication-repair-hardening.md),
ant-node `perf/replication-send-path`, saorsa-core `perf/replication-send-path`,
saorsa-transport `perf/replication-send-path`,
`ant-testnet/state/comparisons/web-support-memory-diag-0921/` (heap profiles
and per-minute live/RSS reports from the diagnosis)

## Context

Under sustained client uploads on a 60-node DigitalOcean testnet, individual
nodes grew from ~100 MiB to 1–2 GiB of resident memory within an hour and
kept growing for as long as writes continued. Heap profiles taken on the
running nodes attributed roughly 80% of live memory to one site: encoded
`FreshReplicationOffer` messages queued by the fresh-write drainer.

The mechanism was structural rather than a leak:

- Every accepted PUT was pushed to the drainer with the full chunk, and
`replicate_fresh` encoded the offer (chunk plus proof, up to ~4–5 MiB)
immediately, before any send permit was held.
- One send task per close-group peer then waited for one of
`MAX_CONCURRENT_REPLICATION_SENDS` (3) permits while pinning that encoded
buffer. Nothing bounded how many chunks could be waiting in that state.
- On a real network each send holds its permit for seconds (QUIC delivery
acknowledgement, retries, unreachable NAT peers), so a write rate above
the send rate grew the queue without limit. Loopback devnets never showed
it because sends complete instantly.

Once the backlog was bounded, the profile showed the remaining cost per
in-flight send: the same frame existed as the caller's serialized message
*and* as the QUIC stream's copy for the whole transfer, plus transient
copies made while framing (payload clone per channel attempt, owned wire
message for signing, and a doubling `Vec` that left chunk-sized frames with
up to twice their length in capacity).

## Decision Drivers

- The chunk-sized buffers held for replication must stay bounded under any
client write rate, and what queues ahead of them must be small;
replication may be delayed by backpressure but must not be dropped by it.
- The change must not alter the wire format, storage format or payment
logic, so it can ship as a behavioural fix.
- Existing callers of the send APIs in saorsa-core and saorsa-transport must
keep compiling and behaving the same.

## Considered Options

1. Bound the fresh-write channel and drop or block PUT handling when full.
Rejected: either silently loses replication or blocks client responses
on network conditions.
2. Raise `MAX_CONCURRENT_REPLICATION_SENDS`. Rejected: only moves the
knee of the curve and increases bandwidth pressure on home links; the
queue behind the permits would still be unbounded.
3. Keep events small and take a bounded permit before materialising an
offer; separately remove the avoidable copies on the send path. Chosen.

## Decision

We will bound the number of encoded fresh offers that can exist at once and
make the send path hand a single owned buffer down to the QUIC stream:

- `FreshWriteEvent` carries only the key and the payment proof, and fresh
replication runs as two stages. The fresh-write drainer never waits for
chunk back-pressure: for every event, at arrival rate, it records the key
in `PaidForList(self)` and sends `PaidNotify` to the paid close group —
the evidence later repair depends on — then forwards the event to the
offer dispatcher. The dispatcher is the only permit-gated stage: it
acquires a `MAX_PENDING_FRESH_OFFERS` (8) permit before it reads the
chunk back from storage and encodes it; the permit lives with the encoded
offer until the last handle to its bytes is dropped, whether a per-peer
send's or the transport's. A backlog therefore waits as
small queued events, and at most ~40 MiB of encoded offers exist per node.
Nothing is dropped by back-pressure: both queues are unbounded and FIFO,
and every offer is dispatched with the same fan-out, retries and delayed
possession check.
- The read-back is the same verified read the fetch path serves from
(`ChunkStore::get`), not a raw one. The bytes were content-checked when
they were stored, but they now come off disk, possibly long after, and
every receiver rejects an offer that does not hash to its key and charges
the sender for it. A chunk that fails verification is quarantined by that
read, as on any serve, so the node stops advertising it and ordinary
repair replaces it; it is never offered. Like every serve, verification
follows the store's `verify_on_read` setting (on by default).
- A failed read-back is retried up to `MAX_FRESH_READ_ATTEMPTS` (7) times
with the permit released in between, the pause doubling from
`FRESH_READ_RETRY_DELAY` (1, 2, 4, 8, 16 and 32 s). The delay runs on a
task of its own, never on the dispatcher, so a chunk that alone cannot be
read does not hold the healthy writes queued behind it. Read faults tend
to be store-wide (exhausted descriptors), though, and then every queued
write fails together: each failure frees its permit for the next write at
once, so a fixed short retry window would spend the whole backlog's
attempts within seconds. The backoff spreads them over about a minute, and
a store-wide fault shorter than that costs no offers. Only a chunk that is
no longer stored is skipped without retry.
- The chunk and its proof move into the offer rather than being copied, and
`ReplicationMessage::encode` serializes into an exactly-sized buffer
(shared with the WebRTC browser path as `codec::encode_exact`). A chunk is
read from disk into a buffer sized from its file, not grown to twice its
size. The
chunk-carrying fields — the offer's data and proof, `PaidNotify`'s proof,
`FetchResponse::Success::data` and a subtree slice's `bao_slice` — are
byte strings (`serde_bytes`), which postcard lays out exactly like a `u8`
sequence: one copy each way instead of a per-byte loop, and an
exactly-sized buffer on decode.
- The encoded offer is shared as `Bytes` (`Bytes::from_owner`, owning the
buffer and the permit); saorsa-core's `send_message`
accepts `impl Into<Bytes>`, frames the payload through a borrowing
`WireMessageRef` (byte-identical to `WireMessage` on the wire) into an
exactly-sized frame, and passes that frame as `Bytes` to
saorsa-transport's new `send_bytes`, where the QUIC stream takes ownership
via `write_chunks` instead of copying it.

## Consequences

### Positive

- Chunk memory under write load is bounded by configuration: pending offers
plus the three in-flight sends, each held once, instead of growing with
the backlog. On the diagnostic fleets peak live memory fell from 1068 MiB to
298 MiB (mimalloc build) and from 674 MiB to 262 MiB (jemalloc build)
after the backpressure change alone.
- Every large send node-wide (chunk GET responses included) stops paying
for a second copy of its frame during the transfer, and a client GET
response is handed to the transport without being copied first. A
replication fetch response's own encoding and decoding are one copy each;
when it is answered over request/response, saorsa-core's envelope around
it still serializes the payload per byte into a growing buffer and
decodes it the same way (a `serde_bytes` payload there would be
wire-identical, and is left for saorsa-core).
- No wire, storage or API break: `send(&[u8])` remains and copies once as
before; `Vec<u8>` callers of `send_message` convert without copying.

### Negative / Trade-offs

- Replication of a burst of writes is spread out in time rather than
encoded eagerly; the delayed possession check is scheduled after each
offer's sends are dispatched, so it shifts by the same amount. A chunk
fetched seconds after its upload can therefore have fewer replicas than
before (the 2026-09-22 comparison measured downloads of just-uploaded
files 8% slower). Paid-list evidence is not affected, and the previous
unbounded fan-out lost that evidence outright under load (2,795
"paid notify dropped at admission" in one hour on the baseline fleet).
- It is not a hard memory bound. The queues ahead of the permit hold a key
and a stripped payment proof per write (about 40 KB for a single-node
proof, about 130 KB for a merkle proof, at most 512 KiB), so a sustained
write rate above the send rate still grows memory, roughly a hundred times
more slowly than one encoded chunk per write did. Capping those queues
would mean dropping offers, which is left to a later decision if testnets
show a sustained backlog.
- The dispatcher re-reads and re-verifies each chunk when its permit
arrives: one extra read and one BLAKE3 pass per accepted write, the same
as serving it once.
- The pending-offer budget bounds the sender's memory, not what receivers
admit. Offers are one-way and the sender does not read the answer, so a
burst of small chunks from one sender can exceed a receiver's per-source
fresh-offer admission cap, and that receiver refuses the excess. Chunk-
sized sends are slow enough that this is rare in practice, and neighbor
sync fills the gap, as it does for any refused offer.

### Neutral / Operational

- `MAX_PENDING_FRESH_OFFERS` and `MAX_CONCURRENT_REPLICATION_SENDS` are
the two knobs; raising the first trades memory for burst absorption.
- A write whose read-back fails `MAX_FRESH_READ_ATTEMPTS` times is not
offered, and neither is its possession check scheduled. A store-wide
fault lasting longer than the retry window does this to every write
queued at the time. Each failed read also leaves the key marked suspect,
so the node stops advertising it. That is self-correcting: the store
clears the mark on the next read that succeeds — another holder's
possession probe, a late duplicate offer or client PUT checking what it
holds, a client GET, or this node's own neighbor sync re-fetching a key it
no longer claims — and on a repair or re-put. The paid-list evidence went
out before the read was attempted, and the client stored the chunk
directly on a majority of the close group, each of which fans it out, so
one node's lost offers cost a replica for a while, not the data.
- The signing step still serializes the payload once to produce the signed
bytes; changing that would alter the signature input and is out of scope.

## Validation

- Unit tests: exact-capacity encoding of chunk-sized offers and
exact-capacity decoding of fetched chunks; exact-capacity chunk reads;
the read-back retry schedule; the pending-offer permit outliving every
handle to the offer; wire equivalence of every byte-string field with its
`u8`-sequence layout, both directions, across the varint length
boundaries (ant-node); byte-for-byte equivalence of `WireMessageRef` with
`WireMessage` (saorsa-core).
- E2E tests over the real harness, whose nodes wire the PUT handler to the
fresh-write pipeline as a node does: a PUT through the handler replicates
and a missing chunk queued ahead of it is skipped; with the send stage
held, a burst three times the budget encodes exactly
`MAX_PENDING_FRESH_OFFERS` offers, then all of them once sends resume,
and returns every permit; a failed read-back releases its permit, does
not delay a healthy write queued behind it (timed from the fault), and is
offered once the fault clears; a chunk corrupted on disk is quarantined
and never offered. The capacity driver feeds the same fresh-write channel,
so it measures receiver admission under the production drainer's pacing.
Each pipeline test fails against a build with its fix reverted.
- Testnet evidence (2026-09-21): with the backpressure change, the node
that had reached 1051 MiB live memory stayed flat at 0.0 MiB/min with a
150 MiB peak, and the worst bootstrap's queued offers dropped from 101
(463 MiB) to 6 (17.7 MiB) in heap profiles.
- Review trigger: any change to fresh replication fan-out, send permits, or
the wire-message framing must re-run the memory diagnostics under
sustained uploads and confirm live memory stays bounded.

## Notes for AI-assisted work

AI tools may help draft this ADR, but **must not mark it Accepted without human review**. Accepted ADRs are immutable: create a new superseding ADR rather than editing an Accepted ADR.
1 change: 1 addition & 0 deletions docs/adr/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,3 +36,4 @@ See [`TOOLING.md`](./TOOLING.md) for `adrs`, `adr-kit`, and AI harness setup.
- [ADR-0013: Settlement version and pre-payment compatibility](./ADR-0013-settlement-version-and-pre-payment-compatibility.md)
- [ADR-0015: Direct browser clients over WebRTC Direct](./ADR-0015-direct-browser-clients-over-webrtc-direct.md)
- [ADR-0016: Pointers — paid mutable references with an immutable owner](./ADR-0016-pointers-immutable-owner.md)
- [ADR-0017: Bounded fresh-replication offers and copy-free message sends](./ADR-0017-bounded-fresh-offers-and-copy-free-sends.md)
57 changes: 57 additions & 0 deletions src/codec.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
//! Exactly-sized postcard encoding shared by the node's wire formats.

use serde::Serialize;

/// Why [`encode_exact`] produced no bytes.
#[derive(Debug)]
pub enum ExactEncodeError {
/// Postcard could not serialize the value.
Serialize(postcard::Error),
/// The encoding would be `size` bytes, over `limit`; nothing was allocated.
TooLarge { size: usize, limit: usize },
}

/// Encode `value` into a buffer of exactly its serialized size, refusing
/// before anything is allocated if that size is over `limit`.
///
/// A growing `Vec` would otherwise keep up to twice the needed capacity for
/// as long as the bytes are held, and chunk-carrying messages are held while
/// they are sent. Sizing first costs next to nothing: postcard's sizing pass
/// adds a byte string's length in one step, and compiles a plain `u8`
/// sequence's per-byte count down to a length too.
pub fn encode_exact<T: Serialize + ?Sized>(
value: &T,
limit: usize,
) -> Result<Vec<u8>, ExactEncodeError> {
let size =
postcard::experimental::serialized_size(value).map_err(ExactEncodeError::Serialize)?;
if size > limit {
return Err(ExactEncodeError::TooLarge { size, limit });
}
postcard::to_extend(value, Vec::with_capacity(size)).map_err(ExactEncodeError::Serialize)
}

#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;

#[test]
fn encodes_into_exactly_the_serialized_size() {
let value = vec![0xABu8; 70_000];
let bytes = encode_exact(&value, usize::MAX).expect("encode");
assert_eq!(bytes, postcard::to_stdvec(&value).expect("reference"));
assert_eq!(bytes.capacity(), bytes.len());
}

#[test]
fn refuses_an_encoding_over_the_limit() {
let value = vec![1u8; 100];
let size = postcard::to_stdvec(&value).expect("reference").len();
assert!(encode_exact(&value, size).is_ok());
assert!(matches!(
encode_exact(&value, size - 1),
Err(ExactEncodeError::TooLarge { size: s, limit }) if s == size && limit == size - 1
));
}
}
1 change: 1 addition & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
pub mod ant_protocol;
pub mod browser;
pub mod client;
mod codec;
pub mod config;
pub mod devnet;
pub mod error;
Expand Down
11 changes: 7 additions & 4 deletions src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1108,24 +1108,27 @@ impl RunningNode {
let traffic_key = handled.traffic_key;
match handled.response {
Ok(Some(response)) => {
let response_len = response.len();
let send_started = Instant::now();
// Handed over as the `Bytes` it already is: a chunk GET
// response is up to a whole chunk, and `to_vec` copied it.
let send_result = p2p
.send_message(source, response_topic, response.to_vec(), &[])
.send_message(source, response_topic, response, &[])
.await;
if let Some(telemetry) = telemetry {
telemetry.finish_send(send_started.elapsed(), send_result.is_ok());
}
// V2-834: attribute response bytes only once the send is
// confirmed; failed sends are itemised separately.
match (&send_result, traffic_key) {
(Ok(()), Some(key)) => storage_traffic::record_tx(key, response.len()),
(Ok(()), Some(key)) => storage_traffic::record_tx(key, response_len),
(Ok(()), None) => {
storage_traffic::record_tx(
storage_traffic::ChunkResponseKey::Other,
response.len(),
response_len,
);
}
(Err(_), _) => storage_traffic::record_send_failed(response.len()),
(Err(_), _) => storage_traffic::record_send_failed(response_len),
}
if let Err(e) = send_result {
warn!("Failed to send {data_type} protocol response to {source}: {e}");
Expand Down
Loading
Loading