mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
Resend a bloom filter announce the peer did not receive
A node counted a FilterAnnounce as delivered as soon as the transport accepted it, and announces go out only when a filter changes. A dropped datagram, or a link outage shorter than the dead timeout, therefore left the peer holding the old filter until something else changed, so destinations could stay missing from its discovery indefinitely. Each tick now checks every outstanding announce against the link's existing receiver reports before sending pending announces, and marks the peer again when a report covering the announce shows fewer frames arrived than were sent, or when the announce cannot be confirmed: once, as soon as a usable report arrives, or after 30 s. A report that shows neither delivery nor a loss, as when a reordered frame could be a late one from before the announce, waits for the 30 s fallback. The resend goes out through the ordinary debounced send path in the same tick. Reports from another session are ignored, at most one unchecked and three loss resends are sent per announce per session, and a peer connection sees at most six resends a minute, and one a minute while losses persist. Nothing on the wire changes. The design notes describe the resend and its limits.
This commit is contained in:
@@ -201,7 +201,7 @@ there is no incremental update.
|
||||
|
||||
### Update Triggers
|
||||
|
||||
Filter updates are event-driven, not periodic:
|
||||
New filter content is sent only on these events, never periodically:
|
||||
|
||||
- Peer connects (new filter includes the new peer's reachability)
|
||||
- Peer disconnects (filter must exclude the departed peer's entries)
|
||||
@@ -209,6 +209,46 @@ Filter updates are event-driven, not periodic:
|
||||
be recomputed)
|
||||
- Local state changes (new identity, leaf-only dependent changes)
|
||||
|
||||
The one timed send is the resend of an announce that was not confirmed
|
||||
delivered. The transport accepting a FilterAnnounce does not mean the peer
|
||||
received it: a datagram can be lost, or a link can go down for less than
|
||||
the dead timeout without the peer being removed. Each announce therefore
|
||||
stays outstanding until the link's ordinary MMP ReceiverReports show that
|
||||
every link frame up to and including the announce's counter arrived. No
|
||||
new message or field is involved.
|
||||
|
||||
A report covering the announce is compared with the report the announce
|
||||
was sent after, its base. Between them the peer counted `got` frames, of
|
||||
which `sure` arrived in counter order. A frame that arrived after a higher
|
||||
counter is a reorder, and it can be one of the window's frames or a late
|
||||
frame from before the base, so the frames of the window number between
|
||||
`sure` and `got`. The announce is lost when `got` is below the number of
|
||||
counters the window holds, and delivered when `sure` equals it, or when
|
||||
`got` equals it and the base has no holes. A base has no holes when
|
||||
nothing had been reported yet in the peer's first session, or when it is a
|
||||
first-session report that counted exactly its highest counter plus one
|
||||
frame. After a rekey the peer's cumulative count includes earlier
|
||||
sessions, so no base can be shown to have no holes. Any other pair is no
|
||||
evidence, and the 30 s fallback below covers the announce.
|
||||
|
||||
An unconfirmed announce is resent:
|
||||
|
||||
- on a report that covers it and shows fewer frames arrived than were
|
||||
sent;
|
||||
- once, when the reports cannot check it (sent before the session's first
|
||||
usable report, carried over a rekey, or a report left over from the
|
||||
previous session), as soon as a usable report arrives, or after 30 s if
|
||||
none does. A usable report that does not yet cover an announce sent in
|
||||
the current session becomes the base it is checked against instead.
|
||||
|
||||
A report whose highest counter is at or above the next counter the current
|
||||
session will use describes another session and is ignored. Each announce
|
||||
gets at most one unchecked resend and three loss resends per session, and
|
||||
a per-peer backoff spaces resends 1, 2, 4 ... s apart up to 60 s until
|
||||
120 s pass with none, so a peer connection sees at most six resends in a
|
||||
minute and one a minute while losses persist. A resend of content the
|
||||
peer already holds changes nothing there, so it propagates no further.
|
||||
|
||||
### Rate Limiting
|
||||
|
||||
Updates are rate-limited at a 500ms minimum interval per peer
|
||||
@@ -241,7 +281,9 @@ property of the data structure). Entries are expired through:
|
||||
- **Implicit timeout**: If a peer becomes unresponsive, the MMP link
|
||||
liveness detector eventually declares the link dead and removes the
|
||||
peer, which triggers filter cleanup as a side effect of peer removal.
|
||||
There is no independent filter staleness timer.
|
||||
There is no independent filter staleness timer; the only timer is the
|
||||
30 s fallback that resends an announce the receiver reports cannot
|
||||
confirm (see Update Triggers).
|
||||
|
||||
## Membership Test
|
||||
|
||||
|
||||
+67
-1
@@ -6,6 +6,7 @@
|
||||
use crate::NodeAddr;
|
||||
use crate::proto::bloom::BloomFilter;
|
||||
use crate::proto::bloom::FilterAnnounce;
|
||||
use crate::proto::bloom::{LinkEvidence, RrCounters};
|
||||
|
||||
use super::reject::BloomReject;
|
||||
use super::{Node, NodeError};
|
||||
@@ -71,6 +72,20 @@ impl Node {
|
||||
|
||||
self.metrics().bloom.sent.inc();
|
||||
|
||||
// Read after the send: anything else taking a counter in between only
|
||||
// makes the recorded counter higher, which delays confirmation rather
|
||||
// than confirming a frame that was never covered.
|
||||
if let Some(link) = self.link_evidence(peer_addr) {
|
||||
let counter = link.next_counter.saturating_sub(1);
|
||||
self.bloom_state.record_announce(
|
||||
*peer_addr,
|
||||
&sent_filter,
|
||||
counter,
|
||||
&link,
|
||||
crate::time::mono_ms(),
|
||||
);
|
||||
}
|
||||
|
||||
// Self-plausibility check: WARN if our own outgoing filter is
|
||||
// above the antipoison cap. Independent detection signal if
|
||||
// aggregation drift or an ingress-check bypass pushes us over
|
||||
@@ -271,10 +286,61 @@ impl Node {
|
||||
.mark_changed_peers(from, &peer_addrs, &peer_filters);
|
||||
}
|
||||
|
||||
/// Read what `peer_addr`'s link shows about delivery of our frames: the
|
||||
/// session's identity and next send counter, and the last ReceiverReport
|
||||
/// accepted on it. `None` when the peer has no session.
|
||||
fn link_evidence(&self, peer_addr: &NodeAddr) -> Option<LinkEvidence> {
|
||||
let peer = self.peers.get(peer_addr)?;
|
||||
let session = peer.noise_session()?;
|
||||
let mut epoch = [0u8; 8];
|
||||
epoch.copy_from_slice(&session.handshake_hash()[..8]);
|
||||
let rr = peer.mmp().and_then(|mmp| mmp.metrics.rr_counters()).map(
|
||||
|(highest, received, reordered)| RrCounters {
|
||||
highest,
|
||||
received,
|
||||
reordered,
|
||||
},
|
||||
);
|
||||
Some(LinkEvidence {
|
||||
epoch: u64::from_le_bytes(epoch),
|
||||
next_counter: session.current_send_counter(),
|
||||
rr,
|
||||
})
|
||||
}
|
||||
|
||||
/// Mark for resend every peer whose outstanding announce the receiver
|
||||
/// reports show was lost, or could not confirm in time.
|
||||
fn check_announces(&mut self) {
|
||||
let now_ms = crate::time::mono_ms();
|
||||
let waiting: Vec<NodeAddr> = self
|
||||
.peers
|
||||
.keys()
|
||||
.filter(|addr| self.bloom_state.announce_outstanding(addr))
|
||||
.copied()
|
||||
.collect();
|
||||
for addr in waiting {
|
||||
let Some(link) = self.link_evidence(&addr) else {
|
||||
continue;
|
||||
};
|
||||
let counter = self.bloom_state.outstanding_counter(&addr);
|
||||
if let Some(reason) = self.bloom_state.check_announce(&addr, &link, now_ms) {
|
||||
debug!(
|
||||
peer = %self.peer_display_name(&addr),
|
||||
reason = ?reason,
|
||||
counter = ?counter,
|
||||
"Resending unconfirmed FilterAnnounce"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Check bloom filter state on tick (called from event loop).
|
||||
///
|
||||
/// Sends any pending debounced filter announces.
|
||||
/// Marks peers whose last announce was not confirmed delivered, then sends
|
||||
/// any pending debounced filter announces, so a resend goes out in the
|
||||
/// same tick through the ordinary send path.
|
||||
pub(super) async fn check_bloom_state(&mut self) {
|
||||
self.check_announces();
|
||||
self.send_pending_filter_announces().await;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -988,3 +988,589 @@ async fn test_bloom_tree_announce_without_tree_peer_flip_marks_no_peer() {
|
||||
assert!(!bloom.needs_update(&c), "C must not be marked");
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
// ===== Resend of a filter announce the peer did not receive =====
|
||||
//
|
||||
// A lost datagram is made by taking M's frame out of P's receive channel
|
||||
// without processing it: the transport returned `Ok` and the receiver never
|
||||
// saw the frame, which is the shape of a real loss on UDP or Ethernet. Final
|
||||
// assertions read the filter P stores for M, never what M believes it sent.
|
||||
|
||||
/// Index of P in `FlipFixture::nodes`.
|
||||
const P: usize = 0;
|
||||
|
||||
/// Set every link report interval on `node` to zero, so the next
|
||||
/// `check_mmp_reports` sends each report that has interval data.
|
||||
///
|
||||
/// Processing a ReceiverReport re-derives the intervals from SRTT, so callers
|
||||
/// re-apply this before every `check_mmp_reports`.
|
||||
fn zero_intervals(node: &mut Node) {
|
||||
for peer in node.peers.values_mut() {
|
||||
if let Some(mmp) = peer.mmp_mut() {
|
||||
mmp.sender.update_report_interval_with_bounds(1_000, 0, 0);
|
||||
mmp.receiver.update_report_interval_with_bounds(1_000, 0, 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// One MMP exchange between M and P: M reports, P processes, P reports, M
|
||||
/// processes. C is never asked to report.
|
||||
async fn mmp_round(nodes: &mut [TestNode]) {
|
||||
zero_intervals(&mut nodes[M].node);
|
||||
zero_intervals(&mut nodes[P].node);
|
||||
nodes[M].node.check_mmp_reports().await;
|
||||
process_available_packets(nodes).await;
|
||||
zero_intervals(&mut nodes[M].node);
|
||||
zero_intervals(&mut nodes[P].node);
|
||||
nodes[P].node.check_mmp_reports().await;
|
||||
process_available_packets(nodes).await;
|
||||
}
|
||||
|
||||
/// Process packets on every node until a pass handles none, at most 50 passes.
|
||||
async fn drain_quiet(nodes: &mut [TestNode]) {
|
||||
for _ in 0..50 {
|
||||
if process_available_packets(nodes).await == 0 {
|
||||
return;
|
||||
}
|
||||
}
|
||||
panic!("setup: packets still flowing after 50 passes");
|
||||
}
|
||||
|
||||
/// Wait at most 1 s for `tn` to hold a queued frame, then take every queued
|
||||
/// frame without processing it. Returns how many were taken.
|
||||
async fn drop_queued(tn: &mut TestNode) -> usize {
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(1);
|
||||
while tn.packet_rx.is_empty() && std::time::Instant::now() < deadline {
|
||||
tokio::time::sleep(Duration::from_millis(5)).await;
|
||||
}
|
||||
let mut dropped = 0;
|
||||
while tn.packet_rx.try_recv().is_ok() {
|
||||
dropped += 1;
|
||||
}
|
||||
dropped
|
||||
}
|
||||
|
||||
/// ReceiverReports M has seen from P, including stale and duplicate ones.
|
||||
fn reports_seen(fx: &FlipFixture) -> u64 {
|
||||
fx.nodes[M]
|
||||
.node
|
||||
.get_peer(&fx.p)
|
||||
.and_then(|peer| peer.mmp())
|
||||
.map_or(0, |mmp| mmp.metrics.reports_seen())
|
||||
}
|
||||
|
||||
/// Whether the filter P stores for M contains the marker.
|
||||
fn holds_marker(fx: &FlipFixture) -> bool {
|
||||
fx.nodes[P]
|
||||
.node
|
||||
.get_peer(&fx.m)
|
||||
.and_then(|peer| peer.inbound_filter())
|
||||
.is_some_and(|filter| filter.contains(&marker()))
|
||||
}
|
||||
|
||||
/// Parent-switch counts at M and P before a run of MMP rounds.
|
||||
///
|
||||
/// A first RTT sample can re-evaluate the parent, and a switch marks every
|
||||
/// peer, which would pass a resend test for a reason unrelated to the resend.
|
||||
struct SwitchGuard {
|
||||
m: u64,
|
||||
p: u64,
|
||||
}
|
||||
|
||||
/// Snapshot M's and P's parent-switch counters.
|
||||
fn switch_guard(fx: &FlipFixture) -> SwitchGuard {
|
||||
SwitchGuard {
|
||||
m: fx.nodes[M].node.metrics().tree.parent_switches.get(),
|
||||
p: fx.nodes[P].node.metrics().tree.parent_switches.get(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Assert neither M nor P switched parent since `guard`, and M's parent is P.
|
||||
fn assert_unswitched(fx: &FlipFixture, guard: &SwitchGuard) {
|
||||
assert_eq!(
|
||||
fx.nodes[M].node.metrics().tree.parent_switches.get(),
|
||||
guard.m,
|
||||
"setup: M must not switch parent during the MMP rounds"
|
||||
);
|
||||
assert_eq!(
|
||||
fx.nodes[P].node.metrics().tree.parent_switches.get(),
|
||||
guard.p,
|
||||
"setup: P must not switch parent during the MMP rounds"
|
||||
);
|
||||
assert_eq!(
|
||||
fx.nodes[M].node.tree_state().my_declaration().parent_id(),
|
||||
&fx.p,
|
||||
"setup: M's parent must still be P"
|
||||
);
|
||||
}
|
||||
|
||||
/// FilterAnnounces M has sent.
|
||||
fn sent_count(fx: &FlipFixture) -> u64 {
|
||||
fx.nodes[M].node.metrics().bloom.sent.get()
|
||||
}
|
||||
|
||||
/// Drain the fixture, then run MMP rounds until M has seen a report from P.
|
||||
async fn start_reports(fx: &mut FlipFixture) {
|
||||
drain_quiet(&mut fx.nodes).await;
|
||||
for _ in 0..10 {
|
||||
if reports_seen(fx) >= 1 {
|
||||
break;
|
||||
}
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
}
|
||||
assert!(reports_seen(fx) >= 1, "setup: P must report to M");
|
||||
}
|
||||
|
||||
/// Deliver C's filter carrying the marker and send M's announce of it to P.
|
||||
/// Returns M's sent count after the send.
|
||||
async fn send_marker(fx: &mut FlipFixture) -> u64 {
|
||||
let c = fx.c;
|
||||
deliver_filter(fx, &[c, marker()]).await;
|
||||
let before = sent_count(fx);
|
||||
fx.nodes[M].node.send_pending_filter_announces().await;
|
||||
let after = sent_count(fx);
|
||||
assert_eq!(after, before + 1, "setup: M must send exactly one announce");
|
||||
after
|
||||
}
|
||||
|
||||
/// Lose the announce just sent to P, and check the loss took.
|
||||
async fn lose_announce(fx: &mut FlipFixture) {
|
||||
assert_eq!(
|
||||
drop_queued(&mut fx.nodes[P]).await,
|
||||
1,
|
||||
"setup: exactly the one announce frame must be lost"
|
||||
);
|
||||
assert!(
|
||||
!holds_marker(fx),
|
||||
"control: P must not hold the lost announce's content"
|
||||
);
|
||||
assert!(
|
||||
sent_to_parent(fx).contains(&marker()),
|
||||
"control: M must record the lost announce as sent"
|
||||
);
|
||||
}
|
||||
|
||||
/// Get reports flowing from P to M, then send M's announce carrying the
|
||||
/// marker and lose it on the way to P. Returns M's sent count after the send.
|
||||
async fn lose_marker(fx: &mut FlipFixture) -> u64 {
|
||||
start_reports(fx).await;
|
||||
let sent = send_marker(fx).await;
|
||||
lose_announce(fx).await;
|
||||
sent
|
||||
}
|
||||
|
||||
/// The first eight bytes of the handshake hash of M's current session with P.
|
||||
fn link_epoch(fx: &FlipFixture) -> [u8; 8] {
|
||||
let hash = fx.nodes[M]
|
||||
.node
|
||||
.get_peer(&fx.p)
|
||||
.and_then(|peer| peer.noise_session())
|
||||
.expect("M has a session with P")
|
||||
.handshake_hash();
|
||||
let mut epoch = [0u8; 8];
|
||||
epoch.copy_from_slice(&hash[..8]);
|
||||
epoch
|
||||
}
|
||||
|
||||
/// A FilterAnnounce lost in transit is resent once a receiver report shows
|
||||
/// the loss, so the peer ends up holding the filter.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_filter_announce_lost_in_transit_reaches_the_peer_after_a_receiver_report() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
lose_marker(&mut fx).await;
|
||||
|
||||
let seen = reports_seen(&fx);
|
||||
let guard = switch_guard(&fx);
|
||||
for _ in 0..3 {
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
fx.nodes[M].node.check_bloom_state().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
}
|
||||
assert!(
|
||||
reports_seen(&fx) > seen,
|
||||
"setup: a receiver report must arrive after the loss"
|
||||
);
|
||||
assert_unswitched(&fx, &guard);
|
||||
|
||||
assert!(
|
||||
holds_marker(&fx),
|
||||
"P must hold the filter whose announce was lost"
|
||||
);
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
/// Receiving the same filter again under a newer sequence changes no outgoing
|
||||
/// filter, so it marks no peer and cannot cascade.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_unchanged_filter_with_newer_sequence_marks_no_peer() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
let (c, p) = (fx.c, fx.p);
|
||||
|
||||
deliver_filter(&mut fx, &[c, marker()]).await;
|
||||
fx.nodes[M].node.send_pending_filter_announces().await;
|
||||
let bloom = &fx.nodes[M].node.bloom_state;
|
||||
assert!(
|
||||
!bloom.needs_update(&p) && !bloom.needs_update(&c),
|
||||
"control: the first delivery must be fully sent"
|
||||
);
|
||||
|
||||
deliver_filter(&mut fx, &[c, marker()]).await;
|
||||
let bloom = &fx.nodes[M].node.bloom_state;
|
||||
assert!(!bloom.needs_update(&p), "P must not be marked");
|
||||
assert!(!bloom.needs_update(&c), "C must not be marked");
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
/// Make M rekey on its next check: one message on a session is enough, time
|
||||
/// never triggers it, and both ends of M's links are aged past the
|
||||
/// responder's rekey-acceptance gate so both rekeys are ordinary ones.
|
||||
fn arm_rekey(fx: &mut FlipFixture) {
|
||||
fx.nodes[M].node.replace_context(|ctx| {
|
||||
let mut cfg = (*ctx.config).clone();
|
||||
cfg.node.rekey.enabled = true;
|
||||
cfg.node.rekey.after_messages = 1;
|
||||
cfg.node.rekey.after_secs = u64::MAX;
|
||||
ctx.config = std::sync::Arc::new(cfg);
|
||||
});
|
||||
let (m, p, c) = (fx.m, fx.p, fx.c);
|
||||
let age = Duration::from_secs(31);
|
||||
for (i, remote) in [(M, p), (P, m), (M, c), (C, m)] {
|
||||
fx.nodes[i]
|
||||
.node
|
||||
.get_peer_mut(&remote)
|
||||
.expect("setup: link peer present")
|
||||
.test_backdate_session_established(age);
|
||||
}
|
||||
}
|
||||
|
||||
/// Drive the real rekey handshake until M's session with P is cut over.
|
||||
async fn rekey_cutover(fx: &mut FlipFixture) {
|
||||
let before = link_epoch(fx);
|
||||
for _ in 0..6 {
|
||||
fx.nodes[M].node.check_rekey().await;
|
||||
fx.nodes[P].node.check_rekey().await;
|
||||
for _ in 0..3 {
|
||||
tokio::time::sleep(Duration::from_millis(5)).await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
}
|
||||
if link_epoch(fx) != before {
|
||||
break;
|
||||
}
|
||||
}
|
||||
assert_ne!(link_epoch(fx), before, "setup: M's link to P must rekey");
|
||||
let (m, p) = (fx.m, fx.p);
|
||||
assert!(
|
||||
!fx.nodes[M].node.get_peer(&p).unwrap().rekey_in_progress(),
|
||||
"setup: M's rekey with P must be complete"
|
||||
);
|
||||
assert!(
|
||||
!fx.nodes[P].node.get_peer(&m).unwrap().rekey_in_progress(),
|
||||
"setup: P's rekey with M must be complete"
|
||||
);
|
||||
}
|
||||
|
||||
/// An announce lost just before a link rekey is resent on the new session.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_announce_lost_before_a_link_rekey_reaches_the_peer_after_the_cutover() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
let sent = lose_marker(&mut fx).await;
|
||||
|
||||
arm_rekey(&mut fx);
|
||||
rekey_cutover(&mut fx).await;
|
||||
|
||||
let guard = switch_guard(&fx);
|
||||
for _ in 0..5 {
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
fx.nodes[M].node.check_bloom_state().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
}
|
||||
assert_unswitched(&fx, &guard);
|
||||
|
||||
assert_eq!(
|
||||
sent_count(&fx),
|
||||
sent + 1,
|
||||
"M must resend to P exactly once after the loss"
|
||||
);
|
||||
assert!(
|
||||
holds_marker(&fx),
|
||||
"P must hold the filter whose announce was lost before the rekey"
|
||||
);
|
||||
assert!(
|
||||
!fx.nodes[M].node.bloom_state.announce_outstanding(&fx.p),
|
||||
"M's resend on the new session must be confirmed"
|
||||
);
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
/// An announce that arrives is confirmed from the receiver reports and never
|
||||
/// resent.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_delivered_announce_is_confirmed_without_a_resend() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
let p = fx.p;
|
||||
start_reports(&mut fx).await;
|
||||
let sent = send_marker(&mut fx).await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
assert!(holds_marker(&fx), "control: P must hold the announce");
|
||||
|
||||
let guard = switch_guard(&fx);
|
||||
for _ in 0..3 {
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
fx.nodes[M].node.check_bloom_state().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
}
|
||||
assert_unswitched(&fx, &guard);
|
||||
|
||||
let bloom = &fx.nodes[M].node.bloom_state;
|
||||
assert_eq!(
|
||||
sent_count(&fx),
|
||||
sent,
|
||||
"M must not resend a delivered announce"
|
||||
);
|
||||
assert!(!bloom.needs_update(&p), "P must not be marked");
|
||||
assert!(
|
||||
!bloom.announce_outstanding(&p),
|
||||
"the delivered announce must be confirmed"
|
||||
);
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
/// On a converged mesh with clean links, every announce is confirmed from the
|
||||
/// receiver reports and none is resent. The convergence announces go out
|
||||
/// before any report, so this checks the zero baseline on real counters.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_clean_links_confirm_every_announce_without_a_resend() {
|
||||
let mut nodes = run_tree_test(3, &[(0, 1), (1, 2)], false).await;
|
||||
drain_quiet(&mut nodes).await;
|
||||
let snapshot = |nodes: &[TestNode]| -> Vec<(u64, u64, NodeAddr)> {
|
||||
nodes
|
||||
.iter()
|
||||
.map(|tn| {
|
||||
(
|
||||
tn.node.metrics().bloom.sent.get(),
|
||||
tn.node.metrics().tree.parent_switches.get(),
|
||||
*tn.node.tree_state().my_declaration().parent_id(),
|
||||
)
|
||||
})
|
||||
.collect()
|
||||
};
|
||||
let before = snapshot(&nodes);
|
||||
|
||||
for _ in 0..5 {
|
||||
for i in 0..nodes.len() {
|
||||
zero_intervals(&mut nodes[i].node);
|
||||
nodes[i].node.check_mmp_reports().await;
|
||||
process_available_packets(&mut nodes).await;
|
||||
}
|
||||
for tn in nodes.iter_mut() {
|
||||
tn.node.check_bloom_state().await;
|
||||
}
|
||||
process_available_packets(&mut nodes).await;
|
||||
}
|
||||
|
||||
assert_eq!(
|
||||
snapshot(&nodes),
|
||||
before,
|
||||
"no node may send an announce, switch parent or change parent"
|
||||
);
|
||||
for (i, tn) in nodes.iter().enumerate() {
|
||||
for peer in tn.node.peers.keys() {
|
||||
assert!(
|
||||
!tn.node.bloom_state.announce_outstanding(peer),
|
||||
"node {i} must have confirmed its announce to every peer"
|
||||
);
|
||||
}
|
||||
}
|
||||
cleanup_nodes(&mut nodes).await;
|
||||
}
|
||||
|
||||
/// With no receiver report at all, a lost announce is still resent once the
|
||||
/// fallback interval passes.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_lost_announce_is_resent_after_the_fallback_when_no_receiver_report_arrives() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
drain_quiet(&mut fx.nodes).await;
|
||||
fx.nodes[M].node.bloom_state.set_fallback(0);
|
||||
let seen = reports_seen(&fx);
|
||||
send_marker(&mut fx).await;
|
||||
lose_announce(&mut fx).await;
|
||||
|
||||
fx.nodes[M].node.check_bloom_state().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
|
||||
assert_eq!(reports_seen(&fx), seen, "setup: P must send no report");
|
||||
assert!(
|
||||
holds_marker(&fx),
|
||||
"P must hold the filter once the fallback resends it"
|
||||
);
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
/// The cumulative counters of the last report `node` accepted from `peer`.
|
||||
fn rr_counters(node: &Node, peer: &NodeAddr) -> Option<(u64, u64, u32)> {
|
||||
node.get_peer(peer)?.mmp()?.metrics.rr_counters()
|
||||
}
|
||||
|
||||
/// The next send counter of `node`'s current session with `peer`.
|
||||
fn next_counter(node: &Node, peer: &NodeAddr) -> u64 {
|
||||
node.get_peer(peer)
|
||||
.and_then(|p| p.noise_session())
|
||||
.expect("setup: session present")
|
||||
.current_send_counter()
|
||||
}
|
||||
|
||||
/// Around a rekey, reports that describe the previous session reach both
|
||||
/// ends of the link: the initiator accepts one the responder built before it
|
||||
/// switched, and the responder's frames from the old session pollute the
|
||||
/// initiator's receiver. Neither kind of report may trigger a resend.
|
||||
#[tokio::test]
|
||||
async fn test_bloom_reports_from_the_previous_session_do_not_trigger_resends() {
|
||||
let mut fx = flip_fixture(true).await;
|
||||
let (m, p) = (fx.m, fx.p);
|
||||
fx.nodes[P].node.bloom_state.set_update_debounce_ms(0);
|
||||
start_reports(&mut fx).await;
|
||||
arm_rekey(&mut fx);
|
||||
|
||||
// A session reaching its rekey has carried many frames. Reserve counters
|
||||
// on both old sessions so their counters stay above the new sessions'
|
||||
// for the whole test, as they do in the field; with a short history the
|
||||
// new counters pass them within a few rounds, the reports become usable,
|
||||
// and each announce spends its one unchecked resend.
|
||||
for (i, remote) in [(M, p), (P, m)] {
|
||||
let session = fx.nodes[i]
|
||||
.node
|
||||
.get_peer_mut(&remote)
|
||||
.and_then(|peer| peer.noise_session_mut())
|
||||
.expect("setup: session present");
|
||||
for _ in 0..1000 {
|
||||
session
|
||||
.take_send_counter()
|
||||
.expect("setup: counter available");
|
||||
}
|
||||
}
|
||||
|
||||
// Reports both ways, then M reports alone so P holds interval data.
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
zero_intervals(&mut fx.nodes[M].node);
|
||||
zero_intervals(&mut fx.nodes[P].node);
|
||||
fx.nodes[M].node.check_mmp_reports().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
|
||||
// M starts the rekey and holds the new session, not yet cut over.
|
||||
let before = link_epoch(&fx);
|
||||
fx.nodes[M].node.check_rekey().await;
|
||||
for _ in 0..10 {
|
||||
if fx.nodes[M]
|
||||
.node
|
||||
.get_peer(&p)
|
||||
.is_some_and(|peer| peer.pending_new_session().is_some())
|
||||
{
|
||||
break;
|
||||
}
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
}
|
||||
assert!(
|
||||
fx.nodes[M]
|
||||
.node
|
||||
.get_peer(&p)
|
||||
.is_some_and(|peer| peer.pending_new_session().is_some()),
|
||||
"setup: M must hold P's new session"
|
||||
);
|
||||
assert_eq!(link_epoch(&fx), before, "setup: M must not have cut over");
|
||||
|
||||
// P reports on the old session; hold its frames back from M.
|
||||
assert!(
|
||||
fx.nodes[M].packet_rx.is_empty(),
|
||||
"setup: M's queue is empty"
|
||||
);
|
||||
zero_intervals(&mut fx.nodes[P].node);
|
||||
fx.nodes[P].node.check_mmp_reports().await;
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(1);
|
||||
while fx.nodes[M].packet_rx.is_empty() && std::time::Instant::now() < deadline {
|
||||
tokio::time::sleep(Duration::from_millis(5)).await;
|
||||
}
|
||||
let mut held = Vec::new();
|
||||
while let Ok(packet) = fx.nodes[M].packet_rx.try_recv() {
|
||||
held.push(packet);
|
||||
}
|
||||
assert!(!held.is_empty(), "setup: P must queue its reports for M");
|
||||
|
||||
// M cuts over, then receives P's old-session frames.
|
||||
fx.nodes[M].node.check_rekey().await;
|
||||
assert_ne!(link_epoch(&fx), before, "setup: M must cut over");
|
||||
for packet in held {
|
||||
fx.nodes[M].node.handle_encrypted_frame(packet).await;
|
||||
}
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
|
||||
let m_rr = rr_counters(&fx.nodes[M].node, &p).expect("setup: M holds a report");
|
||||
let p_rr = rr_counters(&fx.nodes[P].node, &m).expect("setup: P holds a report");
|
||||
assert!(
|
||||
m_rr.0 >= next_counter(&fx.nodes[M].node, &p),
|
||||
"setup: M's report must describe M's previous session"
|
||||
);
|
||||
assert!(
|
||||
p_rr.0 >= next_counter(&fx.nodes[P].node, &m),
|
||||
"setup: P's report must carry P's previous-session counter"
|
||||
);
|
||||
|
||||
fx.nodes[M].node.bloom_state.mark_update_needed(p);
|
||||
fx.nodes[P].node.bloom_state.mark_update_needed(m);
|
||||
fx.nodes[M].node.send_pending_filter_announces().await;
|
||||
fx.nodes[P].node.send_pending_filter_announces().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
assert!(
|
||||
fx.nodes[M].node.bloom_state.announce_outstanding(&p)
|
||||
&& fx.nodes[P].node.bloom_state.announce_outstanding(&m),
|
||||
"setup: both ends must have an announce outstanding"
|
||||
);
|
||||
|
||||
let sent_m = fx.nodes[M].node.metrics().bloom.sent.get();
|
||||
let sent_p = fx.nodes[P].node.metrics().bloom.sent.get();
|
||||
let guard = switch_guard(&fx);
|
||||
let mut last = rr_counters(&fx.nodes[P].node, &m);
|
||||
let mut changes = 0;
|
||||
for _ in 0..5 {
|
||||
mmp_round(&mut fx.nodes).await;
|
||||
fx.nodes[M].node.check_bloom_state().await;
|
||||
fx.nodes[P].node.check_bloom_state().await;
|
||||
process_available_packets(&mut fx.nodes).await;
|
||||
let now = rr_counters(&fx.nodes[P].node, &m);
|
||||
if now != last {
|
||||
changes += 1;
|
||||
}
|
||||
last = now;
|
||||
}
|
||||
assert!(
|
||||
changes >= 2,
|
||||
"setup: P must accept at least two reports from M, saw {changes}"
|
||||
);
|
||||
assert_unswitched(&fx, &guard);
|
||||
let m_rr = rr_counters(&fx.nodes[M].node, &p).expect("setup: M holds a report");
|
||||
let p_rr = rr_counters(&fx.nodes[P].node, &m).expect("setup: P holds a report");
|
||||
assert!(
|
||||
m_rr.0 >= next_counter(&fx.nodes[M].node, &p)
|
||||
&& p_rr.0 >= next_counter(&fx.nodes[P].node, &m),
|
||||
"setup: both reports must still describe a previous session"
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
fx.nodes[M].node.metrics().bloom.sent.get(),
|
||||
sent_m,
|
||||
"M must not resend on reports from the previous session"
|
||||
);
|
||||
assert_eq!(
|
||||
fx.nodes[P].node.metrics().bloom.sent.get(),
|
||||
sent_p,
|
||||
"P must not resend on reports carrying its previous-session counter"
|
||||
);
|
||||
assert!(
|
||||
fx.nodes[M].node.bloom_state.announce_outstanding(&p),
|
||||
"M must still hold its announce to P"
|
||||
);
|
||||
assert!(
|
||||
fx.nodes[P].node.bloom_state.announce_outstanding(&m),
|
||||
"P must still hold its announce to M"
|
||||
);
|
||||
cleanup_nodes(&mut fx.nodes).await;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user