From 2c9b79e34eb470f6427633540844f43df65b8b72 Mon Sep 17 00:00:00 2001 From: grumbach Date: Tue, 29 Sep 2026 18:55:02 +0900 Subject: [PATCH 1/7] 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); From 971673adb685bb24273fe582c28f2a7ea26aaf70 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 11:21:45 +0900 Subject: [PATCH 2/7] fix(pointer): charge one invalid record once in a possession check A possession check fetched the record from a peer, and a record that did not verify was penalised by the fetch and then again by the check as a record not held. The fetch now says whether it penalised an invalid record, and the check does not add a second penalty for it. --- src/replication/pointer.rs | 57 ++++++++++++++++++++++++++++---------- 1 file changed, 42 insertions(+), 15 deletions(-) diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index f5fc0b32..3470727a 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -234,6 +234,16 @@ pub(crate) fn evaluate( } } +/// What asking a peer for a pointer record came to. +enum Fetched { + /// A record that verifies and belongs at the address. + Record(Pointer), + /// A record that does not, which the peer has already been penalised for. + Invalid, + /// No record: no answer, an empty one, or one about another address. + Nothing, +} + /// 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 @@ -433,7 +443,16 @@ impl PointerReplication { /// back: a valid signature and the right address. Anything else counts /// against the peer. async fn fetch_record(&self, peer: &PeerId, address: &XorName) -> Option { - let body = self + match self.fetch_outcome(peer, address).await { + Fetched::Record(record) => Some(record), + Fetched::Invalid | Fetched::Nothing => None, + } + } + + /// [`Self::fetch_record`], saying whether a failure was a record that did + /// not verify, which has already cost the peer, or no record at all. + async fn fetch_outcome(&self, peer: &PeerId, address: &XorName) -> Fetched { + let Some(body) = self .request( peer, ReplicationMessageBody::PointerFetchRequest(PointerFetchRequest { @@ -441,29 +460,34 @@ impl PointerReplication { }), self.config.fetch_request_timeout, ) - .await?; + .await + else { + return Fetched::Nothing; + }; let ReplicationMessageBody::PointerFetchResponse(response) = body else { - return None; + return Fetched::Nothing; }; if response.address != *address { - return None; + return Fetched::Nothing; } - let bytes = response.record?; + let Some(bytes) = response.record else { + return Fetched::Nothing; + }; // 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()?; + let Ok(parsed) = spawn_blocking(move || Pointer::from_bytes(&bytes)).await else { + return Fetched::Nothing; + }; match parsed { - Ok(record) if record.address() == *address => Some(record), + Ok(record) if record.address() == *address => Fetched::Record(record), Ok(_) | Err(_) => { debug!( "Peer {peer} served an invalid pointer record for {}", hex::encode(address) ); self.penalise(peer).await; - None + Fetched::Invalid } } } @@ -1170,13 +1194,16 @@ impl PointerReplication { if *peer == self_id || !group.contains(peer) || !self.is_capable(peer) { continue; } - let holds = self - .fetch_record(peer, &fresh.address) - .await - .is_some_and(|record| { + let holds = match self.fetch_outcome(peer, &fresh.address).await { + Fetched::Record(record) => { let state = record.state(); state.state_id == fresh.state_id || state.replaces(&fresh) - }); + } + // Serving a record that does not verify has already been + // charged to the peer, once; it is not charged again here. + Fetched::Invalid => continue, + Fetched::Nothing => false, + }; if !holds { warn!( "Peer {peer} does not hold pointer {} it was offered", From 152452314267c7d2097f366ccce04be764f43ba3 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 11:52:25 +0900 Subject: [PATCH 3/7] fix(pointer): judge a possession check on what a peer owes when it answers A possession check looked up the close group once and then asked each peer in turn. Asking waits on that peer's other requests and on the peers before it, so a peer could leave the group before it was asked and still be penalised for not holding the record. Membership and capability are now checked again before a penalty. A record this node could not check, because the verification task failed here, was treated as no record and charged to the peer. It is now its own outcome and charges nobody. The judgement of one answer is a pure function with a unit test, which fails if an invalid record, already charged by the fetch, is charged again as missing. --- src/replication/pointer.rs | 119 ++++++++++++++++++++++++++++--------- 1 file changed, 92 insertions(+), 27 deletions(-) diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 3470727a..6d6799b8 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -242,6 +242,37 @@ enum Fetched { Invalid, /// No record: no answer, an empty one, or one about another address. Nothing, + /// A record this node could not check, for a reason of its own. + LocalFailure, +} + +/// What a possession check makes of one peer's answer about `fresh`. +#[derive(Debug, PartialEq, Eq)] +enum Possession { + /// The peer served `fresh` or a state that replaces it. + Holds, + /// The peer could not produce `fresh` or anything newer. + Missing, + /// Nothing to judge the peer on: it served a record that does not verify, + /// which the fetch has already charged it for, or this node could not + /// check what it served. + NotJudged, +} + +/// Judge one peer's answer in a possession check for `fresh`. +fn judge_possession(fetched: Fetched, fresh: &PointerState) -> Possession { + match fetched { + Fetched::Record(record) => { + let state = record.state(); + if state.state_id == fresh.state_id || state.replaces(fresh) { + Possession::Holds + } else { + Possession::Missing + } + } + Fetched::Nothing => Possession::Missing, + Fetched::Invalid | Fetched::LocalFailure => Possession::NotJudged, + } } /// Whether nothing but the map holds a peer's outbound permits. @@ -445,7 +476,7 @@ impl PointerReplication { async fn fetch_record(&self, peer: &PeerId, address: &XorName) -> Option { match self.fetch_outcome(peer, address).await { Fetched::Record(record) => Some(record), - Fetched::Invalid | Fetched::Nothing => None, + Fetched::Invalid | Fetched::Nothing | Fetched::LocalFailure => None, } } @@ -477,7 +508,12 @@ impl PointerReplication { // can demand once per record it serves: off the async executor, as // every other pointer signature check is (ADR-0016). let Ok(parsed) = spawn_blocking(move || Pointer::from_bytes(&bytes)).await else { - return Fetched::Nothing; + // The check itself failed here, which says nothing about the peer. + warn!( + "Could not verify the pointer record {peer} served for {}", + hex::encode(address) + ); + return Fetched::LocalFailure; }; match parsed { Ok(record) if record.address() == *address => Fetched::Record(record), @@ -1180,40 +1216,43 @@ impl PointerReplication { /// that cannot produce `fresh` or a state that replaces it. pub async fn check_possession(&self, fresh: PointerState, peers: &[PeerId]) { let self_id = *self.p2p.peer_id(); - let group: HashSet = self - .p2p - .dht_manager() - .find_closest_nodes_local_with_self(&fresh.address, self.config.close_group_size) - .await - .into_iter() - .map(|node| node.peer_id) - .collect(); for peer in peers { // A peer that has left the group owes nothing, and one that cannot // be asked is never judged on its silence. - if *peer == self_id || !group.contains(peer) || !self.is_capable(peer) { + if *peer == self_id || !self.owes(peer, &fresh.address).await { continue; } - let holds = match self.fetch_outcome(peer, &fresh.address).await { - Fetched::Record(record) => { - let state = record.state(); - state.state_id == fresh.state_id || state.replaces(&fresh) - } - // Serving a record that does not verify has already been - // charged to the peer, once; it is not charged again here. - Fetched::Invalid => continue, - Fetched::Nothing => false, - }; - if !holds { - warn!( - "Peer {peer} does not hold pointer {} it was offered", - hex::encode(fresh.address) - ); - self.penalise(peer).await; + let fetched = self.fetch_outcome(peer, &fresh.address).await; + if judge_possession(fetched, &fresh) != Possession::Missing { + continue; + } + // Asking can wait on this peer's other requests and on the peers + // before it, and it may have left the group meanwhile. It is + // judged on what it owes now, not on what it owed when this began. + if !self.owes(peer, &fresh.address).await { + continue; } + warn!( + "Peer {peer} does not hold pointer {} it was offered", + hex::encode(fresh.address) + ); + self.penalise(peer).await; } } + /// Whether `peer` is responsible for `address` in this node's view, and + /// can be asked about it. + async fn owes(&self, peer: &PeerId, address: &XorName) -> bool { + self.is_capable(peer) + && self + .p2p + .dht_manager() + .find_closest_nodes_local_with_self(address, self.config.close_group_size) + .await + .iter() + .any(|node| node.peer_id == *peer) + } + // ----------------------------------------------------------------------- // Pruning // ----------------------------------------------------------------------- @@ -1350,6 +1389,7 @@ impl PointerReplication { mod tests { use super::*; use ant_protocol::pointer::{PointerTarget, PointerTargetKind}; + use saorsa_pqc::api::sig::ml_dsa_65; fn state(counter: u64, target: u8, id: u8) -> PointerState { PointerState { @@ -1370,6 +1410,31 @@ mod tests { /// A peer's outbound permits are dropped only once nothing holds them, so /// a request in flight keeps the one set later requests share. + fn record(counter: u64, target: u8) -> Pointer { + let (pk, sk) = ml_dsa_65().generate_keypair_from_seed(&[3; 32]); + let target = PointerTarget::new(PointerTargetKind::Chunk, [target; 32]); + Pointer::sign(&sk, &pk, counter, target).expect("sign") + } + + #[test] + fn a_possession_check_judges_each_answer_once() { + let fresh = record(5, 1); + let judge = |fetched| judge_possession(fetched, &fresh.state()); + assert_eq!(judge(Fetched::Record(fresh.clone())), Possession::Holds); + assert_eq!( + judge(Fetched::Record(record(6, 9))), + Possession::Holds, + "a newer state is as good" + ); + assert_eq!(judge(Fetched::Record(record(4, 1))), Possession::Missing); + assert_eq!(judge(Fetched::Nothing), Possession::Missing); + // The fetch already charged the peer for an invalid record; charging + // it again as missing would count one bad answer twice. + assert_eq!(judge(Fetched::Invalid), Possession::NotJudged); + // A failure on this node says nothing about the peer. + assert_eq!(judge(Fetched::LocalFailure), Possession::NotJudged); + } + #[tokio::test] async fn outbound_permits_in_use_are_never_dropped() { let permits = Arc::new(Semaphore::new(MAX_REQUESTS_PER_PEER)); From 651b10c58c6eb02e70539bd0358e2ed9424603bf Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 12:58:16 +0900 Subject: [PATCH 4/7] fix(pointer): push hints only for admitted syncs, once per peer per interval Every neighbour-sync request started a hint push before admission, so a peer sending small sync requests, even refused ones, made this node scan every pointer it holds, with a routing lookup per record, as often as it liked. Hints now answer only an admitted request, at most once per peer per shortest sync interval, which an honest peer never syncs faster than, and at most eight answers run at once. A request that finds them busy gets no hints this time; the next sync round sends them anyway. The traffic summary counted the six pointer messages but never logged them. It now does, in a fourth summary line. --- src/replication/mod.rs | 13 +++-- src/replication/pointer.rs | 106 ++++++++++++++++++++++++++++++++++++ src/replication/protocol.rs | 17 ++++++ 3 files changed, 130 insertions(+), 6 deletions(-) diff --git a/src/replication/mod.rs b/src/replication/mod.rs index cf731c34..6b4016de 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -6540,12 +6540,6 @@ async fn dispatch_neighbor_sync_request( received_at: Instant, rr_message_id: Option<&str>, ) -> Result<()> { - // A peer syncing with us gets our pointer hints as well — including a - // node that is bootstrapping, which is how it learns the pointers it - // should hold. - if let Some(pointers) = &ctx.pointers { - pointers.push_hints_detached(vec![source]); - } let guard = match admit_bounded_responder( &ctx.neighbor_sync_responder_admission_semaphore, &ctx.neighbor_sync_responder_inflight, @@ -6569,6 +6563,13 @@ async fn dispatch_neighbor_sync_request( return Ok(()); } }; + // A peer whose sync request is admitted gets our pointer hints as well, + // including a node that is bootstrapping, which is how it learns the + // pointers it should hold. Only an admitted request: each answer scans + // every pointer held, and a refused one must cost nothing. + if let Some(pointers) = &ctx.pointers { + pointers.answer_sync_with_hints(source); + } let worker_semaphore = Arc::clone(&ctx.neighbor_sync_responder_worker_semaphore); let p2p_node = Arc::clone(&ctx.p2p_node); diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 6d6799b8..557c7230 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -111,6 +111,16 @@ const _: () = assert!( /// is learned again from the peer's next hint push. const MAX_CAPABLE_PEERS: usize = 4096; +/// Most hint pushes answering peers' sync requests at once. Each scans every +/// pointer this node holds, so a peer's sync request that finds them all +/// running gets no hints this time, and the next sync round delivers them. +const MAX_CONCURRENT_SYNC_ANSWERS: usize = 8; + +/// Most peers remembered as having been answered with hints. An honest peer +/// syncs once per interval, so this only bounds what a flood of identities +/// could leave behind. +const MAX_ANSWERED_SYNC_PEERS: usize = 4096; + /// How often the pending hints are looked at. const VERIFICATION_TICK: Duration = Duration::from_millis(500); @@ -319,10 +329,40 @@ pub struct PointerReplication { pending: Mutex>, /// When each held address was first seen continuously out of range. out_of_range: Mutex>, + /// Hint pushes answering sync requests that may run at once. + sync_answers: Arc, + /// When each peer was last answered with hints. + answered_syncs: Mutex, shutdown: CancellationToken, tracker: TaskTracker, } +/// When each peer was last sent hints in answer to its own sync request. +#[derive(Default)] +struct AnsweredSyncs { + at: HashMap, +} + +impl AnsweredSyncs { + /// Whether `peer` may be answered at `now`, and if so note it: not if it + /// was answered less than `spacing` ago, nor while the map is full of + /// peers answered that recently. + fn admit(&mut self, peer: PeerId, now: Instant, spacing: Duration) -> bool { + let recent = |at: &Instant| now.saturating_duration_since(*at) < spacing; + if self.at.get(&peer).is_some_and(recent) { + return false; + } + if self.at.len() >= MAX_ANSWERED_SYNC_PEERS { + self.at.retain(|_, at| recent(at)); + if self.at.len() >= MAX_ANSWERED_SYNC_PEERS { + return false; + } + } + self.at.insert(peer, now); + true + } +} + impl PointerReplication { /// Pointer replication over `store`, sharing the engine's resources. #[allow(clippy::too_many_arguments)] @@ -353,6 +393,8 @@ impl PointerReplication { capable: Mutex::new(HashSet::new()), pending: Mutex::new(HashMap::new()), out_of_range: Mutex::new(HashMap::new()), + sync_answers: Arc::new(Semaphore::new(MAX_CONCURRENT_SYNC_ANSWERS)), + answered_syncs: Mutex::new(AnsweredSyncs::default()), shutdown, tracker, } @@ -546,6 +588,34 @@ impl PointerReplication { .spawn(async move { this.push_hints(&peers).await }); } + /// Answer a peer's admitted sync request with this node's pointer hints, + /// in the background: how a peer, including one still bootstrapping, + /// learns the pointers it should hold. + /// + /// Each answer scans every pointer held, so a peer is answered at most + /// once per shortest sync interval, which an honest peer never syncs + /// faster than, and at most [`MAX_CONCURRENT_SYNC_ANSWERS`] answers run at + /// once. A request that finds them all busy gets no hints this time; the + /// next sync round sends them anyway. + pub(crate) fn answer_sync_with_hints(self: &Arc, peer: PeerId) { + let Ok(permit) = Arc::clone(&self.sync_answers).try_acquire_owned() else { + debug!("Not answering {peer}'s sync with pointer hints: too many answers running"); + return; + }; + if !self.answered_syncs.lock().admit( + peer, + Instant::now(), + self.config.neighbor_sync_interval_min, + ) { + return; + } + let this = Arc::clone(self); + self.tracker.spawn(async move { + let _permit = permit; + this.push_hints(&[peer]).await; + }); + } + /// Push hints to `peers` and wait until they are sent. This is what each /// neighbour-sync round runs in the background; tests call it to drive a /// round. @@ -1435,6 +1505,42 @@ mod tests { assert_eq!(judge(Fetched::LocalFailure), Possession::NotJudged); } + #[test] + fn a_peer_is_answered_with_hints_once_per_sync_interval() { + let spacing = Duration::from_secs(600); + let start = Instant::now(); + let mut answered = AnsweredSyncs::default(); + assert!(answered.admit(peer(1), start, spacing)); + assert!( + !answered.admit(peer(1), start + Duration::from_secs(1), spacing), + "a second request inside the interval gets no second scan" + ); + assert!( + answered.admit(peer(2), start, spacing), + "other peers are not held up" + ); + assert!(answered.admit(peer(1), start + spacing, spacing)); + + // A flood of identities fills the map, and then no one new is + // answered until the oldest are old enough to forget. + let mut full = AnsweredSyncs::default(); + for i in 0..MAX_ANSWERED_SYNC_PEERS { + let id = u32::try_from(i).expect("fits"); + let mut bytes = [0u8; 32]; + if let Some(prefix) = bytes.get_mut(..4) { + prefix.copy_from_slice(&id.to_be_bytes()); + } + assert!(full.admit(PeerId::from_bytes(bytes), start, spacing)); + } + assert!(!full.admit(peer(0xEE), start, spacing)); + assert!(full.admit(peer(0xEE), start + spacing, spacing)); + assert_eq!( + full.at.len(), + 1, + "everyone older than the interval was forgotten" + ); + } + #[tokio::test] async fn outbound_permits_in_use_are_never_dropped() { let permits = Arc::new(Semaphore::new(MAX_REQUESTS_PER_PEER)); diff --git a/src/replication/protocol.rs b/src/replication/protocol.rs index b70ae9a7..430c7e6c 100644 --- a/src/replication/protocol.rs +++ b/src/replication/protocol.rs @@ -637,6 +637,23 @@ pub(crate) fn log_traffic_summary() { get_commitment_by_pin_response_rx_count = rc(16), "replication traffic summary (cumulative)" ); + crate::logging::info!( + target: "ant_node::replication::traffic", + group = 4, + pointer_fresh_offer_tx_bytes = tb(17), pointer_fresh_offer_tx_count = tc(17), + pointer_fresh_offer_rx_bytes = rb(17), pointer_fresh_offer_rx_count = rc(17), + pointer_hints_tx_bytes = tb(18), pointer_hints_tx_count = tc(18), + pointer_hints_rx_bytes = rb(18), pointer_hints_rx_count = rc(18), + pointer_fetch_request_tx_bytes = tb(19), pointer_fetch_request_tx_count = tc(19), + pointer_fetch_request_rx_bytes = rb(19), pointer_fetch_request_rx_count = rc(19), + pointer_fetch_response_tx_bytes = tb(20), pointer_fetch_response_tx_count = tc(20), + pointer_fetch_response_rx_bytes = rb(20), pointer_fetch_response_rx_count = rc(20), + pointer_state_request_tx_bytes = tb(21), pointer_state_request_tx_count = tc(21), + pointer_state_request_rx_bytes = rb(21), pointer_state_request_rx_count = rc(21), + pointer_state_response_tx_bytes = tb(22), pointer_state_response_tx_count = tc(22), + pointer_state_response_rx_bytes = rb(22), pointer_state_response_rx_count = rc(22), + "replication traffic summary (cumulative)" + ); } // --------------------------------------------------------------------------- From f59bec59f080b74032359335c48c0ee5ea88dcab Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 13:49:19 +0900 Subject: [PATCH 5/7] fix(pointer): answer only fresh syncs from routed peers with hints, and never refuse a new one Hint answers ran before the neighbour-sync worker's freshness check, so a request later shed as stale had already cost a scan of every pointer, and any sender could be answered. A full map of answered peers also refused every new peer, an honest bootstrapping one included, for a whole sync interval. Answers now run only for a request that is admitted and still fresh, only for a routing-table peer, and a full map forgets the peer answered longest ago instead of refusing anyone. --- src/replication/mod.rs | 17 ++++---- src/replication/pointer.rs | 83 +++++++++++++++++++++++--------------- 2 files changed, 59 insertions(+), 41 deletions(-) diff --git a/src/replication/mod.rs b/src/replication/mod.rs index 6b4016de..c2f7b91a 100644 --- a/src/replication/mod.rs +++ b/src/replication/mod.rs @@ -6563,14 +6563,7 @@ async fn dispatch_neighbor_sync_request( return Ok(()); } }; - // A peer whose sync request is admitted gets our pointer hints as well, - // including a node that is bootstrapping, which is how it learns the - // pointers it should hold. Only an admitted request: each answer scans - // every pointer held, and a refused one must cost nothing. - if let Some(pointers) = &ctx.pointers { - pointers.answer_sync_with_hints(source); - } - + let pointers = ctx.pointers.clone(); let worker_semaphore = Arc::clone(&ctx.neighbor_sync_responder_worker_semaphore); let p2p_node = Arc::clone(&ctx.p2p_node); let storage = Arc::clone(&ctx.storage); @@ -6603,6 +6596,14 @@ async fn dispatch_neighbor_sync_request( ); return; } + // A peer whose sync request is admitted and fresh gets our pointer + // hints as well, including a node that is bootstrapping, which is how + // it learns the pointers it should hold. Only such a request: each + // answer scans every pointer held, and a refused or stale one must + // cost nothing. + if let Some(pointers) = &pointers { + pointers.answer_sync_with_hints(source); + } if let Err(e) = handle_neighbor_sync_request( &source, &request, diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 557c7230..333412d5 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -116,9 +116,10 @@ const MAX_CAPABLE_PEERS: usize = 4096; /// running gets no hints this time, and the next sync round delivers them. const MAX_CONCURRENT_SYNC_ANSWERS: usize = 8; -/// Most peers remembered as having been answered with hints. An honest peer -/// syncs once per interval, so this only bounds what a flood of identities -/// could leave behind. +/// Most peers remembered as having been answered with hints. Only routing-table +/// peers are answered, so this only bounds what a lagging table could leave +/// behind; past it the peer answered longest ago is forgotten, never a new one +/// refused. const MAX_ANSWERED_SYNC_PEERS: usize = 4096; /// How often the pending hints are looked at. @@ -345,17 +346,24 @@ struct AnsweredSyncs { impl AnsweredSyncs { /// Whether `peer` may be answered at `now`, and if so note it: not if it - /// was answered less than `spacing` ago, nor while the map is full of - /// peers answered that recently. + /// was answered less than `spacing` ago. A full map forgets the peers + /// answered longest ago rather than refuse anyone new. fn admit(&mut self, peer: PeerId, now: Instant, spacing: Duration) -> bool { let recent = |at: &Instant| now.saturating_duration_since(*at) < spacing; if self.at.get(&peer).is_some_and(recent) { return false; } - if self.at.len() >= MAX_ANSWERED_SYNC_PEERS { + if self.at.len() >= MAX_ANSWERED_SYNC_PEERS && !self.at.contains_key(&peer) { self.at.retain(|_, at| recent(at)); - if self.at.len() >= MAX_ANSWERED_SYNC_PEERS { - return false; + } + if self.at.len() >= MAX_ANSWERED_SYNC_PEERS && !self.at.contains_key(&peer) { + let oldest = self + .at + .iter() + .min_by_key(|(_, at)| **at) + .map(|(peer, _)| *peer); + if let Some(oldest) = oldest { + self.at.remove(&oldest); } } self.at.insert(peer, now); @@ -588,31 +596,35 @@ impl PointerReplication { .spawn(async move { this.push_hints(&peers).await }); } - /// Answer a peer's admitted sync request with this node's pointer hints, - /// in the background: how a peer, including one still bootstrapping, - /// learns the pointers it should hold. + /// Answer a peer's sync request, admitted and still fresh, with this + /// node's pointer hints, in the background: how a peer, including one + /// still bootstrapping, learns the pointers it should hold. /// - /// Each answer scans every pointer held, so a peer is answered at most - /// once per shortest sync interval, which an honest peer never syncs - /// faster than, and at most [`MAX_CONCURRENT_SYNC_ANSWERS`] answers run at - /// once. A request that finds them all busy gets no hints this time; the - /// next sync round sends them anyway. + /// Each answer scans every pointer held, so only a routing-table peer is + /// answered, at most once per shortest sync interval, which an honest + /// peer never syncs faster than, and at most + /// [`MAX_CONCURRENT_SYNC_ANSWERS`] answers run at once. A request that + /// finds them all busy gets no hints this time; the next sync round sends + /// them anyway. pub(crate) fn answer_sync_with_hints(self: &Arc, peer: PeerId) { let Ok(permit) = Arc::clone(&self.sync_answers).try_acquire_owned() else { debug!("Not answering {peer}'s sync with pointer hints: too many answers running"); return; }; - if !self.answered_syncs.lock().admit( - peer, - Instant::now(), - self.config.neighbor_sync_interval_min, - ) { - return; - } let this = Arc::clone(self); self.tracker.spawn(async move { let _permit = permit; - this.push_hints(&[peer]).await; + if !this.p2p.dht_manager().is_in_routing_table(&peer).await { + return; + } + let admitted = this.answered_syncs.lock().admit( + peer, + Instant::now(), + this.config.neighbor_sync_interval_min, + ); + if admitted { + this.push_hints(&[peer]).await; + } }); } @@ -1521,24 +1533,29 @@ mod tests { ); assert!(answered.admit(peer(1), start + spacing, spacing)); - // A flood of identities fills the map, and then no one new is - // answered until the oldest are old enough to forget. + // A full map forgets the peer answered longest ago; a new peer is + // never refused for it. let mut full = AnsweredSyncs::default(); + let mut first = None; for i in 0..MAX_ANSWERED_SYNC_PEERS { let id = u32::try_from(i).expect("fits"); let mut bytes = [0u8; 32]; if let Some(prefix) = bytes.get_mut(..4) { prefix.copy_from_slice(&id.to_be_bytes()); } - assert!(full.admit(PeerId::from_bytes(bytes), start, spacing)); + let at = start + Duration::from_millis(u64::from(id)); + let each = PeerId::from_bytes(bytes); + first.get_or_insert(each); + assert!(full.admit(each, at, spacing)); } - assert!(!full.admit(peer(0xEE), start, spacing)); - assert!(full.admit(peer(0xEE), start + spacing, spacing)); - assert_eq!( - full.at.len(), - 1, - "everyone older than the interval was forgotten" + let late = start + Duration::from_secs(10); + assert!( + full.admit(peer(0xEE), late, spacing), + "a new peer is answered" ); + assert_eq!(full.at.len(), MAX_ANSWERED_SYNC_PEERS); + let first = first.expect("one was admitted"); + assert!(!full.at.contains_key(&first), "the oldest was forgotten"); } #[tokio::test] From 7132e3c3000e65aca2596129684e93324d903645 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 15:07:37 +0900 Subject: [PATCH 6/7] fix(pointer): give up a request that cannot be sent in time, charge nobody for it, and stop serving at shutdown This node's own requests to one peer waited for a permit with no limit and past shutdown, outside the request's timeout, so repeated paid updates to one close peer could pile up possession checks waiting forever and hold shutdown open. A request now waits for its permit no longer than its own timeout and not past shutdown, and one never sent is its own outcome: a possession check judges nobody on it. Serve tasks waiting for a serve permit ignored shutdown, so requests admitted just before it kept reading disk and answering after cancellation. The wait now ends at shutdown, and the work behind it. --- src/replication/pointer.rs | 113 ++++++++++++++++++++++++++++++++----- 1 file changed, 99 insertions(+), 14 deletions(-) diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 333412d5..8546fd99 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -43,7 +43,7 @@ use parking_lot::Mutex; use rand::Rng; use saorsa_core::identity::PeerId; use saorsa_core::{P2PNode, TrustEvent}; -use tokio::sync::{mpsc, RwLock, Semaphore}; +use tokio::sync::{mpsc, OwnedSemaphorePermit, RwLock, Semaphore}; use tokio::task::{spawn_blocking, JoinHandle}; use tokio_util::sync::CancellationToken; use tokio_util::task::TaskTracker; @@ -253,10 +253,35 @@ enum Fetched { Invalid, /// No record: no answer, an empty one, or one about another address. Nothing, - /// A record this node could not check, for a reason of its own. + /// A record this node could not check, or a request it never sent, for a + /// reason of its own. LocalFailure, } +/// What asking a peer came to (see [`PointerReplication::ask`]). +enum Asked { + /// The peer's answer. + Answered(ReplicationMessageBody), + /// Sent, and no usable answer came back. + Silent, + /// Never sent: this node is shutting down, or its own requests to the + /// peer stayed busy for the whole timeout. + NotSent, +} + +/// One of `permits`, if one comes free within `timeout` and before +/// `shutdown`. +async fn acquire_or_give_up( + permits: Arc, + timeout: Duration, + shutdown: &CancellationToken, +) -> Option { + tokio::select! { + permit = tokio::time::timeout(timeout, permits.acquire_owned()) => permit.ok()?.ok(), + () = shutdown.cancelled() => None, + } +} + /// What a possession check makes of one peer's answer about `fresh`. #[derive(Debug, PartialEq, Eq)] enum Possession { @@ -495,20 +520,41 @@ impl PointerReplication { body: ReplicationMessageBody, timeout: Duration, ) -> Option { - let _permit = self.outbound_permits(peer).acquire_owned().await.ok()?; + match self.ask(peer, body, timeout).await { + Asked::Answered(body) => Some(body), + Asked::Silent | Asked::NotSent => None, + } + } + + /// Ask `peer`, saying whether a missing answer was the peer's silence or + /// this node never sending the request at all. + /// + /// The request waits for one of this node's own permits for `peer` first, + /// but no longer than `timeout` and not past shutdown: a request that + /// could not be sent in that time is dropped as not sent, never left + /// waiting, and never held against the peer. + async fn ask(&self, peer: &PeerId, body: ReplicationMessageBody, timeout: Duration) -> Asked { + let Some(_permit) = + acquire_or_give_up(self.outbound_permits(peer), timeout, &self.shutdown).await + else { + return Asked::NotSent; + }; let msg = ReplicationMessage { request_id: rand::thread_rng().gen::(), body, }; - let bytes = msg.encode().ok()?; - let response = self + let Ok(bytes) = msg.encode() else { + return Asked::NotSent; + }; + let Ok(response) = self .p2p .send_request(peer, REPLICATION_PROTOCOL_ID, bytes, timeout) .await - .ok()?; + else { + return Asked::Silent; + }; ReplicationMessage::decode(&response.data) - .ok() - .map(|msg| msg.body) + .map_or(Asked::Silent, |msg| Asked::Answered(msg.body)) } async fn penalise(&self, peer: &PeerId) { @@ -533,8 +579,8 @@ impl PointerReplication { /// [`Self::fetch_record`], saying whether a failure was a record that did /// not verify, which has already cost the peer, or no record at all. async fn fetch_outcome(&self, peer: &PeerId, address: &XorName) -> Fetched { - let Some(body) = self - .request( + let body = match self + .ask( peer, ReplicationMessageBody::PointerFetchRequest(PointerFetchRequest { address: *address, @@ -542,8 +588,11 @@ impl PointerReplication { self.config.fetch_request_timeout, ) .await - else { - return Fetched::Nothing; + { + Asked::Answered(body) => body, + Asked::Silent => return Fetched::Nothing, + // Never asked: nothing to judge the peer on. + Asked::NotSent => return Fetched::LocalFailure, }; let ReplicationMessageBody::PointerFetchResponse(response) = body else { return Fetched::Nothing; @@ -1012,7 +1061,12 @@ impl PointerReplication { let this = Arc::clone(self); self.tracker.spawn(async move { let _guard = guard; - let Ok(_permit) = this.serve_permits.acquire().await else { + // Shutting down ends the wait, and the work behind it. + let permit = tokio::select! { + permit = this.serve_permits.acquire() => permit, + () = this.shutdown.cancelled() => return, + }; + let Ok(_permit) = permit else { return; }; this.mark_capable(&source).await; @@ -1060,7 +1114,11 @@ impl PointerReplication { let this = Arc::clone(self); self.tracker.spawn(async move { let _guard = guard; - let Ok(_permit) = this.serve_permits.acquire().await else { + let permit = tokio::select! { + permit = this.serve_permits.acquire() => permit, + () = this.shutdown.cancelled() => return, + }; + let Ok(_permit) = permit else { return; }; this.mark_capable(&source).await; @@ -1498,6 +1556,33 @@ mod tests { Pointer::sign(&sk, &pk, counter, target).expect("sign") } + #[tokio::test(start_paused = true)] + async fn a_request_that_cannot_be_sent_in_time_is_given_up_not_queued() { + let permits = Arc::new(Semaphore::new(1)); + let shutdown = CancellationToken::new(); + let held = acquire_or_give_up(Arc::clone(&permits), Duration::from_secs(1), &shutdown) + .await + .expect("a free permit is taken"); + + // Busy past the timeout: given up, not left waiting. + assert!( + acquire_or_give_up(Arc::clone(&permits), Duration::from_secs(1), &shutdown) + .await + .is_none() + ); + + // Shutdown ends a wait at once, however long it could still run. + shutdown.cancel(); + let waited = tokio::time::Instant::now(); + assert!( + acquire_or_give_up(Arc::clone(&permits), Duration::from_secs(3600), &shutdown) + .await + .is_none() + ); + assert!(waited.elapsed() < Duration::from_secs(1)); + drop(held); + } + #[test] fn a_possession_check_judges_each_answer_once() { let fresh = record(5, 1); From 5f5b457cf48454b8f77a1c416323e4bc97fcd883 Mon Sep 17 00:00:00 2001 From: grumbach Date: Wed, 30 Sep 2026 15:38:49 +0900 Subject: [PATCH 7/7] fix(pointer): let shutdown win when it and a permit are ready together The waits added for shutdown picked either branch at random when a permit was free and shutdown had already begun, so a request could still be sent, or a stored record read and served, after cancellation. Shutdown is now checked first and again once a permit is held. A test asks repeatedly with a free permit after shutdown and requires nothing be sent each time. --- src/replication/pointer.rs | 32 ++++++++++++++++++++++++++------ 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/src/replication/pointer.rs b/src/replication/pointer.rs index 8546fd99..bbd43a74 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -270,16 +270,18 @@ enum Asked { } /// One of `permits`, if one comes free within `timeout` and before -/// `shutdown`. +/// `shutdown`. Shutdown wins when both are ready at once. async fn acquire_or_give_up( permits: Arc, timeout: Duration, shutdown: &CancellationToken, ) -> Option { - tokio::select! { - permit = tokio::time::timeout(timeout, permits.acquire_owned()) => permit.ok()?.ok(), + let permit = tokio::select! { + biased; () = shutdown.cancelled() => None, - } + permit = tokio::time::timeout(timeout, permits.acquire_owned()) => permit.ok()?.ok(), + }?; + (!shutdown.is_cancelled()).then_some(permit) } /// What a possession check makes of one peer's answer about `fresh`. @@ -1063,12 +1065,16 @@ impl PointerReplication { let _guard = guard; // Shutting down ends the wait, and the work behind it. let permit = tokio::select! { - permit = this.serve_permits.acquire() => permit, + biased; () = this.shutdown.cancelled() => return, + permit = this.serve_permits.acquire() => permit, }; let Ok(_permit) = permit else { return; }; + if this.shutdown.is_cancelled() { + return; + } this.mark_capable(&source).await; // `get` verifies the signature before serving, so a record damaged // on this disk is never handed on. @@ -1115,12 +1121,16 @@ impl PointerReplication { self.tracker.spawn(async move { let _guard = guard; let permit = tokio::select! { - permit = this.serve_permits.acquire() => permit, + biased; () = this.shutdown.cancelled() => return, + permit = this.serve_permits.acquire() => permit, }; let Ok(_permit) = permit else { return; }; + if this.shutdown.is_cancelled() { + return; + } this.mark_capable(&source).await; let states = request .addresses @@ -1581,6 +1591,16 @@ mod tests { ); assert!(waited.elapsed() < Duration::from_secs(1)); drop(held); + + // With a permit free and shutdown already under way, nothing is sent, + // every time. + for _ in 0..64 { + assert!( + acquire_or_give_up(Arc::clone(&permits), Duration::from_secs(1), &shutdown) + .await + .is_none() + ); + } } #[test]