iuna

iuna - experimental devnet protocol
git clone https://iuna.jhx.app/git/iuna.git
Log | Files | Refs | README | LICENSE

commit 42c882aec6d434bbabf5bb10b0e9f16bea0e05c3
parent c9bc44019715361bc1f5b8f594ed4399647f2999
Author: Joris Hartog <jorishartog@hotmail.com>
Date:   Tue, 28 Jul 2026 14:12:03 +0200

Expose transaction gossip rejections

Diffstat:
Msrc/adapters/http.rs | 38+++++++++++++++++++++++++++++++++++++-
Msrc/adapters/p2p.rs | 148++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Msrc/app.rs | 35++++++++++++++++++++++++++++++++++-
3 files changed, 207 insertions(+), 14 deletions(-)

diff --git a/src/adapters/http.rs b/src/adapters/http.rs @@ -822,6 +822,7 @@ fn network_health_at( let last_error = peers.iter().rev().find_map(|peer| { peer.last_error .as_ref() + .or(peer.last_transaction_rejection.as_ref()) .map(|error| format!("{}: {error}", peer.address)) }); @@ -2182,7 +2183,7 @@ const INDEX_HTML: &str = r#"<!doctype html> <td><code x-text="short(peer.last_known_tip_hash)"></code></td> <td x-text="peer.messages_sent"></td> <td x-text="peer.messages_received"></td> - <td x-text="peer.last_error || ''"></td> + <td x-text="peer.last_error || peer.last_transaction_rejection || ''"></td> <td><div class="peer-actions"><button class="peer-remove" type="button" x-show="canRemovePeer(peer)" @click="removePeer(peer)">Remove</button><span class="muted" x-show="!canRemovePeer(peer)">Observed</span></div></td> </tr> </template> @@ -2785,9 +2786,11 @@ mod tests { last_known_height: Some(3), last_known_tip_hash: Some("remote-tip".to_string()), last_error: None, + last_transaction_rejection: None, last_contact_ms: Some(10_000), last_success_ms: Some(10_000), last_error_ms: None, + last_transaction_rejection_ms: None, misbehavior_score: 0, banned_until_ms: None, ban_reason: None, @@ -2808,9 +2811,11 @@ mod tests { last_known_height: None, last_known_tip_hash: None, last_error: Some("connection refused".to_string()), + last_transaction_rejection: None, last_contact_ms: Some(10_000), last_success_ms: None, last_error_ms: Some(10_000), + last_transaction_rejection_ms: None, misbehavior_score: 1, banned_until_ms: None, ban_reason: Some("connection refused".to_string()), @@ -2823,6 +2828,33 @@ mod tests { Some("127.0.0.1:9446: connection refused") ); + let tx_rejection = super::network_health( + &status, + &[PeerInfo { + address: "127.0.0.1:9449".to_string(), + direction: PeerDirection::Outbound, + messages_sent: 1, + messages_received: 1, + last_known_height: Some(0), + last_known_tip_hash: Some("tip".to_string()), + last_error: None, + last_transaction_rejection: Some( + "peer rejected transaction abc: conflict".to_string(), + ), + last_contact_ms: Some(10_000), + last_success_ms: Some(10_000), + last_error_ms: None, + last_transaction_rejection_ms: Some(10_000), + misbehavior_score: 0, + banned_until_ms: None, + ban_reason: None, + }], + ); + assert_eq!( + tx_rejection.last_error.as_deref(), + Some("127.0.0.1:9449: peer rejected transaction abc: conflict") + ); + let stale = super::network_health_at( &status, &[PeerInfo { @@ -2833,9 +2865,11 @@ mod tests { last_known_height: Some(0), last_known_tip_hash: Some("tip".to_string()), last_error: None, + last_transaction_rejection: None, last_contact_ms: Some(1), last_success_ms: Some(1), last_error_ms: None, + last_transaction_rejection_ms: None, misbehavior_score: 0, banned_until_ms: None, ban_reason: None, @@ -2856,9 +2890,11 @@ mod tests { last_known_height: None, last_known_tip_hash: None, last_error: Some("invalid block".to_string()), + last_transaction_rejection: None, last_contact_ms: Some(10), last_success_ms: None, last_error_ms: Some(10), + last_transaction_rejection_ms: None, misbehavior_score: 3, banned_until_ms: Some(1_000), ban_reason: Some("invalid block".to_string()), diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs @@ -44,6 +44,7 @@ const JOIN_RESPONSE_TIMEOUT: Duration = Duration::from_secs(5); const MAX_JOIN_RESPONSE_ENVELOPES: usize = 16; const INITIAL_RECONNECT_DELAY: Duration = Duration::from_secs(1); const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30); +static NEXT_NODE_ID: AtomicU64 = AtomicU64::new(1); #[derive(Clone, Debug, Eq, PartialEq)] struct PeerStatus { @@ -93,6 +94,7 @@ struct GossipNetworkInner { node: SharedNode, peers: SharedPeerBook, listen_addr: SocketAddr, + node_id: String, sessions: Mutex<BTreeMap<String, mpsc::Sender<OutboundBatch>>>, tx_delivery: Mutex<BTreeMap<String, PeerTransactionDelivery>>, metrics: P2pMetricsCounters, @@ -259,6 +261,7 @@ impl GossipNetwork { node, peers, listen_addr: "127.0.0.1:0".parse().unwrap(), + node_id: new_node_id(), sessions: Mutex::new(BTreeMap::new()), tx_delivery: Mutex::new(BTreeMap::new()), metrics: P2pMetricsCounters::default(), @@ -275,6 +278,7 @@ impl GossipNetwork { node, peers, listen_addr: addr, + node_id: new_node_id(), sessions: Mutex::new(BTreeMap::new()), tx_delivery: Mutex::new(BTreeMap::new()), metrics: P2pMetricsCounters::default(), @@ -376,6 +380,21 @@ impl GossipNetwork { peer_delivery.accepted.remove(&rejection.signature); peer_delivery.sent.remove(&rejection.signature); } + let last_rejection = rejected.last().map(|rejection| { + format!( + "peer rejected transaction {}: {}", + short_signature(&rejection.signature), + rejection.reason + ) + }); + drop(delivery); + if let Some(reason) = last_rejection { + self.inner + .peers + .lock() + .await + .record_transaction_rejection(peer, reason); + } } async fn pending_transactions_for_retry(&self, peer: &str) -> Vec<Transaction> { @@ -697,12 +716,10 @@ async fn session_loop( .as_ref() .map(|peer| format!("outbound {peer}")) .unwrap_or_else(|| format!("inbound {remote_addr}")); - let hello = network - .inner - .node - .lock() - .await - .hello(Some(network.inner.listen_addr.to_string())); + let hello = network.inner.node.lock().await.hello( + Some(network.inner.listen_addr.to_string()), + Some(network.inner.node_id.clone()), + ); write_envelope(&mut writer, &hello).await?; let mut reader = LimitedLineReader::new(reader); let mut sync_tick = interval_at( @@ -1010,25 +1027,31 @@ async fn process_transactions( ); } - for rejection in rejected { + for rejection in &rejected { let mut peers = network.inner.peers.lock().await; match known_peer.as_deref() { Some(peer) if transaction_rejection_counts_as_misbehavior(&rejection.reason) => { - peers.record_misbehavior(peer, rejection.reason); + peers.record_misbehavior(peer, rejection.reason.clone()); } Some(peer) => { - peers.record_inbound_error(peer, rejection.reason); + peers.record_inbound_transaction_rejection(peer, rejection.reason.clone()); } None if transaction_rejection_counts_as_misbehavior(&rejection.reason) => { - peers.record_inbound_misbehavior(&remote_addr.to_string(), rejection.reason); + peers + .record_inbound_misbehavior(&remote_addr.to_string(), rejection.reason.clone()); } None => { - peers.record_inbound_error(&remote_addr.to_string(), rejection.reason); + peers.record_inbound_transaction_rejection( + &remote_addr.to_string(), + rejection.reason.clone(), + ); } } } network.forward_outbox().await; - record_inbound_result(network, known_peer, remote_addr, Ok(())).await; + if rejected.is_empty() { + record_inbound_result(network, known_peer, remote_addr, Ok(())).await; + } Ok(()) } @@ -1183,6 +1206,10 @@ fn transactions_in_envelopes(envelopes: &[GossipEnvelope]) -> Vec<(String, Trans transactions } +fn short_signature(signature: &str) -> String { + signature.chars().take(12).collect() +} + fn transaction_rejection_counts_as_misbehavior(reason: &str) -> bool { let reason = reason.to_ascii_lowercase(); if transaction_rejection_is_state_dependent(&reason) { @@ -1727,6 +1754,15 @@ async fn process_hello( NETWORK_ID ); } + if hello + .node_id + .as_deref() + .is_some_and(|node_id| node_id == network.inner.node_id) + { + P2pMetricsCounters::inc(&network.inner.metrics.self_peer_rejections); + forget_stale_self_peer(network, known_peer).await; + return Ok(PeerStatus::new(hello.height, hello.tip_hash)); + } let (local_genesis, local_accepts_remote_genesis) = { let node = network.inner.node.lock().await; ( @@ -1779,6 +1815,19 @@ fn setup_placeholder_genesis_hash() -> String { Ledger::new(BTreeMap::new(), 1).genesis_hash().to_string() } +fn new_node_id() -> String { + let now_nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|duration| duration.as_nanos()) + .unwrap_or_default(); + format!( + "{}-{}-{}", + std::process::id(), + now_nanos, + NEXT_NODE_ID.fetch_add(1, Ordering::Relaxed) + ) +} + async fn remember_discoverable_advertised_peer( network: &GossipNetwork, remote_addr: SocketAddr, @@ -2238,6 +2287,7 @@ mod tests { node: Arc::clone(&node), peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2296,6 +2346,7 @@ mod tests { node: Arc::clone(&node), peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2338,6 +2389,7 @@ mod tests { node: Arc::clone(&node), peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2548,6 +2600,7 @@ mod tests { node, peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2566,6 +2619,7 @@ mod tests { .genesis_hash() .to_string(), listen_addr: Some("127.0.0.1:9545".to_string()), + node_id: None, height: 0, tip_hash: "tip".to_string(), }; @@ -2587,6 +2641,7 @@ mod tests { network_id: NETWORK_ID.to_string(), genesis_hash: "not-local-genesis".to_string(), listen_addr: Some("127.0.0.1:9545".to_string()), + node_id: None, height: 0, tip_hash: "tip".to_string(), }; @@ -2615,6 +2670,7 @@ mod tests { .genesis_hash() .to_string(), listen_addr: Some("127.0.0.1:9545".to_string()), + node_id: None, height: 0, tip_hash: "tip".to_string(), }; @@ -2650,6 +2706,7 @@ mod tests { "iuna.jhx.app:9444".to_string(), ]))), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2669,6 +2726,7 @@ mod tests { network_id: NETWORK_ID.to_string(), genesis_hash: remote_genesis.clone(), listen_addr: Some("142.132.164.59:9444".to_string()), + node_id: None, height: 5, tip_hash: "remote-tip".to_string(), }; @@ -2726,6 +2784,7 @@ mod tests { node: Arc::clone(&node), peers: Arc::new(tokio::sync::Mutex::new(PeerBook::default())), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2737,6 +2796,7 @@ mod tests { network_id: NETWORK_ID.to_string(), genesis_hash: setup_ledger.genesis_hash().to_string(), listen_addr: Some("127.0.0.1:9545".to_string()), + node_id: None, height: 0, tip_hash: setup_ledger.status().tip_hash, }; @@ -2770,6 +2830,7 @@ mod tests { node: Arc::clone(&node), peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2781,6 +2842,7 @@ mod tests { network_id: NETWORK_ID.to_string(), genesis_hash: node.lock().await.ledger().genesis_hash().to_string(), listen_addr: Some("10.42.1.1:12138".to_string()), + node_id: None, height: status.height, tip_hash: status.tip_hash, }; @@ -2810,6 +2872,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2875,6 +2938,7 @@ mod tests { node, peers, listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2901,6 +2965,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "127.0.0.1:9544".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2931,6 +2996,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2964,6 +3030,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -2998,6 +3065,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -3022,6 +3090,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9545".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -3053,6 +3122,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9545".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -3070,6 +3140,7 @@ mod tests { .genesis_hash() .to_string(), listen_addr: Some("127.0.0.1:9545".to_string()), + node_id: None, height: 0, tip_hash: "tip".to_string(), }; @@ -3103,6 +3174,7 @@ mod tests { node, peers: Arc::clone(&peers), listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), sessions: tokio::sync::Mutex::new(BTreeMap::new()), tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), metrics: super::P2pMetricsCounters::default(), @@ -3120,6 +3192,7 @@ mod tests { .genesis_hash() .to_string(), listen_addr: Some("127.0.0.1:9444".to_string()), + node_id: None, height: 0, tip_hash: "tip".to_string(), }; @@ -3139,6 +3212,57 @@ mod tests { assert!(peers.lock().await.addresses().is_empty()); } + #[tokio::test] + async fn hello_removes_outbound_peer_with_same_node_id() { + let alice = Wallet::from_seed("hello-self-node-id-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::from_addresses(vec![ + "142.132.164.59:9444".to_string(), + ]))); + let network = super::GossipNetwork { + inner: Arc::new(super::GossipNetworkInner { + node, + peers: Arc::clone(&peers), + listen_addr: "0.0.0.0:9444".parse().unwrap(), + node_id: super::new_node_id(), + sessions: tokio::sync::Mutex::new(BTreeMap::new()), + tx_delivery: tokio::sync::Mutex::new(BTreeMap::new()), + metrics: super::P2pMetricsCounters::default(), + }), + }; + let hello = ProtocolHello { + protocol_version: PROTOCOL_VERSION, + network_id: NETWORK_ID.to_string(), + genesis_hash: network + .inner + .node + .lock() + .await + .ledger() + .genesis_hash() + .to_string(), + listen_addr: Some("0.0.0.0:9444".to_string()), + node_id: Some(network.inner.node_id.clone()), + height: 0, + tip_hash: "tip".to_string(), + }; + let mut known_peer = Some("142.132.164.59:9444".to_string()); + + super::process_hello( + &network, + "142.132.164.59:52144".parse().unwrap(), + &mut known_peer, + hello, + ) + .await + .unwrap(); + + assert_eq!(network.metrics().self_peer_rejections, 1); + assert!(known_peer.is_none()); + assert!(peers.lock().await.addresses().is_empty()); + } + fn node(_network_key: &str, wallet: Wallet, allocations: BTreeMap<String, Amount>) -> NodeCore { let ledger = Ledger::new_with_genesis_burns( allocations, diff --git a/src/app.rs b/src/app.rs @@ -122,6 +122,8 @@ pub struct ProtocolHello { pub network_id: String, pub genesis_hash: String, pub listen_addr: Option<String>, + #[serde(default)] + pub node_id: Option<String>, pub height: u64, pub tip_hash: String, } @@ -358,13 +360,14 @@ impl NodeCore { self.ledger.snapshot() } - pub fn hello(&self, listen_addr: Option<String>) -> GossipEnvelope { + pub fn hello(&self, listen_addr: Option<String>, node_id: Option<String>) -> GossipEnvelope { let status = self.ledger.status(); GossipEnvelope::Hello(ProtocolHello { protocol_version: PROTOCOL_VERSION, network_id: NETWORK_ID.to_string(), genesis_hash: self.ledger.genesis_hash().to_string(), listen_addr, + node_id, height: status.height, tip_hash: status.tip_hash, }) @@ -1168,6 +1171,12 @@ impl PeerBook { if to_peer.last_error.is_none() { to_peer.last_error = from_peer.last_error; } + if to_peer.last_transaction_rejection.is_none() { + to_peer.last_transaction_rejection = from_peer.last_transaction_rejection; + } + to_peer.last_transaction_rejection_ms = to_peer + .last_transaction_rejection_ms + .max(from_peer.last_transaction_rejection_ms); to_peer.misbehavior_score = to_peer .misbehavior_score .saturating_add(from_peer.misbehavior_score); @@ -1265,6 +1274,26 @@ impl PeerBook { peer.last_error = Some(error.into()); } + pub fn record_transaction_rejection(&mut self, address: &str, reason: impl Into<String>) { + let now = now_ms(); + let peer = self.ensure(address, PeerDirection::Outbound); + peer.last_contact_ms = Some(now); + peer.last_transaction_rejection_ms = Some(now); + peer.last_transaction_rejection = Some(reason.into()); + } + + pub fn record_inbound_transaction_rejection( + &mut self, + address: &str, + reason: impl Into<String>, + ) { + let now = now_ms(); + let peer = self.ensure(address, PeerDirection::Inbound); + peer.last_contact_ms = Some(now); + peer.last_transaction_rejection_ms = Some(now); + peer.last_transaction_rejection = Some(reason.into()); + } + pub fn record_received(&mut self, address: &str, count: u64) { let now = now_ms(); let peer = self.ensure(address, PeerDirection::Inbound); @@ -1334,9 +1363,11 @@ pub struct PeerInfo { pub last_known_height: Option<u64>, pub last_known_tip_hash: Option<String>, pub last_error: Option<String>, + pub last_transaction_rejection: Option<String>, pub last_contact_ms: Option<u64>, pub last_success_ms: Option<u64>, pub last_error_ms: Option<u64>, + pub last_transaction_rejection_ms: Option<u64>, pub misbehavior_score: u32, pub banned_until_ms: Option<u64>, pub ban_reason: Option<String>, @@ -1352,9 +1383,11 @@ impl PeerInfo { last_known_height: None, last_known_tip_hash: None, last_error: None, + last_transaction_rejection: None, last_contact_ms: None, last_success_ms: None, last_error_ms: None, + last_transaction_rejection_ms: None, misbehavior_score: 0, banned_until_ms: None, ban_reason: None,