Resend a tree announce the peer did not receive

A tree announce was counted as delivered once the transport accepted it, so
a lost datagram left the peer on our old tree position until the periodic
re-broadcast up to a minute later, and indefinitely on a node with a single
peer, which has no periodic re-broadcast. Meanwhile the peer could leave
destinations out of discovery or route toward them by stale coordinates.

The tree state now keeps, per peer, the tree announce still awaiting
confirmation from the link's receiver reports, using the same delivery check
as filter announces. The declaration sequence is the lineage: a new sequence
refills the resend budgets, while a resend or periodic re-broadcast of the
same declaration spends from them. The tracking is dropped with the peer.

The node resends on a reported loss, and once after 30 s when the reports
cannot confirm the announce. Per declaration and session it resends at most
once unchecked and three times on loss, and at most six times a minute per
peer. The resend is an ordinary announce of the current declaration through
the rate-limited send path, and the "Sent TreeAnnounce" trace line now
carries the declaration sequence.

The node tests run on a converged line with every role read from the
converged tree, and share the loopback MMP, loss and rekey helpers of the
filter announce tests, now generalised to node indices.

The Ethernet mesh scenario's 60 s delivery window no longer sits at the
ordinary recovery bound, since tree and filter announces lost to a link flap
are now resent within seconds of the link returning, or after about 30 s
when the reports cannot confirm them. A comment beside the assertion says
so, so a red there is investigated rather than expected.
This commit is contained in:
Johnathan Corgan
2026-09-27 19:35:48 +00:00
parent 633450ba4e
commit 32bf476660
8 changed files with 1002 additions and 84 deletions
+13
View File
@@ -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
+3 -2
View File
@@ -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<LinkEvidence> {
/// 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<LinkEvidence> {
let peer = self.peers.get(peer_addr)?;
let session = peer.noise_session()?;
let mut epoch = [0u8; 8];
+602 -76
View File
@@ -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<NodeAddr> = 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]
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(&remote)
.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<TestNode>,
/// 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<usize> = [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<u64> {
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<u64>) {
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;
}
+61 -3
View File
@@ -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<NodeAddr> = 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
+70
View File
@@ -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<NodeAddr, ParentDeclaration>,
/// Each peer's full ancestry to root.
peer_ancestry: BTreeMap<NodeAddr, TreeCoordinate>,
/// Per-peer delivery tracking for sent TreeAnnounces.
acks: Acks,
/// The declaration sequence last recorded as sent to each peer.
announced: BTreeMap<NodeAddr, u64>,
/// 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<ResendReason> {
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<u64> {
self.acks.outstanding(peer)
}
/// The declaration sequence of the TreeAnnounce to `peer` awaiting
/// confirmation.
pub fn outstanding_seq(&self, peer: &NodeAddr) -> Option<u64> {
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.
+241
View File
@@ -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<RrCounters> {
Some(RrCounters {
highest,
received,
reordered,
})
}
/// Link evidence for session `epoch`.
fn link(epoch: u64, next_counter: u64, rr: Option<RrCounters>) -> LinkEvidence {
LinkEvidence {
epoch,
next_counter,
rr,
}
}
/// 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<ResendReason> {
self.state.check_announce(&self.peer, &link, now_ms)
}
/// Whether the announce is still unconfirmed.
fn outstanding(&self) -> bool {
self.state.announce_outstanding(&self.peer)
}
/// Check every 1,000 ms from `from_ms` to `to_ms` inclusive, 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);
}
+1
View File
@@ -1,5 +1,6 @@
//! STP primitive unit tests. Shared helpers live in `util`.
mod acks;
mod coordinate;
mod limits;
mod state;
@@ -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