diff --git a/docs/design/fips-bloom-filters.md b/docs/design/fips-bloom-filters.md index d0fce166..7b394367 100644 --- a/docs/design/fips-bloom-filters.md +++ b/docs/design/fips-bloom-filters.md @@ -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 diff --git a/src/node/bloom.rs b/src/node/bloom.rs index 5c6586af..e85ea4c9 100644 --- a/src/node/bloom.rs +++ b/src/node/bloom.rs @@ -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 { + 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 = 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; } } diff --git a/src/node/tests/bloom.rs b/src/node/tests/bloom.rs index 0542b0a2..978b57cd 100644 --- a/src/node/tests/bloom.rs +++ b/src/node/tests/bloom.rs @@ -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; +}