Track whether each bloom filter announce reached its peer

A filter announce is recorded as sent once the transport accepts the
frame, which on UDP and Ethernet says nothing about delivery. The bloom
state now keeps each peer's last announce outstanding, with its link
counter and session, and decides from the peer's ordinary receiver
reports whether it arrived, by comparing the report that covers the
announce with the report taken before the send, its base.

A loss is concluded only when fewer frames arrived between the two
reports than the window holds. Delivery is confirmed when the frames that
arrived in counter order cover the window, or when all arrivals do and
the base had no holes: nothing reported yet in the peer's first session,
or a first-session report that counted every counter up to its highest.
A frame counted as a reorder may be one of the window's frames or a late
frame from before the base, so any other pair is no evidence and the
announce stays outstanding. A report from another session, recognisable
because its highest counter is at or above the next counter the current
session will use, is ignored, and an inconsistent pair of reports proves
nothing. An announce the reports cannot check is resent once, as soon as
a usable report arrives or after 30 s; a usable report that does not yet
cover an announce sent in the current session becomes its base instead.
After a rekey the peer's cumulative count includes earlier sessions, so
no base can be shown to have no holes, and an in-window reorder there
costs one fallback resend.

Resends are bounded per announce and per peer: at most one unchecked
and three loss resends per announce per session, and a per-peer backoff
of 1, 2, 4 ... s up to 60 s that resets only after 120 s without a
resend. MMP metrics gain a read-only accessor for the last accepted
report's cumulative counters.
This commit is contained in:
Johnathan Corgan
2026-09-26 18:59:43 +00:00
parent 16c4d42b2f
commit 26be2b7283
4 changed files with 958 additions and 1 deletions
+1 -1
View File
@@ -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.
+375
View File
@@ -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<RrCounters>,
}
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<RrCounters> {
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<SentAnnounce>,
/// Backoff level: the number of resends in the current run, capped.
level: u8,
/// When the last resend was triggered.
resent_ms: Option<u64>,
/// 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<bool> {
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<NodeAddr, BloomFilter>,
/// How long an unchecked announce waits for its fallback resend (ms).
fallback_ms: u64,
/// Per-peer delivery tracking for sent announces.
acks: BTreeMap<NodeAddr, AckState>,
}
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<ResendReason> {
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<u64> {
self.acks.get(peer)?.sent.map(|sent| sent.counter)
}
/// Mark only peers whose outgoing filter has actually changed.
+571
View File
@@ -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<RrCounters> {
Some(RrCounters {
highest,
received,
reordered,
})
}
/// Link evidence for session `epoch`.
fn link(epoch: u64, next_counter: u64, rr: Option<RrCounters>) -> 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<ResendReason> {
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<u64> {
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<u64> = [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)]);
}
+11
View File
@@ -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