diff --git a/Cargo.lock b/Cargo.lock index 7e146a36..2c85fc5f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -862,6 +862,7 @@ dependencies = [ "self_encryption", "semver 1.0.28", "serde", + "serde_bytes", "serde_json", "serial_test", "sha2", @@ -5305,8 +5306,7 @@ dependencies = [ [[package]] name = "saorsa-core" version = "0.28.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "84923dc4955c42dfad8e06cfc34ef75b8e89d4d2fed048b18b81dc9ebadf790b" +source = "git+https://github.com/WithAutonomi/saorsa-core?rev=7f815141ff8d5e07f443c25c99607586fa7099cf#7f815141ff8d5e07f443c25c99607586fa7099cf" dependencies = [ "anyhow", "async-trait", @@ -5374,8 +5374,7 @@ dependencies = [ [[package]] name = "saorsa-transport" version = "0.37.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "261ed20ee1ed8555594fef412b428fa20c21d2cbeebfe3f11868c7df9489a410" +source = "git+https://github.com/WithAutonomi/saorsa-transport?rev=6a772bd3cf806e381cf5706cfe34492bede3141e#6a772bd3cf806e381cf5706cfe34492bede3141e" dependencies = [ "anyhow", "async-trait", diff --git a/Cargo.toml b/Cargo.toml index a129dd58..11b23ac7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" @@ -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 @@ -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`, 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 diff --git a/docs/adr/ADR-0017-bounded-fresh-offers-and-copy-free-sends.md b/docs/adr/ADR-0017-bounded-fresh-offers-and-copy-free-sends.md new file mode 100644 index 00000000..ee76da8a --- /dev/null +++ b/docs/adr/ADR-0017-bounded-fresh-offers-and-copy-free-sends.md @@ -0,0 +1,218 @@ +# ADR-0017: Bounded fresh-replication offers and copy-free message sends + +- **Status:** Proposed +- **Date:** 2026-09-22 +- **Decision owners:** +- **Reviewers:** +- **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`, 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` 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. diff --git a/docs/adr/README.md b/docs/adr/README.md index 5061984e..cb17f2f3 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -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) diff --git a/src/codec.rs b/src/codec.rs new file mode 100644 index 00000000..3f401ee0 --- /dev/null +++ b/src/codec.rs @@ -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( + value: &T, + limit: usize, +) -> Result, 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 + )); + } +} diff --git a/src/lib.rs b/src/lib.rs index 852c837a..3adbaace 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -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; diff --git a/src/node.rs b/src/node.rs index 56813244..b926fd85 100644 --- a/src/node.rs +++ b/src/node.rs @@ -1108,9 +1108,12 @@ 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()); @@ -1118,14 +1121,14 @@ impl RunningNode { // 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}"); diff --git a/src/replication/config.rs b/src/replication/config.rs index 276ba1bb..7f35d2fe 100644 --- a/src/replication/config.rs +++ b/src/replication/config.rs @@ -170,6 +170,47 @@ pub const SELF_LOOKUP_INTERVAL_MAX: Duration = Duration::from_secs(SELF_LOOKUP_I /// at most ~12 MB queued for the upload link at any instant. pub const MAX_CONCURRENT_REPLICATION_SENDS: usize = 3; +/// Maximum number of encoded fresh-replication offers held in memory. +/// +/// Each accepted write is encoded once (chunk plus proof, up to ~4 MB) and +/// that buffer stays alive until the last of its per-peer sends completes. +/// With only `MAX_CONCURRENT_REPLICATION_SENDS` transfers in flight, a write +/// rate above the network's send rate would otherwise queue an unbounded +/// number of encoded offers behind the send permits. The offer dispatcher +/// takes one of these permits before it reads and encodes a chunk, so the +/// backlog waits as small queued events instead of chunk-sized buffers. +pub const MAX_PENDING_FRESH_OFFERS: usize = 8; + +/// How many times the offer dispatcher tries to read an accepted chunk back +/// from storage before giving up on its fresh offers. +/// +/// The chunk was stored moments earlier, so a failed read is a transient +/// fault (exhausted descriptors, an I/O hiccup) far more often than a lost +/// chunk, and such a fault is usually store-wide: every queued write fails +/// at once. The retries back off (see [`fresh_read_retry_delay`]), so these +/// attempts span about a minute and a fault shorter than that costs no +/// offers. A lost chunk reports `None` and is skipped without retry. +pub const MAX_FRESH_READ_ATTEMPTS: u32 = 7; + +/// Pause before the first retry of a failed chunk read-back in the offer +/// dispatcher; each later retry waits [`FRESH_READ_RETRY_BACKOFF_FACTOR`] +/// times longer than the one before. +pub const FRESH_READ_RETRY_DELAY: Duration = Duration::from_secs(1); + +/// Growth of the pause between successive read-back retries. +pub const FRESH_READ_RETRY_BACKOFF_FACTOR: u32 = 2; + +/// Pause before retrying a read-back that has failed `failed_attempts` times. +/// +/// [`FRESH_READ_RETRY_DELAY`], grown by [`FRESH_READ_RETRY_BACKOFF_FACTOR`] +/// per earlier failure. With [`MAX_FRESH_READ_ATTEMPTS`] the retries come 1, +/// 2, 4, 8, 16 and 32 seconds apart. +#[must_use] +pub fn fresh_read_retry_delay(failed_attempts: u32) -> Duration { + let growth = FRESH_READ_RETRY_BACKOFF_FACTOR.saturating_pow(failed_attempts.saturating_sub(1)); + FRESH_READ_RETRY_DELAY.saturating_mul(growth) +} + /// Maximum number of concurrent in-flight audit-responder tasks. /// /// The LIGHT audit-responder handlers — responsible-chunk audits and subtree @@ -2079,4 +2120,26 @@ mod tests { "a capped retry must still get several looks inside one entry lifetime" ); } + + #[test] + fn fresh_read_retries_back_off_over_about_a_minute() { + let delays: Vec = (1..MAX_FRESH_READ_ATTEMPTS) + .map(fresh_read_retry_delay) + .collect(); + let expected: Vec = [1, 2, 4, 8, 16, 32] + .into_iter() + .map(Duration::from_secs) + .collect(); + assert_eq!(delays, expected); + assert_eq!( + delays.iter().sum::(), + Duration::from_secs(63), + "a store-wide read fault shorter than this costs no fresh offers" + ); + // Out-of-range inputs saturate rather than panic. + assert_eq!(fresh_read_retry_delay(0), FRESH_READ_RETRY_DELAY); + assert!( + fresh_read_retry_delay(u32::MAX) >= fresh_read_retry_delay(MAX_FRESH_READ_ATTEMPTS) + ); + } } diff --git a/src/replication/fresh.rs b/src/replication/fresh.rs index 80e9766e..af656cb2 100644 --- a/src/replication/fresh.rs +++ b/src/replication/fresh.rs @@ -2,22 +2,28 @@ //! //! When a node accepts a newly written record with valid `PoP`: //! 1. Store locally (already done by chunk handler). -//! 2. Send fresh offers to `CLOSE_GROUP_SIZE` nearest peers (excluding self). -//! 3. Send `PaidNotify` to all peers in `PaidCloseGroup(K)`. +//! 2. Record the key in `PaidForList(self)` and send `PaidNotify` to every +//! peer in `PaidCloseGroup(K)` — immediately, never behind back-pressure. +//! 3. Send fresh offers to `CLOSE_GROUP_SIZE` nearest peers (excluding self), +//! bounded by the pending-offer permits so a write burst cannot pile up +//! chunk-sized buffers. +use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use crate::logging::{debug, warn}; +use bytes::Bytes; use rand::Rng; use saorsa_core::identity::PeerId; use saorsa_core::P2PNode; -use tokio::sync::Semaphore; +use tokio::sync::{mpsc, OwnedSemaphorePermit, Semaphore}; use crate::ant_protocol::XorName; use crate::replication::config::{ ReplicationConfig, FRESH_REPLICATION_DELIVERY_MAX_RETRIES, REPLICATION_PROTOCOL_ID, }; use crate::replication::paid_list::PaidList; +use crate::replication::possession::PossessionCheckEvent; use crate::replication::protocol::{ FreshReplicationOffer, PaidNotify, ReplicationMessage, ReplicationMessageBody, }; @@ -26,48 +32,101 @@ use crate::replication::protocol::{ /// /// Sent from the chunk PUT handler to the replication engine via an /// unbounded channel so that the PUT response is not blocked by -/// replication fan-out. +/// replication fan-out. The event deliberately carries no chunk bytes: the +/// chunk is already on disk, and the offer dispatcher reads it back only +/// once it holds a pending-offer permit, so a replication backlog queues as +/// small events rather than chunk-sized buffers. pub struct FreshWriteEvent { /// Content-address of the stored chunk. pub key: XorName, - /// The chunk data. - pub data: Vec, /// Serialized proof-of-payment. pub payment_proof: Vec, } -/// Execute fresh replication for a newly accepted record. +/// A write whose paid-list evidence has been announced and whose chunk offer +/// is waiting for a pending-offer permit. Carries no chunk bytes. +pub(crate) struct FreshOfferEvent { + pub(crate) key: XorName, + pub(crate) payment_proof: Vec, + /// Storage read-backs attempted so far; see `MAX_FRESH_READ_ATTEMPTS`. + pub(crate) read_attempts: u32, +} + +/// Handles shared by everything that sends fresh offers, so the offer +/// dispatcher task and the direct entry point run one pipeline. +pub(crate) struct FreshOfferContext { + pub(crate) p2p_node: Arc, + pub(crate) config: Arc, + /// Limits concurrent outbound chunk transfers across the engine. + pub(crate) send_semaphore: Arc, + /// Delayed possession checks (ADR-0003) are scheduled here once an + /// offer's sends are dispatched. + pub(crate) possession_check_tx: mpsc::UnboundedSender, + /// Offers encoded and handed to at least one per-peer send, so tests can + /// tell a write that went out from one lost to back-pressure. + pub(crate) dispatched: Arc, +} + +/// An encoded fresh offer and the pending-offer permit it holds. /// -/// Sends fresh offers to close group members (with bounded delivery retries, -/// ADR-0003) and `PaidNotify` to `PaidCloseGroup`. Returns the close-group -/// peers responsible for the key (excluding self) so the caller can schedule -/// the delayed possession check; `PaidNotify` remains fire-and-forget. +/// Wrapped with [`Bytes::from_owner`], so the per-peer send tasks and the +/// transport share the one buffer through reference-counted handles, and the +/// permit is released only when the last handle anywhere is dropped. That is +/// what caps the encoded offers alive at once at `MAX_PENDING_FRESH_OFFERS`. +struct EncodedOffer { + bytes: Vec, + _pending: OwnedSemaphorePermit, +} + +impl AsRef<[u8]> for EncodedOffer { + fn as_ref(&self) -> &[u8] { + &self.bytes + } +} + +/// Rules 6-8: record the paid key locally and announce it to +/// `PaidCloseGroup(K)`. /// -/// The `send_semaphore` limits how many outbound chunk transfers can be -/// in-flight concurrently across the entire replication engine, preventing -/// bandwidth saturation on home broadband connections. -pub async fn replicate_fresh( +/// This is the evidence peers need to repair the key later, so it runs the +/// moment a write is accepted and is never gated by the pending-offer permit +/// or the send semaphore; both messages are small metadata. +pub(crate) async fn announce_paid_write( key: &XorName, - data: &[u8], proof_of_payment: &[u8], + paid_list: &PaidList, p2p_node: &Arc, - paid_list: &Arc, config: &ReplicationConfig, - send_semaphore: &Arc, -) -> Vec { - let self_id = *p2p_node.peer_id(); - +) { // Rule 6: Node that validates PoP adds K to PaidForList(self). if let Err(e) = paid_list.insert(key).await { warn!("Failed to add key {} to PaidForList: {e}", hex::encode(key)); } + // Rules 7-8: PaidNotify to every member of PaidCloseGroup(K). + send_paid_notify(key, proof_of_payment, p2p_node, config).await; +} + +/// Rules 2-3: send fresh offers to the close group and schedule the delayed +/// possession check (ADR-0003) for the responsible peers. +/// +/// `pending_offer` is the caller's permit from the pending-offer semaphore; +/// it is held with the encoded offer until the last handle to it is dropped. +/// `data` and `proof_of_payment` are taken by value so they move into the +/// offer instead of being copied. +pub(crate) async fn send_fresh_offers( + ctx: &FreshOfferContext, + key: &XorName, + data: Vec, + proof_of_payment: Vec, + pending_offer: OwnedSemaphorePermit, +) { + let self_id = *ctx.p2p_node.peer_id(); - // Rule 2-3: Send fresh offers to CLOSE_GROUP_SIZE nearest peers - // (excluding self). Use self-inclusive query to get the true close group, - // then filter self out. - let closest = p2p_node + // Use the self-inclusive query to get the true close group, then filter + // self out. + let closest = ctx + .p2p_node .dht_manager() - .find_closest_nodes_local_with_self(key, config.close_group_size) + .find_closest_nodes_local_with_self(key, ctx.config.close_group_size) .await; let target_peers: Vec = closest .iter() @@ -77,8 +136,8 @@ pub async fn replicate_fresh( let offer = FreshReplicationOffer { key: *key, - data: data.to_vec(), - proof_of_payment: proof_of_payment.to_vec(), + data, + proof_of_payment, }; let request_id = rand::thread_rng().gen::(); let offer_msg = ReplicationMessage { @@ -91,17 +150,20 @@ pub async fn replicate_fresh( "Failed to encode FreshReplicationOffer for {}", hex::encode(key), ); - return Vec::new(); + return; }; - // Share one encoded copy across the per-peer send tasks so a retry only - // re-materialises the buffer for the (consuming) send call, keeping the - // common single-attempt path at one clone per peer. - let encoded = Arc::new(encoded); + // One encoded copy serves every per-peer send task and every retry; the + // transport borrows it through `Bytes` instead of taking a copy. The + // pending-offer permit travels with the buffer. + let encoded = Bytes::from_owner(EncodedOffer { + bytes: encoded, + _pending: pending_offer, + }); for peer in &target_peers { - let p2p = Arc::clone(p2p_node); - let data = Arc::clone(&encoded); + let p2p = Arc::clone(&ctx.p2p_node); + let offer = encoded.clone(); let peer_id = *peer; - let sem = Arc::clone(send_semaphore); + let sem = Arc::clone(&ctx.send_semaphore); tokio::spawn(async move { // Acquire a permit before sending — this caps the number of // concurrent outbound replication transfers across the engine. @@ -117,12 +179,7 @@ pub async fn replicate_fresh( let mut attempt = 0u32; loop { match p2p - .send_message( - &peer_id, - REPLICATION_PROTOCOL_ID, - data.as_ref().clone(), - &[], - ) + .send_message(&peer_id, REPLICATION_PROTOCOL_ID, offer.clone(), &[]) .await { Ok(()) => break, @@ -145,23 +202,28 @@ pub async fn replicate_fresh( }); } - // Rule 7-8: Send PaidNotify to every member of PaidCloseGroup(K). - // PaidNotify messages are small metadata (no chunk data), so they don't - // need semaphore gating. - send_paid_notify(key, proof_of_payment, p2p_node, config).await; - debug!( - "Fresh replication initiated for {} to {} peers + PaidNotify", + "Fresh replication initiated for {} to {} peers", hex::encode(key), target_peers.len() ); - target_peers + // Schedule the delayed possession check (ADR-0003) for the responsible + // close-group peers. A closed receiver (engine shutting down) is ignored. + if !target_peers.is_empty() { + ctx.dispatched.fetch_add(1, Ordering::Relaxed); + let _ = ctx.possession_check_tx.send(PossessionCheckEvent { + key: *key, + peers: target_peers, + }); + } } /// Send `PaidNotify(K)` to every peer in `PaidCloseGroup(K)` (fire-and-forget). /// -/// Per Invariant 16: sender MUST attempt delivery to every member. +/// Per Invariant 16: sender MUST attempt delivery to every member. The +/// message is small metadata (no chunk data), so it is neither gated by the +/// send semaphore nor by the pending-offer permit. async fn send_paid_notify( key: &XorName, proof_of_payment: &[u8], @@ -188,7 +250,8 @@ async fn send_paid_notify( warn!("Failed to encode PaidNotify for {}", hex::encode(key)); return; }; - + // One buffer for every recipient; the sends only take handles. + let encoded = Bytes::from(encoded); for node in &paid_group { if node.peer_id == self_id { continue; @@ -206,3 +269,29 @@ async fn send_paid_notify( }); } } + +#[cfg(test)] +#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] +mod tests { + use super::*; + + /// The permit is released with the last handle to the encoded offer, not + /// the first: a slice the transport still holds keeps the offer counted. + #[test] + fn the_pending_offer_permit_outlives_every_handle() { + let permits = Arc::new(Semaphore::new(1)); + let permit = Arc::clone(&permits) + .try_acquire_owned() + .expect("free permit"); + let offer = Bytes::from_owner(EncodedOffer { + bytes: vec![7; 16], + _pending: permit, + }); + let held_elsewhere = offer.slice(4..8); + + drop(offer); + assert_eq!(permits.available_permits(), 0); + drop(held_elsewhere); + assert_eq!(permits.available_permits(), 1); + } +} diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 73d91f45..1beaa8c8 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -79,12 +79,12 @@ use crate::replication::commitment_state::{ PeerCommitmentRecord, PersistedRetention, ResponderCommitmentState, GOSSIP_ANSWERABILITY_TTL, }; use crate::replication::config::{ - max_parallel_fetch, storage_admission_width, ReplicationConfig, MAX_AUDIT_RESPONSES_PER_PEER, - MAX_CONCURRENT_AUDIT_RESPONSES, MAX_CONCURRENT_REPLICATION_SENDS, - MAX_DIGEST_AUDIT_RESPONSES_PER_PEER, MAX_INCOMING_VERIFICATION_KEYS, - MAX_SUBTREE_ROUND1_PER_PEER, MAX_SUBTREE_SESSIONS, MAX_VERIFICATION_KEYS_PER_CYCLE, - REPLICATION_PROTOCOL_ID, SUBTREE_AUDIT_PROTOCOL_ID, SUBTREE_ROUND1_WORK_BURST_BYTES, - SUBTREE_ROUND1_WORK_REFILL_BYTES_PER_SEC, SUBTREE_SESSION_TTL, + fresh_read_retry_delay, max_parallel_fetch, storage_admission_width, ReplicationConfig, + MAX_AUDIT_RESPONSES_PER_PEER, MAX_CONCURRENT_AUDIT_RESPONSES, MAX_CONCURRENT_REPLICATION_SENDS, + MAX_DIGEST_AUDIT_RESPONSES_PER_PEER, MAX_FRESH_READ_ATTEMPTS, MAX_INCOMING_VERIFICATION_KEYS, + MAX_PENDING_FRESH_OFFERS, MAX_SUBTREE_ROUND1_PER_PEER, MAX_SUBTREE_SESSIONS, + MAX_VERIFICATION_KEYS_PER_CYCLE, REPLICATION_PROTOCOL_ID, SUBTREE_AUDIT_PROTOCOL_ID, + SUBTREE_ROUND1_WORK_BURST_BYTES, SUBTREE_ROUND1_WORK_REFILL_BYTES_PER_SEC, SUBTREE_SESSION_TTL, }; use crate::replication::paid_list::PaidList; use crate::replication::protocol::{ @@ -1748,6 +1748,11 @@ pub struct ReplicationEngine { /// Limits concurrent outbound replication sends to prevent bandwidth /// saturation on home broadband connections. send_semaphore: Arc, + /// Bounds how many encoded fresh offers can wait behind `send_semaphore`; + /// see [`MAX_PENDING_FRESH_OFFERS`]. + pending_offer_semaphore: Arc, + /// Fresh offers encoded and handed to their per-peer sends. + fresh_offers_dispatched: Arc, /// Bounds concurrent IN-FLIGHT LIGHT audit-responder tasks (responsible-chunk /// audits + subtree slice round 2). The heavy subtree round 1 has its own /// tighter pool ([`SubtreeRound1Limiter`]). Those are spawned off the serial @@ -1814,16 +1819,18 @@ pub struct ReplicationEngine { subtree_round1: SubtreeRound1Limiter, /// Receiver for fresh-write events from the chunk PUT handler. /// - /// When present, `start()` spawns a drainer task that calls - /// `replicate_fresh` for each event. + /// When present, `start()` spawns the fresh-write drainer, which records + /// paid-list evidence for each event immediately and forwards it to the + /// offer dispatcher. fresh_write_rx: Option>, /// Pointer replication (ADR-0016), when this node stores pointers. pointers: Option>, /// Receiver for fresh pointer writes, taken by `start()`. pointer_fresh_rx: Option>, - /// Sender for delayed possession-check events (ADR-0003). The fresh-write - /// drainer pushes the responsible close-group peers here after each fresh - /// replication; the possession-check scheduler drains the paired receiver. + /// Sender for delayed possession-check events (ADR-0003). Sending a write's + /// fresh offers (`fresh::send_fresh_offers`) pushes its responsible + /// close-group peers here; the possession-check scheduler drains the paired + /// receiver. possession_check_tx: mpsc::UnboundedSender, /// Receiver paired with `possession_check_tx`; taken by the scheduler task. possession_check_rx: Option>, @@ -1917,6 +1924,8 @@ impl ReplicationEngine { recent_provers: Arc::new(RwLock::new(RecentProvers::new())), sig_verify_attempts: Arc::new(RwLock::new(HashMap::new())), send_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_REPLICATION_SENDS)), + pending_offer_semaphore: Arc::new(Semaphore::new(MAX_PENDING_FRESH_OFFERS)), + fresh_offers_dispatched: Arc::new(AtomicU64::new(0)), audit_responder_semaphore: Arc::new(Semaphore::new(MAX_CONCURRENT_AUDIT_RESPONSES)), audit_responder_inflight: Arc::new(RwLock::new(HashMap::new())), audit_responder_metrics: Arc::new(AuditResponderMetrics::default()), @@ -2244,6 +2253,38 @@ impl ReplicationEngine { self.pointers.as_ref() } + /// Test-only: pending-offer permits not currently held by an encoded + /// fresh offer. Equals [`MAX_PENDING_FRESH_OFFERS`] when no fresh + /// replication is in flight, which is how tests prove a burst of writes + /// drained without leaking a permit. + #[cfg(any(test, feature = "test-utils"))] + #[must_use] + pub fn pending_offer_permits_available(&self) -> usize { + self.pending_offer_semaphore.available_permits() + } + + /// Test-only: fresh offers encoded and handed to their per-peer sends + /// since the engine was created. A write skipped because its chunk is gone + /// is not counted. + #[cfg(any(test, feature = "test-utils"))] + #[must_use] + pub fn fresh_offers_dispatched(&self) -> u64 { + self.fresh_offers_dispatched.load(Ordering::Relaxed) + } + + /// Test-only: take every outbound replication send permit. Until the + /// returned permit is dropped, encoded fresh offers wait behind the send + /// stage, so a test can fill the pending-offer budget deterministically. + /// The send semaphore is never closed, so this returns `Some`. + #[cfg(any(test, feature = "test-utils"))] + pub async fn hold_replication_sends(&self) -> Option { + let all = u32::try_from(MAX_CONCURRENT_REPLICATION_SENDS).ok()?; + Arc::clone(&self.send_semaphore) + .acquire_many_owned(all) + .await + .ok() + } + /// Start all background tasks. /// /// `dht_events` must be subscribed **before** `P2PNode::start()` so that @@ -2269,7 +2310,7 @@ impl ReplicationEngine { self.start_fetch_worker(); self.start_verification_worker(); self.start_bootstrap_sync(dht_events); - self.start_fresh_write_drainer(); + self.start_fresh_replication(); self.start_possession_check_scheduler(); if let Some(pointers) = &self.pointers { self.task_handles.push(pointers.start_verification_loop()); @@ -2408,26 +2449,47 @@ impl ReplicationEngine { self.sync_trigger.notify_one(); } - /// Execute fresh replication for a newly stored record, then schedule the - /// delayed possession check for the responsible close-group peers - /// (ADR-0003). The production PUT path schedules via the fresh-write - /// drainer; this direct entry point schedules here so callers (and tests) - /// that drive replication directly still get the possession check. + /// Execute fresh replication for a newly stored record: announce it, then + /// send its offers and schedule the delayed possession check for the + /// responsible close-group peers (ADR-0003), as the PUT path does. + /// + /// Unlike the PUT path, this waits for a pending-offer permit before + /// returning, and it offers the caller's bytes rather than reading them + /// back from storage. pub async fn replicate_fresh(&self, key: &XorName, data: &[u8], proof_of_payment: &[u8]) { - let peers = fresh::replicate_fresh( + fresh::announce_paid_write( key, - data, proof_of_payment, - &self.p2p_node, &self.paid_list, + &self.p2p_node, &self.config, - &self.send_semaphore, ) .await; - if !peers.is_empty() { - let _ = self - .possession_check_tx - .send(possession::PossessionCheckEvent { key: *key, peers }); + // Never closed, so this cannot fail; the arm keeps the call panic-free. + let Ok(pending_offer) = Arc::clone(&self.pending_offer_semaphore) + .acquire_owned() + .await + else { + return; + }; + fresh::send_fresh_offers( + &self.fresh_offer_context(), + key, + data.to_vec(), + proof_of_payment.to_vec(), + pending_offer, + ) + .await; + } + + /// Handles the offer dispatcher and the direct entry point share. + fn fresh_offer_context(&self) -> fresh::FreshOfferContext { + fresh::FreshOfferContext { + p2p_node: Arc::clone(&self.p2p_node), + config: Arc::clone(&self.config), + send_semaphore: Arc::clone(&self.send_semaphore), + possession_check_tx: self.possession_check_tx.clone(), + dispatched: Arc::clone(&self.fresh_offers_dispatched), } } @@ -2435,45 +2497,59 @@ impl ReplicationEngine { // Background task launchers // ======================================================================= - /// Spawn a task that drains the fresh-write channel and triggers - /// replication for each newly-stored chunk. - fn start_fresh_write_drainer(&mut self) { - let Some(mut rx) = self.fresh_write_rx.take() else { + /// Spawn fresh replication's two stages, joined by an unbounded FIFO of + /// key + proof events: the drainer, which announces every write at arrival + /// rate, and the offer dispatcher, the only permit-gated stage. + fn start_fresh_replication(&mut self) { + let Some(writes) = self.fresh_write_rx.take() else { return; }; + let (offer_tx, offer_rx) = mpsc::unbounded_channel(); + self.start_fresh_write_drainer(writes, offer_tx.clone()); + self.start_fresh_offer_dispatcher(offer_rx, offer_tx); + } + + /// Spawn a task that drains the fresh-write channel: it announces each + /// newly-stored chunk and hands it to the offer dispatcher. + fn start_fresh_write_drainer( + &mut self, + mut rx: mpsc::UnboundedReceiver, + offer_tx: mpsc::UnboundedSender, + ) { let p2p = Arc::clone(&self.p2p_node); let paid_list = Arc::clone(&self.paid_list); let config = Arc::clone(&self.config); - let send_semaphore = Arc::clone(&self.send_semaphore); - let possession_tx = self.possession_check_tx.clone(); let shutdown = self.shutdown.clone(); let handle = tokio::spawn(async move { loop { - tokio::select! { + let event = tokio::select! { () = shutdown.cancelled() => break, event = rx.recv() => { let Some(event) = event else { break }; - let peers = fresh::replicate_fresh( - &event.key, - &event.data, - &event.payment_proof, - &p2p, - &paid_list, - &config, - &send_semaphore, - ) - .await; - // Schedule the delayed possession check (ADR-0003) for - // the responsible close-group peers. A closed receiver - // (engine shutting down) is ignored. - if !peers.is_empty() { - let _ = possession_tx.send(possession::PossessionCheckEvent { - key: event.key, - peers, - }); - } + event } + }; + // Stage one never waits for chunk back-pressure: the paid-list + // entry and PaidNotify are what let peers repair the key later, + // so every queued write gets them at arrival rate. + fresh::announce_paid_write( + &event.key, + &event.payment_proof, + &paid_list, + &p2p, + &config, + ) + .await; + if offer_tx + .send(fresh::FreshOfferEvent { + key: event.key, + payment_proof: event.payment_proof, + read_attempts: 0, + }) + .is_err() + { + break; } } debug!("Fresh-write drainer shut down"); @@ -2481,6 +2557,115 @@ impl ReplicationEngine { self.task_handles.push(handle); } + /// Spawn the fresh-offer dispatcher: the only stage of fresh replication + /// that waits for a pending-offer permit. + /// + /// For each forwarded write it acquires a permit, reads the chunk back + /// from storage and dispatches the offers. A missing chunk is skipped; a + /// failed read is retried up to `MAX_FRESH_READ_ATTEMPTS` times, backing + /// off over about a minute with the permit released in between, so a + /// transient I/O fault does not lose the write's replication. The retry + /// waits on its own task, never on the dispatcher. + fn start_fresh_offer_dispatcher( + &mut self, + mut rx: mpsc::UnboundedReceiver, + offer_tx: mpsc::UnboundedSender, + ) { + let storage = Arc::clone(&self.storage); + let pending_offer_semaphore = Arc::clone(&self.pending_offer_semaphore); + let ctx = self.fresh_offer_context(); + let shutdown = self.shutdown.clone(); + let retries = self.detached_task_tracker.clone(); + + let handle = tokio::spawn(async move { + loop { + let event = tokio::select! { + () = shutdown.cancelled() => break, + event = rx.recv() => { + let Some(event) = event else { break }; + event + } + }; + // Wait for a pending-offer permit before touching the chunk so a + // send backlog holds queued events, not encoded chunk buffers. + let pending_offer = tokio::select! { + () = shutdown.cancelled() => break, + permit = Arc::clone(&pending_offer_semaphore).acquire_owned() => { + let Ok(permit) = permit else { break }; + permit + } + }; + let key_hex = hex::encode(event.key); + // The same verified read the fetch path serves from. The bytes + // were content-checked when they were stored, but they come + // off disk now, possibly long after, and every receiver + // charges the sender for an offer that does not hash to its + // key. A chunk that fails verification is quarantined by this + // read, so its retry finds it gone and skips it. + let data = match storage.get(&event.key).await { + Ok(Some(data)) => data, + Ok(None) => { + debug!("Chunk {key_hex} no longer stored, skipping fresh replication"); + continue; + } + Err(e) => { + // The permit is released as this iteration ends, well + // before any retry is due. + let attempts = event.read_attempts + 1; + if attempts >= MAX_FRESH_READ_ATTEMPTS { + warn!( + "Giving up fresh replication of {key_hex} after {attempts} failed reads: {e}" + ); + continue; + } + // One warning per write: a store-wide fault fails the + // whole backlog, once per retry. + if attempts == 1 { + warn!( + "Failed to read chunk {key_hex} for fresh replication; will retry: {e}" + ); + } else { + debug!( + "Failed to read chunk {key_hex} for fresh replication (attempt {attempts}): {e}" + ); + } + // Wait on a task of its own: a read fault is often + // shared (exhausted descriptors), and sleeping here + // would hold every healthy write queued behind this + // one for the whole delay, once per failed read. The + // delay grows, so a store-wide fault spends a write's + // attempts over about a minute rather than seconds. + let retry_tx = offer_tx.clone(); + let retry_shutdown = shutdown.clone(); + let delay = fresh_read_retry_delay(attempts); + retries.spawn(async move { + tokio::select! { + () = retry_shutdown.cancelled() => {} + () = tokio::time::sleep(delay) => { + let _ = retry_tx.send(fresh::FreshOfferEvent { + read_attempts: attempts, + ..event + }); + } + } + }); + continue; + } + }; + fresh::send_fresh_offers( + &ctx, + &event.key, + data, + event.payment_proof, + pending_offer, + ) + .await; + } + debug!("Fresh-offer dispatcher shut down"); + }); + self.task_handles.push(handle); + } + /// Spawn the possession-check scheduler (ADR-0003). /// /// Drains scheduled possession-check events and, for each, waits a @@ -5735,7 +5920,7 @@ fn fresh_offer_structural_rejection( /// Tell `source` its offer for `key` was not taken. /// -/// Note the sender does not currently read this: `fresh::replicate_fresh` uses +/// Note the sender does not currently read this: `fresh::send_fresh_offers` uses /// one-way `send_message`, so the refusal is observed only as a later absence by /// the delayed possession check. Recorded in ADR-0005 as a known gap. async fn refuse_fresh_offer( diff --git a/src/replication/protocol.rs b/src/replication/protocol.rs index b2ca7962..9d617340 100644 --- a/src/replication/protocol.rs +++ b/src/replication/protocol.rs @@ -10,6 +10,7 @@ use saorsa_core::identity::PeerId; use serde::{Deserialize, Serialize}; use crate::ant_protocol::XorName; +use crate::codec::{encode_exact, ExactEncodeError}; use super::types::AuditFailureReason; @@ -46,9 +47,6 @@ impl ReplicationMessage { /// Returns [`ReplicationProtocolError::SerializationFailed`] if postcard /// serialization fails. pub fn encode(&self) -> Result, ReplicationProtocolError> { - let bytes = postcard::to_stdvec(self) - .map_err(|e| ReplicationProtocolError::SerializationFailed(e.to_string()))?; - // The same family ceiling the decoder applies, from the same table and // with the same arms, including the unclassified case. Every receiver // drops a subtree-audit body over that ceiling before decoding it, so @@ -66,13 +64,22 @@ impl ReplicationMessage { // the largest is a round-1 proof at the commitment // key-count cap, pinned under it with headroom by // `max_round1_proof_fits_the_audit_family_ceiling`. + // + // The buffer is sized exactly, and an oversized body is refused + // before anything is allocated. Chunk-carrying bodies run to several + // MiB and are held while they are sent. let max_size = ceiling_for(family_of_variant(self.body.variant_index())); - if bytes.len() > max_size { - return Err(ReplicationProtocolError::MessageTooLarge { - size: bytes.len(), - max_size, - }); - } + let bytes = encode_exact(self, max_size).map_err(|e| match e { + ExactEncodeError::Serialize(e) => { + ReplicationProtocolError::SerializationFailed(e.to_string()) + } + ExactEncodeError::TooLarge { size, limit } => { + ReplicationProtocolError::MessageTooLarge { + size, + max_size: limit, + } + } + })?; // V2-623: cumulative per-variant tx accounting. Every replication send // funnels through here, so this is the single tx choke point. @@ -822,9 +829,13 @@ pub(crate) fn log_served_peers_summary() { pub struct FreshReplicationOffer { /// The record key. pub key: XorName, - /// The record data. + /// The record data. Encoded as a byte string, which postcard lays out + /// exactly like a `u8` sequence, so the wire format is unchanged while + /// serialization and sizing copy the payload in one pass. + #[serde(with = "serde_bytes")] pub data: Vec, /// Proof of Payment (required, validated by receiver). + #[serde(with = "serde_bytes")] pub proof_of_payment: Vec, } @@ -854,6 +865,7 @@ pub struct PaidNotify { /// The record key. pub key: XorName, /// Proof of Payment for receiver-side verification. + #[serde(with = "serde_bytes")] pub proof_of_payment: Vec, } @@ -957,7 +969,10 @@ pub enum FetchResponse { Success { /// The record key. key: XorName, - /// The record data. + /// The record data. Encoded as a byte string, which postcard lays out + /// exactly like a `u8` sequence: one copy of the chunk each way instead + /// of a per-byte loop, and an exactly-sized buffer on decode. + #[serde(with = "serde_bytes")] data: Vec, }, /// Record not found on this peer. @@ -1350,7 +1365,9 @@ pub enum SubtreeSliceItem { block_index: u32, /// Bao verified slice: the block bytes plus the BLAKE3 parent hashes that /// authenticate them against the chunk address. The auditor decodes this - /// to recover the verified block bytes. + /// to recover the verified block bytes. Encoded as a byte string, laid + /// out exactly like a `u8` sequence on the wire. + #[serde(with = "serde_bytes")] bao_slice: Vec, /// Sibling hashes on the path from this block up to the committed /// `nonced_root`, bottom-up. The auditor folds the block leaf with these @@ -1510,6 +1527,17 @@ impl std::error::Error for ReplicationProtocolError {} mod tests { use super::*; + /// Payload lengths on either side of postcard's one- and two-byte varint + /// length boundaries, plus one past 64 KiB. + const BYTE_STRING_WIRE_LENGTHS: [usize; 6] = [0, 127, 128, 16_383, 16_384, 70_000]; + /// Proof length used alongside each payload in the byte-string wire checks. + const BYTE_STRING_WIRE_PROOF_LEN: usize = 300; + /// Block index carried by the subtree-slice item in the byte-string checks. + const BYTE_STRING_WIRE_BLOCK_INDEX: u32 = 7; + /// A payload that is not a power of two, so a doubling buffer would show + /// slack: three MiB and change. + const CHUNK_SIZED_TEST_PAYLOAD_LEN: usize = 3 * 1024 * 1024 + 123; + // === Audit outcome counters === /// No other test mutates the audit outcome counters, so per-slot deltas are @@ -2595,6 +2623,149 @@ mod tests { assert_eq!(decoded.request_id, 7); } + /// Two types that must be interchangeable on the wire: identical bytes, and + /// each decodes what the other encodes. + fn assert_same_wire(annotated: &A, plain: &P) + where + A: Serialize + serde::de::DeserializeOwned, + P: Serialize + serde::de::DeserializeOwned, + { + let wire = postcard::to_stdvec(annotated).unwrap(); + assert_eq!(wire, postcard::to_stdvec(plain).unwrap()); + let as_plain: P = postcard::from_bytes(&wire).unwrap(); + assert_eq!(postcard::to_stdvec(&as_plain).unwrap(), wire); + let as_annotated: A = postcard::from_bytes(&wire).unwrap(); + assert_eq!(postcard::to_stdvec(&as_annotated).unwrap(), wire); + } + + /// `serde_bytes` must not change the wire layout: postcard encodes a byte + /// string and a `u8` sequence identically (varint length + raw bytes). Each + /// annotated field is checked against its pre-annotation layout in both + /// directions, on either side of the varint length boundaries. + #[test] + fn byte_string_fields_encode_like_u8_sequences() { + #[derive(Serialize, Deserialize)] + struct PlainOffer { + key: XorName, + data: Vec, + proof_of_payment: Vec, + } + #[derive(Serialize, Deserialize)] + struct PlainNotify { + key: XorName, + proof_of_payment: Vec, + } + // Only the first variant is mirrored: postcard tags it 0 in both. + #[derive(Serialize, Deserialize)] + enum PlainFetchResponse { + Success { key: XorName, data: Vec }, + } + #[derive(Serialize, Deserialize)] + enum PlainSliceItem { + Present { + key: XorName, + block_index: u32, + bao_slice: Vec, + nonced_siblings: Vec<[u8; 32]>, + }, + } + + let key = [5; 32]; + let proof = vec![9u8; BYTE_STRING_WIRE_PROOF_LEN]; + let siblings = vec![[1u8; 32]]; + for len in BYTE_STRING_WIRE_LENGTHS { + let data: Vec = (0..=u8::MAX).cycle().take(len).collect(); + assert_same_wire( + &FreshReplicationOffer { + key, + data: data.clone(), + proof_of_payment: proof.clone(), + }, + &PlainOffer { + key, + data: data.clone(), + proof_of_payment: proof.clone(), + }, + ); + assert_same_wire( + &PaidNotify { + key, + proof_of_payment: data.clone(), + }, + &PlainNotify { + key, + proof_of_payment: data.clone(), + }, + ); + assert_same_wire( + &FetchResponse::Success { + key, + data: data.clone(), + }, + &PlainFetchResponse::Success { + key, + data: data.clone(), + }, + ); + assert_same_wire( + &SubtreeSliceItem::Present { + key, + block_index: BYTE_STRING_WIRE_BLOCK_INDEX, + bao_slice: data.clone(), + nonced_siblings: siblings.clone(), + }, + &PlainSliceItem::Present { + key, + block_index: BYTE_STRING_WIRE_BLOCK_INDEX, + bao_slice: data, + nonced_siblings: siblings.clone(), + }, + ); + } + } + + #[test] + fn fetch_response_decodes_into_an_exactly_sized_buffer() { + // A fetched chunk is held by the receiver until it is stored; a + // byte-string field decodes in one copy with no growth slack. + let data = vec![0xCD; CHUNK_SIZED_TEST_PAYLOAD_LEN]; + let msg = ReplicationMessage { + request_id: 11, + body: ReplicationMessageBody::FetchResponse(FetchResponse::Success { + key: [4; 32], + data, + }), + }; + let encoded = msg.encode().unwrap(); + assert_eq!(encoded.capacity(), encoded.len()); + let decoded = ReplicationMessage::decode(&encoded).unwrap(); + let ReplicationMessageBody::FetchResponse(FetchResponse::Success { data, .. }) = + decoded.body + else { + panic!("decoded a different body"); + }; + assert_eq!(data.len(), CHUNK_SIZED_TEST_PAYLOAD_LEN); + assert_eq!(data.capacity(), CHUNK_SIZED_TEST_PAYLOAD_LEN); + } + + #[test] + fn encode_allocates_exactly_the_serialized_size() { + // A chunk-sized offer must not carry growth slack: the encoded buffer is + // shared by every per-peer send task for as long as it is queued. + let msg = ReplicationMessage { + request_id: 7, + body: ReplicationMessageBody::FreshReplicationOffer(FreshReplicationOffer { + key: [3; 32], + data: vec![0xAB; CHUNK_SIZED_TEST_PAYLOAD_LEN], + proof_of_payment: vec![1, 2, 3], + }), + }; + let encoded = msg.encode().unwrap(); + assert_eq!(encoded.capacity(), encoded.len()); + let decoded = ReplicationMessage::decode(&encoded).unwrap(); + assert_eq!(decoded.request_id, 7); + } + #[test] fn encode_rejects_oversized_message() { // Build a message whose serialized form exceeds the limit. diff --git a/src/storage/chunk_store.rs b/src/storage/chunk_store.rs index fe297705..e6da0c4f 100644 --- a/src/storage/chunk_store.rs +++ b/src/storage/chunk_store.rs @@ -1104,6 +1104,14 @@ impl ChunkStore { std::fs::metadata(self.legacy_env_dir.join(LEGACY_DATA_FILE)).map_or(0, |m| m.len()) } + /// Test-only path of the file a chunk is (or would be) stored in, so tests + /// can damage or remove it without re-deriving the store's layout. + #[cfg(any(test, feature = "test-utils"))] + #[must_use] + pub fn test_chunk_path(&self, address: &XorName) -> PathBuf { + self.files.chunk_path(address) + } + /// Test-only handle to the file store's put gate. #[cfg(any(test, feature = "test-utils"))] #[must_use] diff --git a/src/storage/file_store.rs b/src/storage/file_store.rs index 3172bdda..b4b087d7 100644 --- a/src/storage/file_store.rs +++ b/src/storage/file_store.rs @@ -1622,7 +1622,7 @@ impl FileStore { } /// Absolute path of a chunk file. - fn chunk_path(&self, address: &XorName) -> PathBuf { + pub(crate) fn chunk_path(&self, address: &XorName) -> PathBuf { self.chunks_dir .join(shard_name(address)) .join(hex::encode(address)) @@ -2462,7 +2462,12 @@ fn open_regular(path: &Path) -> Result> { /// during an ordinary GET or an audit response. fn read_bounded(file: File, path: &Path) -> Result> { let ceiling = MAX_CHUNK_SIZE as u64; - let mut buf = Vec::new(); + // Sized from the file, so a chunk lands in an exactly-sized buffer: grown from + // empty, a full chunk ended in twice its size, held for as long as the bytes are + // served or offered. Capped at the ceiling, so a planted oversized file still + // cannot make this allocate more than a chunk up front. + let expected = file.metadata().map_or(0, |m| m.len()).min(ceiling); + let mut buf = Vec::with_capacity(usize::try_from(expected).unwrap_or(0)); let read = file.take(ceiling + 1).read_to_end(&mut buf).map_err(|e| { Error::Storage(format!("Failed to read chunk file {}: {e}", path.display())) })?; @@ -3189,6 +3194,21 @@ mod tests { assert_eq!(raw, b"tampered"); } + /// A full-size chunk reads into an exactly-sized buffer, verified or raw. + #[tokio::test] + async fn reads_allocate_exactly_the_chunk() { + let (store, _dir) = test_store().await; + let content: Vec = (0..=u8::MAX).cycle().take(MAX_CHUNK_SIZE).collect(); + let addr = crate::client::compute_address(&content); + store.put(&addr, &content).await.expect("put"); + + let verified = store.get(&addr).await.expect("get").expect("bytes"); + assert_eq!(verified, content); + assert_eq!(verified.capacity(), MAX_CHUNK_SIZE); + let raw = store.get_raw(&addr).await.expect("get_raw").expect("bytes"); + assert_eq!(raw.capacity(), MAX_CHUNK_SIZE); + } + /// A put whose caller goes away does not admit a key on bytes nothing has read. /// /// The blocking half of a put outlives the future that started it, deliberately, so diff --git a/src/storage/handler.rs b/src/storage/handler.rs index 840a5ef9..a712bd5a 100644 --- a/src/storage/handler.rs +++ b/src/storage/handler.rs @@ -430,6 +430,13 @@ impl AntProtocol { self.fresh_write_tx = Some(tx); } + /// The fresh-write sender, if one was set. Lets tests drive the + /// replication engine's fresh-write pipeline exactly as a PUT would. + #[cfg(any(test, feature = "test-utils"))] + pub fn fresh_write_sender(&self) -> Option> { + self.fresh_write_tx.clone() + } + /// Get the protocol identifier. #[must_use] pub fn protocol_id(&self) -> &'static str { @@ -871,14 +878,11 @@ impl AntProtocol { // fall back to the original proof rather than dropping the // replication entirely. let proof = Self::strip_commitment_sidecars(proof); - // `request.content` is now `bytes::Bytes`; FreshWriteEvent - // still carries the chunk as `Vec` for compatibility - // with the replication wire format, so materialise once - // here. Done only on the success path, where storage has - // already accepted the chunk. + // Storage has already accepted the chunk on this path, so + // the event carries only the key; the offer dispatcher + // reads the chunk back when it is ready to send it. let event = FreshWriteEvent { key: address, - data: request.content.to_vec(), payment_proof: proof, }; if tx.send(event).is_err() { diff --git a/src/web_rtc.rs b/src/web_rtc.rs index 6bb3e8e0..b682aba1 100644 --- a/src/web_rtc.rs +++ b/src/web_rtc.rs @@ -13,6 +13,7 @@ use crate::ant_protocol::{ ChunkQuoteResponse, MAX_CHUNK_SIZE, }; use crate::browser::{browser_payment_network, BrowserEndpoint, BrowserPaymentNetwork}; +use crate::codec::{encode_exact, ExactEncodeError}; use crate::config::WebRtcDirectConfig; use crate::error::{Error, Result}; use crate::logging::{debug, info, warn}; @@ -1762,16 +1763,12 @@ fn binary_response_limit(body: &ChunkMessageBody) -> ServerResult { fn encode_binary_response(message: &ChunkMessage, limit: usize) -> ServerResult> { // Size without allocating, then encode into exactly the charged space. - // This also prevents Vec growth from retaining an oversized capacity. - let length = postcard::experimental::serialized_size(message) - .map_err(|error| public_error("invalid_response", error))?; - if length > limit { - return Err("chunk protocol response exceeds operation limit".to_string()); - } - let mut bytes = vec![0; length]; - postcard::to_slice(message, &mut bytes) - .map_err(|error| public_error("invalid_response", error))?; - Ok(bytes) + encode_exact(message, limit).map_err(|error| match error { + ExactEncodeError::Serialize(error) => public_error("invalid_response", error), + ExactEncodeError::TooLarge { .. } => { + "chunk protocol response exceeds operation limit".to_string() + } + }) } fn hello_response(request_id: u64, state: &ServerState) -> Response { diff --git a/tests/e2e/fresh_offer_capacity.rs b/tests/e2e/fresh_offer_capacity.rs index 5774c2a0..66db3358 100644 --- a/tests/e2e/fresh_offer_capacity.rs +++ b/tests/e2e/fresh_offer_capacity.rs @@ -26,6 +26,7 @@ use super::TestHarness; use ant_node::client::compute_address; +use ant_node::replication::fresh::FreshWriteEvent; use ant_node::replication::{ fresh_offer_admission_refusals, fresh_offer_refusals_global_pool, fresh_offer_refusals_per_peer_share, paid_notify_admission_refusals, @@ -92,25 +93,15 @@ async fn normal_upload_never_reaches_fresh_offer_capacity() { // Every node must accept the offers without an on-chain proof; Anvil is not // running in this suite. for (_, address) in &chunks { - for index in 0..harness.node_count() { - if let Some(node) = harness.test_node(index) { - if let Some(protocol) = node.ant_protocol.as_ref() { - protocol.payment_verifier().cache_insert(*address); - } - } - } + harness.prepopulate_payment_cache_everywhere(address); } let source = harness.test_node(UPLOAD_SOURCE_INDEX).expect("source node"); - let source_storage = source - .ant_protocol - .as_ref() - .expect("source protocol") - .storage(); - let engine = source - .replication_engine - .as_ref() - .expect("source replication engine"); + let source_protocol = source.ant_protocol.as_ref().expect("source protocol"); + let source_storage = source_protocol.storage(); + let fresh_writes = source_protocol + .fresh_write_sender() + .expect("fresh-write sender wired by the harness"); // Counters are process-global and cumulative, and other tests share this // binary, so measure a delta rather than an absolute. @@ -119,14 +110,21 @@ async fn normal_upload_never_reaches_fresh_offer_capacity() { let global_before = fresh_offer_refusals_global_pool(); let share_before = fresh_offer_refusals_per_peer_share(); - // Drive the upload the way the fresh-write drainer does: store, then hand - // off to replication, moving to the next chunk immediately. No pacing — - // pacing here would be the test quietly avoiding the very pressure it - // exists to apply. + // Drive the upload through the fresh-write channel, exactly as the PUT + // handler does: store, then hand the write to the drainer, moving to the + // next chunk immediately. The drainer announces every write at arrival + // rate and only the offers wait for pending-offer permits, as in + // production. No pacing here — pacing would be the test quietly avoiding + // the very pressure it exists to apply. let started = Instant::now(); for (content, address) in &chunks { source_storage.put(address, content).await.expect("put"); - engine.replicate_fresh(address, content, &DUMMY_POP).await; + fresh_writes + .send(FreshWriteEvent { + key: *address, + payment_proof: DUMMY_POP.to_vec(), + }) + .expect("queue fresh write"); } let dispatch_elapsed = started.elapsed(); diff --git a/tests/e2e/harness.rs b/tests/e2e/harness.rs index eaa560f9..786b18c5 100644 --- a/tests/e2e/harness.rs +++ b/tests/e2e/harness.rs @@ -379,6 +379,19 @@ impl TestHarness { self.network.node_mut(index) } + /// Pre-populate the payment cache on every node. + /// + /// The source of a write and every receiver of its fresh offers and + /// `PaidNotify` then accept a dummy proof for `address`, as if it had + /// been paid for on chain (no Anvil runs in most suites). + pub fn prepopulate_payment_cache_everywhere(&self, address: &XorName) { + for node in self.network.nodes() { + if let Some(ref protocol) = node.ant_protocol { + protocol.payment_verifier().cache_insert(*address); + } + } + } + /// Pre-populate the payment cache on the node matching `peer_id`. /// /// Inserts `address` into the target node's payment verifier cache so diff --git a/tests/e2e/replication.rs b/tests/e2e/replication.rs index abf5f92b..876e42c9 100644 --- a/tests/e2e/replication.rs +++ b/tests/e2e/replication.rs @@ -5,14 +5,17 @@ #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] -use super::testnet::TestNetworkConfig; +use super::testnet::{TestNetworkConfig, DEFAULT_CHUNK_OPERATION_TIMEOUT_SECS}; use super::TestHarness; +use ant_node::ant_protocol::{ChunkMessage, ChunkMessageBody, ChunkPutRequest, ChunkPutResponse}; use ant_node::client::compute_address; use ant_node::replication::audit_coordinator::AuditChallengeCoordinator; use ant_node::replication::commitment_state::{BuiltCommitment, ResponderCommitmentState}; use ant_node::replication::config::{ - storage_admission_width, K_BUCKET_SIZE, REPLICATION_PROTOCOL_ID, + storage_admission_width, FRESH_READ_RETRY_DELAY, K_BUCKET_SIZE, MAX_PENDING_FRESH_OFFERS, + REPLICATION_PROTOCOL_ID, }; +use ant_node::replication::fresh::FreshWriteEvent; use ant_node::replication::protocol::{ compute_audit_digest, AuditChallenge, AuditResponse, FetchRequest, FetchResponse, FreshReplicationOffer, FreshReplicationResponse, NeighborSyncRequest, ReplicationMessage, @@ -21,11 +24,14 @@ use ant_node::replication::protocol::{ use ant_node::replication::pruning; use ant_node::replication::scheduling::ReplicationQueues; use ant_node::replication::types::{NeighborSyncState, RepairProofs}; +use ant_node::storage::XorName; use ant_node::ReplicationConfig; +use bytes::Bytes; use saorsa_core::identity::PeerId; use saorsa_core::{P2PNode, TrustEvent}; use serial_test::serial; use std::collections::HashSet; +use std::fs; use std::sync::Arc; use std::time::Duration; use tokio::sync::RwLock; @@ -48,6 +54,26 @@ const FULL_NODE_SHUN_POSSESSION_DELAY_MAX: Duration = Duration::from_millis(500) const DUMMY_PAYMENT_PROOF_LEN: usize = 64; /// Dummy proof byte used when a test only needs to reach pre-payment gates. const DUMMY_PAYMENT_PROOF_BYTE: u8 = 0x01; +/// A regular (non-bootstrap) node of the minimal harness; source of the fresh +/// replication tests. +const FRESH_PIPELINE_SOURCE_INDEX: usize = 3; +/// Writes queued at once by the saturation test: three times the pending-offer +/// budget, so the dispatcher must block on and recycle permits to drain it. +const FRESH_BURST_WRITES: usize = 3 * MAX_PENDING_FRESH_OFFERS; +/// How long the saturation test watches a full budget for offers encoded past +/// it. The dispatcher encodes a small chunk in well under a millisecond. +const FRESH_BURST_SETTLE: Duration = Duration::from_millis(500); +/// Wait budget for every burst write to be offered once sends resume. +const FRESH_BURST_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); +/// Wait budget for pending-offer permits to come back once the offers holding +/// them have finished sending. +const PERMIT_RELEASE_TIMEOUT: Duration = Duration::from_secs(10); +/// Poll interval for timing the dispatcher against the retry delay: fine +/// enough that the poll cannot hide a stall of that length. +const DISPATCH_POLL_INTERVAL: Duration = Duration::from_millis(10); +/// Extension a chunk file is moved aside under while a directory stands in +/// for it, injecting a read fault. +const FAULT_SET_ASIDE_EXTENSION: &str = "set-aside"; /// Minimal paid-list repair close group used by the deterministic repair e2e. const PAID_REPAIR_GROUP_SIZE: usize = 5; /// Storage threshold configured above majority so one holder is below quorum. @@ -200,62 +226,453 @@ async fn test_fresh_replication_propagates_to_close_group() { let harness = TestHarness::setup_minimal().await.expect("setup"); harness.warmup_dht().await.expect("warmup"); - // Pick a non-bootstrap node with replication engine - let source_idx = 3; // first regular node - let source = harness.test_node(source_idx).expect("source node"); - let source_protocol = source.ant_protocol.as_ref().expect("protocol"); - let source_storage = source_protocol.storage(); - - // Create and store a chunk + let source_idx = FRESH_PIPELINE_SOURCE_INDEX; let content = b"hello replication world"; - let address = compute_address(content); - source_storage.put(&address, content).await.expect("put"); + // Paid on every node, so receivers accept the offer without Anvil. + let address = store_paid_chunk(&harness, source_idx, content).await; + harness + .test_node(source_idx) + .and_then(|node| node.replication_engine.as_ref()) + .expect("source replication engine") + .replicate_fresh(&address, content, &dummy_payment_proof()) + .await; - // Pre-populate payment cache on ALL nodes so receivers accept the offer - // (bypasses EVM verification, which is unavailable without Anvil). - for i in 0..harness.node_count() { - if let Some(node) = harness.test_node(i) { - if let Some(protocol) = &node.ant_protocol { - protocol.payment_verifier().cache_insert(address); - } - } - } + assert!( + wait_until_replicated(&harness, source_idx, &address, PROPAGATION_TIMEOUT).await, + "Chunk should have replicated to at least one other node" + ); - // Trigger fresh replication with a dummy PoP - let dummy_pop = [0x01u8; 64]; - if let Some(ref engine) = source.replication_engine { - engine.replicate_fresh(&address, content, &dummy_pop).await; - } + harness.teardown().await.expect("teardown"); +} - // Poll until replication propagates (or timeout). - let deadline = tokio::time::Instant::now() + PROPAGATION_TIMEOUT; - let mut found_on_other = false; - while tokio::time::Instant::now() < deadline { - for i in 0..harness.node_count() { - if i == source_idx { - continue; - } - if let Some(node) = harness.test_node(i) { - if let Some(protocol) = &node.ant_protocol { - if protocol.storage().exists(&address).unwrap_or(false) { - found_on_other = true; - } - } - } - } - if found_on_other { - break; - } - tokio::time::sleep(PROPAGATION_POLL_INTERVAL).await; +/// The whole PUT path: a chunk PUT through the handler emits the fresh-write +/// event itself, and the drainer → dispatcher pipeline replicates it. A write +/// whose chunk is no longer stored, queued ahead of it, is skipped without +/// being offered, without stalling the pipeline and without keeping its +/// pending-offer permit. +#[tokio::test] +#[serial] +async fn fresh_write_pipeline_replicates_a_put_and_skips_missing_chunks() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let source = harness + .test_node(FRESH_PIPELINE_SOURCE_INDEX) + .expect("source node"); + let engine = source + .replication_engine + .as_ref() + .expect("replication engine"); + let fresh_tx = fresh_write_sender(&harness, FRESH_PIPELINE_SOURCE_INDEX); + + // A write whose chunk was never stored goes first: the dispatcher must + // skip it and carry on with the PUT queued behind it. + let missing = compute_address(b"never stored anywhere"); + fresh_tx + .send(FreshWriteEvent { + key: missing, + payment_proof: dummy_payment_proof(), + }) + .expect("queue missing write"); + let address = put_paid_chunk( + &harness, + FRESH_PIPELINE_SOURCE_INDEX, + b"chunk PUT through the handler", + ) + .await; + + assert!( + wait_until_replicated( + &harness, + FRESH_PIPELINE_SOURCE_INDEX, + &address, + PROPAGATION_TIMEOUT + ) + .await, + "the PUT should have replicated through the fresh-write pipeline" + ); + assert!( + wait_until( + || engine.pending_offer_permits_available() == MAX_PENDING_FRESH_OFFERS, + PERMIT_RELEASE_TIMEOUT + ) + .await, + "a pending-offer permit was not released" + ); + assert_eq!( + engine.fresh_offers_dispatched(), + 1, + "only the stored chunk may be offered; the missing one is skipped" + ); + + harness.teardown().await.expect("teardown"); +} + +/// Saturation: with the send stage held, a burst three times the pending-offer +/// budget encodes exactly `MAX_PENDING_FRESH_OFFERS` offers and the dispatcher +/// then waits for a permit instead of encoding more. Once sends resume, every +/// write is offered — none is lost to back-pressure — and every permit comes +/// back. +#[tokio::test] +#[serial] +async fn fresh_write_pipeline_holds_a_burst_at_the_offer_budget() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let source = harness + .test_node(FRESH_PIPELINE_SOURCE_INDEX) + .expect("source node"); + let engine = source + .replication_engine + .as_ref() + .expect("replication engine"); + let budget = u64::try_from(MAX_PENDING_FRESH_OFFERS).expect("budget fits u64"); + let burst = u64::try_from(FRESH_BURST_WRITES).expect("burst fits u64"); + + let held_sends = engine + .hold_replication_sends() + .await + .expect("replication send permits"); + for i in 0..FRESH_BURST_WRITES { + let content = format!("fresh-write burst chunk {i}"); + put_paid_chunk(&harness, FRESH_PIPELINE_SOURCE_INDEX, content.as_bytes()).await; } + assert!( - found_on_other, - "Chunk should have replicated to at least one other node" + wait_until( + || engine.fresh_offers_dispatched() >= budget, + PROPAGATION_TIMEOUT + ) + .await, + "the burst never filled the pending-offer budget ({} offers encoded)", + engine.fresh_offers_dispatched() + ); + // Give a dispatcher that ignored the budget time to overshoot it. + tokio::time::sleep(FRESH_BURST_SETTLE).await; + assert_eq!( + engine.fresh_offers_dispatched(), + budget, + "offers were encoded past the pending-offer budget" + ); + assert_eq!(engine.pending_offer_permits_available(), 0); + + drop(held_sends); + assert!( + wait_until( + || engine.fresh_offers_dispatched() == burst, + FRESH_BURST_DRAIN_TIMEOUT + ) + .await, + "only {} of {burst} burst writes were offered", + engine.fresh_offers_dispatched() + ); + assert!( + wait_until( + || engine.pending_offer_permits_available() == MAX_PENDING_FRESH_OFFERS, + PERMIT_RELEASE_TIMEOUT + ) + .await, + "the burst leaked a pending-offer permit" ); harness.teardown().await.expect("teardown"); } +/// Retry: a chunk whose read-back fails is retried after +/// `FRESH_READ_RETRY_DELAY` without holding its permit or the dispatcher. A +/// healthy write queued behind it is offered within half the delay of the +/// failed read, which a dispatcher sleeping the delay out cannot do, and once +/// the fault clears the failed write is offered too. The fault is a directory +/// standing where the chunk file should be, which the store refuses to read +/// whatever the platform or user. +#[tokio::test] +#[serial] +async fn fresh_write_pipeline_retries_a_failed_read_off_the_dispatcher() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let source = harness + .test_node(FRESH_PIPELINE_SOURCE_INDEX) + .expect("source node"); + let engine = source + .replication_engine + .as_ref() + .expect("replication engine"); + let storage = source.ant_protocol.as_ref().expect("protocol").storage(); + let fresh_tx = fresh_write_sender(&harness, FRESH_PIPELINE_SOURCE_INDEX); + let faulty = store_paid_chunk( + &harness, + FRESH_PIPELINE_SOURCE_INDEX, + b"chunk whose first read-back fails", + ) + .await; + // Stored up front, so queueing it later takes no time out of the window. + let healthy = store_paid_chunk( + &harness, + FRESH_PIPELINE_SOURCE_INDEX, + b"healthy write queued behind a failed read", + ) + .await; + + let chunk_file = storage.test_chunk_path(&faulty); + assert!(chunk_file.is_file(), "chunk file on disk"); + let set_aside = chunk_file.with_extension(FAULT_SET_ASIDE_EXTENSION); + fs::rename(&chunk_file, &set_aside).expect("move the chunk file aside"); + fs::create_dir(&chunk_file).expect("put a directory in its place"); + + fresh_tx + .send(FreshWriteEvent { + key: faulty, + payment_proof: dummy_payment_proof(), + }) + .expect("queue write"); + // A failed read marks the chunk suspect, which hides it from `exists`: + // the observable proof that the dispatcher's first attempt hit the fault. + // Polled finely, so the window below starts within a poll of the fault. + assert!( + wait_until_every( + || !storage.exists(&faulty).unwrap_or(true), + PROPAGATION_TIMEOUT, + DISPATCH_POLL_INTERVAL + ) + .await, + "the dispatcher never attempted the faulty read" + ); + let fault_seen = tokio::time::Instant::now(); + assert!( + wait_until_every( + || engine.pending_offer_permits_available() == MAX_PENDING_FRESH_OFFERS, + PERMIT_RELEASE_TIMEOUT, + DISPATCH_POLL_INTERVAL + ) + .await, + "the failed read kept its pending-offer permit while waiting to retry" + ); + + // A dispatcher sleeping out the retry delay would hold this write until + // the delay ended; one that is not offers it straight away. The window is + // half the delay, counted from the failed read. + fresh_tx + .send(FreshWriteEvent { + key: healthy, + payment_proof: dummy_payment_proof(), + }) + .expect("queue healthy write"); + let window = (FRESH_READ_RETRY_DELAY / 2).saturating_sub(fault_seen.elapsed()); + assert!( + wait_until_every( + || engine.fresh_offers_dispatched() > 0, + window, + DISPATCH_POLL_INTERVAL + ) + .await, + "a healthy write waited behind the failed read's retry delay" + ); + + fs::remove_dir(&chunk_file).expect("remove the directory"); + fs::rename(&set_aside, &chunk_file).expect("restore the chunk file"); + assert!( + wait_until( + || engine.fresh_offers_dispatched() == 2, + FRESH_READ_RETRY_DELAY + PROPAGATION_TIMEOUT + ) + .await, + "the failed write was not offered once its read succeeded" + ); + assert!( + storage.exists(&faulty).unwrap_or(false), + "the successful retry clears the suspect mark on the source" + ); + assert!( + wait_until_replicated( + &harness, + FRESH_PIPELINE_SOURCE_INDEX, + &faulty, + PROPAGATION_TIMEOUT + ) + .await, + "the retried write should have replicated" + ); + + harness.teardown().await.expect("teardown"); +} + +/// A chunk whose bytes rotted on disk after it was stored is never offered: +/// every receiver would reject it and charge the sender. The read-back +/// verifies the chunk and quarantines it on the mismatch, and nothing is +/// offered for it by the time its first retry is due. +#[tokio::test] +#[serial] +async fn fresh_write_pipeline_never_offers_a_corrupt_chunk() { + let harness = TestHarness::setup_minimal().await.expect("setup"); + harness.warmup_dht().await.expect("warmup"); + + let source = harness + .test_node(FRESH_PIPELINE_SOURCE_INDEX) + .expect("source node"); + let engine = source + .replication_engine + .as_ref() + .expect("replication engine"); + let storage = source.ant_protocol.as_ref().expect("protocol").storage(); + let fresh_tx = fresh_write_sender(&harness, FRESH_PIPELINE_SOURCE_INDEX); + let content = b"chunk that rots on disk before its offer"; + let address = store_paid_chunk(&harness, FRESH_PIPELINE_SOURCE_INDEX, content).await; + + let chunk_file = storage.test_chunk_path(&address); + assert!(chunk_file.is_file(), "chunk file on disk"); + let mut rotted = content.to_vec(); + rotted.reverse(); + fs::write(&chunk_file, &rotted).expect("corrupt the chunk file"); + + fresh_tx + .send(FreshWriteEvent { + key: address, + payment_proof: dummy_payment_proof(), + }) + .expect("queue write"); + assert!( + wait_until(|| !chunk_file.is_file(), PROPAGATION_TIMEOUT).await, + "the read-back never quarantined the corrupt chunk" + ); + // Past the first retry, which must not offer anything either. + tokio::time::sleep(FRESH_READ_RETRY_DELAY + FRESH_BURST_SETTLE).await; + assert_eq!( + engine.fresh_offers_dispatched(), + 0, + "a chunk that does not match its address was offered" + ); + assert!(!storage.exists(&address).unwrap_or(true)); + assert_eq!( + engine.pending_offer_permits_available(), + MAX_PENDING_FRESH_OFFERS + ); + + harness.teardown().await.expect("teardown"); +} + +/// Proof bytes for tests that only need to reach the pre-payment gates. +fn dummy_payment_proof() -> Vec { + vec![DUMMY_PAYMENT_PROOF_BYTE; DUMMY_PAYMENT_PROOF_LEN] +} + +/// Store a chunk on the source directly, bypassing the handler, so no +/// fresh-write event is emitted for it. +async fn store_paid_chunk(harness: &TestHarness, source_idx: usize, content: &[u8]) -> XorName { + let address = compute_address(content); + harness + .test_node(source_idx) + .expect("source node") + .ant_protocol + .as_ref() + .expect("protocol") + .storage() + .put(&address, content) + .await + .expect("put"); + harness.prepopulate_payment_cache_everywhere(&address); + address +} + +/// PUT a chunk through the source node's handler, as a client would. The +/// handler stores it and emits the fresh-write event itself, carrying the +/// dummy proof it was paid with. +async fn put_paid_chunk(harness: &TestHarness, source_idx: usize, content: &[u8]) -> XorName { + let address = compute_address(content); + harness.prepopulate_payment_cache_everywhere(&address); + let request = ChunkMessage { + request_id: rand::random(), + body: ChunkMessageBody::PutRequest(ChunkPutRequest::with_payment( + address, + Bytes::copy_from_slice(content), + dummy_payment_proof(), + )), + }; + let protocol = harness + .test_node(source_idx) + .expect("source node") + .ant_protocol + .as_ref() + .expect("protocol"); + let response = tokio::time::timeout( + Duration::from_secs(DEFAULT_CHUNK_OPERATION_TIMEOUT_SECS), + protocol.try_handle_request(&request.encode().expect("encode PUT")), + ) + .await + .expect("PUT through the handler timed out") + .expect("handle PUT") + .expect("PUT response"); + match ChunkMessage::decode(&response) + .expect("decode PUT response") + .body + { + ChunkMessageBody::PutResponse(ChunkPutResponse::Success { .. }) => address, + other => panic!("PUT through the handler failed: {other:?}"), + } +} + +/// The sender the source's PUT handler feeds its fresh-write pipeline from, +/// for tests that queue a write the handler would not emit. +fn fresh_write_sender( + harness: &TestHarness, + source_idx: usize, +) -> tokio::sync::mpsc::UnboundedSender { + harness + .test_node(source_idx) + .expect("source node") + .ant_protocol + .as_ref() + .expect("protocol") + .fresh_write_sender() + .expect("fresh-write sender wired by the harness") +} + +/// Whether any node other than `source_idx` currently stores `address`. +fn stored_on_another_node(harness: &TestHarness, source_idx: usize, address: &XorName) -> bool { + (0..harness.node_count()) + .filter(|&i| i != source_idx) + .filter_map(|i| harness.test_node(i)) + .filter_map(|node| node.ant_protocol.as_ref()) + .any(|protocol| protocol.storage().exists(address).unwrap_or(false)) +} + +/// Poll until `address` is stored on a node other than `source_idx`, or +/// `budget` runs out. +async fn wait_until_replicated( + harness: &TestHarness, + source_idx: usize, + address: &XorName, + budget: Duration, +) -> bool { + wait_until( + || stored_on_another_node(harness, source_idx, address), + budget, + ) + .await +} + +/// Poll `condition` every `PROPAGATION_POLL_INTERVAL` until it holds or +/// `budget` runs out. +async fn wait_until(condition: impl FnMut() -> bool, budget: Duration) -> bool { + wait_until_every(condition, budget, PROPAGATION_POLL_INTERVAL).await +} + +/// Poll `condition` every `interval` until it holds or `budget` runs out. +async fn wait_until_every( + mut condition: impl FnMut() -> bool, + budget: Duration, + interval: Duration, +) -> bool { + let deadline = tokio::time::Instant::now() + budget; + while tokio::time::Instant::now() < deadline { + if condition() { + return true; + } + tokio::time::sleep(interval).await; + } + false +} + /// ADR-0003: the delayed possession check penalises a responsible peer that /// does NOT hold the chunk, and leaves a peer that DOES hold it unpenalised. /// @@ -434,7 +851,7 @@ async fn possession_scheduler_penalises_absent_close_peer_after_delay() { // Trigger fresh replication; the engine enqueues the possession check, which // fires ~200-500 ms later and penalises the absent close peers. - let dummy_pop = [0x01u8; 64]; + let dummy_pop = dummy_payment_proof(); engine_a .replicate_fresh(&address, content, &dummy_pop) .await; @@ -526,20 +943,13 @@ async fn full_close_group_node_rejects_replica_and_is_penalised_as_absent() { let (content, address) = candidate.expect("find key where full node is a responsible close-group peer"); - for idx in 0..harness.node_count() { - if let Some(protocol) = harness - .test_node(idx) - .and_then(|node| node.ant_protocol.as_ref()) - { - protocol.payment_verifier().cache_insert(address); - } - } + harness.prepopulate_payment_cache_everywhere(&address); - let dummy_payment_proof = vec![DUMMY_PAYMENT_PROOF_BYTE; DUMMY_PAYMENT_PROOF_LEN]; + let dummy_proof = dummy_payment_proof(); let offer = FreshReplicationOffer { key: address, data: content.clone(), - proof_of_payment: dummy_payment_proof.clone(), + proof_of_payment: dummy_proof.clone(), }; let response = send_replication_request( checker_p2p, @@ -595,7 +1005,7 @@ async fn full_close_group_node_rejects_replica_and_is_penalised_as_absent() { let trust_before = checker_p2p.peer_trust(&full_peer); checker_engine - .replicate_fresh(&address, &content, &dummy_payment_proof) + .replicate_fresh(&address, &content, &dummy_proof) .await; let deadline = tokio::time::Instant::now() + PROPAGATION_TIMEOUT; @@ -1807,7 +2217,7 @@ async fn test_fresh_offer_with_mismatched_content_address_rejected() { let offer = FreshReplicationOffer { key: wrong_address, data: content.to_vec(), - proof_of_payment: vec![0x01; 64], + proof_of_payment: dummy_payment_proof(), }; let msg = ReplicationMessage { request_id: 1001, @@ -2033,16 +2443,10 @@ async fn scenario_1_and_24_fresh_replication_stores_and_propagates_paid_list() { // Pre-populate payment cache on ALL nodes so receivers accept the offer // (bypasses EVM verification, which is unavailable without Anvil). - for i in 0..harness.node_count() { - if let Some(node) = harness.test_node(i) { - if let Some(p) = &node.ant_protocol { - p.payment_verifier().cache_insert(address); - } - } - } + harness.prepopulate_payment_cache_everywhere(&address); // Trigger fresh replication (sends FreshReplicationOffer + PaidNotify) - let dummy_pop = [0x01u8; 64]; + let dummy_pop = dummy_payment_proof(); if let Some(ref engine) = source.replication_engine { engine.replicate_fresh(&address, content, &dummy_pop).await; } @@ -2595,16 +2999,10 @@ async fn scenario_24_fresh_replication_propagates_paid_notify() { // Pre-populate payment cache on ALL nodes so receivers accept the offer // and PaidNotify (bypasses EVM verification, unavailable without Anvil). - for i in 0..harness.node_count() { - if let Some(node) = harness.test_node(i) { - if let Some(p) = &node.ant_protocol { - p.payment_verifier().cache_insert(address); - } - } - } + harness.prepopulate_payment_cache_everywhere(&address); // Trigger fresh replication (includes PaidNotify to PaidCloseGroup) - let dummy_pop = [0x01u8; 64]; + let dummy_pop = dummy_payment_proof(); if let Some(ref engine) = source.replication_engine { engine.replicate_fresh(&address, content, &dummy_pop).await; } @@ -2783,14 +3181,7 @@ async fn scenario_26_paid_list_majority_repairs_missing_replica_below_storage_qu .await .expect("put source record"); - for idx in 0..harness.node_count() { - if let Some(protocol) = harness - .test_node(idx) - .and_then(|node| node.ant_protocol.as_ref()) - { - protocol.payment_verifier().cache_insert(address); - } - } + harness.prepopulate_payment_cache_everywhere(&address); for idx in 0..PAID_REPAIR_CONFIRMING_NODES { let engine = harness @@ -3029,17 +3420,11 @@ async fn test_late_joiner_replicates_responsible_chunks() { } for (address, _) in &chunks { - for i in 0..harness.node_count() { - if let Some(node) = harness.test_node(i) { - if let Some(protocol) = &node.ant_protocol { - protocol.payment_verifier().cache_insert(*address); - } - } - } + harness.prepopulate_payment_cache_everywhere(address); } // Trigger fresh replication for each chunk so they spread to close groups. - let dummy_pop = [0x01u8; 64]; + let dummy_pop = dummy_payment_proof(); { let source = harness.test_node(source_idx).expect("source node"); if let Some(ref engine) = source.replication_engine { diff --git a/tests/e2e/testnet.rs b/tests/e2e/testnet.rs index 6e878709..b97a6c0c 100644 --- a/tests/e2e/testnet.rs +++ b/tests/e2e/testnet.rs @@ -23,6 +23,7 @@ use ant_node::payment::{ QuotingMetricsTracker, }; use ant_node::replication::config::MAX_REPLICATION_MESSAGE_SIZE; +use ant_node::replication::fresh::FreshWriteEvent; use ant_node::storage::{AntProtocol, ChunkStore, ChunkStoreConfig}; use ant_node::{ReplicationConfig, ReplicationEngine}; use bytes::Bytes; @@ -101,7 +102,7 @@ const SMALL_STABILIZATION_TIMEOUT_SECS: u64 = 60; /// conservative; the happy path completes in well under a second on /// loopback, so the larger budget only shows up on flakes. Test-only — /// no production code path reads this constant. -const DEFAULT_CHUNK_OPERATION_TIMEOUT_SECS: u64 = 90; +pub const DEFAULT_CHUNK_OPERATION_TIMEOUT_SECS: u64 = 90; /// Short node-level network timeout for E2E test harness. /// @@ -439,6 +440,11 @@ pub struct TestNode { /// Shutdown token for the replication engine. pub replication_shutdown: Option, + + /// Fresh-write events from this node's PUT handler, waiting for the + /// replication engine that `start_node` creates to take them. The sender + /// half lives in `ant_protocol`, as it does in a real node. + pub fresh_write_rx: Option>, } impl TestNode { @@ -1082,13 +1088,17 @@ impl TestNetwork { .get(&index) .copied() .unwrap_or_default(); - let ant_protocol = Self::create_ant_protocol_with_disk_reserve( + let mut ant_protocol = Self::create_ant_protocol_with_disk_reserve( &data_dir, self.config.evm_network.clone(), storage_disk_reserve, &identity, ) .await?; + // Wired as a real node wires it, so a PUT through the handler feeds + // the replication engine's fresh-write pipeline. + let (fresh_write_tx, fresh_write_rx) = tokio::sync::mpsc::unbounded_channel(); + ant_protocol.set_fresh_write_sender(fresh_write_tx); Ok(TestNode { index, @@ -1104,6 +1114,7 @@ impl TestNetwork { protocol_task: None, replication_engine: None, replication_shutdown: None, + fresh_write_rx: Some(fresh_write_rx), }) } @@ -1348,15 +1359,21 @@ impl TestNetwork { } // Start replication engine for this node. A node without an identity - // skips ONLY the engine (no early return — the node must still be - // tracked in `self.nodes` below, or its already-started P2P/protocol - // tasks would keep running untracked by the harness). - if let (Some(ref p2p), Some(ref protocol), Some(ref id)) = - (&node.p2p_node, &node.ant_protocol, &node.node_identity) - { + // or a fresh-write channel skips ONLY the engine (no early return — + // the node must still be tracked in `self.nodes` below, or its + // already-started P2P/protocol tasks would keep running untracked by + // the harness). `create_node` fills the channel and each node is + // started once, so it is missing only if that invariant breaks. + let fresh_write_rx = node.fresh_write_rx.take(); + let has_fresh_writes = fresh_write_rx.is_some(); + if let (Some(ref p2p), Some(ref protocol), Some(ref id), Some(fresh_rx)) = ( + &node.p2p_node, + &node.ant_protocol, + &node.node_identity, + fresh_write_rx, + ) { let shutdown = CancellationToken::new(); let repl_config = self.config.replication_config.clone().unwrap_or_default(); - let (_fresh_tx, fresh_rx) = tokio::sync::mpsc::unbounded_channel(); let node_identity = Arc::clone(id); match ReplicationEngine::new( repl_config, @@ -1395,6 +1412,12 @@ impl TestNetwork { "Node {} has no identity; skipping replication engine", node.index ); + } else if !has_fresh_writes { + warn!( + "Node {} has no fresh-write channel (started twice?); skipping \ + replication engine", + node.index + ); } debug!("Node {} started successfully", node.index);