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..c2f7b91a 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(()) } @@ -6536,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, @@ -6565,7 +6563,7 @@ async fn dispatch_neighbor_sync_request( return Ok(()); } }; - + 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); @@ -6598,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 df62a0f4..bbd43a74 100644 --- a/src/replication/pointer.rs +++ b/src/replication/pointer.rs @@ -43,8 +43,8 @@ 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::task::JoinHandle; +use tokio::sync::{mpsc, OwnedSemaphorePermit, RwLock, Semaphore}; +use tokio::task::{spawn_blocking, JoinHandle}; use tokio_util::sync::CancellationToken; use tokio_util::task::TaskTracker; @@ -79,6 +79,49 @@ 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; + +/// 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. 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. const VERIFICATION_TICK: Duration = Duration::from_millis(500); @@ -202,6 +245,93 @@ 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, + /// 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`. Shutdown wins when both are ready at once. +async fn acquire_or_give_up( + permits: Arc, + timeout: Duration, + shutdown: &CancellationToken, +) -> Option { + 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`. +#[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. +/// +/// 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,16 +344,60 @@ 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. 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. 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 && !self.at.contains_key(&peer) { + self.at.retain(|_, at| recent(at)); + } + 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); + true + } +} + impl PointerReplication { /// Pointer replication over `store`, sharing the engine's resources. #[allow(clippy::too_many_arguments)] @@ -248,9 +422,14 @@ 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()), + sync_answers: Arc::new(Semaphore::new(MAX_CONCURRENT_SYNC_ANSWERS)), + answered_syncs: Mutex::new(AnsweredSyncs::default()), shutdown, tracker, } @@ -265,8 +444,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 +470,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,19 +522,41 @@ impl PointerReplication { body: ReplicationMessageBody, timeout: Duration, ) -> Option { + 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) { @@ -339,31 +572,59 @@ 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 - .request( + match self.fetch_outcome(peer, address).await { + Fetched::Record(record) => Some(record), + Fetched::Invalid | Fetched::Nothing | Fetched::LocalFailure => 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 body = match self + .ask( peer, ReplicationMessageBody::PointerFetchRequest(PointerFetchRequest { address: *address, }), self.config.fetch_request_timeout, ) - .await?; + .await + { + 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 None; + return Fetched::Nothing; }; if response.address != *address { - return None; + return Fetched::Nothing; } - let bytes = response.record?; - match Pointer::from_bytes(&bytes) { - Ok(record) if record.address() == *address => Some(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 Ok(parsed) = spawn_blocking(move || Pointer::from_bytes(&bytes)).await else { + // 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), Ok(_) | Err(_) => { debug!( "Peer {peer} served an invalid pointer record for {}", hex::encode(address) ); self.penalise(peer).await; - None + Fetched::Invalid } } } @@ -386,6 +647,38 @@ impl PointerReplication { .spawn(async move { this.push_hints(&peers).await }); } + /// 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 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; + }; + let this = Arc::clone(self); + self.tracker.spawn(async move { + let _permit = permit; + 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; + } + }); + } + /// 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. @@ -433,7 +726,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 +799,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 +864,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 +1025,57 @@ 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 Ok(_permit) = this.serve_permits.acquire().await else { + let _guard = guard; + // Shutting down ends the wait, and the work behind it. + let permit = tokio::select! { + 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. let record = match this.store.get(&request.address).await { @@ -751,19 +1097,41 @@ 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 Ok(_permit) = this.serve_permits.acquire().await else { + let _guard = guard; + let permit = tokio::select! { + 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 .iter() @@ -875,7 +1243,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 +1250,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; }); } @@ -998,37 +1366,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 = self - .fetch_record(peer, &fresh.address) - .await - .is_some_and(|record| { - let state = record.state(); - state.state_id == fresh.state_id || state.replaces(&fresh) - }); - 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 // ----------------------------------------------------------------------- @@ -1165,6 +1539,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 { @@ -1183,6 +1558,149 @@ 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. + 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") + } + + #[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); + + // 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] + 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); + } + + #[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 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()); + } + 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)); + } + 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] + 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); 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)" + ); } // ---------------------------------------------------------------------------