From ba698b512b9f88d7e45374b3ae33004b12010c1b Mon Sep 17 00:00:00 2001 From: Bronek Kozicki Date: Wed, 16 Sep 2026 13:44:38 +0100 Subject: [PATCH 1/4] List and count peers by identity Deduplicate peer API responses while retaining each connection until its disconnect event. Use the latest observed connection for each peer's address and direction, tracking arrival order independently of reusable connection handles. Cover duplicate identities, decreasing handles, and both disconnect orders. Check peer response JSON values and invalid-filter HTTP status without requiring exact body bytes or header order. Document the connected-only inventory, null ENRs, and missed connection events before the API's first read in ADR 0004. Assisted-by: Claude:claude-fable-5-1 Assisted-by: Codex:gpt-6-astra --- crates/beacon_api/src/peers.rs | 34 ++++++-- crates/beacon_api/src/routes.rs | 105 +++++++++++++++++++------ docs/adr/0004-sync-materialized-api.md | 19 +++++ 3 files changed, 129 insertions(+), 29 deletions(-) diff --git a/crates/beacon_api/src/peers.rs b/crates/beacon_api/src/peers.rs index e00a3c3c..0f491e90 100644 --- a/crates/beacon_api/src/peers.rs +++ b/crates/beacon_api/src/peers.rs @@ -29,27 +29,47 @@ impl Peer { } } +/// Connection handles are reused and do not establish arrival order. #[derive(Default)] -pub(crate) struct PeerTable(FxHashMap); +pub(crate) struct PeerTable { + connections: FxHashMap, + arrivals: u64, +} impl PeerTable { pub(crate) fn new() -> Self { - Self(FxHashMap::with_capacity_and_hasher(256, Default::default())) + Self { + connections: FxHashMap::with_capacity_and_hasher(256, Default::default()), + arrivals: 0, + } } + pub(crate) fn insert(&mut self, connection: usize, peer: Peer) { - self.0.insert(connection, peer); + self.arrivals += 1; + self.connections.insert(connection, (self.arrivals, peer)); } pub(crate) fn remove(&mut self, connection: usize) { - self.0.remove(&connection); + self.connections.remove(&connection); + } + + fn by_identity(&self) -> Vec<&Peer> { + let mut latest: FxHashMap = FxHashMap::default(); + for (arrival, peer) in self.connections.values() { + let entry = latest.entry(peer.id).or_insert((*arrival, peer)); + if *arrival > entry.0 { + *entry = (*arrival, peer); + } + } + latest.into_values().map(|(_, peer)| peer).collect() } - pub(crate) fn len(&self) -> usize { - self.0.len() + pub(crate) fn connected(&self) -> usize { + self.by_identity().len() } pub(crate) fn matching<'a>(&'a self, filter: &'a PeerFilter) -> impl Iterator { - self.0.values().filter(move |peer| filter.admits(peer)) + self.by_identity().into_iter().filter(move |peer| filter.admits(peer)) } } diff --git a/crates/beacon_api/src/routes.rs b/crates/beacon_api/src/routes.rs index cedd1421..930a69c9 100644 --- a/crates/beacon_api/src/routes.rs +++ b/crates/beacon_api/src/routes.rs @@ -274,7 +274,7 @@ fn peers(req: &Request<'_>, ctx: &ApiCtx, resp: &mut Response<'_>) { } fn peer_count(_req: &Request<'_>, ctx: &ApiCtx, resp: &mut Response<'_>) { - let connected = ctx.peers.len() as u64; + let connected = ctx.peers.connected() as u64; resp.json_body(|json| json.data_envelope(|json| json.peer_count(connected))); } @@ -606,30 +606,55 @@ mod tests { } } + fn peer(connection: usize, secret: u8, inbound: bool) -> Peer { + Peer { + id: Keypair::from_secret(&[secret; 32]).unwrap().peer_id(), + ip: IpBytes::V4([10, 0, 0, connection as u8]), + port: 9000, + inbound, + } + } + fn two_peer_ctx() -> ApiCtx { let mut ctx = anchor_ctx(); for (connection, inbound) in [(1, true), (2, false)] { - ctx.peers.insert(connection, Peer { - id: Keypair::from_secret(&[connection as u8; 32]).unwrap().peer_id(), - ip: IpBytes::V4([10, 0, 0, connection as u8]), - port: 9000, - inbound, - }); + ctx.peers.insert(connection, peer(connection, connection as u8, inbound)); } ctx } - fn peers_json(query: &str) -> serde_json::Value { - let resp = query_get(&Router::new(ROUTES), &two_peer_ctx(), "/eth/v1/node/peers", query); - assert!(resp.starts_with(b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n")); - serde_json::from_slice(body(&resp)).unwrap() + fn peers_response(ctx: &ApiCtx, query: &str) -> Vec { + query_get(&Router::new(ROUTES), ctx, "/eth/v1/node/peers", query) + } + + fn json_ok(resp: &[u8]) -> serde_json::Value { + assert!(resp.starts_with(b"HTTP/1.1 200 OK\r\n")); + let head_end = resp.windows(4).position(|w| w == b"\r\n\r\n").unwrap(); + let headers = std::str::from_utf8(&resp[..head_end]).unwrap(); + assert!(headers.lines().any(|l| l == "Content-Type: application/json"), "{headers}"); + serde_json::from_slice(body(resp)).unwrap() + } + + fn peers_json(ctx: &ApiCtx, query: &str) -> serde_json::Value { + json_ok(&peers_response(ctx, query)) + } + + fn peer_count_json(ctx: &ApiCtx) -> serde_json::Value { + json_ok(&get(&Router::new(ROUTES), ctx, "/eth/v1/node/peer_count")) + } + + fn address(listed: &serde_json::Value) -> &str { + listed["data"][0]["last_seen_p2p_address"].as_str().unwrap() } #[test] fn peers_list_every_connection_and_honour_the_filters() { - let all = peers_json(""); + let ctx = two_peer_ctx(); + let all = peers_json(&ctx, ""); assert_eq!(all["meta"]["count"], 2); - let inbound = peers_json("direction=inbound&state=connected&state=connecting"); + assert_eq!(all["data"].as_array().unwrap().len(), 2); + + let inbound = peers_json(&ctx, "direction=inbound&state=connected&state=connecting"); assert_eq!(inbound["meta"]["count"], 1); let peer = &inbound["data"][0]; let id = peer["peer_id"].as_str().unwrap(); @@ -640,20 +665,56 @@ mod tests { ); assert_eq!(peer["state"], "connected"); assert_eq!(peer["direction"], "inbound"); - assert_eq!(peers_json("state=disconnected")["meta"]["count"], 0); + assert_eq!(peers_json(&ctx, "state=disconnected")["meta"]["count"], 0); + } - let resp = - query_get(&Router::new(ROUTES), &two_peer_ctx(), "/eth/v1/node/peers", "state=x"); - assert_eq!(body(&resp), br#"{"code":400,"message":"invalid state or direction"}"#); + #[test] + fn unknown_state_or_direction_is_a_400() { + let resp = peers_response(&two_peer_ctx(), "state=x"); + assert!(resp.starts_with(b"HTTP/1.1 400 Bad Request\r\n")); + let error: serde_json::Value = serde_json::from_slice(body(&resp)).unwrap(); + assert_eq!(error["code"], 400); + assert!(error["message"].is_string()); } #[test] fn peer_count_reports_connected_peers_only() { - let resp = get(&Router::new(ROUTES), &two_peer_ctx(), "/eth/v1/node/peer_count"); - assert_eq!( - body(&resp), - br#"{"data":{"disconnected":"0","connecting":"0","connected":"2","disconnecting":"0"}}"# - ); + let count = peer_count_json(&two_peer_ctx()); + assert_eq!(count["data"]["connected"], "2"); + for state in ["disconnected", "connecting", "disconnecting"] { + assert_eq!(count["data"][state], "0", "{state}"); + } + } + + /// A reused handle can be smaller than an older connection's handle. + #[test] + fn duplicate_connections_count_one_peer() { + let mut ctx = anchor_ctx(); + ctx.peers.insert(20, peer(20, 7, true)); + ctx.peers.insert(3, peer(3, 7, false)); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "1"); + + let listed = peers_json(&ctx, ""); + assert_eq!(listed["meta"]["count"], 1); + assert!(address(&listed).starts_with("/ip4/10.0.0.3/"), "{}", address(&listed)); + assert_eq!(listed["data"][0]["direction"], "outbound"); + } + + #[test] + fn closing_one_duplicate_connection_keeps_the_peer() { + for (closed, kept) in [(20, "/ip4/10.0.0.3/"), (3, "/ip4/10.0.0.20/")] { + let mut ctx = anchor_ctx(); + ctx.peers.insert(20, peer(20, 7, true)); + ctx.peers.insert(3, peer(3, 7, false)); + + ctx.peers.remove(closed); + let listed = peers_json(&ctx, ""); + assert_eq!(listed["meta"]["count"], 1, "closed {closed}"); + assert!(address(&listed).starts_with(kept), "closed {closed}: {}", address(&listed)); + + ctx.peers.remove(23 - closed); + assert_eq!(peers_json(&ctx, "")["meta"]["count"], 0, "closed both"); + } } #[test] diff --git a/docs/adr/0004-sync-materialized-api.md b/docs/adr/0004-sync-materialized-api.md index 4cb11d3c..7181d9e9 100644 --- a/docs/adr/0004-sync-materialized-api.md +++ b/docs/adr/0004-sync-materialized-api.md @@ -170,6 +170,25 @@ when no clients subscribe to `block_gossip`. If the cache has overwritten a block's bytes, the API logs a warning and emits no event for that request. This does not cancel the publication request. +## Peer inventory + +`/eth/v1/node/peers` and `/eth/v1/node/peer_count` list and count unique peer +identities from the API's observed connection events. Every listed peer is +`connected`; the other state counts are zero. The inventory does not track +disconnected, connecting or disconnecting peers. The `enr` field is `null`, +as permitted by the [peer schema](https://github.com/ethereum/beacon-APIs/blob/master/types/p2p.yaml). + +Multiple connections can share an identity while QUIC drains a duplicate. +The API retains each connection until its disconnect event. Removing one +connection leaves the peer listed while another remains. Among retained +connections, the most recently observed supplies the peer's direction and +`last_seen_p2p_address`, a QUIC multiaddr. Arrival order is tracked separately +because connection handles can be reused. + +Connection events published before the API's first read are missed. A peer +absent for this reason remains absent until the API observes another +connection for that identity. + ## Node status API construction requires a published anchor and seeds node status from it. From 4a8aabd81e45516e1fe3b85e19dc8f3e4469b1d0 Mon Sep 17 00:00:00 2001 From: Bronek Kozicki Date: Thu, 17 Sep 2026 16:39:31 +0100 Subject: [PATCH 2/4] Group peer connections by identity on insertion Keep each identity's connections in arrival order and select the latest survivor for API output. Count identities directly and list peers through an iterator, removing per-query deduplication and its temporary map and vector. Use the identity and handle in disconnect events to remove the closed connection. Preserve the surviving connection's address and direction, including when the newest connection closes first. Extend tests for direction filters, three overlapping connections, and unknown or repeated events. Update ADR 0004's description of connection storage. Assisted-by: Codex:gpt-6-astra --- crates/beacon_api/src/peers.rs | 49 +++++++------- crates/beacon_api/src/routes.rs | 94 ++++++++++++++++++++++++-- crates/beacon_api/src/server.rs | 4 +- docs/adr/0004-sync-materialized-api.md | 4 +- 4 files changed, 117 insertions(+), 34 deletions(-) diff --git a/crates/beacon_api/src/peers.rs b/crates/beacon_api/src/peers.rs index 0f491e90..e27a5c66 100644 --- a/crates/beacon_api/src/peers.rs +++ b/crates/beacon_api/src/peers.rs @@ -1,4 +1,7 @@ -use std::net::{Ipv4Addr, Ipv6Addr}; +use std::{ + collections::hash_map::Entry, + net::{Ipv4Addr, Ipv6Addr}, +}; use rustc_hash::FxHashMap; use silver_common::{Eth2Addr, IpBytes, PeerId}; @@ -29,47 +32,45 @@ impl Peer { } } -/// Connection handles are reused and do not establish arrival order. +struct Connection { + handle: usize, + peer: Peer, +} + #[derive(Default)] pub(crate) struct PeerTable { - connections: FxHashMap, - arrivals: u64, + peers: FxHashMap>, } impl PeerTable { pub(crate) fn new() -> Self { - Self { - connections: FxHashMap::with_capacity_and_hasher(256, Default::default()), - arrivals: 0, - } + Self { peers: FxHashMap::with_capacity_and_hasher(256, Default::default()) } } pub(crate) fn insert(&mut self, connection: usize, peer: Peer) { - self.arrivals += 1; - self.connections.insert(connection, (self.arrivals, peer)); + let connections = self.peers.entry(peer.id).or_default(); + connections.retain(|c| c.handle != connection); + connections.push(Connection { handle: connection, peer }); } - pub(crate) fn remove(&mut self, connection: usize) { - self.connections.remove(&connection); - } - - fn by_identity(&self) -> Vec<&Peer> { - let mut latest: FxHashMap = FxHashMap::default(); - for (arrival, peer) in self.connections.values() { - let entry = latest.entry(peer.id).or_insert((*arrival, peer)); - if *arrival > entry.0 { - *entry = (*arrival, peer); - } + pub(crate) fn remove(&mut self, peer_id: PeerId, connection: usize) { + let Entry::Occupied(mut entry) = self.peers.entry(peer_id) else { return }; + entry.get_mut().retain(|c| c.handle != connection); + if entry.get().is_empty() { + entry.remove(); } - latest.into_values().map(|(_, peer)| peer).collect() } pub(crate) fn connected(&self) -> usize { - self.by_identity().len() + self.peers.len() } pub(crate) fn matching<'a>(&'a self, filter: &'a PeerFilter) -> impl Iterator { - self.by_identity().into_iter().filter(move |peer| filter.admits(peer)) + self.peers + .values() + .filter_map(|connections| connections.last()) + .map(|connection| &connection.peer) + .filter(move |peer| filter.admits(peer)) } } diff --git a/crates/beacon_api/src/routes.rs b/crates/beacon_api/src/routes.rs index 930a69c9..4cc725b0 100644 --- a/crates/beacon_api/src/routes.rs +++ b/crates/beacon_api/src/routes.rs @@ -702,21 +702,101 @@ mod tests { #[test] fn closing_one_duplicate_connection_keeps_the_peer() { - for (closed, kept) in [(20, "/ip4/10.0.0.3/"), (3, "/ip4/10.0.0.20/")] { + for closed in [20, 3] { let mut ctx = anchor_ctx(); - ctx.peers.insert(20, peer(20, 7, true)); - ctx.peers.insert(3, peer(3, 7, false)); - - ctx.peers.remove(closed); + let older = peer(20, 7, true); + let newer = peer(3, 7, false); + let survivor = if closed == 20 { &newer } else { &older }; + let expected_address = survivor.multiaddr(); + let expected_direction = survivor.direction(); + let peer_id = survivor.id; + ctx.peers.insert(20, older); + ctx.peers.insert(3, newer); + + ctx.peers.remove(peer_id, closed); let listed = peers_json(&ctx, ""); assert_eq!(listed["meta"]["count"], 1, "closed {closed}"); - assert!(address(&listed).starts_with(kept), "closed {closed}: {}", address(&listed)); + assert_eq!(address(&listed), expected_address, "closed {closed}"); + assert_eq!(listed["data"][0]["direction"], expected_direction, "closed {closed}"); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "1"); + for direction in ["inbound", "outbound"] { + let count = peers_json(&ctx, &format!("direction={direction}"))["meta"]["count"] + .as_u64() + .unwrap(); + assert_eq!(count, u64::from(direction == expected_direction), "closed {closed}"); + } - ctx.peers.remove(23 - closed); + ctx.peers.remove(peer_id, 23 - closed); assert_eq!(peers_json(&ctx, "")["meta"]["count"], 0, "closed both"); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "0"); } } + #[test] + fn removing_older_connections_preserves_arrival_order() { + let mut ctx = anchor_ctx(); + let middle = peer(3, 7, false); + let newest = peer(11, 7, true); + let peer_id = newest.id; + let middle_address = middle.multiaddr(); + let newest_address = newest.multiaddr(); + ctx.peers.insert(20, peer(20, 7, true)); + ctx.peers.insert(3, middle); + ctx.peers.insert(11, newest); + + ctx.peers.remove(peer_id, 20); + let listed = peers_json(&ctx, ""); + assert_eq!(listed["meta"]["count"], 1); + assert_eq!(address(&listed), newest_address); + assert_eq!(listed["data"][0]["direction"], "inbound"); + + ctx.peers.remove(peer_id, 11); + let listed = peers_json(&ctx, ""); + assert_eq!(listed["meta"]["count"], 1); + assert_eq!(address(&listed), middle_address); + assert_eq!(listed["data"][0]["direction"], "outbound"); + } + + #[test] + fn unknown_or_repeated_disconnect_keeps_surviving_connections() { + let mut ctx = anchor_ctx(); + let survivor = peer(3, 7, false); + let peer_id = survivor.id; + let expected_address = survivor.multiaddr(); + ctx.peers.insert(3, survivor); + + ctx.peers.remove(peer_id, 20); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "1"); + assert_eq!(address(&peers_json(&ctx, "")), expected_address); + + ctx.peers.insert(20, peer(20, 7, true)); + ctx.peers.remove(peer_id, 20); + ctx.peers.remove(peer_id, 20); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "1"); + assert_eq!(address(&peers_json(&ctx, "")), expected_address); + + ctx.peers.remove(peer_id, 3); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "0"); + } + + #[test] + fn repeated_connection_replaces_its_details() { + let mut ctx = anchor_ctx(); + let replacement = Peer { port: 9001, ..peer(3, 7, false) }; + let peer_id = replacement.id; + let expected_address = replacement.multiaddr(); + ctx.peers.insert(3, peer(3, 7, true)); + ctx.peers.insert(3, replacement); + let listed = peers_json(&ctx, ""); + assert_eq!(listed["meta"]["count"], 1); + assert_eq!(address(&listed), expected_address); + assert_eq!(listed["data"][0]["direction"], "outbound"); + + ctx.peers.remove(peer_id, 3); + assert_eq!(peer_count_json(&ctx)["data"]["connected"], "0"); + assert_eq!(peers_json(&ctx, "")["meta"]["count"], 0); + } + #[test] fn metrics_response_valid_prometheus_format() { let router = Router::new(ROUTES); diff --git a/crates/beacon_api/src/server.rs b/crates/beacon_api/src/server.rs index 1ac912a0..dd9f0da5 100644 --- a/crates/beacon_api/src/server.rs +++ b/crates/beacon_api/src/server.rs @@ -499,7 +499,9 @@ impl BeaconApi { let peer = Peer { id: peer_id_full, ip, port, inbound: !local_dial }; self.ctx.peers.insert(p2p_peer_id, peer); } - PeerEvent::P2pDisconnect { p2p_peer, .. } => self.ctx.peers.remove(p2p_peer), + PeerEvent::P2pDisconnect { p2p_peer, peer_id } => { + self.ctx.peers.remove(peer_id, p2p_peer) + } PeerEvent::SendGossip { topic: GossipTopic::BeaconBlock, ssz, .. } => { self.publish_relayed_block(ssz) } diff --git a/docs/adr/0004-sync-materialized-api.md b/docs/adr/0004-sync-materialized-api.md index 7181d9e9..116a5c92 100644 --- a/docs/adr/0004-sync-materialized-api.md +++ b/docs/adr/0004-sync-materialized-api.md @@ -182,8 +182,8 @@ Multiple connections can share an identity while QUIC drains a duplicate. The API retains each connection until its disconnect event. Removing one connection leaves the peer listed while another remains. Among retained connections, the most recently observed supplies the peer's direction and -`last_seen_p2p_address`, a QUIC multiaddr. Arrival order is tracked separately -because connection handles can be reused. +`last_seen_p2p_address`, a QUIC multiaddr. Connections are grouped by identity +and kept in arrival order; handle values do not establish that order. Connection events published before the API's first read are missed. A peer absent for this reason remains absent until the API observes another From 502fd43691e61c1db9ce3f35e89048380f1d6d07 Mon Sep 17 00:00:00 2001 From: Bronek Kozicki Date: Fri, 18 Sep 2026 10:13:37 +0100 Subject: [PATCH 3/4] Switch to peer connections to SmallVec --- Cargo.lock | 1 + Cargo.toml | 1 + crates/beacon_api/Cargo.toml | 1 + crates/beacon_api/src/peers.rs | 3 ++- 4 files changed, 5 insertions(+), 1 deletion(-) diff --git a/Cargo.lock b/Cargo.lock index 7c76eddc..7f53b3bb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4789,6 +4789,7 @@ dependencies = [ "silver_beacon_state_data", "silver_common", "silver_httpcore", + "smallvec", "tempfile", "toml", "tracing", diff --git a/Cargo.toml b/Cargo.toml index e2ea8cac..3ea731d2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -144,6 +144,7 @@ rcgen = "0.14.7" ring = "0.17" rustc-hash = "2" secp256k1 = { version = "0.30", features = ["global-context", "rand", "hashes"] } +smallvec = "1.15.2" rustls = "0.23.37" serde = { version = "1.0.228", features = ["derive"] } clap = { version = "4", features = ["derive"] } diff --git a/crates/beacon_api/Cargo.toml b/crates/beacon_api/Cargo.toml index 7c80e6ad..2bbc5c6e 100644 --- a/crates/beacon_api/Cargo.toml +++ b/crates/beacon_api/Cargo.toml @@ -14,6 +14,7 @@ silver_common.workspace = true silver_httpcore.workspace = true serde.workspace = true serde_json.workspace = true +smallvec.workspace = true tracing.workspace = true [dev-dependencies] diff --git a/crates/beacon_api/src/peers.rs b/crates/beacon_api/src/peers.rs index e27a5c66..53724840 100644 --- a/crates/beacon_api/src/peers.rs +++ b/crates/beacon_api/src/peers.rs @@ -6,6 +6,7 @@ use std::{ use rustc_hash::FxHashMap; use silver_common::{Eth2Addr, IpBytes, PeerId}; use silver_httpcore::Query; +use smallvec::SmallVec; pub(crate) struct Peer { pub(crate) id: PeerId, @@ -39,7 +40,7 @@ struct Connection { #[derive(Default)] pub(crate) struct PeerTable { - peers: FxHashMap>, + peers: FxHashMap>, } impl PeerTable { From 739d7e0c44355a0c6dafcb268c7acb65c708e03d Mon Sep 17 00:00:00 2001 From: Bronek Kozicki Date: Fri, 18 Sep 2026 12:00:48 +0100 Subject: [PATCH 4/4] Bump pre-allocated PeerTable to 600 --- crates/beacon_api/src/peers.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/beacon_api/src/peers.rs b/crates/beacon_api/src/peers.rs index 53724840..4abe0e43 100644 --- a/crates/beacon_api/src/peers.rs +++ b/crates/beacon_api/src/peers.rs @@ -45,7 +45,8 @@ pub(crate) struct PeerTable { impl PeerTable { pub(crate) fn new() -> Self { - Self { peers: FxHashMap::with_capacity_and_hasher(256, Default::default()) } + // 600 is the default Config::max_connections; preallocate to avoid rehashing. + Self { peers: FxHashMap::with_capacity_and_hasher(600, Default::default()) } } pub(crate) fn insert(&mut self, connection: usize, peer: Peer) {