diff --git a/docs/design/fips-spanning-tree.md b/docs/design/fips-spanning-tree.md index c7b5d117..5e3e9cf1 100644 --- a/docs/design/fips-spanning-tree.md +++ b/docs/design/fips-spanning-tree.md @@ -207,6 +207,19 @@ branch, and each only re-announces to its peers. - **Convergence time**: A tree of depth D reconverges in roughly D × 0.5s to D × 1.0s +The transport accepting a TreeAnnounce is not delivery. Each announce stays +outstanding until the link's receiver reports confirm it, by the same rule +as a FilterAnnounce (see "Update Triggers" in +[fips-bloom-filters.md](fips-bloom-filters.md)). It is resent when the +reports show a loss, and once after 30 s when they cannot confirm it, with +the same budgets and backoff: per declaration sequence per session, one +unchecked resend and three on loss, spaced by a per-peer backoff of 1, 2, +4 ... up to 60 s. A resend carries the current declaration, so a peer that +missed an older one gets the newer position, and one that already holds it +ignores it as not fresher. The periodic re-evaluation under Stability +Mechanisms, which re-broadcasts an unchanged declaration, stays as the +backstop on a node with two or more peers. + ### Transitive Trust (v1) In the v1 protocol, only the sender's outer signature on the TreeAnnounce diff --git a/src/node/bloom.rs b/src/node/bloom.rs index c4877d59..48c08415 100644 --- a/src/node/bloom.rs +++ b/src/node/bloom.rs @@ -288,8 +288,9 @@ impl Node { /// 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 { + /// accepted on it. `None` when the peer has no session. Both the filter + /// and the tree announce resends read it. + pub(super) 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]; diff --git a/src/node/tests/bloom.rs b/src/node/tests/bloom.rs index 978b57cd..70104193 100644 --- a/src/node/tests/bloom.rs +++ b/src/node/tests/bloom.rs @@ -1013,17 +1013,24 @@ fn zero_intervals(node: &mut Node) { } } +/// One MMP exchange between nodes `a` and `b`: `a` reports, everyone +/// processes, `b` reports, everyone processes. No other node is asked to +/// report. +async fn mmp_between(nodes: &mut [TestNode], a: usize, b: usize) { + zero_intervals(&mut nodes[a].node); + zero_intervals(&mut nodes[b].node); + nodes[a].node.check_mmp_reports().await; + process_available_packets(nodes).await; + zero_intervals(&mut nodes[a].node); + zero_intervals(&mut nodes[b].node); + nodes[b].node.check_mmp_reports().await; + process_available_packets(nodes).await; +} + /// 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; + mmp_between(nodes, M, P).await; } /// Process packets on every node until a pass handles none, at most 50 passes. @@ -1050,15 +1057,22 @@ async fn drop_queued(tn: &mut TestNode) -> usize { dropped } -/// ReceiverReports M has seen from P, including stale and duplicate ones. -fn reports_seen(fx: &FlipFixture) -> u64 { - fx.nodes[M] +/// ReceiverReports node `a` has seen from node `b`, including stale and +/// duplicate ones. +fn seen_by(nodes: &[TestNode], a: usize, b: usize) -> u64 { + let remote = *nodes[b].node.node_addr(); + nodes[a] .node - .get_peer(&fx.p) + .get_peer(&remote) .and_then(|peer| peer.mmp()) .map_or(0, |mmp| mmp.metrics.reports_seen()) } +/// ReceiverReports M has seen from P, including stale and duplicate ones. +fn reports_seen(fx: &FlipFixture) -> u64 { + seen_by(&fx.nodes, M, P) +} + /// Whether the filter P stores for M contains the marker. fn holds_marker(fx: &FlipFixture) -> bool { fx.nodes[P] @@ -1068,40 +1082,52 @@ fn holds_marker(fx: &FlipFixture) -> bool { .is_some_and(|filter| filter.contains(&marker())) } -/// Parent-switch counts at M and P before a run of MMP rounds. +/// Parent-switch counts at the two ends of a link 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, + a: u64, + b: u64, +} + +/// Snapshot the parent-switch counters of nodes `a` and `b`. +fn guard_of(nodes: &[TestNode], a: usize, b: usize) -> SwitchGuard { + SwitchGuard { + a: nodes[a].node.metrics().tree.parent_switches.get(), + b: nodes[b].node.metrics().tree.parent_switches.get(), + } +} + +/// Assert neither node `a` nor node `b` switched parent since `guard`, and +/// `a`'s parent is still `parent`. +fn assert_steady(nodes: &[TestNode], a: usize, b: usize, guard: &SwitchGuard, parent: &NodeAddr) { + assert_eq!( + nodes[a].node.metrics().tree.parent_switches.get(), + guard.a, + "setup: node {a} must not switch parent during the MMP rounds" + ); + assert_eq!( + nodes[b].node.metrics().tree.parent_switches.get(), + guard.b, + "setup: node {b} must not switch parent during the MMP rounds" + ); + assert_eq!( + nodes[a].node.tree_state().my_declaration().parent_id(), + parent, + "setup: node {a}'s parent must not change" + ); } /// 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(), - } + guard_of(&fx.nodes, M, P) } /// 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" - ); + assert_steady(&fx.nodes, M, P, guard, &fx.p); } /// FilterAnnounces M has sent. @@ -1109,16 +1135,25 @@ 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; +/// Drain `nodes`, then run MMP exchanges between nodes `a` and `b` until `a` +/// has seen a report from `b`. +async fn await_report(nodes: &mut [TestNode], a: usize, b: usize) { + drain_quiet(nodes).await; for _ in 0..10 { - if reports_seen(fx) >= 1 { + if seen_by(nodes, a, b) >= 1 { break; } - mmp_round(&mut fx.nodes).await; + mmp_between(nodes, a, b).await; } - assert!(reports_seen(fx) >= 1, "setup: P must report to M"); + assert!( + seen_by(nodes, a, b) >= 1, + "setup: node {b} must report to node {a}" + ); +} + +/// Drain the fixture, then run MMP rounds until M has seen a report from P. +async fn start_reports(fx: &mut FlipFixture) { + await_report(&mut fx.nodes, M, P).await; } /// Deliver C's filter carrying the marker and send M's announce of it to P. @@ -1159,19 +1194,26 @@ async fn lose_marker(fx: &mut FlipFixture) -> u64 { 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] +/// The first eight bytes of the handshake hash of node `a`'s current session +/// with node `b`. +fn epoch_of(nodes: &[TestNode], a: usize, b: usize) -> [u8; 8] { + let remote = *nodes[b].node.node_addr(); + let hash = nodes[a] .node - .get_peer(&fx.p) + .get_peer(&remote) .and_then(|peer| peer.noise_session()) - .expect("M has a session with P") + .expect("setup: the link has a session") .handshake_hash(); let mut epoch = [0u8; 8]; epoch.copy_from_slice(&hash[..8]); epoch } +/// The first eight bytes of the handshake hash of M's current session with P. +fn link_epoch(fx: &FlipFixture) -> [u8; 8] { + epoch_of(&fx.nodes, M, P) +} + /// A FilterAnnounce lost in transit is resent once a receiver report shows /// the loss, so the peer ends up holding the filter. #[tokio::test] @@ -1221,52 +1263,76 @@ async fn test_bloom_unchanged_filter_with_newer_sequence_marks_no_peer() { 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| { +/// Make node `a` rekey on its next check: one message on a session is +/// enough, time never triggers it, and both ends of every link of `a` are +/// aged past the responder's rekey-acceptance gate so every rekey is an +/// ordinary one. +fn arm_rekeys(nodes: &mut [TestNode], a: usize) { + nodes[a].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 local = *nodes[a].node.node_addr(); + let remotes: Vec = nodes[a].node.peers.keys().copied().collect(); 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); + for remote in remotes { + let b = nodes + .iter() + .position(|tn| *tn.node.node_addr() == remote) + .expect("setup: every peer is a test node"); + for (i, addr) in [(a, remote), (b, local)] { + nodes[i] + .node + .get_peer_mut(&addr) + .expect("setup: link peer present") + .test_backdate_session_established(age); + } } } +/// Make M rekey on its next check, with both ends of both of M's links aged +/// past the responder's rekey-acceptance gate. +fn arm_rekey(fx: &mut FlipFixture) { + arm_rekeys(&mut fx.nodes, M); +} + +/// Drive the real rekey handshake until node `a`'s session with node `b` is +/// cut over. +async fn cutover(nodes: &mut [TestNode], a: usize, b: usize) { + let before = epoch_of(nodes, a, b); + for _ in 0..6 { + nodes[a].node.check_rekey().await; + nodes[b].node.check_rekey().await; + for _ in 0..3 { + tokio::time::sleep(Duration::from_millis(5)).await; + process_available_packets(nodes).await; + } + if epoch_of(nodes, a, b) != before { + break; + } + } + assert_ne!( + epoch_of(nodes, a, b), + before, + "setup: node {a}'s link to node {b} must rekey" + ); + let (addr_a, addr_b) = (*nodes[a].node.node_addr(), *nodes[b].node.node_addr()); + assert!( + !nodes[a].node.get_peer(&addr_b).unwrap().rekey_in_progress(), + "setup: node {a}'s rekey with node {b} must be complete" + ); + assert!( + !nodes[b].node.get_peer(&addr_a).unwrap().rekey_in_progress(), + "setup: node {b}'s rekey with node {a} must be complete" + ); +} + /// 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" - ); + cutover(&mut fx.nodes, M, P).await; } /// An announce lost just before a link rekey is resent on the new session. @@ -1574,3 +1640,463 @@ async fn test_bloom_reports_from_the_previous_session_do_not_trigger_resends() { ); cleanup_nodes(&mut fx.nodes).await; } + +// ===== Resend of a tree announce the peer did not receive ===== +// +// Tree announces are confirmed and resent by the same receiver-report check +// as filter announces, so their node tests share this file's MMP, rekey and +// loss helpers. They run on real converged state only: every role is read +// from the converged tree, and nothing in any node's tree state is forged. +// Final assertions read the sequence the receiver stores for the sender, +// never what the sender believes it sent. + +/// A converged line of loopback nodes with a sender S that is not root and +/// its parent R. +/// +/// In a line every neighbour of S other than its parent is its child, whose +/// ancestry contains S, so S has no alternative parent and cannot switch. +struct TreeLine { + nodes: Vec, + /// Index of the sender. + s: usize, + /// Index of the sender's parent. + r: usize, +} + +/// Index of the node whose address is `addr`. +fn index_of(nodes: &[TestNode], addr: &NodeAddr) -> usize { + nodes + .iter() + .position(|tn| tn.node.node_addr() == addr) + .expect("setup: address belongs to a test node") +} + +/// Converge a line of `n` nodes (2 or 4) and pick S and R from the result. +/// +/// With 4 nodes, S is the first of indices 1 and 2 that is not root; with 2, +/// S is the node that is not root. R is S's parent. Every node's per-peer +/// tree announce rate limit is set to 0 and whatever convergence left +/// pending is flushed, so a resend is never held back by the rate limit. +async fn tree_line(n: usize) -> TreeLine { + let edges: Vec<(usize, usize)> = (1..n).map(|i| (i - 1, i)).collect(); + let mut nodes = run_tree_test(n, &edges, false).await; + + let root = (0..n) + .find(|&i| nodes[i].node.tree_state().is_root()) + .expect("setup: the line must have a root"); + let candidates: &[usize] = if n == 2 { &[0, 1] } else { &[1, 2] }; + let s = *candidates + .iter() + .find(|&&i| !nodes[i].node.tree_state().is_root()) + .expect("setup: one candidate sender is not root"); + let s_addr = *nodes[s].node.node_addr(); + let r = index_of( + &nodes, + nodes[s].node.tree_state().my_declaration().parent_id(), + ); + let neighbours: Vec = [s.checked_sub(1), Some(s + 1)] + .into_iter() + .flatten() + .filter(|&i| i < n) + .collect(); + eprintln!( + "tree_line({n}): root index {root}, S {s}, R {r}, shape: {}", + if r == root { + "R is root" + } else { + "R is not root" + } + ); + + assert!( + !nodes[s].node.tree_state().is_root(), + "setup: S is not root" + ); + assert_eq!( + nodes[s].node.peers.len(), + neighbours.len(), + "setup: S has exactly its line neighbours as peers" + ); + assert!(neighbours.contains(&r), "setup: R is S's neighbour"); + for &c in neighbours.iter().filter(|&&i| i != r) { + assert_eq!( + nodes[c].node.tree_state().my_declaration().parent_id(), + &s_addr, + "setup: S's other neighbour declares S as its parent" + ); + } + + for tn in nodes.iter_mut() { + for peer in tn.node.peers.values_mut() { + peer.set_tree_announce_min_interval_ms(0); + } + tn.node.send_pending_tree_announces().await; + } + drain_quiet(&mut nodes).await; + for (i, tn) in nodes.iter().enumerate() { + for peer in tn.node.peers.values() { + assert!( + !peer.has_pending_tree_announce(), + "setup: node {i} must have no pending tree announce" + ); + } + } + TreeLine { nodes, s, r } +} + +/// Turn off the periodic parent re-evaluation on `node`, so its periodic +/// re-broadcast cannot be what delivers a lost announce. +fn no_reeval(node: &mut Node) { + node.replace_context(|ctx| { + let mut cfg = (*ctx.config).clone(); + cfg.node.tree.reeval_interval_secs = 0; + ctx.config = std::sync::Arc::new(cfg); + }); +} + +/// Hold off S's fallback resend. S's announces to a child that never +/// reports stay outstanding, and on a loaded host a test running past the +/// fallback would add an unchecked resend to the child and break the exact +/// send counts. +fn hold_fallback(line: &mut TreeLine) { + line.nodes[line.s] + .node + .tree_state_mut() + .set_fallback(u64::MAX); +} + +/// Whether S's announce to R still awaits confirmation. +fn parent_outstanding(line: &TreeLine) -> bool { + let parent = *line.nodes[line.r].node.node_addr(); + line.nodes[line.s] + .node + .tree_state() + .announce_outstanding(&parent) +} + +/// TreeAnnounces node `i` has sent. +fn tree_sent(nodes: &[TestNode], i: usize) -> u64 { + nodes[i].node.metrics().tree.sent.get() +} + +/// The declaration sequence node `r` stores for node `s`. +fn held_seq(nodes: &[TestNode], r: usize, s: usize) -> Option { + let sender = *nodes[s].node.node_addr(); + nodes[r] + .node + .tree_state() + .peer_declaration(&sender) + .map(|decl| decl.sequence()) +} + +/// Give S a new declaration sequence with the same parent, signed, as a +/// position change would. Returns the new sequence. +fn bump(line: &mut TreeLine) -> u64 { + let (s, r) = (line.s, line.r); + let parent = *line.nodes[r].node.node_addr(); + let identity = line.nodes[s].node.identity().clone(); + let ts = line.nodes[s].node.tree_state_mut(); + let seq = ts.my_declaration().sequence() + 1; + let timestamp = ts.my_declaration().timestamp() + 1; + ts.set_parent(parent, seq, timestamp, crate::time::mono_ms()); + ts.recompute_coords(); + sign_declaration(ts.my_declaration_mut(), &identity).unwrap(); + + let ts = line.nodes[s].node.tree_state(); + assert!(!ts.is_root(), "setup: S is still not root after the bump"); + assert_eq!( + ts.my_declaration().parent_id(), + &parent, + "setup: S's parent is still R after the bump" + ); + assert!( + held_seq(&line.nodes, r, s).is_some_and(|held| seq > held), + "setup: the new sequence is fresher than the one R holds for S" + ); + seq +} + +/// Send S's current announce to R and check exactly one was sent. Returns +/// S's sent count after the send. +async fn send_up(line: &mut TreeLine) -> u64 { + let (s, r) = (line.s, line.r); + let parent = *line.nodes[r].node.node_addr(); + let before = tree_sent(&line.nodes, s); + line.nodes[s] + .node + .send_tree_announce_to_peer(&parent) + .await + .expect("setup: the send must succeed"); + let after = tree_sent(&line.nodes, s); + assert_eq!(after, before + 1, "setup: S must send exactly one announce"); + after +} + +/// Lose the announce just sent to R, and check the loss took. +async fn lose_up(line: &mut TreeLine, old: Option) { + let (s, r) = (line.s, line.r); + assert_eq!( + drop_queued(&mut line.nodes[r]).await, + 1, + "setup: exactly the one announce frame must be lost" + ); + assert_eq!( + held_seq(&line.nodes, r, s), + old, + "control: R must still hold S's old sequence" + ); +} + +/// Bump S's declaration, send it to R and lose it on the way. Returns the new +/// sequence and S's sent count after the send. +async fn lose_bump(line: &mut TreeLine) -> (u64, u64) { + let old = held_seq(&line.nodes, line.r, line.s); + let seq = bump(line); + let sent = send_up(line).await; + lose_up(line, old).await; + (seq, sent) +} + +/// `rounds` rounds of an MMP exchange between S and R followed by S's tree +/// tick, with checks that a report arrived and nothing switched parent. +async fn tree_rounds(line: &mut TreeLine, rounds: usize) { + let (s, r) = (line.s, line.r); + let parent = *line.nodes[r].node.node_addr(); + let seen = seen_by(&line.nodes, s, r); + let guard = guard_of(&line.nodes, s, r); + for _ in 0..rounds { + mmp_between(&mut line.nodes, s, r).await; + line.nodes[s].node.check_tree_state().await; + process_available_packets(&mut line.nodes).await; + } + assert!( + seen_by(&line.nodes, s, r) > seen, + "setup: a receiver report must arrive after the send" + ); + assert_steady(&line.nodes, s, r, &guard, &parent); +} + +/// Whether S has a tree announce pending for R. +fn parent_pending(line: &TreeLine) -> bool { + let parent = *line.nodes[line.r].node.node_addr(); + line.nodes[line.s] + .node + .get_peer(&parent) + .expect("setup: R is S's peer") + .has_pending_tree_announce() +} + +/// A TreeAnnounce lost in transit is resent once a receiver report shows the +/// loss, so the parent ends up holding the new declaration. +#[tokio::test] +async fn test_tree_announce_lost_in_transit_reaches_the_peer_after_a_receiver_report() { + let mut line = tree_line(4).await; + let (s, r) = (line.s, line.r); + await_report(&mut line.nodes, s, r).await; + no_reeval(&mut line.nodes[s].node); + hold_fallback(&mut line); + + let (seq, sent) = lose_bump(&mut line).await; + tree_rounds(&mut line, 3).await; + + assert_eq!( + tree_sent(&line.nodes, s), + sent + 1, + "S must resend to R exactly once after the loss" + ); + assert!( + !parent_pending(&line), + "the rate limit must not be holding the resend" + ); + assert_eq!( + held_seq(&line.nodes, r, s), + Some(seq), + "R must hold the declaration whose announce was lost" + ); + cleanup_nodes(&mut line.nodes).await; +} + +/// A node with a single peer has no periodic re-broadcast, so without a +/// resend a lost announce is never recovered. +#[tokio::test] +async fn test_tree_announce_lost_by_a_node_with_one_peer_reaches_the_peer_after_a_receiver_report() +{ + let mut line = tree_line(2).await; + let (s, r) = (line.s, line.r); + await_report(&mut line.nodes, s, r).await; + + let (seq, _) = lose_bump(&mut line).await; + tree_rounds(&mut line, 3).await; + + assert_eq!( + held_seq(&line.nodes, r, s), + Some(seq), + "R must hold the declaration whose announce was lost" + ); + cleanup_nodes(&mut line.nodes).await; +} + +/// An announce lost just before a link rekey is resent on the new session. +#[tokio::test] +async fn test_tree_announce_lost_before_a_link_rekey_reaches_the_peer_after_the_cutover() { + let mut line = tree_line(4).await; + let (s, r) = (line.s, line.r); + await_report(&mut line.nodes, s, r).await; + no_reeval(&mut line.nodes[s].node); + hold_fallback(&mut line); + + let (seq, sent) = lose_bump(&mut line).await; + arm_rekeys(&mut line.nodes, s); + cutover(&mut line.nodes, s, r).await; + tree_rounds(&mut line, 5).await; + + assert_eq!( + tree_sent(&line.nodes, s), + sent + 1, + "S must resend to R exactly once after the loss" + ); + assert_eq!( + held_seq(&line.nodes, r, s), + Some(seq), + "R must hold the declaration whose announce was lost before the rekey" + ); + assert!( + !parent_outstanding(&line), + "S's resend on the new session must be confirmed" + ); + cleanup_nodes(&mut line.nodes).await; +} + +/// An announce that arrives is confirmed from the receiver reports and never +/// resent. +#[tokio::test] +async fn test_tree_delivered_announce_is_confirmed_without_a_resend() { + let mut line = tree_line(4).await; + let (s, r) = (line.s, line.r); + await_report(&mut line.nodes, s, r).await; + no_reeval(&mut line.nodes[s].node); + hold_fallback(&mut line); + + let seq = bump(&mut line); + let sent = send_up(&mut line).await; + assert!( + parent_outstanding(&line), + "control: the tracker must hold the announce straight after the send" + ); + process_available_packets(&mut line.nodes).await; + assert_eq!( + held_seq(&line.nodes, r, s), + Some(seq), + "control: R must hold the announce" + ); + tree_rounds(&mut line, 3).await; + + assert_eq!( + tree_sent(&line.nodes, s), + sent, + "S must not resend a delivered announce" + ); + assert!(!parent_pending(&line), "R must not be marked"); + assert!( + !parent_outstanding(&line), + "the delivered announce must be confirmed" + ); + cleanup_nodes(&mut line.nodes).await; +} + +/// On a converged mesh with clean links, every tree 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_tree_clean_links_confirm_every_announce_without_a_resend() { + let mut nodes = run_tree_test(3, &[(0, 1), (1, 2)], false).await; + for tn in nodes.iter_mut() { + for peer in tn.node.peers.values_mut() { + peer.set_tree_announce_min_interval_ms(0); + } + tn.node.send_pending_tree_announces().await; + no_reeval(&mut tn.node); + } + drain_quiet(&mut nodes).await; + for (i, tn) in nodes.iter().enumerate() { + for peer in tn.node.peers.values() { + assert!( + !peer.has_pending_tree_announce(), + "setup: node {i} must have no pending tree announce" + ); + } + } + let snapshot = |nodes: &[TestNode]| -> Vec<(u64, u64, NodeAddr)> { + nodes + .iter() + .map(|tn| { + ( + tn.node.metrics().tree.sent.get(), + tn.node.metrics().tree.parent_switches.get(), + *tn.node.tree_state().my_declaration().parent_id(), + ) + }) + .collect() + }; + let before = snapshot(&nodes); + assert!( + (0..nodes.len()).all(|i| tree_sent(&nodes, i) > 0), + "setup: every node must have sent tree announces while converging" + ); + + 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_tree_state().await; + } + process_available_packets(&mut nodes).await; + } + + assert_eq!( + snapshot(&nodes), + before, + "no node may send a tree announce, switch parent or change parent" + ); + for (i, tn) in nodes.iter().enumerate() { + for peer in tn.node.peers.keys() { + assert!( + !tn.node.tree_state().announce_outstanding(peer), + "node {i} must have confirmed its tree announce to every peer" + ); + } + } + cleanup_nodes(&mut nodes).await; +} + +/// With no receiver report at all, a lost tree announce is still resent once +/// the fallback interval passes. +#[tokio::test] +async fn test_tree_lost_announce_is_resent_after_the_fallback_when_no_receiver_report_arrives() { + let mut line = tree_line(4).await; + let (s, r) = (line.s, line.r); + drain_quiet(&mut line.nodes).await; + no_reeval(&mut line.nodes[s].node); + line.nodes[s].node.tree_state_mut().set_fallback(0); + let seen = seen_by(&line.nodes, s, r); + + let (seq, _) = lose_bump(&mut line).await; + line.nodes[s].node.check_tree_state().await; + process_available_packets(&mut line.nodes).await; + + assert_eq!( + seen_by(&line.nodes, s, r), + seen, + "setup: R must send no report" + ); + assert_eq!( + held_seq(&line.nodes, r, s), + Some(seq), + "R must hold the declaration once the fallback resends it" + ); + cleanup_nodes(&mut line.nodes).await; +} diff --git a/src/node/tree.rs b/src/node/tree.rs index 100e28d8..21dadad6 100644 --- a/src/node/tree.rs +++ b/src/node/tree.rs @@ -139,7 +139,26 @@ impl Node { peer.record_tree_announce_sent(now_ms); } - trace!(peer = %self.peer_display_name(peer_addr), "Sent TreeAnnounce"); + // Keep the announce outstanding until the receiver reports confirm + // it. 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. + let seq = announce.declaration.sequence(); + if let Some(link) = self.link_evidence(peer_addr) { + self.tree_state.record_announce( + *peer_addr, + seq, + link.next_counter.saturating_sub(1), + &link, + crate::time::mono_ms(), + ); + } + + trace!( + peer = %self.peer_display_name(peer_addr), + seq = seq, + "Sent TreeAnnounce" + ); Ok(()) } @@ -584,13 +603,52 @@ impl Node { /// Periodic tree maintenance, called from the tick handler. /// - /// Sends pending rate-limited announces and checks for periodic - /// parent re-evaluation based on current MMP link costs. + /// Marks peers whose last announce was not confirmed delivered, sends + /// pending rate-limited announces (so a resend goes out in the same tick + /// through the ordinary send path), and checks for periodic parent + /// re-evaluation based on current MMP link costs. pub(super) async fn check_tree_state(&mut self) { + self.tree_resend(); self.send_pending_tree_announces().await; self.check_periodic_parent_reeval().await; } + /// Mark for resend every peer whose outstanding TreeAnnounce the receiver + /// reports show was lost, or could not confirm in time. + /// + /// The resend is an ordinary announce built from the current declaration, + /// so if our position moved on since the lost one, the peer gets the + /// newer position. + fn tree_resend(&mut self) { + let now_ms = crate::time::mono_ms(); + let waiting: Vec = self + .peers + .keys() + .filter(|addr| self.tree_state.announce_outstanding(addr)) + .copied() + .collect(); + for addr in waiting { + let Some(link) = self.link_evidence(&addr) else { + continue; + }; + let counter = self.tree_state.outstanding_counter(&addr); + let seq = self.tree_state.outstanding_seq(&addr); + let Some(reason) = self.tree_state.check_announce(&addr, &link, now_ms) else { + continue; + }; + if let Some(peer) = self.peers.get_mut(&addr) { + peer.mark_tree_announce_pending(); + } + debug!( + peer = %self.peer_display_name(&addr), + reason = ?reason, + counter = ?counter, + seq = ?seq, + "Resending unconfirmed TreeAnnounce" + ); + } + } + /// Periodic parent re-evaluation based on current MMP link costs. /// /// Self-paces using `last_parent_reeval` and the configured diff --git a/src/proto/stp/state.rs b/src/proto/stp/state.rs index 0c99a665..e84efd46 100644 --- a/src/proto/stp/state.rs +++ b/src/proto/stp/state.rs @@ -7,6 +7,7 @@ use super::core::ParentEval; use super::limits::FlapDampener; use super::{CoordEntry, ParentDeclaration, TreeCoordinate}; use crate::NodeAddr; +use crate::proto::mmp::delivery::{Acks, LinkEvidence, ResendReason}; /// What a parent-loss recovery did: whether the tree state changed, and /// whether the recovery switch was the one that armed a dampening episode. @@ -40,6 +41,10 @@ pub struct TreeState { peer_declarations: BTreeMap, /// Each peer's full ancestry to root. peer_ancestry: BTreeMap, + /// Per-peer delivery tracking for sent TreeAnnounces. + acks: Acks, + /// The declaration sequence last recorded as sent to each peer. + announced: BTreeMap, /// Hysteresis factor for cost-based parent re-selection (0.0-1.0). parent_hysteresis: f64, /// Flap-dampening / hold-down state machine. @@ -64,6 +69,8 @@ impl TreeState { root: my_node_addr, peer_declarations: BTreeMap::new(), peer_ancestry: BTreeMap::new(), + acks: Acks::new(), + announced: BTreeMap::new(), parent_hysteresis: 0.0, flap: FlapDampener::new(), } @@ -149,6 +156,69 @@ impl TreeState { pub fn remove_peer(&mut self, peer_id: &NodeAddr) { self.peer_declarations.remove(peer_id); self.peer_ancestry.remove(peer_id); + self.acks.remove(peer_id); + self.announced.remove(peer_id); + } + + /// Record a TreeAnnounce the transport accepted for `peer`, carrying our + /// declaration sequence `seq` and 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 lost announce would + /// leave the peer on our old tree position until the next announce, which + /// on a node with one peer may never come. The announce's content changes + /// only with the declaration sequence, so a sequence other than the one + /// last recorded for `peer` starts a new lineage with fresh resend + /// budgets, and a resend or periodic re-broadcast of the same declaration + /// spends from its lineage's budget. + pub fn record_announce( + &mut self, + peer: NodeAddr, + seq: u64, + counter: u64, + link: &LinkEvidence, + now_ms: u64, + ) { + let fresh = self.announced.insert(peer, seq) != Some(seq); + self.acks.record(peer, fresh, counter, link, now_ms); + } + + /// Decide whether the outstanding TreeAnnounce to `peer` must be resent, + /// by the shared receiver-report rule ([`Acks::check`]). + /// + /// Marks nothing: on `Some`, the caller marks the peer's announce pending + /// so the ordinary send path delivers the current declaration. + pub fn check_announce( + &mut self, + peer: &NodeAddr, + link: &LinkEvidence, + now_ms: u64, + ) -> Option { + self.acks.check(peer, link, now_ms) + } + + /// Whether a TreeAnnounce 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 TreeAnnounce to `peer` awaiting confirmation. + pub fn outstanding_counter(&self, peer: &NodeAddr) -> Option { + self.acks.outstanding(peer) + } + + /// The declaration sequence of the TreeAnnounce to `peer` awaiting + /// confirmation. + pub fn outstanding_seq(&self, peer: &NodeAddr) -> Option { + self.outstanding_counter(peer)?; + self.announced.get(peer).copied() + } + + /// Set how long a TreeAnnounce the receiver reports cannot check waits + /// before its fallback resend. Defaults to + /// [`FALLBACK_MS`](crate::proto::mmp::delivery::FALLBACK_MS). + pub fn set_fallback(&mut self, ms: u64) { + self.acks.set_fallback(ms); } /// Update this node's parent selection. diff --git a/src/proto/stp/tests/acks.rs b/src/proto/stp/tests/acks.rs new file mode 100644 index 00000000..21be4db7 --- /dev/null +++ b/src/proto/stp/tests/acks.rs @@ -0,0 +1,241 @@ +//! Delivery tracking for sent tree 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. The shared rule itself (bases, sessions, budgets, +//! backoff, rekey artifacts) is covered by the filter announce tests; these +//! cover what the tree adds: the declaration sequence as the lineage, and +//! removal with the peer. + +use super::util::make_node_addr; +use crate::NodeAddr; +use crate::proto::mmp::delivery::{FALLBACK_MS, LinkEvidence, ResendReason, RrCounters}; +use crate::proto::stp::TreeState; + +/// First session. +const E1: u64 = 0x0e01; +/// Second session. +const E2: u64 = 0x0e02; + +/// A report's cumulative counters. +fn rr(highest: u64, received: u64, reordered: u32) -> Option { + Some(RrCounters { + highest, + received, + reordered, + }) +} + +/// Link evidence for session `epoch`. +fn link(epoch: u64, next_counter: u64, rr: Option) -> LinkEvidence { + LinkEvidence { + epoch, + next_counter, + rr, + } +} + +/// One peer's tree announces, driven as the shell drives them. +struct Track { + state: TreeState, + peer: NodeAddr, + counter: u64, +} + +impl Track { + /// A tracker with nothing sent yet. + fn new() -> Self { + Self { + state: TreeState::new(make_node_addr(0), 1000), + peer: make_node_addr(1), + counter: 0, + } + } + + /// Record a send of declaration sequence `seq` at `counter`. + fn send(&mut self, seq: u64, counter: u64, link: LinkEvidence, now_ms: u64) { + self.state + .record_announce(self.peer, seq, counter, &link, now_ms); + self.counter = counter; + } + + /// One tick of the tracker. + fn check(&mut self, link: LinkEvidence, now_ms: u64) -> Option { + self.state.check_announce(&self.peer, &link, now_ms) + } + + /// Whether the announce is still unconfirmed. + fn outstanding(&self) -> bool { + self.state.announce_outstanding(&self.peer) + } + + /// Check every 1,000 ms from `from_ms` to `to_ms` inclusive, recording + /// each resend as a send of `seq`, as the shell resends the current + /// declaration. `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. Returns the resends. + fn hold( + &mut self, + seq: u64, + from_ms: u64, + to_ms: u64, + model: impl Fn(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(self.counter, false), now) { + resends.push((now, reason)); + let ev = model(self.counter, true); + self.send(seq, 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(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)) + } +} + +/// Send sequence 5 into a lossy link and spend its three loss resends. +/// +/// The resends are checked but not recorded, so spending the budget does not +/// depend on the lineage decision; the one record each test then makes is +/// the only lineage decision it observes. The window ends before the +/// backoff would allow a fourth resend (7 s), so the budget limit itself is +/// observed only by the test's final assertion. +fn spend_budget() -> Track { + let mut t = Track::new(); + t.send(5, 12, lossy(11, true), 0); + let mut resends = Vec::new(); + for now in (0..=6_000).step_by(1_000) { + if let Some(reason) = t.check(lossy(12, false), now) { + resends.push((now, reason)); + } + } + assert_eq!( + resends, + vec![ + (0, ResendReason::Loss), + (1_000, ResendReason::Loss), + (3_000, ResendReason::Loss), + ], + "setup: the lineage spends its three loss resends" + ); + assert!( + t.outstanding(), + "setup: the lost announce is still outstanding" + ); + t +} + +/// A new declaration sequence starts a new lineage, whose loss budget is +/// full again. +#[test] +fn test_tree_ack_new_sequence_starts_a_new_lineage() { + let mut t = spend_budget(); + let c = t.counter; + t.send(6, c + 1, lossy(c, true), 10_000); + assert_eq!( + t.check(lossy(t.counter, false), 11_000), + Some(ResendReason::Loss), + "a new sequence must refill the loss budget" + ); +} + +/// The periodic re-broadcast sends the same sequence again, so it stays in +/// the lineage and does not refill the loss budget; the one unchecked +/// resend still comes, 30 s after the re-broadcast. +#[test] +fn test_tree_ack_periodic_resend_of_the_same_sequence_keeps_the_budget() { + let mut t = spend_budget(); + let c = t.counter; + t.send(5, c + 1, lossy(c, true), 10_000); + let resends = t.hold(5, 11_000, 120_000, lossy); + assert_eq!( + resends, + vec![(10_000 + FALLBACK_MS, ResendReason::Timeout)], + "no fourth loss resend, and exactly one unchecked resend" + ); +} + +/// Removing the peer forgets its announce, so the next send starts a fresh +/// entry whose first session measures from counter zero. +#[test] +fn test_tree_ack_removed_peer_starts_fresh() { + // Control: without the removal, an announce in a later session with no + // report has no base, and a first report already covering it cannot + // check it. + let mut t = Track::new(); + t.send(5, 12, link(E1, 13, rr(9, 10, 0)), 0); + t.send(5, 3, link(E2, 4, None), 1_000); + assert_eq!( + t.check(link(E2, 4, rr(3, 4, 0)), 2_000), + Some(ResendReason::Unverified), + "control: the kept entry cannot check the announce" + ); + + let mut t = Track::new(); + t.send(5, 12, link(E1, 13, rr(9, 10, 0)), 0); + t.state.remove_peer(&t.peer); + assert!(!t.outstanding(), "removal must forget the announce"); + t.send(5, 3, link(E2, 4, None), 1_000); + assert_eq!(t.check(link(E2, 4, rr(3, 4, 0)), 2_000), None); + assert!(!t.outstanding(), "the new entry must measure from zero"); +} + +/// 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. This is the traced shape of a tree announce and a filter +/// announce sent in the same tick arriving swapped. +#[test] +fn test_tree_ack_in_window_reorder_on_the_zero_base_confirms() { + let mut t = Track::new(); + t.send(5, 3, link(E1, 4, None), 0); + // 0..=4 all arrived, two of them after a higher counter. + assert_eq!(t.check(link(E1, 5, rr(4, 5, 2)), 1_000), None); + assert!( + !t.outstanding(), + "every counter arrived, so it must confirm" + ); +} + +/// A covering report whose receipts since the base fall short of the span is +/// a loss, and the announce stays outstanding. +#[test] +fn test_tree_ack_short_upper_bound_is_a_loss() { + let mut t = Track::new(); + // Base (9, 10, 0): all of 0..=9 counted. + t.send(5, 12, link(E1, 13, rr(9, 10, 0)), 0); + // Four of 10..=14 arrived. + assert_eq!( + t.check(link(E1, 15, rr(14, 14, 0)), 1_000), + Some(ResendReason::Loss) + ); + assert!(t.outstanding(), "a lost announce stays outstanding"); +} + +/// With no report ever, the announce gets its one unchecked resend exactly +/// at the default fallback, and none after. +#[test] +fn test_tree_ack_fallback_resends_once_per_lineage() { + let mut t = Track::new(); + t.send(5, 12, link(E1, 13, None), 0); + assert_eq!(t.check(link(E1, 13, None), FALLBACK_MS - 1), None); + assert!(t.outstanding()); + let quiet = |c: u64, recording: bool| link(E1, c + if recording { 2 } else { 1 }, None); + let resends = t.hold(5, FALLBACK_MS, 120_000, quiet); + assert_eq!(resends, vec![(FALLBACK_MS, ResendReason::Timeout)]); + assert_eq!(FALLBACK_MS, 30_000); +} diff --git a/src/proto/stp/tests/mod.rs b/src/proto/stp/tests/mod.rs index 71199998..dc2f9f26 100644 --- a/src/proto/stp/tests/mod.rs +++ b/src/proto/stp/tests/mod.rs @@ -1,5 +1,6 @@ //! STP primitive unit tests. Shared helpers live in `util`. +mod acks; mod coordinate; mod limits; mod state; diff --git a/testing/chaos/scenarios/ethernet-mesh.yaml b/testing/chaos/scenarios/ethernet-mesh.yaml index a980856f..ab48cb5d 100644 --- a/testing/chaos/scenarios/ethernet-mesh.yaml +++ b/testing/chaos/scenarios/ethernet-mesh.yaml @@ -82,6 +82,14 @@ traffic: # happened, not a loss ratio, so a busy host slows it without failing # it. It runs at teardown, after flapped links are restored. # +# Recovery: a tree or bloom announce lost to a flap is resent from the +# link's receiver reports a few seconds after the link returns, or within +# about 30 s when those reports cannot confirm it; only a per-peer backoff +# built up by earlier resends on a lossy link can hold a resend longer, up +# to 60 s after the previous one. The 60 s window therefore no longer sits +# at the ordinary recovery bound, and a red here is investigated, not +# expected. +# # Sizes: payload 0 is the smallest ICMPv6 echo and 1200 is near the # 1280-byte TUN MTU. Neither can reach the receiver's padding trim on a # veth link. Measured on a live run of this scenario (2026-09-19): an