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 e00a3c3c..4abe0e43 100644 --- a/crates/beacon_api/src/peers.rs +++ b/crates/beacon_api/src/peers.rs @@ -1,8 +1,12 @@ -use std::net::{Ipv4Addr, Ipv6Addr}; +use std::{ + collections::hash_map::Entry, + net::{Ipv4Addr, Ipv6Addr}, +}; 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, @@ -29,27 +33,46 @@ impl Peer { } } +struct Connection { + handle: usize, + peer: Peer, +} + #[derive(Default)] -pub(crate) struct PeerTable(FxHashMap); +pub(crate) struct PeerTable { + peers: FxHashMap>, +} impl PeerTable { pub(crate) fn new() -> Self { - Self(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) { - self.0.insert(connection, 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.0.remove(&connection); + 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(); + } } - pub(crate) fn len(&self) -> usize { - self.0.len() + pub(crate) fn connected(&self) -> usize { + self.peers.len() } pub(crate) fn matching<'a>(&'a self, filter: &'a PeerFilter) -> impl Iterator { - self.0.values().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 cedd1421..4cc725b0 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,136 @@ 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 in [20, 3] { + let mut ctx = anchor_ctx(); + 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_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(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] 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 4cb11d3c..116a5c92 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. 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 +connection for that identity. + ## Node status API construction requires a published anchor and seeds node status from it.