diff --git a/src/proto/bloom/mod.rs b/src/proto/bloom/mod.rs index cf4332e0..b58ad251 100644 --- a/src/proto/bloom/mod.rs +++ b/src/proto/bloom/mod.rs @@ -28,7 +28,7 @@ mod tests; pub use core::BloomFilter; pub use limits::{DEFAULT_FILTER_SIZE_BITS, DEFAULT_HASH_COUNT, V1_SIZE_CLASS}; -pub use state::BloomState; +pub use state::{BloomState, LinkEvidence, RrCounters}; pub use wire::FilterAnnounce; /// Errors related to Bloom filter operations. diff --git a/src/proto/bloom/state.rs b/src/proto/bloom/state.rs index 034dd2fe..e827c4f9 100644 --- a/src/proto/bloom/state.rs +++ b/src/proto/bloom/state.rs @@ -5,6 +5,225 @@ use alloc::collections::{BTreeMap, BTreeSet}; use super::BloomFilter; use crate::NodeAddr; +/// How long an announce the receiver reports cannot check waits before its +/// one unchecked resend, in milliseconds. Equal to the default link dead +/// timeout, so an outage that did not remove the peer has ended by then. +pub const FALLBACK_MS: u64 = 30_000; + +/// The largest gap the per-peer resend backoff imposes, in milliseconds. +pub const MAXGAP_MS: u64 = 60_000; + +/// A run of resends with no resend for this long, in milliseconds, resets the +/// backoff. It must exceed [`MAXGAP_MS`], or a sustained trigger resending at +/// the largest gap would reset its own backoff every time. +pub const QUIET_MS: u64 = 120_000; + +/// Unchecked resends (`Unverified`, `SessionChanged` or `Timeout`) allowed per +/// announce lineage per session. +pub const UNVERIFIED_BUDGET: u8 = 1; + +/// Resends on reported loss allowed per announce lineage per session. +pub const LOSS_BUDGET: u8 = 3; + +/// Highest backoff level. `gap` at this level is already capped at +/// [`MAXGAP_MS`], so a higher level would add nothing; the cap keeps the +/// shift in range. +const MAX_LEVEL: u8 = 7; + +/// The cumulative counters of one ReceiverReport the peer sent about our +/// frames on a link. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RrCounters { + /// Highest link counter the peer had received from us. + pub highest: u64, + /// Link frames from us the peer had counted, cumulative. + pub received: u64, + /// Of those, frames that arrived below the highest counter, cumulative. + pub reordered: u32, +} + +/// What the shell reads from one peer's link at one moment, for deciding +/// whether an announce reached that peer. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct LinkEvidence { + /// Identity of the current link session (from its handshake hash). The + /// send counter restarts at 0 in every session, so counters are compared + /// only within one epoch. + pub epoch: u64, + /// The next send counter the current session will use. + pub next_counter: u64, + /// The last ReceiverReport accepted in the current session, if any. + pub rr: Option, +} + +impl LinkEvidence { + /// The report, when it can describe the current session. + /// + /// The peer cannot have received a counter this session has not used yet, + /// so a report whose highest counter is at or above `next_counter` + /// describes another session. That happens briefly around a rekey, when a + /// report or frame of the old session is counted against the new one, and + /// such a report is no evidence either way. + pub fn usable_rr(&self) -> Option { + self.rr.filter(|rr| rr.highest < self.next_counter) + } +} + +/// Why an announce is being resent, for the shell's log line. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ResendReason { + /// A report covering the announce shows fewer frames arrived since its + /// base than were sent. + Loss, + /// The first usable report already covers an announce it cannot check. + Unverified, + /// The announce was sent on an earlier session and cannot be checked. + SessionChanged, + /// No usable report checked the announce within the fallback interval. + Timeout, +} + +/// What an announce's delivery is measured from. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Base { + /// The last usable report before the announce was sent, and whether it + /// has no holes: every counter up to its highest had arrived. + Counted(RrCounters, bool), + /// Nothing received yet in the peer's first session: every counter from 0 + /// must arrive. + Zero, + /// No base: a later report can become one if it does not yet cover the + /// announce. + Unknown, +} + +/// The one announce to a peer still awaiting confirmation. +#[derive(Clone, Copy, Debug)] +struct SentAnnounce { + /// Link counter the announce was sent with. + counter: u64, + /// Session the counter belongs to (or, once orphaned, the session that + /// orphaned it). + epoch: u64, + /// What delivery is measured from. + base: Base, + /// Sent on an earlier session, so it can never be checked. + orphan: bool, + /// When it was sent or orphaned, for the fallback. + at_ms: u64, +} + +/// Per-peer delivery tracking for filter announces. +#[derive(Clone, Debug)] +struct AckState { + /// Session in which this entry was created; only there does a missing + /// report mean the peer has received nothing yet. + first_epoch: u64, + /// The outstanding announce, if one is unconfirmed. + sent: Option, + /// Backoff level: the number of resends in the current run, capped. + level: u8, + /// When the last resend was triggered. + resent_ms: Option, + /// Session the budgets were last refilled for. + budget_epoch: u64, + /// Unchecked resends left for the current lineage in this session. + unverified_left: u8, + /// Loss resends left for the current lineage in this session. + loss_left: u8, +} + +impl AckState { + /// A fresh entry for a peer first sent to in session `epoch`. + fn new(epoch: u64) -> Self { + Self { + first_epoch: epoch, + sent: None, + level: 0, + resent_ms: None, + budget_epoch: epoch, + unverified_left: UNVERIFIED_BUDGET, + loss_left: LOSS_BUDGET, + } + } + + /// Refill both budgets for session `epoch`. + fn refill(&mut self, epoch: u64) { + self.budget_epoch = epoch; + self.unverified_left = UNVERIFIED_BUDGET; + self.loss_left = LOSS_BUDGET; + } + + /// Whether report `rr`, taken in session `epoch`, shows no holes: every + /// counter up to its highest had arrived. Only in the peer's first session + /// does its cumulative count start at 0, so a later session's report is + /// never known to be whole. + fn whole(&self, epoch: u64, rr: RrCounters) -> bool { + epoch == self.first_epoch && rr.highest.checked_add(1) == Some(rr.received) + } + + /// Whether the backoff allows a resend at `now_ms`, resetting the level + /// after a quiet period. + fn backoff_allows(&mut self, now_ms: u64) -> bool { + let Some(last) = self.resent_ms else { + return true; + }; + if now_ms >= last.saturating_add(QUIET_MS) { + self.level = 0; + } + now_ms >= last.saturating_add(gap(self.level)) + } +} + +/// Minimum time after a resend before the next one, at backoff `level`. +fn gap(level: u8) -> u64 { + match level { + 0 => 0, + n => (1000u64 << (n.min(MAX_LEVEL) - 1)).min(MAXGAP_MS), + } +} + +/// Whether every counter in the report's range since `base` arrived. +/// +/// Within one receiver epoch every frame counted between two reports is a +/// distinct counter at or below `h1`. Those in `(h0, h1]` number at most all +/// receipts, `got`, and at least the non-reorder receipts, `sure`: a frame +/// that arrived after a higher counter is a reorder whether its counter lies +/// inside the window or at or below `h0`. +/// +/// - `got` below the span is a loss: fewer frames arrived than the window +/// holds. +/// - `sure` equal to the span is delivery. +/// - `got` equal to the span is delivery when the base has no holes, since +/// then no counter at or below `h0` is left to arrive late. +/// +/// Anything else is ambiguous and proves nothing, as is an inconsistent pair +/// (a counter went backwards), which means the two reports straddle a +/// receiver reset or another session's frame. Frames reserved but never +/// sent and frames dropped before counting only lower the counts, so a lost +/// frame is confirmed only if the peer overcounts. +fn delivered(base: Base, rr: RrCounters) -> Option { + let (r0, o0, span, complete) = match base { + Base::Counted(b, complete) => ( + b.received, + b.reordered, + rr.highest.checked_sub(b.highest)?, + complete, + ), + Base::Zero => (0, 0, rr.highest.checked_add(1)?, true), + Base::Unknown => return None, + }; + let got = rr.received.checked_sub(r0)?; + let sure = got.checked_sub(u64::from(rr.reordered.checked_sub(o0)?))?; + if got < span { + Some(false) + } else if sure == span || (complete && got == span) { + Some(true) + } else { + None + } +} + /// State for managing Bloom filter announcements. /// /// Tracks local filter state and what needs to be sent to peers. @@ -26,6 +245,10 @@ pub struct BloomState { sequence: u64, /// Last outgoing filter sent to each peer (for change detection). last_sent_filters: BTreeMap, + /// How long an unchecked announce waits for its fallback resend (ms). + fallback_ms: u64, + /// Per-peer delivery tracking for sent announces. + acks: BTreeMap, } impl BloomState { @@ -40,6 +263,8 @@ impl BloomState { pending_updates: BTreeSet::new(), sequence: 0, last_sent_filters: BTreeMap::new(), + fallback_ms: FALLBACK_MS, + acks: BTreeMap::new(), } } @@ -81,6 +306,12 @@ impl BloomState { self.update_debounce_ms = ms; } + /// Set how long an announce the receiver reports cannot check waits + /// before its fallback resend. Defaults to [`FALLBACK_MS`]. + pub fn set_fallback(&mut self, ms: u64) { + self.fallback_ms = ms; + } + /// Add a leaf dependent that we'll include in our filter. pub fn add_leaf_dependent(&mut self, node_addr: NodeAddr) { self.leaf_dependents.insert(node_addr); @@ -159,6 +390,150 @@ impl BloomState { self.last_sent_filters.remove(peer_id); self.last_update_sent.remove(peer_id); self.pending_updates.remove(peer_id); + self.acks.remove(peer_id); + } + + /// Record an announce the transport accepted for `peer`, sent with link + /// counter `counter`, so it stays outstanding until the peer's receiver + /// reports show it arrived. + /// + /// The transport accepting a frame is not delivery: a datagram can still + /// be lost, and announces are sent only when a filter changes, so a lost + /// one would otherwise leave the peer's filter stale indefinitely. + /// + /// Call this before [`record_sent_filter`](Self::record_sent_filter) for + /// the same send: it compares `filter` with the last one sent, and new + /// content starts a new lineage with fresh resend budgets. The budgets + /// also refill when the session changes, and never otherwise, so a resend + /// of the same content spends from its lineage's budget. + pub fn record_announce( + &mut self, + peer: NodeAddr, + filter: &BloomFilter, + counter: u64, + link: &LinkEvidence, + now_ms: u64, + ) { + let new_lineage = self.last_sent_filters.get(&peer) != Some(filter); + let ack = self + .acks + .entry(peer) + .or_insert_with(|| AckState::new(link.epoch)); + if new_lineage || link.epoch != ack.budget_epoch { + ack.refill(link.epoch); + } + // With no report yet, the zero baseline holds only in the peer's first + // session: a later session's cumulative count includes earlier ones. + let base = match link.usable_rr() { + Some(rr) if rr.highest < counter => Base::Counted(rr, ack.whole(link.epoch, rr)), + _ if link.rr.is_none() && link.epoch == ack.first_epoch => Base::Zero, + _ => Base::Unknown, + }; + ack.sent = Some(SentAnnounce { + counter, + epoch: link.epoch, + base, + orphan: false, + at_ms: now_ms, + }); + } + + /// Decide whether the outstanding announce to `peer` must be resent. + /// + /// Confirms the announce when a usable report covering its counter shows + /// every counter since its base arrived. Resends on a covering report that + /// shows a loss; once, as soon as a usable report arrives, for an announce + /// the reports cannot check; and once after the fallback interval when no + /// usable report checks it. A usable report that does not yet cover an + /// announce sent in the current session becomes its base instead of + /// triggering a resend. A report from another session, or a pair of + /// reports that is inconsistent or cannot tell a late frame from before + /// the base from one inside the window, is no evidence. Each announce lineage + /// gets [`UNVERIFIED_BUDGET`] unchecked and [`LOSS_BUDGET`] loss resends + /// per session, and a per-peer backoff spaces all resends by 1, 2, 4 ... + /// up to 60 s until [`QUIET_MS`] passes with none. + /// + /// On `Some`, the peer has been marked for an update; the ordinary send + /// path delivers the resend. + pub fn check_announce( + &mut self, + peer: &NodeAddr, + link: &LinkEvidence, + now_ms: u64, + ) -> Option { + let fallback_ms = self.fallback_ms; + let ack = self.acks.get_mut(peer)?; + let mut sent = ack.sent?; + + if link.epoch != ack.budget_epoch { + ack.refill(link.epoch); + } + if sent.epoch != link.epoch { + sent.orphan = true; + sent.epoch = link.epoch; + sent.base = Base::Unknown; + sent.at_ms = now_ms; + } + ack.sent = Some(sent); + + let rr = link.usable_rr(); + let due = now_ms >= sent.at_ms.saturating_add(fallback_ms); + let candidate = match (sent.base, rr) { + (Base::Unknown, Some(rr)) => { + if !sent.orphan && rr.highest < sent.counter { + sent.base = Base::Counted(rr, ack.whole(link.epoch, rr)); + ack.sent = Some(sent); + return None; + } + Some(if sent.orphan { + ResendReason::SessionChanged + } else { + ResendReason::Unverified + }) + } + (Base::Unknown, None) => due.then_some(ResendReason::Timeout), + (base, rr) => { + let covering = rr.filter(|rr| rr.highest >= sent.counter); + let loss = match covering.and_then(|rr| delivered(base, rr)) { + Some(true) => { + ack.sent = None; + return None; + } + Some(false) if ack.loss_left > 0 => Some(ResendReason::Loss), + _ => None, + }; + loss.or(due.then_some(ResendReason::Timeout)) + } + }?; + + let loss = candidate == ResendReason::Loss; + let left = if loss { + ack.loss_left + } else { + ack.unverified_left + }; + if left == 0 || !ack.backoff_allows(now_ms) { + return None; + } + if loss { + ack.loss_left -= 1; + } else { + ack.unverified_left -= 1; + } + ack.level = (ack.level + 1).min(MAX_LEVEL); + ack.resent_ms = Some(now_ms); + self.mark_update_needed(*peer); + Some(candidate) + } + + /// Whether an announce to `peer` is still awaiting confirmation. + pub fn announce_outstanding(&self, peer: &NodeAddr) -> bool { + self.outstanding_counter(peer).is_some() + } + + /// The link counter of the announce to `peer` awaiting confirmation. + pub fn outstanding_counter(&self, peer: &NodeAddr) -> Option { + self.acks.get(peer)?.sent.map(|sent| sent.counter) } /// Mark only peers whose outgoing filter has actually changed. diff --git a/src/proto/bloom/tests/state.rs b/src/proto/bloom/tests/state.rs index 159b2767..c5e66e74 100644 --- a/src/proto/bloom/tests/state.rs +++ b/src/proto/bloom/tests/state.rs @@ -313,3 +313,574 @@ fn test_bloom_state_mark_changed_peers_excludes_source() { assert!(!state.needs_update(&peer1)); } + +// ===== Delivery tracking for sent announces ===== +// +// Synthetic milliseconds, counters and receiver reports, no I/O. A report is +// written `(highest, received, reordered)`. Unless a test says otherwise, a +// send is recorded with `next_counter = counter + 1`, as the shell reads it +// straight after the send. + +use crate::NodeAddr; +use crate::proto::bloom::state::{ + FALLBACK_MS, LOSS_BUDGET, LinkEvidence, QUIET_MS, ResendReason, RrCounters, UNVERIFIED_BUDGET, +}; + +/// First session. +const E1: u64 = 0x0e01; +/// Second session. +const E2: u64 = 0x0e02; +/// Third session. +const E3: u64 = 0x0e03; + +/// A report's cumulative counters. +fn rr(highest: u64, received: u64, reordered: u32) -> Option { + Some(RrCounters { + highest, + received, + reordered, + }) +} + +/// Link evidence for session `epoch`. +fn link(epoch: u64, next_counter: u64, rr: Option) -> LinkEvidence { + LinkEvidence { + epoch, + next_counter, + rr, + } +} + +/// Filter content number `n`: distinct numbers give distinct filters. +fn content(n: u8) -> BloomFilter { + let mut filter = BloomFilter::new(); + filter.insert(&make_node_addr(100u8.wrapping_add(n))); + filter +} + +/// One peer's announces, driven as the shell drives them. +struct Track { + state: BloomState, + peer: NodeAddr, + content: u8, + counter: u64, +} + +impl Track { + /// A tracker with nothing sent yet, sending content 1. + fn new() -> Self { + Self { + state: BloomState::new(make_node_addr(0)), + peer: make_node_addr(1), + content: 1, + counter: 0, + } + } + + /// Record a send of the current content at `counter`, then the sent + /// filter, in the shell's order. + fn send(&mut self, counter: u64, link: LinkEvidence, now_ms: u64) { + let filter = content(self.content); + self.state + .record_announce(self.peer, &filter, counter, &link, now_ms); + self.state.record_sent_filter(self.peer, filter); + self.counter = counter; + } + + /// Send new content at `counter`. + fn send_new(&mut self, counter: u64, link: LinkEvidence, now_ms: u64) { + self.content += 1; + self.send(counter, link, now_ms); + } + + /// One tick of the tracker. + fn check(&mut self, link: LinkEvidence, now_ms: u64) -> Option { + self.state.check_announce(&self.peer, &link, now_ms) + } + + /// Whether the announce is still unconfirmed. + fn outstanding(&self) -> bool { + self.state.announce_outstanding(&self.peer) + } + + /// Check every 1,000 ms from `from_ms` to `to_ms` inclusive. `model` gives + /// the evidence at a time, from the outstanding counter: for a check with + /// `false`, and for recording a resend with `true`, where the resend takes + /// the counter `next_counter - 1` of that evidence. Each resend is + /// recorded, with new content when `renew` is set. Returns the resends. + fn hold( + &mut self, + from_ms: u64, + to_ms: u64, + renew: bool, + model: impl Fn(u64, u64, bool) -> LinkEvidence, + ) -> Vec<(u64, ResendReason)> { + let mut resends = Vec::new(); + let mut now = from_ms; + while now <= to_ms { + if let Some(reason) = self.check(model(now, self.counter, false), now) { + resends.push((now, reason)); + let ev = model(now, self.counter, true); + if renew { + self.content += 1; + } + self.send(ev.next_counter - 1, ev, now); + } + now += 1_000; + } + resends + } +} + +/// Loss on every check: each send is based on a report just below it, and +/// each check sees a report two counters on with one frame missing. +fn lossy(_now: u64, counter: u64, recording: bool) -> LinkEvidence { + if recording { + let n = counter + 1; + link(E1, n + 1, rr(n - 1, n, 0)) + } else { + link(E1, counter + 2, rr(counter + 1, counter + 1, 0)) + } +} + +/// The resend times of `resends`, in ms. +fn times(resends: &[(u64, ResendReason)]) -> Vec { + resends.iter().map(|(t, _)| *t).collect() +} + +/// A covering report with a frame missing since the base is a loss. +#[test] +fn test_bloom_ack_covering_report_with_a_missing_frame_resends_on_loss() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + // Four of 10..=14 arrived; 12 is the missing one. + assert_eq!( + t.check(link(E1, 15, rr(14, 14, 0)), 1_000), + Some(ResendReason::Loss) + ); + assert!(t.state.needs_update(&t.peer), "the peer must be marked"); +} + +/// A covering report with every frame since the base confirms. +#[test] +fn test_bloom_ack_covering_report_with_every_frame_confirms() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + assert_eq!(t.check(link(E1, 15, rr(14, 15, 0)), 1_000), None); + assert!(!t.outstanding(), "the announce must be confirmed"); + assert!( + !t.state.needs_update(&t.peer), + "the peer must not be marked" + ); +} + +/// A late frame from before the base cannot stand in for the lost +/// announce. The base has a hole, so the pair cannot tell that late frame +/// from an in-window reorder: no evidence, and the fallback resends. +#[test] +fn test_bloom_ack_reordered_frame_does_not_mask_a_lost_announce() { + let mut t = Track::new(); + // Base (9, 9, 0): frame 5 missing at the time. + t.send(12, link(E1, 13, rr(9, 9, 0)), 0); + // Late frame 5 plus 10, 11, 13, 14; 12 lost. A naive count gives 5 == 5. + let late = |_, c, recording| link(E1, c + if recording { 2 } else { 3 }, rr(14, 14, 1)); + assert_eq!(t.check(late(1_000, 12, false), 1_000), None); + assert!(t.outstanding(), "a lost announce must not confirm"); + let resends = t.hold(2_000, 120_000, false, late); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// A report that does not yet cover the announce decides nothing, and the +/// next one is measured against the base taken at the send. +#[test] +fn test_bloom_ack_report_below_the_announce_waits_for_a_covering_one() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + assert_eq!(t.check(link(E1, 13, rr(11, 12, 0)), 1_000), None); + assert!(t.outstanding(), "an uncovered announce stays outstanding"); + assert_eq!(t.check(link(E1, 15, rr(14, 15, 0)), 2_000), None); + assert!(!t.outstanding(), "the covering report must confirm"); +} + +/// An announce from an earlier session is resent once the new session can +/// check the resend, or once after the fallback if it never can. +#[test] +fn test_bloom_ack_announce_from_an_earlier_session_is_resent_once() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + assert_eq!(t.check(link(E2, 1, None), 1_000), None); + assert_eq!( + t.check(link(E2, 5, rr(3, 4, 0)), 2_000), + Some(ResendReason::SessionChanged) + ); + t.send(5, link(E2, 6, rr(3, 4, 0)), 2_000); + assert_eq!(t.check(link(E2, 7, rr(6, 7, 0)), 3_000), None); + assert!(!t.outstanding(), "the resend must be confirmed"); + + // No usable report in the new session: one Timeout, 30 s after the change + // was seen, and no second. + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + assert_eq!(t.check(link(E2, 1, None), 1_000), None); + let resends = t.hold(2_000, 121_000, false, |_, c, recording| { + link(E2, c + if recording { 2 } else { 1 }, None) + }); + assert_eq!(resends, vec![(31_000, ResendReason::Timeout)]); +} + +/// In the peer's first session, with no report before the send, every +/// counter from 0 must arrive. +#[test] +fn test_bloom_ack_first_session_measures_from_counter_zero() { + let mut t = Track::new(); + t.send(12, link(E1, 13, None), 0); + assert_eq!(t.check(link(E1, 13, rr(12, 13, 0)), 1_000), None); + assert!(!t.outstanding(), "13 of 0..=12 must confirm"); + + let mut t = Track::new(); + t.send(12, link(E1, 13, None), 0); + assert_eq!( + t.check(link(E1, 13, rr(12, 12, 0)), 1_000), + Some(ResendReason::Loss) + ); +} + +/// In a later session, the peer's cumulative count includes earlier +/// sessions, so no report means no base: a first report below the announce +/// becomes the base, and one already covering it cannot check it. +#[test] +fn test_bloom_ack_later_session_without_a_report_has_no_base() { + // Case A: re-based, then confirmed with no resend. + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + t.send(12, link(E2, 13, None), 1_000); + assert_eq!(t.check(link(E2, 13, rr(8, 509, 0)), 2_000), None); + assert!(t.outstanding()); + assert_eq!(t.check(link(E2, 15, rr(14, 515, 0)), 3_000), None); + assert!(!t.outstanding(), "the re-based announce must confirm"); + + // Case B: the first usable report already covers the announce. + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + t.send(12, link(E2, 13, None), 1_000); + assert_eq!( + t.check(link(E2, 15, rr(14, 515, 0)), 2_000), + Some(ResendReason::Unverified) + ); +} + +/// A report from another session at send time is no evidence, and the +/// announce gets exactly one fallback resend. +#[test] +fn test_bloom_ack_report_from_another_session_at_send_gets_one_timeout() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(20, 30, 0)), 0); + let resends = t.hold(1_000, 120_000, false, |_, c, recording| { + link(E1, c + if recording { 2 } else { 1 }, rr(20, 30, 0)) + }); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// With no report ever, one fallback resend per session. +#[test] +fn test_bloom_ack_no_report_ever_resends_once_per_session() { + let mut t = Track::new(); + t.send(12, link(E1, 13, None), 0); + assert_eq!(t.check(link(E1, 13, None), FALLBACK_MS - 1), None); + let quiet = |_, c, recording| link(E1, c + if recording { 2 } else { 1 }, None); + let resends = t.hold(FALLBACK_MS, 120_000, false, quiet); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); + + assert_eq!(t.check(link(E2, 1, None), 121_000), None); + let later = |_, c, recording| link(E2, c + if recording { 2 } else { 1 }, None); + let resends = t.hold(122_000, 240_000, false, later); + assert_eq!(resends, vec![(151_000, ResendReason::Timeout)]); +} + +/// A trigger that never stops is spaced 1, 2, 4 ... s apart up to 60 s. +#[test] +fn test_bloom_ack_sustained_trigger_backs_off_to_one_resend_a_minute() { + let mut t = Track::new(); + t.send(12, lossy(0, 11, true), 0); + let resends = t.hold(0, 600_000, true, lossy); + let expected: Vec = [0, 1, 3, 7, 15, 31, 63] + .iter() + .map(|s| s * 1_000) + .chain((123_000..=600_000).step_by(60_000)) + .collect(); + assert_eq!(times(&resends), expected); + for (i, &(start, _)) in resends.iter().enumerate() { + let in_window = resends[i..] + .iter() + .take_while(|(t, _)| *t < start + 60_000) + .count(); + assert!(in_window <= 6, "{in_window} resends in 60 s from {start}"); + } +} + +/// 120 s with no resend resets the backoff; 119 s does not. +#[test] +fn test_bloom_ack_backoff_resets_only_after_the_quiet_period() { + let run = |next_ms: u64| { + let mut t = Track::new(); + t.send(12, lossy(0, 11, true), 0); + let first = t.hold(0, 31_000, true, lossy); + assert_eq!(times(&first), vec![0, 1_000, 3_000, 7_000, 15_000, 31_000]); + t.hold(next_ms, next_ms + 70_000, true, lossy) + }; + // Case A: 120 s after the last resend, the level resets to 0. + let a = run(31_000 + QUIET_MS); + assert_eq!(times(&a[..2]), vec![151_000, 152_000]); + // Case B: 119 s after, level 6 still applies and reaches level 7. + let b = run(31_000 + QUIET_MS - 1_000); + assert_eq!(times(&b[..2]), vec![150_000, 210_000]); +} + +/// A report that went backwards within a session is no evidence +/// (defensive: the MMP layer never stores one). +#[test] +fn test_bloom_ack_regressed_report_is_no_evidence() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + let regressed = |_, c, recording| link(E1, c + if recording { 2 } else { 3 }, rr(14, 8, 0)); + assert_eq!(t.check(regressed(1_000, 12, false), 1_000), None); + assert!(t.outstanding()); + let resends = t.hold(2_000, FALLBACK_MS, false, regressed); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// Removing the peer forgets its announce, and the next send starts a +/// fresh entry whose first session measures from counter zero. +#[test] +fn test_bloom_ack_removed_peer_starts_fresh() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + t.state.remove_peer_state(&t.peer); + assert_eq!(t.check(link(E1, 15, rr(14, 14, 0)), 1_000), None); + assert!(!t.outstanding()); + + t.send(3, link(E2, 4, None), 2_000); + assert_eq!(t.check(link(E2, 4, rr(3, 4, 0)), 3_000), None); + assert!(!t.outstanding(), "the new entry must measure from zero"); +} + +/// A confirmation does not reset the backoff. +#[test] +fn test_bloom_ack_confirmation_keeps_the_backoff() { + let mut t = Track::new(); + t.send(12, lossy(0, 11, true), 0); + assert_eq!(t.check(lossy(0, 12, false), 0), Some(ResendReason::Loss)); + t.send(13, lossy(0, 12, true), 0); + assert_eq!(t.check(link(E1, 15, rr(14, 15, 0)), 100), None); + assert!(!t.outstanding(), "setup: the resend is confirmed"); + + t.send_new(15, lossy(200, 14, true), 200); + assert_eq!(t.check(lossy(500, 15, false), 500), None); + assert_eq!( + t.check(lossy(1_000, 15, false), 1_000), + Some(ResendReason::Loss) + ); +} + +/// The initiator holds a report from the responder's view of the old +/// session, frozen until the new session's counter passes it. It never +/// triggers more than each announce's one unchecked resend. +#[test] +fn test_bloom_ack_frozen_report_after_a_rekey_spends_only_the_unchecked_budget() { + // The session sends 100 frames a second from counter 1. + let next = |now: u64| 1 + now / 10; + let frozen = rr(5_000, 90_000, 40); + let accepted = rr(6_100, 96_101, 45); + let report = move |now: u64| if now < 62_000 { frozen } else { accepted }; + let model = move |now: u64, _c: u64, recording: bool| { + link(E2, next(now) + u64::from(recording), report(now)) + }; + + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + t.send(3, link(E2, 4, frozen), 0); + let first = t.hold(1_000, 59_000, false, model); + assert_eq!(first, vec![(FALLBACK_MS, ResendReason::Timeout)]); + + // A new-content announce based on the now-usable frozen report. + t.send_new(6_000, link(E2, 6_001, frozen), 60_000); + let second = t.hold(61_000, 120_000, false, model); + assert_eq!(second, vec![(90_000, ResendReason::Timeout)]); + + let third = t.hold(121_000, 140_000, false, |now, c, recording| { + let n = 10 + (now - 121_000) / 1_000 + u64::from(recording); + link(E3, n.max(c + 1), rr(5, 6, 0)) + }); + assert_eq!(third, vec![(121_000, ResendReason::SessionChanged)]); +} + +/// The responder receives reports whose highest counter comes from its +/// own previous session. They change on every report and are never usable, +/// so they trigger nothing beyond the one fallback resend. +#[test] +fn test_bloom_ack_polluted_report_after_a_rekey_spends_only_the_unchecked_budget() { + let model = |now: u64, c: u64, recording: bool| { + let secs = now / 1_000; + let n = (1 + 10 * secs).max(c + 1) + u64::from(recording); + link(E2, n, rr(9_000, 100 + 5 * secs, 5 * secs as u32)) + }; + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + t.send(5, model(0, 4, true), 0); + let resends = t.hold(1_000, 120_000, false, model); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// More receipts than counters means the reports straddle a reset or +/// another session's frame, which is no evidence. +#[test] +fn test_bloom_ack_surplus_receipts_are_no_evidence() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + let surplus = |_, c, recording| link(E1, c + if recording { 2 } else { 3 }, rr(14, 20, 0)); + assert_eq!(t.check(surplus(1_000, 12, false), 1_000), None); + assert!(t.outstanding()); + let resends = t.hold(2_000, FALLBACK_MS, false, surplus); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// One lineage in one session gets three loss resends and one unchecked +/// resend; new content or a new session refills both. +#[test] +fn test_bloom_ack_budgets_bound_resends_per_lineage_per_session() { + let spend = || { + let mut t = Track::new(); + t.send(12, lossy(0, 11, true), 0); + let resends = t.hold(0, 120_000, false, lossy); + assert_eq!( + resends, + vec![ + (0, ResendReason::Loss), + (1_000, ResendReason::Loss), + (3_000, ResendReason::Loss), + (33_000, ResendReason::Timeout), + ] + ); + assert_eq!(LOSS_BUDGET, 3); + assert_eq!(UNVERIFIED_BUDGET, 1); + t + }; + + let mut t = spend(); + let c = t.counter; + t.send_new(c + 1, lossy(121_000, c, true), 121_000); + assert_eq!( + t.check(lossy(122_000, t.counter, false), 122_000), + Some(ResendReason::Loss), + "new content must refill the loss budget" + ); + + let mut t = spend(); + let c = t.counter; + let rekeyed = |ev: LinkEvidence| link(E2, ev.next_counter, ev.rr); + t.send(c + 1, rekeyed(lossy(121_000, c, true)), 121_000); + assert_eq!( + t.check(rekeyed(lossy(122_000, t.counter, false)), 122_000), + Some(ResendReason::Loss), + "a new session must refill the loss budget" + ); +} + +/// A resend carries the content it repeats, so it spends from the same +/// lineage's budget instead of starting a new one. +#[test] +fn test_bloom_ack_resend_of_the_same_content_is_not_a_new_lineage() { + let mut t = Track::new(); + t.send(12, lossy(0, 11, true), 0); + let resends = t.hold(0, 7_000, false, lossy); + assert_eq!(times(&resends), vec![0, 1_000, 3_000]); + assert_eq!( + t.check(lossy(8_000, t.counter, false), 8_000), + None, + "a fourth loss resend must not be allowed" + ); +} + +/// A report polluted after the base was taken is unusable, not a loss. +#[test] +fn test_bloom_ack_report_polluted_after_the_base_is_not_a_loss() { + let mut t = Track::new(); + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + assert_eq!(t.check(link(E1, 20, rr(9_000, 16, 0)), 1_000), None); + assert!(t.outstanding()); +} + +/// A frame inside the checked window that arrives after a higher counter is +/// counted as a reorder, but it did arrive. On a base with no holes (a +/// first-session report that counted every frame up to its highest), every +/// counter in the window arriving confirms the announce. +#[test] +fn test_bloom_ack_in_window_reorder_on_a_complete_base_confirms() { + let mut t = Track::new(); + // Base (9, 10, 0): all of 0..=9 counted. + t.send(12, link(E1, 13, rr(9, 10, 0)), 0); + // 10..=14 all arrived, 12 after 13. + assert_eq!(t.check(link(E1, 15, rr(14, 15, 1)), 1_000), None); + assert!( + !t.outstanding(), + "every counter arrived, so it must confirm" + ); + assert!( + !t.state.needs_update(&t.peer), + "the peer must not be marked" + ); +} + +/// In the peer's first session with no report before the send, frames that +/// arrive out of order are still every counter from 0, so the announce +/// confirms. +#[test] +fn test_bloom_ack_in_window_reorder_on_the_zero_base_confirms() { + let mut t = Track::new(); + t.send(4, link(E1, 5, None), 0); + // 0..=4 all arrived, as 0, 1, 4, 3, 2. + assert_eq!(t.check(link(E1, 5, rr(4, 5, 2)), 1_000), None); + assert!( + !t.outstanding(), + "every counter arrived, so it must confirm" + ); + assert!( + !t.state.needs_update(&t.peer), + "the peer must not be marked" + ); +} + +/// In a later session the base cannot be shown to have no holes, so a +/// covering report with an in-window reorder cannot tell a late frame from +/// before the base from one inside the window. That is no evidence: no loss +/// resend, and the fallback covers the announce. +#[test] +fn test_bloom_ack_ambiguous_pair_in_a_later_session_waits_for_the_fallback() { + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + // Cumulative counts include 500 frames of the earlier session. + t.send(12, link(E2, 13, rr(9, 510, 0)), 0); + let ambiguous = |_, c, recording| link(E2, c + if recording { 2 } else { 3 }, rr(14, 515, 1)); + assert_eq!(t.check(ambiguous(1_000, 12, false), 1_000), None); + assert!(t.outstanding(), "an ambiguous pair must not confirm"); + let resends = t.hold(2_000, 120_000, false, ambiguous); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} + +/// A later-session base whose counts happen to read as having no holes is +/// still not trusted: the cumulative count includes earlier sessions. A late +/// frame from before the base arriving with the announce lost must not +/// confirm it. +#[test] +fn test_bloom_ack_later_session_base_is_never_complete() { + let mut t = Track::new(); + t.send(3, link(E1, 4, None), 0); + // Base (9, 10, 0) in E2: 3 frames of E1 plus 7 of 0..=9, with 5 missing. + t.send(12, link(E2, 13, rr(9, 10, 0)), 0); + // Late frame 5 plus 10, 11, 13, 14; 12 lost. Received rose by 5 == span. + let late = |_, c, recording| link(E2, c + if recording { 2 } else { 3 }, rr(14, 15, 1)); + assert_eq!(t.check(late(1_000, 12, false), 1_000), None); + assert!(t.outstanding(), "a lost announce must not confirm"); + let resends = t.hold(2_000, 120_000, false, late); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); +} diff --git a/src/proto/mmp/metrics.rs b/src/proto/mmp/metrics.rs index 0803ce2e..a7b7146f 100644 --- a/src/proto/mmp/metrics.rs +++ b/src/proto/mmp/metrics.rs @@ -360,6 +360,17 @@ impl MmpMetrics { self.prev_rr_ecn_ce } + /// Cumulative counters of the last accepted ReceiverReport: highest counter, + /// packets received and reorder count. `None` until a report is accepted in + /// the current session. + pub fn rr_counters(&self) -> Option<(u64, u64, u32)> { + self.has_prev_rr.then_some(( + self.prev_rr_highest_counter, + self.prev_rr_cum_packets, + self.prev_rr_reorder, + )) + } + /// ReceiverReports processed, including stale and duplicate ones. pub fn reports_seen(&self) -> u64 { self.reports_seen