From c04a18782a5c9a4a3120675ca0d8959e94ef28de Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 18:55:02 +0900 Subject: [PATCH] fix(pointer): harden pointer replication against thin views and floods - Repair counts its quorum over the whole configured close group, so a node that sees few members of it cannot adopt an owner-signed state nobody paid for on one peer's word, and it does not repair while it is still bootstrapping. - Fetch and state requests are admitted fairly before a task exists, at most 128 outstanding and 16 per peer, as chunk fetches are. A state request larger than an honest one is dropped before admission. - This node keeps at most 8 pointer requests outstanding at any one peer, across repair, possession checks and pruning, so it never exceeds the allowance a peer gives it and is never dropped as a flood. - A peer is remembered as speaking pointers only while it is in the routing table, and that set is capped. - A fetched record's signature is verified off the async executor, as ADR-0016 says every pointer signature check is. - A commit for a record this node lost takes the lost state or a newer one, never an older one, so a write verified before the loss cannot roll the node back. --- src/pointer/store.rs | 45 +++++++- src/replication/mod.rs | 28 +++-- src/replication/pointer.rs | 226 ++++++++++++++++++++++++++++++++++--- 3 files changed, 270 insertions(+), 29 deletions(-) diff --git a/src/pointer/store.rs b/src/pointer/store.rs index 3628d173..0249268d 100644 --- a/src/pointer/store.rs +++ b/src/pointer/store.rs @@ -845,9 +845,21 @@ impl Inner { // have committed while it ran. let replacing = index.get(&address).is_some_and(|entry| entry.on_disk); let outcome = match index.get(&address) { - // Nothing held, or a record this node lost: either way the - // write must happen, whatever state it carries. - None | Some(IndexEntry { on_disk: false, .. }) => PutOutcome::Changed, + // Nothing held: any state lands. + None => PutOutcome::Changed, + // A record this node lost: the lost state itself lands, to + // restore it, and so does anything that replaces it, but + // nothing older, so an arrival verified before the loss cannot + // roll the node back past it (the rule `admits` states). + Some(entry @ IndexEntry { on_disk: false, .. }) => { + if entry.state.state_id == record.state_id() + || record.state().replaces(&entry.state) + { + PutOutcome::Changed + } else { + PutOutcome::Stale + } + } Some(entry) if entry.state.state_id == record.state_id() => PutOutcome::Unchanged, Some(entry) if record.state().replaces(&entry.state) => PutOutcome::Changed, Some(_) => PutOutcome::Stale, @@ -1754,6 +1766,33 @@ mod tests { assert!(store.get(&record.address()).await.expect("get").is_some()); } + /// The commit itself holds the line `admits` states for a lost record, so + /// an older state that reaches it without that gate, or was admitted + /// before the loss, cannot roll the node back. + #[tokio::test] + async fn a_commit_after_a_loss_takes_nothing_older_than_what_was_lost() { + let (store, _dir) = store().await; + let held = signed(1, 3, 1); + store.put_bytes(&held.to_bytes()).await.expect("put"); + std::fs::remove_file(store.file_for(&held.address())).expect("remove"); + assert!(store.get(&held.address()).await.expect("get").is_none()); + + assert_eq!( + store + .put_bytes(&signed(1, 2, 2).to_bytes()) + .await + .expect("put"), + PutOutcome::Stale, + "an older state must not replace the one that was lost" + ); + assert!(!store.contains(&held.address()), "nothing was written"); + assert_eq!( + store.put_bytes(&held.to_bytes()).await.expect("put"), + PutOutcome::Changed, + "the lost state itself is restored" + ); + } + #[tokio::test] async fn a_lost_record_is_restored_and_cannot_be_rolled_back() { // The node remembers what it lost. Without that the address would look diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 73d91f45..cf731c34 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -5549,23 +5549,27 @@ async fn handle_replication_message( } ReplicationMessageBody::PointerFetchRequest(request) => { if let Some(pointers) = &ctx.pointers { - pointers.serve_fetch_detached( - *source, - request, - msg.request_id, - rr_message_id.map(ToOwned::to_owned), - ); + pointers + .serve_fetch_detached( + *source, + request, + msg.request_id, + rr_message_id.map(ToOwned::to_owned), + ) + .await; } Ok(()) } ReplicationMessageBody::PointerStateRequest(request) => { if let Some(pointers) = &ctx.pointers { - pointers.serve_state_detached( - *source, - request, - msg.request_id, - rr_message_id.map(ToOwned::to_owned), - ); + pointers + .serve_state_detached( + *source, + request, + msg.request_id, + rr_message_id.map(ToOwned::to_owned), + ) + .await; } Ok(()) } diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index df62a0f4..f5fc0b32 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -44,7 +44,7 @@ use rand::Rng; use saorsa_core::identity::PeerId; use saorsa_core::{P2PNode, TrustEvent}; use tokio::sync::{mpsc, RwLock, Semaphore}; -use tokio::task::JoinHandle; +use tokio::task::{spawn_blocking, JoinHandle}; use tokio_util::sync::CancellationToken; use tokio_util::task::TaskTracker; @@ -79,6 +79,38 @@ const MAX_CONCURRENT_OFFERS: usize = 16; /// Fetch and state requests served at once. const MAX_CONCURRENT_SERVES: usize = 32; +/// Fetch and state requests admitted at once, served or waiting to be. Past +/// it a request is dropped at admission rather than queued, as a chunk fetch +/// is, so a flood costs a bounded number of tasks. +const MAX_SERVES_OUTSTANDING: usize = MAX_CONCURRENT_SERVES * 4; + +/// Of those, how many one peer may have: twice the [`MAX_REQUESTS_PER_PEER`] +/// an honest node lets itself have outstanding at any one peer, so its +/// requests are never dropped, even those arriving before the reply to the +/// last has released its place here, while one flooding peer cannot take the +/// room every other peer's requests need. +const MAX_SERVES_OUTSTANDING_PER_PEER: u32 = 16; + +/// Pointer requests this node has outstanding at any one peer, across repair, +/// possession checks and pruning together. As many as a repair asks at once, +/// so repair is not slowed; the rest wait their turn rather than go out and +/// be dropped by the peer's per-peer allowance, which would make an honest +/// peer look as though it had not answered. +const MAX_REQUESTS_PER_PEER: usize = VERIFICATION_CONCURRENCY; + +// An honest node's requests to one peer, repair included, never exceed what +// that peer admits from it, with room for a reply still releasing its place. +const _: () = assert!( + MAX_SERVES_OUTSTANDING_PER_PEER as usize >= 2 * MAX_REQUESTS_PER_PEER, + "a peer's serve allowance must cover what an honest node asks of it" +); + +/// Most peers remembered as speaking pointers. Only routed peers are, and a +/// routing table removal forgets one, so this only bounds what removals missed +/// (a lagging event stream) could leave behind. Capability forgotten this way +/// is learned again from the peer's next hint push. +const MAX_CAPABLE_PEERS: usize = 4096; + /// How often the pending hints are looked at. const VERIFICATION_TICK: Duration = Duration::from_millis(500); @@ -202,6 +234,25 @@ pub(crate) fn evaluate( } } +/// Whether nothing but the map holds a peer's outbound permits. +/// +/// Every holder, a request waiting for a permit or one holding it, has a +/// clone, and clones are only made under the map's lock. So an entry nothing +/// else holds can be dropped without a later request getting a second set of +/// permits beside one still in use, which would let this node exceed +/// [`MAX_REQUESTS_PER_PEER`] at that peer. +fn outbound_idle(permits: &Arc) -> bool { + Arc::strong_count(permits) == 1 +} + +/// How many members a repair counts its quorum over: the whole close group, +/// however few of them this node can see (ADR-0016). A member it cannot see is +/// unanswered, so a thin routing table leaves a repair undecided instead of +/// shrinking the quorum to what little it can see. +fn repair_width(seen: usize, close_group_size: usize) -> usize { + seen.max(close_group_size) +} + /// Pointer replication for one node. See the module documentation. pub struct PointerReplication { store: PointerStore, @@ -214,6 +265,13 @@ pub struct PointerReplication { send_semaphore: Arc, offer_permits: Arc, serve_permits: Arc, + /// Admission ahead of `serve_permits`: see [`MAX_SERVES_OUTSTANDING`]. + serve_admission: Arc, + /// Requests each peer has admitted and not yet had answered. + serve_inflight: Arc>>, + /// This node's own requests outstanding at each peer (see + /// [`MAX_REQUESTS_PER_PEER`]). + outbound: Mutex>>, /// Peers that have sent a pointer message: the only ones asked anything. capable: Mutex>, /// Hinted states awaiting verification, by address. @@ -248,6 +306,9 @@ impl PointerReplication { send_semaphore, offer_permits: Arc::new(Semaphore::new(MAX_CONCURRENT_OFFERS)), serve_permits: Arc::new(Semaphore::new(MAX_CONCURRENT_SERVES)), + serve_admission: Arc::new(Semaphore::new(MAX_SERVES_OUTSTANDING)), + serve_inflight: Arc::new(RwLock::new(HashMap::new())), + outbound: Mutex::new(HashMap::new()), capable: Mutex::new(HashSet::new()), pending: Mutex::new(HashMap::new()), out_of_range: Mutex::new(HashMap::new()), @@ -265,8 +326,22 @@ impl PointerReplication { // Capability // ----------------------------------------------------------------------- - fn mark_capable(&self, peer: &PeerId) { - self.capable.lock().insert(*peer); + /// Remember that `peer` speaks pointers, if it is in the routing table. + /// + /// Only a routed peer is remembered: it is the only kind this node ever + /// asks, and the routing table's removals are what forget it again, so the + /// set cannot outgrow the table however many identities send a message. + async fn mark_capable(&self, peer: &PeerId) { + if !self.p2p.dht_manager().is_in_routing_table(peer).await { + return; + } + let mut capable = self.capable.lock(); + if capable.len() >= MAX_CAPABLE_PEERS && !capable.contains(peer) { + if let Some(evicted) = capable.iter().next().copied() { + capable.remove(&evicted); + } + } + capable.insert(*peer); } /// Whether `peer` has sent a pointer message, and may therefore be asked. @@ -277,6 +352,24 @@ impl PointerReplication { /// Forget a peer that left the routing table. pub(crate) fn forget_peer(&self, peer: &PeerId) { self.capable.lock().remove(peer); + let mut outbound = self.outbound.lock(); + if outbound.get(peer).is_some_and(outbound_idle) { + outbound.remove(peer); + } + } + + /// The permits bounding this node's requests to `peer`. + fn outbound_permits(&self, peer: &PeerId) -> Arc { + let mut outbound = self.outbound.lock(); + if outbound.len() >= MAX_CAPABLE_PEERS && !outbound.contains_key(peer) { + // Peers nothing is outstanding at hold no state worth keeping. + outbound.retain(|_, permits| !outbound_idle(permits)); + } + Arc::clone( + outbound + .entry(*peer) + .or_insert_with(|| Arc::new(Semaphore::new(MAX_REQUESTS_PER_PEER))), + ) } /// Addresses currently awaiting verification. Tests only. @@ -311,6 +404,7 @@ impl PointerReplication { body: ReplicationMessageBody, timeout: Duration, ) -> Option { + let _permit = self.outbound_permits(peer).acquire_owned().await.ok()?; let msg = ReplicationMessage { request_id: rand::thread_rng().gen::(), body, @@ -355,7 +449,13 @@ impl PointerReplication { return None; } let bytes = response.record?; - match Pointer::from_bytes(&bytes) { + // Parsing verifies the ML-DSA signature, milliseconds of CPU a peer + // can demand once per record it serves: off the async executor, as + // every other pointer signature check is (ADR-0016). + let parsed = spawn_blocking(move || Pointer::from_bytes(&bytes)) + .await + .ok()?; + match parsed { Ok(record) if record.address() == *address => Some(record), Ok(_) | Err(_) => { debug!( @@ -433,7 +533,7 @@ impl PointerReplication { /// Take in hints from `source`: queue each hinted state this node should /// hold and lacks, or holds an older state than. pub(crate) async fn handle_hints(&self, source: PeerId, hints: Vec) { - self.mark_capable(&source); + self.mark_capable(&source).await; if hints.is_empty() { return; } @@ -506,12 +606,22 @@ impl PointerReplication { }) } - /// Verify and fetch whatever is due now. Also the tests' way to drive it. - pub async fn verify_due(&self) { + /// Whether this node should look for anything to repair now. + async fn may_repair(&self) -> bool { // As for chunks (ADR-0011): a node that cannot store what it would // fetch does not spend the network's time finding it. The hints wait, // and come back each round if they are dropped meanwhile. if self.chunks.capacity_verdict() == CapacityVerdict::Full { + return false; + } + // A node still bootstrapping sees too little of its close group to + // judge what a quorum of it holds. The hints wait for it. + !*self.is_bootstrapping.read().await + } + + /// Verify and fetch whatever is due now. Also the tests' way to drive it. + pub async fn verify_due(&self) { + if !self.may_repair().await { return; } let now = Instant::now(); @@ -561,11 +671,17 @@ impl PointerReplication { }) .collect(); let held = self.store.state(address); + // The quorum is counted over the whole close group, as + // ADR-0016 has it: a member this node cannot see, because its + // routing table is thin, counts as unanswered, never as a + // smaller group. Otherwise one peer in a thin view could vote + // alone for an owner-signed state nobody paid to store. + let width = repair_width(group.len(), self.config.close_group_size); let verdict = evaluate( held.as_ref(), - group.len(), + width, &answered, - self.config.quorum_needed(group.len()), + self.config.quorum_needed(width), ); (*address, verdict) }) @@ -716,20 +832,48 @@ impl PointerReplication { // Serving // ----------------------------------------------------------------------- + /// Admit a fetch or state request from `source`, or drop it. + /// + /// Admitted before a task exists, fairly across peers, as a chunk fetch + /// is: each admitted request would otherwise be a task waiting on a serve + /// permit with nothing bounding how many, and one peer could take every + /// place. + async fn admit_serve(&self, source: &PeerId, kind: &str) -> Option { + match super::admit_bounded_responder( + &self.serve_admission, + &self.serve_inflight, + source, + MAX_SERVES_OUTSTANDING, + MAX_SERVES_OUTSTANDING_PER_PEER, + ) + .await + { + Ok(guard) => Some(guard), + Err(failure) => { + debug!("Dropping a pointer {kind} request from {source}: {failure}"); + None + } + } + } + /// Answer a fetch request in the background. - pub(crate) fn serve_fetch_detached( + pub(crate) async fn serve_fetch_detached( self: &Arc, source: PeerId, request: PointerFetchRequest, request_id: u64, rr_message_id: Option, ) { - self.mark_capable(&source); + let Some(guard) = self.admit_serve(&source, "fetch").await else { + return; + }; let this = Arc::clone(self); self.tracker.spawn(async move { + let _guard = guard; let Ok(_permit) = this.serve_permits.acquire().await else { return; }; + this.mark_capable(&source).await; // `get` verifies the signature before serving, so a record damaged // on this disk is never handed on. let record = match this.store.get(&request.address).await { @@ -751,19 +895,33 @@ impl PointerReplication { } /// Answer a state request in the background, from the index. - pub(crate) fn serve_state_detached( + pub(crate) async fn serve_state_detached( self: &Arc, source: PeerId, request: PointerStateRequest, request_id: u64, rr_message_id: Option, ) { - self.mark_capable(&source); + // An honest peer never asks about more than this at once. The decoded + // request would otherwise be held by its task for as long as it waits, + // and the wire allows one far larger. + if request.addresses.len() > MAX_POINTER_STATE_REQUEST_ADDRESSES { + debug!( + "Dropping a pointer state request from {source}: {} addresses", + request.addresses.len() + ); + return; + } + let Some(guard) = self.admit_serve(&source, "state").await else { + return; + }; let this = Arc::clone(self); self.tracker.spawn(async move { + let _guard = guard; let Ok(_permit) = this.serve_permits.acquire().await else { return; }; + this.mark_capable(&source).await; let states = request .addresses .iter() @@ -875,7 +1033,6 @@ impl PointerReplication { source: PeerId, offer: PointerFreshOffer, ) { - self.mark_capable(&source); let Ok(permit) = Arc::clone(&self.offer_permits).try_acquire_owned() else { info!("Dropping a fresh pointer offer from {source}: too many in flight"); return; @@ -883,6 +1040,7 @@ impl PointerReplication { let this = Arc::clone(self); self.tracker.spawn(async move { let _permit = permit; + this.mark_capable(&source).await; this.accept_offer(source, offer).await; }); } @@ -1183,6 +1341,46 @@ mod tests { list.iter().map(|(p, s)| (peer(*p), *s)).collect() } + /// A peer's outbound permits are dropped only once nothing holds them, so + /// a request in flight keeps the one set later requests share. + #[tokio::test] + async fn outbound_permits_in_use_are_never_dropped() { + let permits = Arc::new(Semaphore::new(MAX_REQUESTS_PER_PEER)); + assert!(outbound_idle(&permits)); + let held = Arc::clone(&permits) + .acquire_owned() + .await + .expect("a permit"); + assert!(!outbound_idle(&permits), "a held permit keeps the set"); + drop(held); + assert!(outbound_idle(&permits)); + } + + /// A node that sees one member of the close group, as one still filling + /// its routing table does, cannot adopt a state on that member's word: the + /// members it cannot see count as unanswered. + #[test] + fn a_thin_view_of_the_group_cannot_adopt_on_one_vote() { + let config = ReplicationConfig::default(); + let unpaid = state(9, 1, 90); + let width = repair_width(1, config.close_group_size); + let got = evaluate( + None, + width, + &answers(&[(1, Some(unpaid))]), + config.quorum_needed(width), + ); + assert!( + matches!(got, Verdict::Undecided), + "one visible peer must not decide, got {got:?}" + ); + assert_eq!( + repair_width(9, config.close_group_size), + 9, + "a view wider than the configured group is counted as it is" + ); + } + #[test] fn a_state_a_quorum_holds_is_adopted_from_its_holders() { let newer = state(3, 1, 30);