From ea9c7f2d8d37946f45a3d1cd522c7ca3aa4e3284 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Tue, 9 Jun 2026 23:55:48 +0000 Subject: [PATCH 1/2] mesh-size: union all peer filters, not just tree peers Estimate the OR-union cardinality over self plus every connected peer's inbound filter, dropping the parent/child tree gating in compute_mesh_size. Filter propagation is split-horizon, so cross-links advertise near-complete mesh views; unioning all peers yields the same set as the tree-only union in steady state (OR dedups overlap and no filter can over-count) while damping the node-count flap on parent switches, since dropping the parent no longer collapses the upward leg. This also removes the estimate's dependence on tree-declaration cache freshness. Rename the debug-log child_count to contributor_count, adapt the two membership-invariant tests to all-peers semantics, and add a test that the estimate stays stable across a parent drop when a healthy cross-link is present. --- src/node/mod.rs | 47 +++++----- src/node/tests/bloom.rs | 197 ++++++++++++++++++++++++++-------------- 2 files changed, 152 insertions(+), 92 deletions(-) diff --git a/src/node/mod.rs b/src/node/mod.rs index a2e2171..8b183ee 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -1142,17 +1142,27 @@ impl Node { /// of estimated entry counts plus one (self) approximates total network size. pub(crate) fn compute_mesh_size(&mut self) { let my_addr = *self.tree_state.my_node_addr(); - let parent_id = *self.tree_state.my_declaration().parent_id(); - let is_root = self.tree_state.is_root(); let max_fpr = self.config.node.bloom.max_inbound_fpr; - let mut child_count: u32 = 0; + let mut contributor_count: u32 = 0; // OR-union of the contributing filters. Summing per-filter // cardinalities over-counts whenever the filters overlap (a stale // or oversized parent filter, a topology loop); OR is idempotent, // so unioning and estimating once deduplicates the overlap. - // Membership is exactly: self + parent + each child inbound_filter. + // + // Membership is self + every connected peer's inbound_filter. We + // deliberately fold in *all* peers, not just the spanning-tree + // neighborhood (parent + children). Filter propagation is + // split-horizon (BloomState::compute_outgoing_filter excludes the + // peer it routes back to), so every routing peer — including + // cross-links — advertises a near-complete "whole mesh minus my + // subtree" view. Unioning all of them yields the same set as the + // tree-only union in steady state (OR-union dedups overlap, and + // each filter is a subset of the mesh so it cannot over-count) while + // damping the node-count flap on a parent switch: dropping the + // parent leaves the cross-links still carrying the upward coverage. + // It also removes the dependency on tree-declaration cache freshness. let mut union: Option = None; // Helper: fold a contributing filter into the union, starting it @@ -1167,26 +1177,13 @@ impl Node { } }; - // Parent's filter: nodes reachable upward through the tree. - if !is_root - && let Some(parent) = self.peers.get(&parent_id) - && let Some(filter) = parent.inbound_filter() - { - add_to_union(&mut union, filter); - } - - // Children's filters: each child's subtree is (ideally) disjoint. - for (peer_addr, peer) in &self.peers { - if peer_addr == &parent_id { - continue; - } - if let Some(decl) = self.tree_state.peer_declaration(peer_addr) - && *decl.parent_id() == my_addr - { - child_count += 1; - if let Some(filter) = peer.inbound_filter() { - add_to_union(&mut union, filter); - } + // Every connected peer's filter contributes. Honest peers are pure + // redundancy (overlapping bits dedup under OR); cross-links carry + // the upward coverage that would otherwise hinge on the parent alone. + for peer in self.peers.values() { + if let Some(filter) = peer.inbound_filter() { + contributor_count += 1; + add_to_union(&mut union, filter); } } @@ -1225,7 +1222,7 @@ impl Node { tracing::debug!( estimated_mesh_size = union_size, peers = self.peers.len(), - children = child_count, + contributors = contributor_count, "Mesh size estimate" ); self.last_mesh_size_log = Some(now); diff --git a/src/node/tests/bloom.rs b/src/node/tests/bloom.rs index 8f53ab6..a1ac716 100644 --- a/src/node/tests/bloom.rs +++ b/src/node/tests/bloom.rs @@ -443,16 +443,18 @@ async fn test_bloom_filter_split_horizon() { cleanup_nodes(&mut nodes).await; } -/// Regression for `compute_mesh_size` parent double-count. +/// Each peer's inbound filter contributes exactly once, independent of +/// tree-declaration state. /// -/// Reproduces the stale-peer-declaration window that follows a local -/// parent-switch: our `my_declaration().parent_id()` already names P, but -/// the cached `peer_declaration(P)` is still the pre-switch advert in -/// which P names US as its parent. Without the explicit parent-skip in -/// the children loop, P would be iterated as a child and its bloom -/// cardinality added a second time on top of the parent contribution. +/// Under all-peers union semantics there is no parent/child gating and no +/// `peer_declaration` lookup, so a stale or contradictory declaration +/// cache cannot cause a peer's bloom cardinality to be folded twice. This +/// retains the spirit of the old parent-double-count regression: we set up +/// the same stale-cache scenario (the cached `peer_declaration(P)` still +/// names US as P's parent) and assert the estimate reflects each filter +/// counted once, not the double-count fingerprint. #[test] -fn compute_mesh_size_skips_parent_under_stale_peer_declaration() { +fn compute_mesh_size_counts_each_peer_filter_once() { use crate::bloom::BloomFilter; use crate::peer::ActivePeer; use crate::tree::ParentDeclaration; @@ -501,6 +503,8 @@ fn compute_mesh_size_skips_parent_under_stale_peer_declaration() { // Inject the stale-cache scenario: peer_declaration(P) still names // US (M) as P's parent (the pre-switch advert that the cache hasn't // refreshed yet). Q is a legitimate child also naming M as parent. + // Under all-peers semantics these declarations no longer affect the + // estimate at all. let parent_decl_stale = ParentDeclaration::new(parent_addr, my_addr, 1, 1); let child_decl = ParentDeclaration::new(child_addr, my_addr, 1, 1); node.tree_state_mut() @@ -522,8 +526,9 @@ fn compute_mesh_size_skips_parent_under_stale_peer_declaration() { .estimated_mesh_size() .expect("estimator should produce a value with filter data present"); - // Expected (parent counted once + child counted once): 1 + 5 + 3 = 9. - // Unfixed (parent double-counted via children loop): 1 + 2*5 + 3 = 14. + // Each filter counted once: 1 (self) + 5 (P) + 3 (Q) = 9. + // A double-count of P (the old declaration-cache bug fingerprint) + // would land near 1 + 2*5 + 3 = 14. // The estimator's log-based math rounds, so allow +/-1 tolerance. let diff = (estimate as i64 - 9).abs(); assert!( @@ -533,23 +538,22 @@ fn compute_mesh_size_skips_parent_under_stale_peer_declaration() { ); } -/// Overlapping parent and child inbound filters must be OR-unioned, not -/// summed. The parent and the child here share several NodeAddrs (plus a -/// few distinct ones each). The naive sum of per-filter cardinalities -/// would over-count the shared entries; the union estimate must instead -/// approximate the number of *distinct* addresses across both filters -/// (plus self). This asserts the over-count fingerprint of the old -/// summing code is gone, and the assertion holds on the union result -/// (which is what estimated_mesh_size carries), so it survives removal -/// of the temporary dual-estimator instrumentation. +/// Overlapping peer inbound filters must be OR-unioned, not summed. +/// +/// Two connected peers share several NodeAddrs (plus a few distinct ones +/// each). The naive sum of per-filter cardinalities would over-count the +/// shared entries; the union estimate must instead approximate the number +/// of *distinct* addresses across both filters (plus self). Under all-peers +/// semantics the union folds in every connected peer regardless of tree +/// role, so no parent/child wiring is needed — this asserts the +/// overlap-dedup property directly on the union result that +/// `estimated_mesh_size` carries. #[test] fn compute_mesh_size_unions_overlapping_filters() { use crate::bloom::BloomFilter; use crate::peer::ActivePeer; - use crate::tree::ParentDeclaration; let mut node = make_node(); - let my_addr = *node.tree_state().my_node_addr(); // Build the set of addresses. SHARED appear in both filters; the // distinct sets appear in only one each. @@ -560,58 +564,36 @@ fn compute_mesh_size_unions_overlapping_filters() { NodeAddr::from_bytes(bytes) }; let shared: Vec = (0..6u8).map(|i| mk(0x10, i)).collect(); - let parent_only: Vec = (0..3u8).map(|i| mk(0x20, i)).collect(); - let child_only: Vec = (0..3u8).map(|i| mk(0x30, i)).collect(); + let peer_a_only: Vec = (0..3u8).map(|i| mk(0x20, i)).collect(); + let peer_b_only: Vec = (0..3u8).map(|i| mk(0x30, i)).collect(); - // Distinct addresses across the union: shared + parent_only + - // child_only + self = 6 + 3 + 3 + 1 = 13. The naive sum of the two + // Distinct addresses across the union: shared + peer_a_only + + // peer_b_only + self = 6 + 3 + 3 + 1 = 13. The naive sum of the two // filters' cardinalities would be (6+3) + (6+3) + 1 = 19. - let distinct = shared.len() + parent_only.len() + child_only.len() + 1; // 13 - let naive_sum = (shared.len() + parent_only.len()) + (shared.len() + child_only.len()) + 1; // 19 + let distinct = shared.len() + peer_a_only.len() + peer_b_only.len() + 1; // 13 + let naive_sum = (shared.len() + peer_a_only.len()) + (shared.len() + peer_b_only.len()) + 1; // 19 - // Generate a parent identity strictly less than my_addr so the - // tree_state defensive check accepts the extension. - let (parent_identity, parent_addr) = loop { - let candidate = make_peer_identity(); - let addr = *candidate.node_addr(); - if addr < my_addr { - break (candidate, addr); - } - }; - let mut parent_peer = ActivePeer::new(parent_identity, LinkId::new(1), 0); - let mut parent_filter = BloomFilter::new(); - for addr in shared.iter().chain(parent_only.iter()) { - parent_filter.insert(addr); + // Peer A with shared + peer_a_only. + let peer_a_identity = make_peer_identity(); + let peer_a_addr = *peer_a_identity.node_addr(); + let mut peer_a = ActivePeer::new(peer_a_identity, LinkId::new(1), 0); + let mut filter_a = BloomFilter::new(); + for addr in shared.iter().chain(peer_a_only.iter()) { + filter_a.insert(addr); } - parent_peer.update_filter(parent_filter, 1, 0); - node.peers.insert(parent_addr, parent_peer); + peer_a.update_filter(filter_a, 1, 0); + node.peers.insert(peer_a_addr, peer_a); - // Child Q with a filter that overlaps the parent's on `shared`. - let child_identity = make_peer_identity(); - let child_addr = *child_identity.node_addr(); - let mut child_peer = ActivePeer::new(child_identity, LinkId::new(2), 0); - let mut child_filter = BloomFilter::new(); - for addr in shared.iter().chain(child_only.iter()) { - child_filter.insert(addr); + // Peer B with a filter that overlaps A's on `shared`. + let peer_b_identity = make_peer_identity(); + let peer_b_addr = *peer_b_identity.node_addr(); + let mut peer_b = ActivePeer::new(peer_b_identity, LinkId::new(2), 0); + let mut filter_b = BloomFilter::new(); + for addr in shared.iter().chain(peer_b_only.iter()) { + filter_b.insert(addr); } - child_peer.update_filter(child_filter, 1, 0); - node.peers.insert(child_addr, child_peer); - - // Wire up the tree so the child names us as parent and our parent is P. - let parent_ancestry = crate::tree::TreeCoordinate::root_with_meta(parent_addr, 1, 1); - let child_ancestry = crate::tree::TreeCoordinate::root_with_meta(child_addr, 1, 1); - let parent_decl = ParentDeclaration::new(parent_addr, my_addr, 1, 1); - let child_decl = ParentDeclaration::new(child_addr, my_addr, 1, 1); - node.tree_state_mut() - .update_peer(parent_decl, parent_ancestry); - node.tree_state_mut() - .update_peer(child_decl, child_ancestry); - node.tree_state_mut().set_parent(parent_addr, 2, 1); - node.tree_state_mut().recompute_coords(); - assert!( - !node.tree_state().is_root(), - "test setup broken: node should not be its own root after parent switch" - ); + peer_b.update_filter(filter_b, 1, 0); + node.peers.insert(peer_b_addr, peer_b); node.compute_mesh_size(); @@ -637,6 +619,87 @@ fn compute_mesh_size_unions_overlapping_filters() { ); } +/// Flap-damping: the estimate survives a parent switch when a healthy +/// cross-link carries the upward coverage. +/// +/// Under the old tree-only union, the upward leg of the estimate hinged +/// entirely on the current parent's filter; dropping the parent (a parent +/// switch transient, before the new parent's filter converges) collapsed +/// the estimate to self + children. Under all-peers semantics a cross-link +/// peer whose split-horizon `inbound_filter` carries the same upward +/// coverage keeps the union nearly intact across the drop. This builds a +/// node with a parent and one healthy cross-link, records the estimate, +/// removes the parent, and asserts the estimate does not collapse. +#[test] +fn compute_mesh_size_stable_across_parent_drop_with_cross_link() { + use crate::bloom::BloomFilter; + use crate::peer::ActivePeer; + + let mut node = make_node(); + + let mk = |hi: u8, lo: u8| { + let mut bytes = [0u8; 16]; + bytes[0] = hi; + bytes[1] = lo; + NodeAddr::from_bytes(bytes) + }; + + // Upward coverage: the set of addresses reachable through the rest of + // the mesh (everything outside our local subtree). Both the parent and + // a healthy cross-link advertise this under split-horizon propagation. + let upward: Vec = (0..12u8).map(|i| mk(0x40, i)).collect(); + // The cross-link additionally knows a couple of addresses of its own. + let cross_extra: Vec = (0..2u8).map(|i| mk(0x50, i)).collect(); + + // Parent P: carries the upward coverage. + let parent_identity = make_peer_identity(); + let parent_addr = *parent_identity.node_addr(); + let mut parent_peer = ActivePeer::new(parent_identity, LinkId::new(1), 0); + let mut parent_filter = BloomFilter::new(); + for addr in &upward { + parent_filter.insert(addr); + } + parent_peer.update_filter(parent_filter, 1, 0); + node.peers.insert(parent_addr, parent_peer); + + // Cross-link X: a non-tree peer whose split-horizon filter also carries + // the upward coverage (plus a little of its own). + let cross_identity = make_peer_identity(); + let cross_addr = *cross_identity.node_addr(); + let mut cross_peer = ActivePeer::new(cross_identity, LinkId::new(2), 0); + let mut cross_filter = BloomFilter::new(); + for addr in upward.iter().chain(cross_extra.iter()) { + cross_filter.insert(addr); + } + cross_peer.update_filter(cross_filter, 1, 0); + node.peers.insert(cross_addr, cross_peer); + + // Baseline estimate with both peers present. + node.compute_mesh_size(); + let before = + node.estimated_mesh_size() + .expect("estimator should produce a value with filter data present") as i64; + + // Simulate a parent switch transient: the old parent is dropped before + // the new parent's filter has converged. The cross-link remains. + node.peers.remove(&parent_addr); + node.compute_mesh_size(); + let after = node + .estimated_mesh_size() + .expect("estimator should still produce a value via the cross-link") as i64; + + // The estimate must not collapse: the cross-link still holds the upward + // coverage. Old tree-only behavior would have lost the entire upward + // leg (~12 addrs) and dropped to roughly self alone. Allow a small + // tolerance for bloom rounding and the cross-link's couple extra bits. + let diff = (before - after).abs(); + assert!( + diff <= 2, + "mesh-size estimate collapsed across parent drop: before={before}, after={after} \ + (cross-link should preserve upward coverage)" + ); +} + /// 100-node random graph: bloom filter exchange at scale. #[tokio::test] async fn test_bloom_filter_convergence_100_nodes() { From 2eea20a216fc979cb4a379917ce627d7055060f2 Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Tue, 9 Jun 2026 23:55:48 +0000 Subject: [PATCH 2/2] tcp: drive inbound connection cap from node.limits.max_connections The per-transport TCP inbound cap was hardwired to 256 and never read node.limits.max_connections, so raising max_connections was a silent no-op for inbound TCP. Resolve the effective cap with precedence: explicit per-transport max_inbound_connections, then node-wide max_connections, then the built-in default of 256. Established peers remain bounded node-wide by add_connection, so deriving the per-transport raw-accept ceiling from max_connections does not admit more real peers across multiple transports. Add effective_max_inbound on the TCP transport with a node_max_connections setter wired from create_transports, plus a precedence unit test. --- src/node/mod.rs | 8 ++++- src/transport/tcp/mod.rs | 65 +++++++++++++++++++++++++++++++++++++++- 2 files changed, 71 insertions(+), 2 deletions(-) diff --git a/src/node/mod.rs b/src/node/mod.rs index 8b183ee..3e96354 100644 --- a/src/node/mod.rs +++ b/src/node/mod.rs @@ -832,9 +832,15 @@ impl Node { .map(|(name, config)| (name.map(|s| s.to_string()), config.clone())) .collect(); + // Node-wide connection budget — used as the TCP inbound-cap fallback + // when a TCP instance has no explicit `max_inbound_connections`, so + // raising `node.limits.max_connections` actually raises the inbound + // ceiling rather than being silently capped at the transport default. + let node_max_connections = self.config.node.limits.max_connections; for (name, tcp_config) in tcp_instances { let transport_id = self.allocate_transport_id(); - let tcp = TcpTransport::new(transport_id, name, tcp_config, packet_tx.clone()); + let mut tcp = TcpTransport::new(transport_id, name, tcp_config, packet_tx.clone()); + tcp.set_node_max_connections(node_max_connections); transports.push(TransportHandle::Tcp(tcp)); } diff --git a/src/transport/tcp/mod.rs b/src/transport/tcp/mod.rs index f3c20fe..350c879 100644 --- a/src/transport/tcp/mod.rs +++ b/src/transport/tcp/mod.rs @@ -122,6 +122,10 @@ pub struct TcpTransport { accept_task: Option>, /// Local listener address (after start, if bind_addr configured). local_addr: Option, + /// Node-wide `node.limits.max_connections`, used as the inbound cap + /// fallback when this transport has no explicit `max_inbound_connections`. + /// `None` means "not provided" — fall through to the built-in default. + node_max_connections: Option, /// Transport statistics. stats: Arc, } @@ -144,10 +148,45 @@ impl TcpTransport { packet_tx, accept_task: None, local_addr: None, + node_max_connections: None, stats: Arc::new(TcpStats::new()), } } + /// Set the node-wide `node.limits.max_connections` value. + /// + /// Used as the inbound-cap fallback when this transport instance has no + /// explicit `transports.tcp.*.max_inbound_connections` set, so raising + /// `node.limits.max_connections` actually raises the per-transport TCP + /// accept ceiling instead of silently capping at the built-in default. + pub fn set_node_max_connections(&mut self, max: usize) { + self.node_max_connections = Some(max); + } + + /// Resolve the effective inbound connection cap for the accept loop. + /// + /// Precedence: explicit per-transport `max_inbound_connections` > + /// node-wide `node.limits.max_connections` > built-in default. This is a + /// per-transport *raw-accept* ceiling; the true node-wide peer budget is + /// still enforced downstream by the handshake-phase `max_connections` + /// admission check (`Node::add_connection`), so deriving this ceiling + /// from `max_connections` does not let multiple transports exceed the + /// node-wide total — it only stops the transport from rejecting inbound + /// below the configured node budget. + fn effective_max_inbound(&self) -> usize { + match ( + self.config.max_inbound_connections, + self.node_max_connections, + ) { + // Explicit per-transport key always wins. + (Some(explicit), _) => explicit, + // No per-transport key: fall back to the node-wide budget. + (None, Some(node_max)) => node_max, + // Neither set: the transport's built-in default (256). + (None, None) => self.config.max_inbound_connections(), + } + } + /// Get the instance name (if configured as a named instance). pub fn name(&self) -> Option<&str> { self.name.as_deref() @@ -197,7 +236,7 @@ impl TcpTransport { let stats = self.stats.clone(); let cfg = AcceptConfig { mtu: self.config.mtu(), - max_inbound: self.config.max_inbound_connections(), + max_inbound: self.effective_max_inbound(), nodelay: self.config.nodelay(), keepalive_secs: self.config.keepalive_secs(), recv_buf: self.config.recv_buf_size(), @@ -1156,6 +1195,30 @@ mod tests { transport.stop_async().await.unwrap(); } + #[test] + fn effective_max_inbound_precedence() { + let (tx, _rx) = packet_channel(100); + + // Neither set: built-in transport default (256). + let t = TcpTransport::new(TransportId::new(1), None, make_config(), tx.clone()); + assert_eq!(t.effective_max_inbound(), 256); + + // Node-wide max_connections drives the cap when no per-transport key. + let mut t = TcpTransport::new(TransportId::new(1), None, make_config(), tx.clone()); + t.set_node_max_connections(512); + assert_eq!(t.effective_max_inbound(), 512); + + // Explicit per-transport key wins over the node-wide value. + let cfg = TcpConfig { + bind_addr: Some("127.0.0.1:0".to_string()), + max_inbound_connections: Some(64), + ..Default::default() + }; + let mut t = TcpTransport::new(TransportId::new(1), None, cfg, tx); + t.set_node_max_connections(512); + assert_eq!(t.effective_max_inbound(), 64); + } + #[tokio::test] async fn test_double_start_fails() { let (tx, _rx) = packet_channel(100);