Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
Expand Down
1 change: 1 addition & 0 deletions crates/beacon_api/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
41 changes: 32 additions & 9 deletions crates/beacon_api/src/peers.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -29,27 +33,46 @@ impl Peer {
}
}

struct Connection {
handle: usize,
peer: Peer,
}

#[derive(Default)]
pub(crate) struct PeerTable(FxHashMap<usize, Peer>);
pub(crate) struct PeerTable {
peers: FxHashMap<PeerId, SmallVec<[Connection; 2]>>,
}

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<Item = &'a Peer> {
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))
}
}

Expand Down
185 changes: 163 additions & 22 deletions crates/beacon_api/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)));
}

Expand Down Expand Up @@ -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<u8> {
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();
Expand All @@ -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]
Expand Down
4 changes: 3 additions & 1 deletion crates/beacon_api/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
19 changes: 19 additions & 0 deletions docs/adr/0004-sync-materialized-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading