commit 077241db423715c64a608534926c520c4c75ca76
parent f5f8e9aa77dd2fe3ec9937847955a0acfae6a14c
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Mon, 10 Aug 2026 21:23:36 +0200
Improve gossip peer exchange discovery
Diffstat:
3 files changed, 106 insertions(+), 20 deletions(-)
diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs
@@ -46,6 +46,7 @@ const STALE_INBOUND_PEER_RETENTION_MS: u64 = 60 * 60 * 1_000;
const CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5);
const SESSION_SYNC_INTERVAL: Duration = Duration::from_secs(2);
+const PEER_EXCHANGE_INTERVAL: Duration = Duration::from_secs(30);
const JOIN_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5);
const MAX_JOIN_RESPONSE_ENVELOPES: usize = 16;
const MAX_PEER_VERIFICATION_ENVELOPES: usize = 8;
@@ -687,7 +688,7 @@ async fn outbound_session(
let mut reconnect_delay = INITIAL_RECONNECT_DELAY;
loop {
let self_filter_addr = network.self_filter_addr().await;
- if !peer_is_configured_outbound(&network, &peer).await
+ if !peer_is_connectable(&network, &peer).await
|| is_self_peer_address_for(&peer, network.inner.listen_addr, self_filter_addr)
{
network.inner.sessions.lock().await.remove(&peer);
@@ -771,7 +772,7 @@ async fn outbound_session(
let (sender, next_receiver) = mpsc::channel(PEER_QUEUE_SIZE);
receiver = next_receiver;
- if !peer_is_configured_outbound(&network, &peer).await {
+ if !peer_is_connectable(&network, &peer).await {
network.inner.sessions.lock().await.remove(&peer);
return;
}
@@ -809,6 +810,10 @@ async fn session_loop(
Instant::now() + SESSION_SYNC_INTERVAL,
SESSION_SYNC_INTERVAL,
);
+ let mut peer_exchange_tick = interval_at(
+ Instant::now() + PEER_EXCHANGE_INTERVAL,
+ PEER_EXCHANGE_INTERVAL,
+ );
let mut outbound_closed = false;
let mut peer_status: Option<PeerStatus> = None;
let is_outbound_session = stable_peer.is_some();
@@ -838,6 +843,7 @@ async fn session_loop(
return Ok(());
}
maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?;
+ write_peer_exchange(&network, &mut writer, &known_peer).await?;
} else if let GossipEnvelope::PeerStatus {
height,
tip_hash,
@@ -848,6 +854,7 @@ async fn session_loop(
record_peer_status(&network, &known_peer, remote_addr, &status).await;
peer_status = Some(status);
maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?;
+ write_peer_exchange(&network, &mut writer, &known_peer).await?;
} else if respond_to_peer_verification_challenge(&network, &mut writer, &envelope)
.await?
{
@@ -897,6 +904,9 @@ async fn session_loop(
}
}
}
+ _ = peer_exchange_tick.tick() => {
+ write_peer_exchange(&network, &mut writer, &known_peer).await?;
+ }
envelope = read_session_envelope(&network, &connection_label, &mut reader) => {
let Some(envelope) = envelope? else {
return Ok(());
@@ -918,6 +928,7 @@ async fn session_loop(
return Ok(());
}
maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?;
+ write_peer_exchange(&network, &mut writer, &known_peer).await?;
continue;
}
if let GossipEnvelope::PeerStatus {
@@ -930,6 +941,7 @@ async fn session_loop(
record_peer_status(&network, &known_peer, remote_addr, &status).await;
peer_status = Some(status);
maybe_request_catchup(&network, &mut writer, peer_status.as_ref().unwrap()).await?;
+ write_peer_exchange(&network, &mut writer, &known_peer).await?;
continue;
}
@@ -968,13 +980,8 @@ async fn respond_to_peer_verification_challenge(
Ok(true)
}
-async fn peer_is_configured_outbound(network: &GossipNetwork, peer: &str) -> bool {
- network
- .inner
- .peers
- .lock()
- .await
- .is_configured_outbound(peer)
+async fn peer_is_connectable(network: &GossipNetwork, peer: &str) -> bool {
+ network.inner.peers.lock().await.is_connectable_peer(peer)
}
async fn process_envelope(
@@ -1355,6 +1362,25 @@ async fn apply_peer_list(
Ok(())
}
+async fn write_peer_exchange(
+ network: &GossipNetwork,
+ writer: &mut OwnedWriteHalf,
+ known_peer: &Option<String>,
+) -> Result<()> {
+ let envelope = network.peer_exchange().await;
+ let GossipEnvelope::PeerList { peers } = &envelope else {
+ return Ok(());
+ };
+ if peers.is_empty() {
+ return Ok(());
+ }
+ write_envelope(writer, &envelope).await?;
+ if let Some(peer) = known_peer {
+ network.inner.peers.lock().await.record_sent(peer, 1);
+ }
+ Ok(())
+}
+
async fn write_payload(writer: &mut OwnedWriteHalf, payload: &[GossipEnvelope]) -> Result<()> {
for envelope in payload {
write_envelope(writer, envelope).await?;
@@ -2396,7 +2422,7 @@ async fn remember_discoverable_advertised_peer(
.peers
.lock()
.await
- .observe_inbound_peer(peer.clone());
+ .add_discovered_peer(peer.clone());
}
*known_peer = Some(peer);
Ok(true)
@@ -3348,7 +3374,14 @@ mod tests {
.iter()
.find(|peer| peer.address == remote_addr.to_string())
.unwrap();
- assert_eq!(peer.direction, PeerDirection::Inbound);
+ assert_eq!(peer.direction, PeerDirection::Discovered);
+ assert!(
+ peers
+ .lock()
+ .await
+ .addresses()
+ .contains(&remote_addr.to_string())
+ );
}
#[tokio::test]
@@ -3513,7 +3546,7 @@ mod tests {
assert_eq!(listed.len(), 1);
let peer = &listed[0];
assert_eq!(peer.address, "142.132.164.59:9444");
- assert_eq!(peer.direction, PeerDirection::Inbound);
+ assert_eq!(peer.direction, PeerDirection::Discovered);
assert_eq!(peer.last_known_height, Some(7));
assert_eq!(peer.messages_received, 0);
@@ -3529,7 +3562,7 @@ mod tests {
assert!(repeated);
assert_eq!(
peers.lock().await.list()[0].direction,
- PeerDirection::Inbound
+ PeerDirection::Discovered
);
}
@@ -3918,6 +3951,40 @@ mod tests {
}
#[tokio::test]
+ async fn peer_exchange_advertises_discovered_listening_peers() {
+ let alice = Wallet::from_seed("px-discovered-alice");
+ let allocations = allocations(std::slice::from_ref(&alice), 1_000);
+ let node = Arc::new(tokio::sync::Mutex::new(node("alice", alice, allocations)));
+ let peers = Arc::new(tokio::sync::Mutex::new(PeerBook::default()));
+ peers
+ .lock()
+ .await
+ .add_discovered_peer("127.0.0.1:9546".to_string());
+ let network = super::GossipNetwork {
+ inner: Arc::new(super::GossipNetworkInner {
+ node,
+ peers,
+ listen_addr: "127.0.0.1:9544".parse().unwrap(),
+ p2p_announce_addr: tokio::sync::Mutex::new(None),
+ node_id: super::new_node_id(),
+ accept_task: tokio::sync::Mutex::new(None),
+ sessions: tokio::sync::Mutex::new(BTreeMap::new()),
+ inbound_limiter: Arc::new(
+ StdMutex::new(super::InboundConnectionLimiter::default()),
+ ),
+ metrics: super::P2pMetricsCounters::default(),
+ }),
+ };
+
+ match network.peer_exchange().await {
+ GossipEnvelope::PeerList { peers } => {
+ assert!(peers.contains(&"127.0.0.1:9546".to_string()));
+ }
+ other => panic!("expected peer list, got {other:?}"),
+ }
+ }
+
+ #[tokio::test]
async fn peer_exchange_advertises_stable_listen_and_known_peers() {
let alice = Wallet::from_seed("px-alice");
let allocations = allocations(std::slice::from_ref(&alice), 1_000);
diff --git a/src/app.rs b/src/app.rs
@@ -2058,11 +2058,22 @@ impl PeerBook {
.peers
.entry(address.clone())
.or_insert_with(|| PeerInfo::new(address, PeerDirection::Outbound));
- if peer.direction == PeerDirection::Inbound {
+ if peer.direction != PeerDirection::Outbound {
peer.direction = PeerDirection::Outbound;
}
}
+ pub fn add_discovered_peer(&mut self, address: impl Into<String>) {
+ let address = address.into();
+ let peer = self
+ .peers
+ .entry(address.clone())
+ .or_insert_with(|| PeerInfo::new(address, PeerDirection::Discovered));
+ if peer.direction == PeerDirection::Inbound {
+ peer.direction = PeerDirection::Discovered;
+ }
+ }
+
pub fn observe_inbound_peer(&mut self, address: impl Into<String>) {
let address = address.into();
self.peers
@@ -2090,6 +2101,10 @@ impl PeerBook {
.or_insert_with(|| PeerInfo::new(to, from_peer.direction.clone()));
if from_peer.direction == PeerDirection::Outbound {
to_peer.direction = PeerDirection::Outbound;
+ } else if from_peer.direction == PeerDirection::Discovered
+ && to_peer.direction == PeerDirection::Inbound
+ {
+ to_peer.direction = PeerDirection::Discovered;
}
to_peer.messages_sent = to_peer
.messages_sent
@@ -2135,7 +2150,7 @@ impl PeerBook {
}
}
- pub fn is_configured_outbound(&self, address: &str) -> bool {
+ pub fn is_connectable_peer(&self, address: &str) -> bool {
self.peers
.get(address)
.is_some_and(|peer| peer.direction != PeerDirection::Inbound)
@@ -2408,6 +2423,7 @@ fn median_i64(mut values: Vec<i64>) -> Option<i64> {
#[serde(rename_all = "snake_case")]
pub enum PeerDirection {
Outbound,
+ Discovered,
Inbound,
}
diff --git a/tests/iuna.rs b/tests/iuna.rs
@@ -2157,17 +2157,20 @@ fn peer_book_tracks_multiple_peers_without_networking() {
}
#[test]
-fn peer_book_reports_only_configured_outbound_peers_as_outbound() {
+fn peer_book_reports_only_connectable_peers_as_outbound() {
let mut peers = PeerBook::from_addresses(vec!["127.0.0.1:9444".to_string()]);
peers.record_received("127.0.0.1:56666", 1);
+ peers.add_discovered_peer("127.0.0.1:9445");
peers.record_inbound_misbehavior("127.0.0.1:57777", "invalid transaction");
- assert!(peers.is_configured_outbound("127.0.0.1:9444"));
- assert!(!peers.is_configured_outbound("127.0.0.1:56666"));
- assert!(!peers.is_configured_outbound("127.0.0.1:57777"));
+ assert!(peers.is_connectable_peer("127.0.0.1:9444"));
+ assert!(peers.is_connectable_peer("127.0.0.1:9445"));
+ assert!(!peers.is_connectable_peer("127.0.0.1:56666"));
+ assert!(!peers.is_connectable_peer("127.0.0.1:57777"));
+ assert!(peers.addresses().contains(&"127.0.0.1:9445".to_string()));
assert!(peers.remove_peer("127.0.0.1:9444"));
- assert!(!peers.is_configured_outbound("127.0.0.1:9444"));
+ assert!(!peers.is_connectable_peer("127.0.0.1:9444"));
}
#[test]