From 2ee00d5d178b16477b55c7ffc55743914d13cb02 Mon Sep 17 00:00:00 2001 From: einliterflasche Date: Fri, 26 Jun 2026 13:59:38 +0200 Subject: [PATCH 1/4] refactor(asb): vendor libp2p connection-limits behaviour verbatim --- Cargo.lock | 3 + swap/Cargo.toml | 3 + swap/src/asb/event_loop.rs | 2 +- swap/src/asb/network.rs | 3 +- swap/src/network.rs | 1 + swap/src/network/connection_limits.rs | 363 ++++++++++++++++++++++++++ swap/src/network/swarm.rs | 2 +- 7 files changed, 374 insertions(+), 3 deletions(-) create mode 100644 swap/src/network/connection_limits.rs diff --git a/Cargo.lock b/Cargo.lock index 9ea461a5df..ec511eed8f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10612,6 +10612,9 @@ dependencies = [ "hyper-util", "jsonrpsee", "libp2p", + "libp2p-core", + "libp2p-identity", + "libp2p-swarm", "libp2p-tor", "mockito", "moka", diff --git a/swap/Cargo.toml b/swap/Cargo.toml index c69156d094..5fdb923e5f 100644 --- a/swap/Cargo.toml +++ b/swap/Cargo.toml @@ -30,6 +30,9 @@ tor-rtcompat = { workspace = true, features = ["tokio"] } # LibP2P libp2p = { workspace = true, features = ["tcp", "yamux", "dns", "noise", "request-response", "ping", "rendezvous", "identify", "macros", "cbor", "json", "tokio", "serde", "rsa", "websocket", "metrics"] } +libp2p-core = "0.41.3" +libp2p-identity = "0.2.13" +libp2p-swarm = "0.44.2" libp2p-tor = { path = "../libp2p-tor", features = ["listen-onion-service"] } # Error handling diff --git a/swap/src/asb/event_loop.rs b/swap/src/asb/event_loop.rs index 1d9498ac71..b2d1e86fb5 100644 --- a/swap/src/asb/event_loop.rs +++ b/swap/src/asb/event_loop.rs @@ -485,7 +485,7 @@ where } SwarmEvent::IncomingConnectionError { send_back_addr: address, error, .. } => { if let libp2p::swarm::ListenError::Denied { cause } = &error { - if let Some(exceeded) = cause.downcast_ref::() { + if let Some(exceeded) = cause.downcast_ref::() { tracing::warn!(%address, error = %exceeded, "Rejected inbound connection to prevent against denial-of-service"); } else { tracing::trace!(%address, "Failed to set up connection with peer: {:?}", error); diff --git a/swap/src/asb/network.rs b/swap/src/asb/network.rs index 4b8b2e34b4..2500ad4404 100644 --- a/swap/src/asb/network.rs +++ b/swap/src/asb/network.rs @@ -145,10 +145,11 @@ pub mod transport { pub mod behaviour { use std::sync::Arc; - use libp2p::{connection_limits, identify, identity, ping, swarm::behaviour::toggle::Toggle}; + use libp2p::{identify, identity, ping, swarm::behaviour::toggle::Toggle}; use swap_p2p::protocols::metered::RequestResponseMetrics; use swap_p2p::{out_event::alice::OutEvent, patches}; + use crate::network::connection_limits; use crate::network::wormhole; use crate::network::wormhole::PeerTrust; use crate::network::wormhole::alice::transport::WormholeChannels; diff --git a/swap/src/network.rs b/swap/src/network.rs index 7b73b996c4..4406bedfa0 100644 --- a/swap/src/network.rs +++ b/swap/src/network.rs @@ -8,6 +8,7 @@ pub use swap_p2p::protocols::rendezvous; pub use swap_p2p::protocols::swap_setup; pub use swap_p2p::protocols::transfer_proof; +pub mod connection_limits; pub mod swarm; pub mod transport; pub mod wormhole; diff --git a/swap/src/network/connection_limits.rs b/swap/src/network/connection_limits.rs new file mode 100644 index 0000000000..7a98dfeb0a --- /dev/null +++ b/swap/src/network/connection_limits.rs @@ -0,0 +1,363 @@ +// Copyright 2023 Protocol Labs. +// +// Permission is hereby granted, free of charge, to any person obtaining a +// copy of this software and associated documentation files (the "Software"), +// to deal in the Software without restriction, including without limitation +// the rights to use, copy, modify, merge, publish, distribute, sublicense, +// and/or sell copies of the Software, and to permit persons to whom the +// Software is furnished to do so, subject to the following conditions: +// +// The above copyright notice and this permission notice shall be included in +// all copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS +// OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING +// FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER +// DEALINGS IN THE SOFTWARE. + +use libp2p_core::{ConnectedPoint, Endpoint, Multiaddr}; +use libp2p_identity::PeerId; +use libp2p_swarm::{ + behaviour::{ConnectionEstablished, DialFailure, ListenFailure}, + dummy, ConnectionClosed, ConnectionDenied, ConnectionId, FromSwarm, NetworkBehaviour, THandler, + THandlerInEvent, THandlerOutEvent, ToSwarm, +}; +use std::collections::{HashMap, HashSet}; +use std::fmt; +use std::task::{Context, Poll}; +use void::Void; + +/// A [`NetworkBehaviour`] that enforces a set of [`ConnectionLimits`]. +/// +/// For these limits to take effect, this needs to be composed into the behaviour tree of your application. +/// +/// If a connection is denied due to a limit, either a [`SwarmEvent::IncomingConnectionError`](libp2p_swarm::SwarmEvent::IncomingConnectionError) +/// or [`SwarmEvent::OutgoingConnectionError`](libp2p_swarm::SwarmEvent::OutgoingConnectionError) will be emitted. +/// The [`ListenError::Denied`](libp2p_swarm::ListenError::Denied) and respectively the [`DialError::Denied`](libp2p_swarm::DialError::Denied) variant +/// contain a [`ConnectionDenied`] type that can be downcast to [`Exceeded`] error if (and only if) **this** +/// behaviour denied the connection. +/// +/// If you employ multiple [`NetworkBehaviour`]s that manage connections, it may also be a different error. +/// +/// # Example +/// +/// ```rust +/// # use libp2p_identify as identify; +/// # use libp2p_ping as ping; +/// # use libp2p_swarm_derive::NetworkBehaviour; +/// # use libp2p_connection_limits as connection_limits; +/// +/// #[derive(NetworkBehaviour)] +/// # #[behaviour(prelude = "libp2p_swarm::derive_prelude")] +/// struct MyBehaviour { +/// identify: identify::Behaviour, +/// ping: ping::Behaviour, +/// limits: connection_limits::Behaviour +/// } +/// ``` +pub struct Behaviour { + limits: ConnectionLimits, + + pending_inbound_connections: HashSet, + pending_outbound_connections: HashSet, + established_inbound_connections: HashSet, + established_outbound_connections: HashSet, + established_per_peer: HashMap>, +} + +impl Behaviour { + pub fn new(limits: ConnectionLimits) -> Self { + Self { + limits, + pending_inbound_connections: Default::default(), + pending_outbound_connections: Default::default(), + established_inbound_connections: Default::default(), + established_outbound_connections: Default::default(), + established_per_peer: Default::default(), + } + } + + /// Returns a mutable reference to [`ConnectionLimits`]. + /// > **Note**: A new limit will not be enforced against existing connections. + pub fn limits_mut(&mut self) -> &mut ConnectionLimits { + &mut self.limits + } +} + +fn check_limit(limit: Option, current: usize, kind: Kind) -> Result<(), ConnectionDenied> { + let limit = limit.unwrap_or(u32::MAX); + let current = current as u32; + + if current >= limit { + return Err(ConnectionDenied::new(Exceeded { limit, kind })); + } + + Ok(()) +} + +/// A connection limit has been exceeded. +#[derive(Debug, Clone, Copy)] +pub struct Exceeded { + limit: u32, + kind: Kind, +} + +impl Exceeded { + pub fn limit(&self) -> u32 { + self.limit + } +} + +impl fmt::Display for Exceeded { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + f, + "connection limit exceeded: at most {} {} are allowed", + self.limit, self.kind + ) + } +} + +#[derive(Debug, Clone, Copy)] +enum Kind { + PendingIncoming, + PendingOutgoing, + EstablishedIncoming, + EstablishedOutgoing, + EstablishedPerPeer, + EstablishedTotal, +} + +impl fmt::Display for Kind { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Kind::PendingIncoming => write!(f, "pending incoming connections"), + Kind::PendingOutgoing => write!(f, "pending outgoing connections"), + Kind::EstablishedIncoming => write!(f, "established incoming connections"), + Kind::EstablishedOutgoing => write!(f, "established outgoing connections"), + Kind::EstablishedPerPeer => write!(f, "established connections per peer"), + Kind::EstablishedTotal => write!(f, "established connections"), + } + } +} + +impl std::error::Error for Exceeded {} + +/// The configurable connection limits. +#[derive(Debug, Clone, Default)] +pub struct ConnectionLimits { + max_pending_incoming: Option, + max_pending_outgoing: Option, + max_established_incoming: Option, + max_established_outgoing: Option, + max_established_per_peer: Option, + max_established_total: Option, +} + +impl ConnectionLimits { + /// Configures the maximum number of concurrently incoming connections being established. + pub fn with_max_pending_incoming(mut self, limit: Option) -> Self { + self.max_pending_incoming = limit; + self + } + + /// Configures the maximum number of concurrently outgoing connections being established. + pub fn with_max_pending_outgoing(mut self, limit: Option) -> Self { + self.max_pending_outgoing = limit; + self + } + + /// Configures the maximum number of concurrent established inbound connections. + pub fn with_max_established_incoming(mut self, limit: Option) -> Self { + self.max_established_incoming = limit; + self + } + + /// Configures the maximum number of concurrent established outbound connections. + pub fn with_max_established_outgoing(mut self, limit: Option) -> Self { + self.max_established_outgoing = limit; + self + } + + /// Configures the maximum number of concurrent established connections (both + /// inbound and outbound). + /// + /// Note: This should be used in conjunction with + /// [`ConnectionLimits::with_max_established_incoming`] to prevent possible + /// eclipse attacks (all connections being inbound). + pub fn with_max_established(mut self, limit: Option) -> Self { + self.max_established_total = limit; + self + } + + /// Configures the maximum number of concurrent established connections per peer, + /// regardless of direction (incoming or outgoing). + pub fn with_max_established_per_peer(mut self, limit: Option) -> Self { + self.max_established_per_peer = limit; + self + } +} + +impl NetworkBehaviour for Behaviour { + type ConnectionHandler = dummy::ConnectionHandler; + type ToSwarm = Void; + + fn handle_pending_inbound_connection( + &mut self, + connection_id: ConnectionId, + _: &Multiaddr, + _: &Multiaddr, + ) -> Result<(), ConnectionDenied> { + check_limit( + self.limits.max_pending_incoming, + self.pending_inbound_connections.len(), + Kind::PendingIncoming, + )?; + + self.pending_inbound_connections.insert(connection_id); + + Ok(()) + } + + fn handle_established_inbound_connection( + &mut self, + connection_id: ConnectionId, + peer: PeerId, + _: &Multiaddr, + _: &Multiaddr, + ) -> Result, ConnectionDenied> { + self.pending_inbound_connections.remove(&connection_id); + + check_limit( + self.limits.max_established_incoming, + self.established_inbound_connections.len(), + Kind::EstablishedIncoming, + )?; + check_limit( + self.limits.max_established_per_peer, + self.established_per_peer + .get(&peer) + .map(|connections| connections.len()) + .unwrap_or(0), + Kind::EstablishedPerPeer, + )?; + check_limit( + self.limits.max_established_total, + self.established_inbound_connections.len() + + self.established_outbound_connections.len(), + Kind::EstablishedTotal, + )?; + + Ok(dummy::ConnectionHandler) + } + + fn handle_pending_outbound_connection( + &mut self, + connection_id: ConnectionId, + _: Option, + _: &[Multiaddr], + _: Endpoint, + ) -> Result, ConnectionDenied> { + check_limit( + self.limits.max_pending_outgoing, + self.pending_outbound_connections.len(), + Kind::PendingOutgoing, + )?; + + self.pending_outbound_connections.insert(connection_id); + + Ok(vec![]) + } + + fn handle_established_outbound_connection( + &mut self, + connection_id: ConnectionId, + peer: PeerId, + _: &Multiaddr, + _: Endpoint, + ) -> Result, ConnectionDenied> { + self.pending_outbound_connections.remove(&connection_id); + + check_limit( + self.limits.max_established_outgoing, + self.established_outbound_connections.len(), + Kind::EstablishedOutgoing, + )?; + check_limit( + self.limits.max_established_per_peer, + self.established_per_peer + .get(&peer) + .map(|connections| connections.len()) + .unwrap_or(0), + Kind::EstablishedPerPeer, + )?; + check_limit( + self.limits.max_established_total, + self.established_inbound_connections.len() + + self.established_outbound_connections.len(), + Kind::EstablishedTotal, + )?; + + Ok(dummy::ConnectionHandler) + } + + fn on_swarm_event(&mut self, event: FromSwarm) { + match event { + FromSwarm::ConnectionClosed(ConnectionClosed { + peer_id, + connection_id, + .. + }) => { + self.established_inbound_connections.remove(&connection_id); + self.established_outbound_connections.remove(&connection_id); + self.established_per_peer + .entry(peer_id) + .or_default() + .remove(&connection_id); + } + FromSwarm::ConnectionEstablished(ConnectionEstablished { + peer_id, + endpoint, + connection_id, + .. + }) => { + match endpoint { + ConnectedPoint::Listener { .. } => { + self.established_inbound_connections.insert(connection_id); + } + ConnectedPoint::Dialer { .. } => { + self.established_outbound_connections.insert(connection_id); + } + } + + self.established_per_peer + .entry(peer_id) + .or_default() + .insert(connection_id); + } + FromSwarm::DialFailure(DialFailure { connection_id, .. }) => { + self.pending_outbound_connections.remove(&connection_id); + } + FromSwarm::ListenFailure(ListenFailure { connection_id, .. }) => { + self.pending_inbound_connections.remove(&connection_id); + } + _ => {} + } + } + + fn on_connection_handler_event( + &mut self, + _id: PeerId, + _: ConnectionId, + event: THandlerOutEvent, + ) { + void::unreachable(event) + } + + fn poll(&mut self, _: &mut Context<'_>) -> Poll>> { + Poll::Pending + } +} diff --git a/swap/src/network/swarm.rs b/swap/src/network/swarm.rs index 76ba4e8d54..b74421c4e1 100644 --- a/swap/src/network/swarm.rs +++ b/swap/src/network/swarm.rs @@ -5,7 +5,7 @@ use crate::{asb, cli}; use anyhow::Result; use arti_client::TorClient; use libp2p::Transport as _; -use libp2p::connection_limits::ConnectionLimits; +use crate::network::connection_limits::ConnectionLimits; use libp2p::core::muxing::StreamMuxerBox; use libp2p::metrics::{BandwidthTransport, Registry}; use libp2p::swarm::NetworkBehaviour; From c9ca4859c4ae310131dce418b3f747921a20994f Mon Sep 17 00:00:00 2001 From: einliterflasche Date: Fri, 26 Jun 2026 14:06:35 +0200 Subject: [PATCH 2/4] feat(asb): exempt honest peers from incoming connection limit --- CHANGELOG.md | 2 + swap-asb/src/main.rs | 1 + swap-env/src/config.rs | 8 + swap-orchestrator/src/main.rs | 3 + swap/src/asb/network.rs | 12 +- swap/src/network/connection_limits.rs | 207 ++++++++++++++++++++++++-- swap/src/network/swarm.rs | 2 + swap/tests/harness/mod.rs | 1 + 8 files changed, 221 insertions(+), 15 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 722dfa7d72..5f2b801c33 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- ASB: Honest peers (those that committed real funds to a recent swap) are now exempt from the incoming connection limit, so an attacker flooding connections from throwaway peer IDs can no longer lock honest peers out and block quotes/swaps. A new `[network]` config option `connection_limit_exemption_freshness_days` controls how recently a peer must have transacted to stay exempt (default: `30`). + ## [4.11.3] - 2026-06-24 ## [4.11.2] - 2026-06-23 diff --git a/swap-asb/src/main.rs b/swap-asb/src/main.rs index e34376cbc4..5a7346609b 100644 --- a/swap-asb/src/main.rs +++ b/swap-asb/src/main.rs @@ -307,6 +307,7 @@ pub async fn main() -> Result<()> { config.tor.wormhole_num_intro_points, config.tor.wormhole_swap_freshness_hours, db.clone(), + config.network.connection_limit_exemption_freshness_days, metrics_registry.as_mut(), )?; diff --git a/swap-env/src/config.rs b/swap-env/src/config.rs index 074cfc8570..07a327af91 100644 --- a/swap-env/src/config.rs +++ b/swap-env/src/config.rs @@ -51,6 +51,12 @@ pub struct Network { /// `/metrics`. When unset, the metrics endpoint is disabled. #[serde(default)] pub prometheus_port: Option, + #[serde(default = "default_connection_limit_exemption_freshness_days")] + pub connection_limit_exemption_freshness_days: u64, +} + +pub fn default_connection_limit_exemption_freshness_days() -> u64 { + 30 } #[derive(Clone, Debug, Deserialize, PartialEq, Eq, Serialize)] @@ -439,6 +445,8 @@ pub fn query_user_for_initial_config_with_network( rendezvous_point: rendezvous_points, external_addresses: vec![], prometheus_port: None, + connection_limit_exemption_freshness_days: + default_connection_limit_exemption_freshness_days(), }, bitcoin: Bitcoin { electrum_rpc_urls, diff --git a/swap-orchestrator/src/main.rs b/swap-orchestrator/src/main.rs index 8704536002..ebb1b60afb 100644 --- a/swap-orchestrator/src/main.rs +++ b/swap-orchestrator/src/main.rs @@ -18,6 +18,7 @@ use std::path::PathBuf; use std::str::FromStr; use swap_env::config::{ Bitcoin, Config, ConfigNotInitialized, Data, Maker, Monero, Network, TorConf, + default_connection_limit_exemption_freshness_days, default_price_ticker_rest_poll_interval_exolix_secs, default_price_ticker_source_enabled, default_price_ticker_validity_duration_secs, }; @@ -347,6 +348,8 @@ fn main() { rendezvous_point: rendezvous_points, external_addresses: vec![], prometheus_port: None, + connection_limit_exemption_freshness_days: + default_connection_limit_exemption_freshness_days(), }, bitcoin: Bitcoin { electrum_rpc_urls: match electrum_server_type { diff --git a/swap/src/asb/network.rs b/swap/src/asb/network.rs index 2500ad4404..beeb740ab3 100644 --- a/swap/src/asb/network.rs +++ b/swap/src/asb/network.rs @@ -156,6 +156,8 @@ pub mod behaviour { use super::*; + const HONEST_PEERS_REFRESH_INTERVAL: Duration = Duration::from_secs(60); + /// A `NetworkBehaviour` that represents an XMR/BTC swap node as Alice. #[derive(NetworkBehaviour)] #[behaviour(out_event = "OutEvent", event_process = false)] @@ -194,6 +196,7 @@ pub mod behaviour { rendezvous_nodes: Vec, connection_limits: connection_limits::ConnectionLimits, trust_provider: Arc, + connection_limit_exemption_freshness_days: u64, wormhole_channels: Option, wormhole_swap_freshness_hours: u64, request_response_metrics: Option, @@ -210,7 +213,7 @@ pub mod behaviour { let wormhole = wormhole_channels.map(|channels| { wormhole::alice::Behaviour::new( &identity, - trust_provider, + Arc::clone(&trust_provider), channels.service_tx, channels.handle_rx, wormhole::alice::Config { @@ -231,7 +234,12 @@ pub mod behaviour { }; Self { - connection_limits: connection_limits::Behaviour::new(connection_limits), + connection_limits: connection_limits::Behaviour::new( + connection_limits, + trust_provider, + connection_limit_exemption_freshness_days, + HONEST_PEERS_REFRESH_INTERVAL, + ), rendezvous: Toggle::from(behaviour), quote: quote::alice(request_response_metrics.clone()), swap_setup: alice::Behaviour::new( diff --git a/swap/src/network/connection_limits.rs b/swap/src/network/connection_limits.rs index 7a98dfeb0a..7ccc4a3562 100644 --- a/swap/src/network/connection_limits.rs +++ b/swap/src/network/connection_limits.rs @@ -27,9 +27,16 @@ use libp2p_swarm::{ }; use std::collections::{HashMap, HashSet}; use std::fmt; +use std::sync::Arc; use std::task::{Context, Poll}; +use std::time::Duration; use void::Void; +use futures::FutureExt; +use futures::future::{BoxFuture, Fuse, FusedFuture, OptionFuture}; + +use crate::network::wormhole::PeerTrust; + /// A [`NetworkBehaviour`] that enforces a set of [`ConnectionLimits`]. /// /// For these limits to take effect, this needs to be composed into the behaviour tree of your application. @@ -66,10 +73,24 @@ pub struct Behaviour { established_inbound_connections: HashSet, established_outbound_connections: HashSet, established_per_peer: HashMap>, + + honest_peers: HashSet, + trust_provider: Arc, + freshness_days: u64, + poll_interval: tokio::time::Interval, + pending_query: OptionFuture>>>, } impl Behaviour { - pub fn new(limits: ConnectionLimits) -> Self { + pub fn new( + limits: ConnectionLimits, + trust_provider: Arc, + freshness_days: u64, + poll_interval: Duration, + ) -> Self { + let mut poll_interval = tokio::time::interval(poll_interval); + poll_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + Self { limits, pending_inbound_connections: Default::default(), @@ -77,6 +98,11 @@ impl Behaviour { established_inbound_connections: Default::default(), established_outbound_connections: Default::default(), established_per_peer: Default::default(), + honest_peers: Default::default(), + trust_provider, + freshness_days, + poll_interval, + pending_query: OptionFuture::from(None), } } @@ -231,11 +257,20 @@ impl NetworkBehaviour for Behaviour { ) -> Result, ConnectionDenied> { self.pending_inbound_connections.remove(&connection_id); - check_limit( - self.limits.max_established_incoming, - self.established_inbound_connections.len(), - Kind::EstablishedIncoming, - )?; + if !self.honest_peers.contains(&peer) { + check_limit( + self.limits.max_established_incoming, + self.established_inbound_connections.len(), + Kind::EstablishedIncoming, + )?; + check_limit( + self.limits.max_established_total, + self.established_inbound_connections.len() + + self.established_outbound_connections.len(), + Kind::EstablishedTotal, + )?; + } + check_limit( self.limits.max_established_per_peer, self.established_per_peer @@ -244,12 +279,6 @@ impl NetworkBehaviour for Behaviour { .unwrap_or(0), Kind::EstablishedPerPeer, )?; - check_limit( - self.limits.max_established_total, - self.established_inbound_connections.len() - + self.established_outbound_connections.len(), - Kind::EstablishedTotal, - )?; Ok(dummy::ConnectionHandler) } @@ -357,7 +386,159 @@ impl NetworkBehaviour for Behaviour { void::unreachable(event) } - fn poll(&mut self, _: &mut Context<'_>) -> Poll>> { + fn poll(&mut self, cx: &mut Context<'_>) -> Poll>> { + match self.pending_query.poll_unpin(cx) { + Poll::Ready(Some(peers)) => { + self.honest_peers = peers.into_iter().collect(); + } + Poll::Ready(None) | Poll::Pending => {} + } + + if self.pending_query.is_terminated() && self.poll_interval.poll_tick(cx).is_ready() { + let trust_provider = Arc::clone(&self.trust_provider); + let freshness_hours = self.freshness_days.saturating_mul(24); + let fut: Fuse>> = async move { + match trust_provider + .peers_with_financially_relevant_swap(freshness_hours) + .await + { + Ok(peers) => peers, + Err(e) => { + tracing::warn!(error = ?e, "Failed to query honest peers for connection-limit exemption"); + Vec::new() + } + } + } + .boxed() + .fuse(); + self.pending_query = OptionFuture::from(Some(fut)); + + cx.waker().wake_by_ref(); + } + Poll::Pending } } + +#[cfg(test)] +mod tests { + use super::*; + use anyhow::Result; + + struct FixedTrust(Vec); + + #[async_trait::async_trait] + impl PeerTrust for FixedTrust { + async fn peers_with_financially_relevant_swap( + &self, + _freshness_hours: u64, + ) -> Result> { + Ok(self.0.clone()) + } + } + + fn addr() -> Multiaddr { + "/ip4/127.0.0.1/tcp/1".parse().unwrap() + } + + async fn refresh_honest_peers(behaviour: &mut Behaviour) { + let waker = futures::task::noop_waker(); + let mut cx = std::task::Context::from_waker(&waker); + for _ in 0..100 { + let _ = behaviour.poll(&mut cx); + if !behaviour.honest_peers.is_empty() { + return; + } + tokio::time::sleep(Duration::from_millis(1)).await; + } + panic!("honest peers were never refreshed"); + } + + #[tokio::test] + async fn honest_peer_is_exempt_from_incoming_limit() { + let honest = PeerId::random(); + let unknown = PeerId::random(); + + let limits = ConnectionLimits::default().with_max_established_incoming(Some(0)); + let mut behaviour = Behaviour::new( + limits, + Arc::new(FixedTrust(vec![honest])), + 24, + Duration::from_secs(60), + ); + refresh_honest_peers(&mut behaviour).await; + + assert!( + behaviour + .handle_established_inbound_connection( + ConnectionId::new_unchecked(0), + honest, + &addr(), + &addr(), + ) + .is_ok(), + "honest peer must be exempt from the incoming connection limit" + ); + assert!( + behaviour + .handle_established_inbound_connection( + ConnectionId::new_unchecked(1), + unknown, + &addr(), + &addr(), + ) + .is_err(), + "unknown peer must still be subject to the incoming connection limit" + ); + } + + #[tokio::test] + async fn per_peer_limit_applies_even_to_honest_peers() { + let honest = PeerId::random(); + + let limits = ConnectionLimits::default() + .with_max_established_incoming(Some(0)) + .with_max_established_per_peer(Some(1)); + let mut behaviour = Behaviour::new( + limits, + Arc::new(FixedTrust(vec![honest])), + 24, + Duration::from_secs(60), + ); + refresh_honest_peers(&mut behaviour).await; + + assert!( + behaviour + .handle_established_inbound_connection( + ConnectionId::new_unchecked(0), + honest, + &addr(), + &addr(), + ) + .is_ok() + ); + let endpoint = ConnectedPoint::Listener { + local_addr: addr(), + send_back_addr: addr(), + }; + behaviour.on_swarm_event(FromSwarm::ConnectionEstablished(ConnectionEstablished { + peer_id: honest, + connection_id: ConnectionId::new_unchecked(0), + endpoint: &endpoint, + failed_addresses: &[], + other_established: 0, + })); + + assert!( + behaviour + .handle_established_inbound_connection( + ConnectionId::new_unchecked(1), + honest, + &addr(), + &addr(), + ) + .is_err(), + "per-peer limit must apply even to honest peers" + ); + } +} diff --git a/swap/src/network/swarm.rs b/swap/src/network/swarm.rs index b74421c4e1..b71bab26bb 100644 --- a/swap/src/network/swarm.rs +++ b/swap/src/network/swarm.rs @@ -44,6 +44,7 @@ pub fn asb( wormhole_num_intro_points: u8, wormhole_swap_freshness_hours: u64, trust_provider: Arc, + connection_limit_exemption_freshness_days: u64, metrics_registry: Option<&mut Registry>, ) -> Result<( Swarm>, @@ -105,6 +106,7 @@ where rendezvous_nodes, connection_limits, trust_provider, + connection_limit_exemption_freshness_days, // Passing None disables the wormhole behaviour entirely. if wormhole_enabled { wormhole_channels diff --git a/swap/tests/harness/mod.rs b/swap/tests/harness/mod.rs index a078a497a7..6ce72b2d83 100644 --- a/swap/tests/harness/mod.rs +++ b/swap/tests/harness/mod.rs @@ -433,6 +433,7 @@ async fn start_alice( 3, 168, db.clone(), + 30, None, ) .unwrap(); From 8a3569f98d4f56a97b597512c55615fe41f0123a Mon Sep 17 00:00:00 2001 From: einliterflasche Date: Fri, 26 Jun 2026 15:11:32 +0200 Subject: [PATCH 3/4] fix(asb): parse/format entered_at robustly so fresh-swap filtering stops dropping peers --- ...3907f5bdf8b433036c4cca93772a1e3bfa3e2.json | 28 - ...8f5fdc638faa0ddffc9e6022a3cc38985c669.json | 38 ++ swap/Cargo.toml | 2 +- swap/src/database.rs | 2 + swap/src/database/entered_at.rs | 617 ++++++++++++++++++ swap/src/database/sqlite.rs | 82 ++- 6 files changed, 717 insertions(+), 52 deletions(-) delete mode 100644 .sqlx/query-32cafb8a1dc70149d335cd8603e3907f5bdf8b433036c4cca93772a1e3bfa3e2.json create mode 100644 .sqlx/query-df060a1ba7fb32b1f4d0b533e338f5fdc638faa0ddffc9e6022a3cc38985c669.json create mode 100644 swap/src/database/entered_at.rs diff --git a/.sqlx/query-32cafb8a1dc70149d335cd8603e3907f5bdf8b433036c4cca93772a1e3bfa3e2.json b/.sqlx/query-32cafb8a1dc70149d335cd8603e3907f5bdf8b433036c4cca93772a1e3bfa3e2.json deleted file mode 100644 index 00ec43f8da..0000000000 --- a/.sqlx/query-32cafb8a1dc70149d335cd8603e3907f5bdf8b433036c4cca93772a1e3bfa3e2.json +++ /dev/null @@ -1,28 +0,0 @@ -{ - "db_name": "SQLite", - "query": "\n SELECT s.swap_id, s.state, p.peer_id\n FROM (\n SELECT max(id) as id, swap_id, state, entered_at\n FROM swap_states\n GROUP BY swap_id\n ) s\n INNER JOIN peers p ON s.swap_id = p.swap_id\n WHERE CAST(strftime('%s', substr(s.entered_at, 1, 19)) AS INTEGER)\n >= CAST(strftime('%s', 'now') AS INTEGER) - ?\n ", - "describe": { - "columns": [ - { - "name": "swap_id", - "ordinal": 0, - "type_info": "Text" - }, - { - "name": "state", - "ordinal": 1, - "type_info": "Text" - }, - { - "name": "peer_id", - "ordinal": 2, - "type_info": "Text" - } - ], - "parameters": { - "Right": 1 - }, - "nullable": [false, false, false] - }, - "hash": "32cafb8a1dc70149d335cd8603e3907f5bdf8b433036c4cca93772a1e3bfa3e2" -} diff --git a/.sqlx/query-df060a1ba7fb32b1f4d0b533e338f5fdc638faa0ddffc9e6022a3cc38985c669.json b/.sqlx/query-df060a1ba7fb32b1f4d0b533e338f5fdc638faa0ddffc9e6022a3cc38985c669.json new file mode 100644 index 0000000000..7fd57d9e69 --- /dev/null +++ b/.sqlx/query-df060a1ba7fb32b1f4d0b533e338f5fdc638faa0ddffc9e6022a3cc38985c669.json @@ -0,0 +1,38 @@ +{ + "db_name": "SQLite", + "query": "\n SELECT s.swap_id, s.state, s.entered_at, p.peer_id\n FROM (\n SELECT max(id) as id, swap_id, state, entered_at\n FROM swap_states\n GROUP BY swap_id\n ) s\n INNER JOIN peers p ON s.swap_id = p.swap_id\n ORDER BY s.id\n LIMIT ?\n OFFSET ?\n ", + "describe": { + "columns": [ + { + "name": "swap_id", + "ordinal": 0, + "type_info": "Text" + }, + { + "name": "state", + "ordinal": 1, + "type_info": "Text" + }, + { + "name": "entered_at", + "ordinal": 2, + "type_info": "Text" + }, + { + "name": "peer_id", + "ordinal": 3, + "type_info": "Text" + } + ], + "parameters": { + "Right": 2 + }, + "nullable": [ + false, + false, + false, + false + ] + }, + "hash": "df060a1ba7fb32b1f4d0b533e338f5fdc638faa0ddffc9e6022a3cc38985c669" +} diff --git a/swap/Cargo.toml b/swap/Cargo.toml index 5fdb923e5f..be3191110f 100644 --- a/swap/Cargo.toml +++ b/swap/Cargo.toml @@ -84,7 +84,7 @@ swap-p2p = { path = "../swap-p2p" } swap-proptest = { path = "../swap-proptest" } swap-serde = { path = "../swap-serde" } tauri = { version = "2.0", features = ["config-json5"], optional = true, default-features = false } -time = "0.3" +time = { version = "0.3", features = ["formatting", "macros", "parsing"] } toml = { workspace = true } toml_edit = { workspace = true } url = { workspace = true } diff --git a/swap/src/database.rs b/swap/src/database.rs index 116d6f7563..39c64da5e3 100644 --- a/swap/src/database.rs +++ b/swap/src/database.rs @@ -1,5 +1,6 @@ pub use alice::Alice; pub use bob::Bob; +pub use entered_at::{format_entered_at, parse_entered_at}; pub use sqlite::SqliteDatabase; use crate::cli::api::tauri_bindings::TauriHandle; @@ -14,6 +15,7 @@ use swap_fs::ensure_directory_exists; pub use swap_db::alice; pub use swap_db::bob; +mod entered_at; mod sqlite; #[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] diff --git a/swap/src/database/entered_at.rs b/swap/src/database/entered_at.rs new file mode 100644 index 0000000000..2538b0ec50 --- /dev/null +++ b/swap/src/database/entered_at.rs @@ -0,0 +1,617 @@ +//! Reading and writing the timestamps stored alongside swap states. +//! +//! Historically these timestamps were written in a format that does not zero-pad +//! the hour and uses a variable number of subsecond digits, so the stored strings +//! cannot be compared or ordered as text. They must always be parsed back into a +//! point in time first. The parser in this module accepts every variant ever +//! written, while new timestamps are written in a canonical fixed-width form that +//! is lossless, orders correctly even as plain text, and is still understood by +//! all existing readers. + +use anyhow::{Context, Result}; +use time::OffsetDateTime; +use time::format_description::BorrowedFormatItem; +use time::macros::format_description; + +/// Lenient reader format covering every variant of `entered_at` ever written. +const LENIENT_FORMAT: &[BorrowedFormatItem<'_>] = format_description!( + "[year]-[month]-[day] [hour padding:none]:[minute]:[second].[subsecond] \ + [offset_hour padding:none sign:mandatory]:[offset_minute]:[offset_second]" +); + +/// Canonical fixed-width writer format: zero-padded hour, exactly 9 subsecond digits. +/// +/// Unlike the legacy `Display` output this is lexicographically sortable, has a +/// stability guarantee (the `time` crate documents `Display` as unstable) and is +/// understood by SQLite's date functions after stripping the offset trailer. +const CANONICAL_FORMAT: &[BorrowedFormatItem<'_>] = format_description!( + "[year]-[month]-[day] [hour]:[minute]:[second].[subsecond digits:9] \ + [offset_hour sign:mandatory]:[offset_minute]:[offset_second]" +); + +/// Parse an `entered_at` value, whether written in the legacy or canonical format. +pub fn parse_entered_at(entered_at: &str) -> Result { + OffsetDateTime::parse(entered_at, LENIENT_FORMAT) + .with_context(|| format!("Failed to parse entered_at timestamp {entered_at:?}")) +} + +/// Render a timestamp in the canonical `entered_at` format used for new rows. +pub fn format_entered_at(timestamp: OffsetDateTime) -> Result { + timestamp + .format(CANONICAL_FORMAT) + .context("Failed to format entered_at timestamp") +} + +#[cfg(test)] +mod tests { + use super::*; + use time::{Date, Month, PrimitiveDateTime, Time}; + + /// (year, month, day, hour, minute, second, nanosecond) + const CORPUS: &[(i32, u8, u8, u8, u8, u8, u32)] = &[ + // --- lexicographic hour traps (same date, single vs double digit) --- + (2024, 8, 19, 9, 5, 0, 0), + (2024, 8, 19, 13, 5, 0, 0), + (2024, 8, 19, 2, 30, 12, 500_000_000), + (2024, 8, 19, 20, 30, 12, 500_000_000), + (2023, 3, 7, 1, 0, 0, 0), + (2023, 3, 7, 10, 0, 0, 0), + (2023, 3, 7, 11, 0, 0, 0), + (2023, 3, 7, 19, 0, 0, 0), + (2022, 11, 15, 3, 45, 9, 250_000_000), + (2022, 11, 15, 23, 45, 9, 250_000_000), + (2025, 6, 1, 0, 0, 0, 0), + (2025, 6, 1, 10, 0, 0, 0), + (2025, 6, 1, 4, 12, 8, 123_000_000), + (2025, 6, 1, 14, 12, 8, 123_000_000), + (2021, 9, 9, 8, 8, 8, 8_000_000), + (2021, 9, 9, 18, 8, 8, 8_000_000), + (2026, 2, 14, 5, 55, 55, 999_999_999), + (2026, 2, 14, 15, 55, 55, 999_999_999), + (2027, 7, 4, 6, 11, 37, 475_038_000), + (2027, 7, 4, 16, 11, 37, 475_038_000), + (2028, 10, 31, 7, 0, 1, 1), + (2028, 10, 31, 17, 0, 1, 1), + (2029, 5, 20, 9, 59, 59, 90_000_000), + (2029, 5, 20, 21, 59, 59, 90_000_000), + (2030, 12, 25, 2, 22, 22, 22_000_000), + (2030, 12, 25, 12, 22, 22, 22_000_000), + // --- fraction width 1..9 coverage --- + (2024, 1, 10, 5, 0, 0, 500_000_000), + (2024, 1, 10, 5, 0, 0, 250_000_000), + (2024, 1, 10, 5, 0, 0, 123_000_000), + (2024, 1, 10, 5, 0, 0, 123_400_000), + (2024, 1, 10, 5, 0, 0, 123_450_000), + (2024, 1, 10, 5, 0, 0, 123_456_000), + (2024, 1, 10, 5, 0, 0, 123_456_700), + (2024, 1, 10, 5, 0, 0, 123_456_780), + (2024, 1, 10, 5, 0, 0, 123_456_789), + (2024, 1, 10, 5, 0, 0, 0), + (2024, 1, 10, 5, 0, 0, 1), + (2024, 1, 10, 5, 0, 0, 999_999_999), + (2024, 1, 10, 5, 0, 0, 100_000_000), + (2024, 1, 10, 5, 0, 0, 10_000_000), + (2024, 1, 10, 5, 0, 0, 1_000_000), + (2024, 1, 10, 5, 0, 0, 90_000_000), + (2024, 1, 10, 5, 0, 0, 9_000_000), + // --- fraction lexicographic stress (same date+time, differ in nanos) --- + (2025, 4, 18, 13, 13, 13, 500_000_000), + (2025, 4, 18, 13, 13, 13, 450_000_000), + (2025, 4, 18, 13, 13, 13, 499_999_999), + (2025, 4, 18, 13, 13, 13, 100_000_000), + (2025, 4, 18, 13, 13, 13, 90_000_000), + (2025, 4, 18, 13, 13, 13, 0), + (2025, 4, 18, 13, 13, 13, 1), + // --- boundary: year rollover adjacency --- + (2021, 12, 31, 23, 59, 59, 999_999_999), + (2022, 1, 1, 0, 0, 0, 0), + (2022, 12, 31, 23, 59, 59, 999_999_999), + (2023, 1, 1, 0, 0, 0, 0), + (2023, 12, 31, 23, 59, 59, 999_999_999), + (2024, 1, 1, 0, 0, 0, 0), + (2029, 12, 31, 23, 59, 59, 999_999_999), + (2030, 1, 1, 0, 0, 0, 0), + // --- boundary: first instant --- + (2021, 1, 1, 0, 0, 0, 0), + // --- end-of-month -> start-of-next-month adjacency --- + (2023, 4, 30, 23, 59, 59, 999_999_999), + (2023, 5, 1, 0, 0, 0, 0), + (2024, 6, 30, 23, 59, 59, 999_999_999), + (2024, 7, 1, 0, 0, 0, 0), + (2025, 9, 30, 23, 59, 59, 999_999_999), + (2025, 10, 1, 0, 0, 0, 0), + (2026, 11, 30, 23, 59, 59, 999_999_999), + (2026, 12, 1, 0, 0, 0, 0), + (2027, 2, 28, 23, 59, 59, 999_999_999), + (2027, 3, 1, 0, 0, 0, 0), + (2030, 1, 31, 23, 59, 59, 999_999_999), + (2030, 2, 1, 0, 0, 0, 0), + // --- leap days --- + (2024, 2, 29, 3, 0, 0, 0), + (2024, 2, 29, 13, 0, 0, 0), + (2024, 2, 29, 9, 30, 30, 123_456_789), + (2024, 2, 29, 23, 59, 59, 999_999_999), + (2024, 3, 1, 0, 0, 0, 0), + (2028, 2, 29, 1, 15, 0, 0), + (2028, 2, 29, 11, 15, 0, 0), + (2028, 2, 29, 8, 45, 12, 500_000_000), + (2028, 2, 29, 20, 45, 12, 500_000_000), + (2028, 3, 1, 0, 0, 0, 1), + // --- midnight hour 0 vs hour 10 same day --- + (2022, 7, 22, 0, 0, 0, 0), + (2022, 7, 22, 10, 0, 0, 0), + (2022, 7, 22, 0, 30, 0, 250_000_000), + (2022, 7, 22, 10, 30, 0, 250_000_000), + // --- exact duplicates --- + (2025, 8, 8, 8, 8, 8, 808_000_000), + (2025, 8, 8, 8, 8, 8, 808_000_000), + (2023, 6, 6, 6, 6, 6, 600_000_000), + (2023, 6, 6, 6, 6, 6, 600_000_000), + // --- pseudo-random spread --- + (2021, 2, 17, 14, 23, 41, 318_902_000), + (2021, 3, 28, 7, 9, 55, 60_400_000), + (2021, 5, 3, 22, 47, 2, 901_000_000), + (2021, 6, 19, 1, 1, 1, 11_000_000), + (2021, 7, 30, 19, 38, 16, 745_120_000), + (2021, 8, 11, 4, 52, 9, 3_000_000), + (2021, 10, 5, 16, 17, 18, 192_837_465), + (2021, 11, 23, 9, 44, 33, 555_000_000), + (2021, 12, 7, 12, 0, 59, 480_000_000), + (2022, 1, 14, 0, 16, 28, 70_000_000), + (2022, 2, 25, 23, 3, 7, 6_500_000), + (2022, 3, 9, 5, 30, 45, 333_330_000), + (2022, 4, 1, 18, 22, 11, 808_080_000), + (2022, 5, 16, 2, 58, 4, 40_000_000), + (2022, 6, 27, 11, 11, 11, 111_111_111), + (2022, 8, 4, 20, 5, 50, 950_000_000), + (2022, 9, 12, 6, 33, 22, 27_000_000), + (2022, 10, 19, 13, 49, 1, 612_000_000), + (2022, 11, 8, 8, 0, 0, 800_000_000), + (2022, 12, 21, 17, 41, 36, 159_260_000), + (2023, 1, 2, 3, 4, 5, 678_900_000), + (2023, 2, 13, 21, 19, 47, 84_000_000), + (2023, 3, 24, 9, 9, 9, 9_990_000), + (2023, 5, 6, 15, 27, 53, 720_810_000), + (2023, 6, 17, 4, 4, 44, 444_400_000), + (2023, 7, 28, 22, 31, 15, 50_500_000), + (2023, 8, 9, 1, 48, 26, 263_000_000), + (2023, 9, 20, 10, 10, 10, 101_010_101), + (2023, 10, 31, 19, 55, 2, 2_000_000), + (2023, 11, 11, 11, 11, 11, 110_000_000), + (2023, 12, 13, 7, 36, 58, 369_120_000), + (2024, 1, 25, 16, 42, 19, 875_300_000), + (2024, 3, 3, 3, 33, 3, 33_300_000), + (2024, 4, 14, 23, 14, 41, 414_140_000), + (2024, 5, 22, 6, 23, 12, 567_000_000), + (2024, 6, 4, 12, 48, 0, 480_000_000), + (2024, 7, 16, 0, 7, 38, 70_700_000), + (2024, 9, 27, 18, 19, 20, 212_223_000), + (2024, 10, 8, 5, 51, 7, 7_070_000), + (2024, 11, 19, 21, 2, 49, 902_000_000), + (2024, 12, 30, 9, 38, 14, 141_592_653), + (2025, 1, 7, 14, 26, 35, 357_900_000), + (2025, 2, 18, 2, 13, 58, 85_000_000), + (2025, 3, 29, 11, 47, 22, 622_000_000), + (2025, 5, 11, 17, 5, 41, 514_000_000), + (2025, 7, 23, 8, 39, 16, 16_160_000), + (2025, 8, 2, 23, 51, 4, 4_040_000), + (2025, 9, 13, 1, 28, 37, 837_500_000), + (2025, 10, 24, 13, 0, 0, 130_000_000), + (2025, 11, 5, 6, 44, 19, 419_000_000), + (2025, 12, 16, 20, 12, 53, 253_790_000), + (2026, 1, 28, 9, 33, 6, 63_000_000), + (2026, 3, 11, 4, 17, 48, 174_800_000), + (2026, 4, 22, 15, 9, 31, 931_000_000), + (2026, 5, 4, 0, 50, 12, 512_400_000), + (2026, 6, 15, 18, 2, 7, 27_000_000), + (2026, 7, 26, 7, 58, 44, 844_000_000), + (2026, 8, 6, 22, 11, 0, 110_000_000), + (2026, 9, 17, 3, 24, 39, 243_900_000), + (2026, 10, 29, 11, 46, 52, 465_200_000), + (2026, 12, 9, 16, 38, 25, 382_500_000), + (2027, 1, 19, 1, 5, 14, 51_400_000), + (2027, 2, 27, 13, 52, 8, 528_000_000), + (2027, 3, 8, 8, 19, 47, 194_700_000), + (2027, 4, 20, 19, 33, 2, 332_000_000), + (2027, 5, 1, 2, 46, 51, 468_510_000), + (2027, 6, 12, 23, 0, 33, 33_000_000), + (2027, 8, 24, 5, 41, 19, 411_900_000), + (2027, 9, 5, 17, 7, 7, 707_070_000), + (2027, 10, 16, 9, 50, 48, 504_800_000), + (2027, 11, 27, 0, 14, 23, 142_300_000), + (2027, 12, 8, 12, 39, 56, 395_600_000), + (2028, 1, 18, 21, 1, 9, 19_000_000), + (2028, 3, 30, 6, 28, 41, 284_100_000), + (2028, 4, 11, 14, 53, 17, 531_700_000), + (2028, 5, 23, 2, 6, 38, 63_800_000), + (2028, 6, 3, 18, 44, 52, 445_200_000), + (2028, 7, 15, 9, 17, 29, 172_900_000), + (2028, 8, 26, 23, 39, 1, 390_100_000), + (2028, 9, 7, 4, 50, 44, 504_400_000), + (2028, 10, 19, 11, 2, 16, 21_600_000), + (2028, 11, 30, 16, 48, 33, 483_300_000), + (2028, 12, 10, 0, 33, 27, 332_700_000), + (2029, 1, 21, 13, 7, 5, 70_500_000), + (2029, 2, 1, 7, 19, 48, 194_800_000), + (2029, 3, 13, 22, 41, 12, 412_000_000), + (2029, 4, 24, 5, 53, 39, 533_900_000), + (2029, 6, 5, 19, 6, 51, 65_100_000), + (2029, 7, 17, 1, 28, 4, 280_400_000), + (2029, 8, 28, 12, 50, 47, 504_700_000), + (2029, 9, 9, 8, 2, 19, 21_900_000), + (2029, 10, 20, 23, 44, 33, 443_300_000), + (2029, 11, 1, 4, 17, 56, 175_600_000), + (2029, 12, 12, 17, 39, 8, 390_800_000), + (2030, 1, 23, 9, 1, 41, 14_100_000), + (2030, 2, 4, 21, 23, 14, 231_400_000), + (2030, 3, 16, 6, 46, 47, 464_700_000), + (2030, 4, 27, 13, 8, 22, 82_200_000), + (2030, 5, 8, 2, 50, 55, 505_500_000), + (2030, 6, 19, 18, 12, 28, 122_800_000), + (2030, 7, 30, 9, 44, 1, 440_100_000), + (2030, 8, 11, 23, 6, 34, 63_400_000), + (2030, 9, 22, 4, 38, 47, 384_700_000), + (2030, 10, 3, 11, 1, 9, 10_900_000), + (2030, 11, 14, 16, 23, 52, 233_500_000), + (2030, 12, 26, 0, 46, 25, 462_500_000), + // --- more single vs double digit hour traps on shared dates --- + (2021, 4, 10, 1, 0, 0, 0), + (2021, 4, 10, 11, 0, 0, 0), + (2021, 4, 10, 19, 0, 0, 0), + (2022, 5, 14, 2, 30, 0, 0), + (2022, 5, 14, 20, 30, 0, 0), + (2022, 5, 14, 21, 30, 0, 0), + (2023, 8, 3, 3, 15, 0, 0), + (2023, 8, 3, 13, 15, 0, 0), + (2023, 8, 3, 23, 15, 0, 0), + (2024, 9, 17, 4, 0, 0, 0), + (2024, 9, 17, 14, 0, 0, 0), + (2025, 11, 28, 5, 45, 0, 0), + (2025, 11, 28, 15, 45, 0, 0), + (2026, 6, 6, 6, 0, 0, 0), + (2026, 6, 6, 16, 0, 0, 0), + (2027, 10, 10, 7, 7, 0, 0), + (2027, 10, 10, 17, 7, 0, 0), + (2028, 12, 12, 8, 0, 0, 0), + (2028, 12, 12, 18, 0, 0, 0), + (2029, 3, 3, 9, 9, 0, 0), + (2029, 3, 3, 19, 9, 0, 0), + (2030, 7, 7, 0, 0, 1, 0), + (2030, 7, 7, 10, 0, 1, 0), + (2030, 7, 7, 12, 0, 1, 0), + // --- hour 2 vs 20, hour 1 vs 10/11/19 dense pairs --- + (2024, 4, 5, 2, 0, 0, 0), + (2024, 4, 5, 20, 0, 0, 0), + (2024, 4, 5, 1, 0, 0, 0), + (2024, 4, 5, 10, 0, 0, 0), + (2024, 4, 5, 11, 0, 0, 0), + (2024, 4, 5, 19, 0, 0, 0), + (2023, 10, 12, 2, 22, 22, 0), + (2023, 10, 12, 20, 22, 22, 0), + (2023, 10, 12, 1, 11, 11, 0), + (2023, 10, 12, 10, 11, 11, 0), + (2023, 10, 12, 11, 11, 11, 0), + (2023, 10, 12, 19, 11, 11, 0), + // --- minute/second zero-padding traps (single vs double digit) --- + (2025, 7, 9, 8, 5, 5, 0), + (2025, 7, 9, 8, 50, 50, 0), + (2025, 7, 9, 8, 9, 9, 0), + (2025, 7, 9, 8, 19, 19, 0), + (2026, 2, 8, 14, 3, 7, 0), + (2026, 2, 8, 14, 30, 7, 0), + (2026, 2, 8, 14, 3, 7, 70_000_000), + (2026, 2, 8, 14, 3, 7, 700_000_000), + // --- additional fraction edge spread --- + (2027, 1, 1, 0, 0, 0, 999_999_990), + (2027, 1, 1, 0, 0, 0, 999_999_900), + (2027, 1, 1, 0, 0, 0, 999_999_000), + (2027, 1, 1, 0, 0, 0, 999_990_000), + (2027, 1, 1, 0, 0, 0, 999_900_000), + (2027, 1, 1, 0, 0, 0, 999_000_000), + (2027, 1, 1, 0, 0, 0, 990_000_000), + (2027, 1, 1, 0, 0, 0, 900_000_000), + (2028, 1, 1, 0, 0, 0, 50_000_000), + (2028, 1, 1, 0, 0, 0, 5_000_000), + (2028, 1, 1, 0, 0, 0, 500_000), + (2028, 1, 1, 0, 0, 0, 50_000), + (2028, 1, 1, 0, 0, 0, 5_000), + (2028, 1, 1, 0, 0, 0, 500), + (2028, 1, 1, 0, 0, 0, 50), + (2028, 1, 1, 0, 0, 0, 5), + // --- Feb 28 in non-leap years --- + (2021, 2, 28, 9, 0, 0, 0), + (2021, 2, 28, 19, 0, 0, 0), + (2022, 2, 28, 23, 59, 59, 999_999_999), + (2023, 2, 28, 0, 0, 0, 1), + (2025, 2, 28, 12, 30, 30, 303_000_000), + (2026, 2, 28, 6, 6, 6, 60_000_000), + (2027, 2, 28, 18, 18, 18, 180_000_000), + (2029, 2, 28, 3, 33, 33, 333_000_000), + (2030, 2, 28, 21, 21, 21, 210_000_000), + // --- 30-day month last-day validity --- + (2024, 4, 30, 9, 0, 0, 0), + (2024, 6, 30, 13, 0, 0, 0), + (2024, 9, 30, 2, 0, 0, 0), + (2024, 11, 30, 20, 0, 0, 0), + (2025, 4, 30, 0, 0, 0, 0), + (2025, 6, 30, 23, 59, 59, 500_000_000), + (2025, 9, 30, 11, 11, 11, 11_000_000), + (2025, 11, 30, 22, 22, 22, 220_000_000), + // --- 31-day month last-day validity --- + (2024, 1, 31, 8, 0, 0, 0), + (2024, 3, 31, 18, 0, 0, 0), + (2024, 5, 31, 4, 0, 0, 0), + (2024, 7, 31, 14, 0, 0, 0), + (2024, 8, 31, 9, 9, 9, 90_000_000), + (2024, 10, 31, 21, 0, 0, 0), + (2024, 12, 31, 1, 0, 0, 0), + (2024, 12, 31, 23, 0, 0, 0), + // --- more random spread --- + (2021, 1, 15, 13, 24, 35, 246_800_000), + (2021, 6, 8, 7, 7, 7, 770_000_000), + (2021, 9, 21, 2, 14, 9, 14_900_000), + (2022, 3, 17, 19, 3, 28, 328_000_000), + (2022, 7, 4, 11, 11, 0, 0), + (2022, 10, 31, 23, 23, 23, 232_300_000), + (2023, 4, 1, 1, 1, 1, 100_100_000), + (2023, 5, 19, 16, 40, 8, 408_000_000), + (2023, 9, 1, 0, 0, 1, 1_000), + (2023, 11, 22, 22, 22, 22, 222_222_222), + (2024, 2, 14, 14, 14, 14, 141_414_141), + (2024, 5, 5, 5, 5, 5, 555_550_000), + (2024, 8, 19, 6, 11, 37, 475_038_000), + (2024, 11, 11, 11, 11, 11, 111_000_000), + (2025, 3, 7, 7, 17, 27, 717_270_000), + (2025, 6, 30, 6, 3, 9, 639_000_000), + (2025, 10, 13, 20, 31, 42, 314_000_000), + (2026, 1, 1, 1, 1, 1, 11_110_000), + (2026, 4, 4, 4, 44, 4, 44_440_000), + (2026, 8, 18, 18, 1, 8, 818_000_000), + (2026, 12, 31, 23, 59, 0, 590_000_000), + (2027, 5, 25, 5, 25, 27, 527_000_000), + (2027, 9, 9, 9, 19, 29, 919_290_000), + (2027, 11, 3, 13, 3, 31, 133_100_000), + (2028, 6, 16, 16, 6, 6, 166_000_000), + (2028, 8, 8, 18, 8, 8, 188_080_000), + (2028, 10, 2, 2, 20, 2, 220_200_000), + (2029, 1, 9, 19, 9, 1, 191_000_000), + (2029, 7, 7, 17, 7, 17, 717_000_000), + (2029, 11, 11, 1, 11, 1, 110_100_000), + (2030, 2, 28, 18, 28, 2, 282_000_000), + (2030, 5, 15, 15, 5, 15, 515_150_000), + (2030, 9, 1, 9, 0, 9, 909_000_000), + (2030, 12, 31, 23, 59, 59, 999_999_999), + (2021, 7, 12, 12, 21, 7, 127_210_000), + (2022, 8, 23, 3, 32, 8, 332_800_000), + (2023, 12, 24, 18, 6, 12, 186_120_000), + (2024, 10, 6, 9, 49, 16, 949_160_000), + (2025, 1, 31, 13, 31, 1, 133_110_000), + (2026, 3, 22, 2, 23, 22, 223_220_000), + (2027, 6, 18, 16, 8, 18, 168_180_000), + (2028, 4, 9, 4, 49, 4, 494_900_000), + (2029, 5, 27, 17, 25, 27, 175_270_000), + (2030, 6, 14, 6, 41, 14, 641_140_000), + ]; + + /// Raw historical strings paired with their exact components. Includes the + /// hypothetical `+0:00:00` offset variant and zero-padded-hour canonical + /// renderings to prove the lenient parser covers all generations. + const RAW_LEGACY: &[(&str, (i32, u8, u8, u8, u8, u8, u32))] = &[ + ( + "2024-08-19 6:11:37.475038 +00:00:00", + (2024, 8, 19, 6, 11, 37, 475_038_000), + ), + ( + "2023-05-09 7:30:00.5 +00:00:00", + (2023, 5, 9, 7, 30, 0, 500_000_000), + ), + ( + "2022-11-15 3:45:09.25 +00:00:00", + (2022, 11, 15, 3, 45, 9, 250_000_000), + ), + ( + "2025-06-01 4:12:08.123 +00:00:00", + (2025, 6, 1, 4, 12, 8, 123_000_000), + ), + ( + "2024-01-10 5:00:00.1234 +00:00:00", + (2024, 1, 10, 5, 0, 0, 123_400_000), + ), + ( + "2024-01-10 5:00:00.12345 +00:00:00", + (2024, 1, 10, 5, 0, 0, 123_450_000), + ), + ( + "2024-01-10 5:00:00.1234567 +00:00:00", + (2024, 1, 10, 5, 0, 0, 123_456_700), + ), + ( + "2024-01-10 5:00:00.12345678 +00:00:00", + (2024, 1, 10, 5, 0, 0, 123_456_780), + ), + ( + "2024-01-10 5:00:00.123456789 +00:00:00", + (2024, 1, 10, 5, 0, 0, 123_456_789), + ), + ("2023-03-07 19:00:00.0 +00:00:00", (2023, 3, 7, 19, 0, 0, 0)), + ( + "2026-02-14 15:55:55.999999999 +00:00:00", + (2026, 2, 14, 15, 55, 55, 999_999_999), + ), + ("2024-08-19 9:05:00.0 +00:00:00", (2024, 8, 19, 9, 5, 0, 0)), + ( + "2021-09-09 8:08:08.009 +00:00:00", + (2021, 9, 9, 8, 8, 8, 9_000_000), + ), + ( + "2029-05-20 9:59:59.09 +00:00:00", + (2029, 5, 20, 9, 59, 59, 90_000_000), + ), + ( + "2028-10-31 7:00:01.000000001 +00:00:00", + (2028, 10, 31, 7, 0, 1, 1), + ), + ("2021-01-01 0:00:00.0 +00:00:00", (2021, 1, 1, 0, 0, 0, 0)), + ( + "2030-12-31 23:59:59.999999999 +00:00:00", + (2030, 12, 31, 23, 59, 59, 999_999_999), + ), + ( + "2025-04-18 13:13:13.45 +0:00:00", + (2025, 4, 18, 13, 13, 13, 450_000_000), + ), + ( + "2027-07-04 6:11:37.475038 +0:00:00", + (2027, 7, 4, 6, 11, 37, 475_038_000), + ), + ( + "2022-07-22 0:30:00.25 +0:00:00", + (2022, 7, 22, 0, 30, 0, 250_000_000), + ), + ( + "2024-09-17 09:05:01.123400000 +00:00:00", + (2024, 9, 17, 9, 5, 1, 123_400_000), + ), + ( + "2023-08-03 03:15:00.500000 +00:00:00", + (2023, 8, 3, 3, 15, 0, 500_000_000), + ), + ( + "2028-02-29 08:45:12.500000000 +00:00:00", + (2028, 2, 29, 8, 45, 12, 500_000_000), + ), + ( + "2024-02-29 13:00:00.0 +00:00:00", + (2024, 2, 29, 13, 0, 0, 0), + ), + ]; + + fn dt(c: &(i32, u8, u8, u8, u8, u8, u32)) -> OffsetDateTime { + let &(year, month, day, hour, minute, second, nano) = c; + let date = Date::from_calendar_date(year, Month::try_from(month).unwrap(), day) + .unwrap_or_else(|e| panic!("invalid corpus date {c:?}: {e}")); + let time = Time::from_hms_nano(hour, minute, second, nano) + .unwrap_or_else(|e| panic!("invalid corpus time {c:?}: {e}")); + PrimitiveDateTime::new(date, time).assume_utc() + } + + /// Renders each corpus entry exactly the way production wrote it historically. + fn legacy(c: &(i32, u8, u8, u8, u8, u8, u32)) -> String { + dt(c).to_string() + } + + /// Pins the `time` crate's `Display` output so an upgrade that changes it + /// (and thereby the legacy assumptions in this module) fails loudly. + #[test] + fn legacy_display_format_is_what_we_assume() { + assert_eq!( + legacy(&(2024, 8, 19, 6, 11, 37, 475_038_000)), + "2024-08-19 6:11:37.475038 +00:00:00" + ); + assert_eq!( + legacy(&(2026, 6, 9, 13, 5, 1, 0)), + "2026-06-09 13:05:01.0 +00:00:00" + ); + assert_eq!( + legacy(&(2026, 6, 9, 9, 5, 1, 123_400_000)), + "2026-06-09 9:05:01.1234 +00:00:00" + ); + } + + #[test] + fn parses_legacy_and_canonical_renderings_losslessly() { + for c in CORPUS { + let expected = dt(c); + + let legacy = legacy(c); + assert_eq!( + parse_entered_at(&legacy).unwrap(), + expected, + "legacy {legacy:?}" + ); + + let canonical = format_entered_at(expected).unwrap(); + assert_eq!( + parse_entered_at(&canonical).unwrap(), + expected, + "canonical {canonical:?}" + ); + } + } + + #[test] + fn parses_raw_historical_samples() { + for (raw, c) in RAW_LEGACY { + assert_eq!(parse_entered_at(raw).unwrap(), dt(c), "{raw:?}"); + } + } + + /// The core property: sorting `entered_at` strings of *mixed* generations by + /// their parsed value yields chronological order of the underlying instants. + #[test] + fn sorting_by_parsed_value_is_chronological_across_mixed_formats() { + let mut entries: Vec<(String, OffsetDateTime)> = CORPUS + .iter() + .enumerate() + .map(|(i, c)| { + let instant = dt(c); + let rendered = if i % 2 == 0 { + instant.to_string() + } else { + format_entered_at(instant).unwrap() + }; + (rendered, instant) + }) + .chain(RAW_LEGACY.iter().map(|(raw, c)| (raw.to_string(), dt(c)))) + .collect(); + + // Deterministic Fisher-Yates shuffle (no `rand` dependency needed). + let mut state: u64 = 0x9E37_79B9_7F4A_7C15; + for i in (1..entries.len()).rev() { + state = state + .wrapping_mul(6_364_136_223_846_793_005) + .wrapping_add(1_442_695_040_888_963_407); + entries.swap(i, (state >> 33) as usize % (i + 1)); + } + + entries.sort_by_key(|(rendered, _)| parse_entered_at(rendered).unwrap()); + + for window in entries.windows(2) { + let (ref a_str, a) = window[0]; + let (ref b_str, b) = window[1]; + assert!(a <= b, "out of order: {a_str:?} sorted before {b_str:?}"); + } + } + + /// Proves the corpus actually exercises the bug this module fixes: for legacy + /// renderings, lexicographic order disagrees with chronological order. + #[test] + fn corpus_contains_lexicographic_chronological_disagreements() { + let mut disagreements = 0; + for a in CORPUS { + for b in CORPUS { + if (dt(a) < dt(b)) != (legacy(a) < legacy(b)) && dt(a) != dt(b) { + disagreements += 1; + } + } + } + assert!(disagreements > 0, "corpus has no lexicographic traps"); + } + + /// The canonical format is fixed-width, so its lexicographic order *does* + /// match chronological order (useful for SQL-side `ORDER BY` on new rows). + #[test] + fn canonical_renderings_sort_lexicographically() { + let mut instants: Vec = CORPUS.iter().map(dt).collect(); + instants.sort(); + + let rendered: Vec = instants + .iter() + .map(|instant| format_entered_at(*instant).unwrap()) + .collect(); + let mut sorted = rendered.clone(); + sorted.sort(); + + assert_eq!(rendered, sorted); + } +} diff --git a/swap/src/database/sqlite.rs b/swap/src/database/sqlite.rs index 40063c389e..b6e3960a7d 100644 --- a/swap/src/database/sqlite.rs +++ b/swap/src/database/sqlite.rs @@ -313,10 +313,9 @@ impl Database for SqliteDatabase { } async fn insert_latest_state(&self, swap_id: Uuid, state: State) -> Result<()> { - let entered_at = OffsetDateTime::now_utc(); + let entered_at = super::format_entered_at(OffsetDateTime::now_utc())?; let swap = serde_json::to_string(&Swap::from(state))?; - let entered_at = entered_at.to_string(); let swap_id_str = swap_id.to_string(); sqlx::query!( @@ -596,34 +595,71 @@ impl Database for SqliteDatabase { impl SqliteDatabase { /// Like [`Database::all`] but only returns swaps whose latest state /// update happened within the last `freshness_hours`. - /// - /// `entered_at` is stored as a Rust-formatted string - /// (`YYYY-MM-DD HH:MM:SS.fff +00:00:00`). SQLite's date functions don't - /// accept the offset trailer, so we `substr` down to the first 19 - /// characters before parsing with `strftime('%s', ...)`. All writes use - /// `OffsetDateTime::now_utc()`, so dropping the offset is safe. pub async fn all_fresh(&self, freshness_hours: u64) -> Result> { - let freshness_seconds = (freshness_hours as i64).saturating_mul(3600); + const PAGE_SIZE: u32 = 100; + + let cutoff = OffsetDateTime::now_utc() + - time::Duration::seconds((freshness_hours as i64).saturating_mul(3600)); + + let mut fresh = Vec::new(); + let mut offset = 0; + + loop { + let page = self.latest_states_paginated(PAGE_SIZE, offset).await?; + let page_len = page.len(); + + for (peer_id, swap_id, state_json, entered_at) in page { + if entered_at < cutoff { + continue; + } + let state = match serde_json::from_str::(&state_json) { + Ok(a) => State::from(a), + Err(e) => { + tracing::error!(%swap_id, error = ?e, "Failed to deserialize state"); + continue; + } + }; + fresh.push((peer_id, swap_id, state)); + } + + if page_len < PAGE_SIZE as usize { + break; + } + offset += PAGE_SIZE; + } + + Ok(fresh) + } + + async fn latest_states_paginated( + &self, + limit: u32, + offset: u32, + ) -> Result> { + let limit = i32::try_from(limit)?; + let offset = i32::try_from(offset)?; let rows = sqlx::query!( r#" - SELECT s.swap_id, s.state, p.peer_id + SELECT s.swap_id, s.state, s.entered_at, p.peer_id FROM ( SELECT max(id) as id, swap_id, state, entered_at FROM swap_states GROUP BY swap_id ) s INNER JOIN peers p ON s.swap_id = p.swap_id - WHERE CAST(strftime('%s', substr(s.entered_at, 1, 19)) AS INTEGER) - >= CAST(strftime('%s', 'now') AS INTEGER) - ? + ORDER BY s.id + LIMIT ? + OFFSET ? "#, - freshness_seconds, + limit, + offset, ) .fetch_all(&self.pool) .await?; - let result = rows - .iter() + let page = rows + .into_iter() .filter_map(|row| { let swap_id = match Uuid::from_str(&row.swap_id) { Ok(id) => id, @@ -632,26 +668,26 @@ impl SqliteDatabase { return None; } }; - let peer_id = match PeerId::from_str(&row.peer_id) { - Ok(id) => id, + let entered_at = match super::parse_entered_at(&row.entered_at) { + Ok(t) => t, Err(e) => { - tracing::error!(%swap_id, error = ?e, "Failed to parse PeerId"); + tracing::error!(%swap_id, entered_at = %row.entered_at, error = ?e, "Failed to parse entered_at"); return None; } }; - let state = match serde_json::from_str::(&row.state) { - Ok(a) => State::from(a), + let peer_id = match PeerId::from_str(&row.peer_id) { + Ok(id) => id, Err(e) => { - tracing::error!(%swap_id, error = ?e, "Failed to deserialize state"); + tracing::error!(%swap_id, error = ?e, "Failed to parse PeerId"); return None; } }; - Some((peer_id, swap_id, state)) + Some((peer_id, swap_id, row.state, entered_at)) }) .collect(); - Ok(result) + Ok(page) } } From 9651737753b23a886afc2620c586bdfc878049e1 Mon Sep 17 00:00:00 2001 From: einliterflasche Date: Fri, 26 Jun 2026 16:18:44 +0200 Subject: [PATCH 4/4] fix(asb): robust honest-peer pagination & refresh-on-error, configurable exemption cap --- swap-asb/src/main.rs | 1 + swap-env/src/config.rs | 7 +++ swap-orchestrator/src/main.rs | 4 +- swap/src/asb/network.rs | 2 + swap/src/database/sqlite.rs | 33 +++++----- swap/src/network/connection_limits.rs | 87 +++++++++++++++++++++++--- swap/src/network/swarm.rs | 2 + swap/src/network/wormhole/alice/mod.rs | 17 +++-- swap/tests/harness/mod.rs | 1 + 9 files changed, 119 insertions(+), 35 deletions(-) diff --git a/swap-asb/src/main.rs b/swap-asb/src/main.rs index 5a7346609b..a310d00ca2 100644 --- a/swap-asb/src/main.rs +++ b/swap-asb/src/main.rs @@ -308,6 +308,7 @@ pub async fn main() -> Result<()> { config.tor.wormhole_swap_freshness_hours, db.clone(), config.network.connection_limit_exemption_freshness_days, + config.network.connection_limit_exemption_max_peers, metrics_registry.as_mut(), )?; diff --git a/swap-env/src/config.rs b/swap-env/src/config.rs index 07a327af91..3934828ff5 100644 --- a/swap-env/src/config.rs +++ b/swap-env/src/config.rs @@ -53,12 +53,18 @@ pub struct Network { pub prometheus_port: Option, #[serde(default = "default_connection_limit_exemption_freshness_days")] pub connection_limit_exemption_freshness_days: u64, + #[serde(default = "default_connection_limit_exemption_max_peers")] + pub connection_limit_exemption_max_peers: u64, } pub fn default_connection_limit_exemption_freshness_days() -> u64 { 30 } +pub fn default_connection_limit_exemption_max_peers() -> u64 { + 500 +} + #[derive(Clone, Debug, Deserialize, PartialEq, Eq, Serialize)] #[serde(deny_unknown_fields)] pub struct Bitcoin { @@ -447,6 +453,7 @@ pub fn query_user_for_initial_config_with_network( prometheus_port: None, connection_limit_exemption_freshness_days: default_connection_limit_exemption_freshness_days(), + connection_limit_exemption_max_peers: default_connection_limit_exemption_max_peers(), }, bitcoin: Bitcoin { electrum_rpc_urls, diff --git a/swap-orchestrator/src/main.rs b/swap-orchestrator/src/main.rs index ebb1b60afb..df2934525e 100644 --- a/swap-orchestrator/src/main.rs +++ b/swap-orchestrator/src/main.rs @@ -18,7 +18,7 @@ use std::path::PathBuf; use std::str::FromStr; use swap_env::config::{ Bitcoin, Config, ConfigNotInitialized, Data, Maker, Monero, Network, TorConf, - default_connection_limit_exemption_freshness_days, + default_connection_limit_exemption_freshness_days, default_connection_limit_exemption_max_peers, default_price_ticker_rest_poll_interval_exolix_secs, default_price_ticker_source_enabled, default_price_ticker_validity_duration_secs, }; @@ -350,6 +350,8 @@ fn main() { prometheus_port: None, connection_limit_exemption_freshness_days: default_connection_limit_exemption_freshness_days(), + connection_limit_exemption_max_peers: + default_connection_limit_exemption_max_peers(), }, bitcoin: Bitcoin { electrum_rpc_urls: match electrum_server_type { diff --git a/swap/src/asb/network.rs b/swap/src/asb/network.rs index beeb740ab3..a073ee5ded 100644 --- a/swap/src/asb/network.rs +++ b/swap/src/asb/network.rs @@ -197,6 +197,7 @@ pub mod behaviour { connection_limits: connection_limits::ConnectionLimits, trust_provider: Arc, connection_limit_exemption_freshness_days: u64, + connection_limit_exemption_max_peers: u64, wormhole_channels: Option, wormhole_swap_freshness_hours: u64, request_response_metrics: Option, @@ -238,6 +239,7 @@ pub mod behaviour { connection_limits, trust_provider, connection_limit_exemption_freshness_days, + connection_limit_exemption_max_peers as usize, HONEST_PEERS_REFRESH_INTERVAL, ), rendezvous: Toggle::from(behaviour), diff --git a/swap/src/database/sqlite.rs b/swap/src/database/sqlite.rs index b6e3960a7d..f297941b70 100644 --- a/swap/src/database/sqlite.rs +++ b/swap/src/database/sqlite.rs @@ -605,8 +605,7 @@ impl SqliteDatabase { let mut offset = 0; loop { - let page = self.latest_states_paginated(PAGE_SIZE, offset).await?; - let page_len = page.len(); + let (fetched, page) = self.latest_states_paginated(PAGE_SIZE, offset).await?; for (peer_id, swap_id, state_json, entered_at) in page { if entered_at < cutoff { @@ -622,7 +621,7 @@ impl SqliteDatabase { fresh.push((peer_id, swap_id, state)); } - if page_len < PAGE_SIZE as usize { + if fetched < PAGE_SIZE as usize { break; } offset += PAGE_SIZE; @@ -635,7 +634,7 @@ impl SqliteDatabase { &self, limit: u32, offset: u32, - ) -> Result> { + ) -> Result<(usize, Vec<(PeerId, Uuid, String, OffsetDateTime)>)> { let limit = i32::try_from(limit)?; let offset = i32::try_from(offset)?; @@ -658,6 +657,8 @@ impl SqliteDatabase { .fetch_all(&self.pool) .await?; + let fetched = rows.len(); + let page = rows .into_iter() .filter_map(|row| { @@ -687,7 +688,7 @@ impl SqliteDatabase { }) .collect(); - Ok(page) + Ok((fetched, page)) } } @@ -700,17 +701,19 @@ impl crate::network::wormhole::PeerTrust for SqliteDatabase { use std::collections::HashSet; let swaps = self.all_fresh(freshness_hours).await?; - let peers = swaps - .into_iter() - .filter_map(|(peer_id, _, state)| { - let State::Alice(alice_state) = state else { - return None; - }; - alice_state.is_at_or_past_btc_locked().then_some(peer_id) - }) - .collect::>(); - Ok(peers.into_iter().collect()) + let mut seen = HashSet::new(); + let mut peers = Vec::new(); + for (peer_id, _, state) in swaps.into_iter().rev() { + let State::Alice(alice_state) = state else { + continue; + }; + if alice_state.is_at_or_past_btc_locked() && seen.insert(peer_id) { + peers.push(peer_id); + } + } + + Ok(peers) } } diff --git a/swap/src/network/connection_limits.rs b/swap/src/network/connection_limits.rs index 7ccc4a3562..173990f18b 100644 --- a/swap/src/network/connection_limits.rs +++ b/swap/src/network/connection_limits.rs @@ -77,8 +77,9 @@ pub struct Behaviour { honest_peers: HashSet, trust_provider: Arc, freshness_days: u64, + max_honest_peers: usize, poll_interval: tokio::time::Interval, - pending_query: OptionFuture>>>, + pending_query: OptionFuture>>>>, } impl Behaviour { @@ -86,6 +87,7 @@ impl Behaviour { limits: ConnectionLimits, trust_provider: Arc, freshness_days: u64, + max_honest_peers: usize, poll_interval: Duration, ) -> Self { let mut poll_interval = tokio::time::interval(poll_interval); @@ -101,6 +103,7 @@ impl Behaviour { honest_peers: Default::default(), trust_provider, freshness_days, + max_honest_peers, poll_interval, pending_query: OptionFuture::from(None), } @@ -387,25 +390,23 @@ impl NetworkBehaviour for Behaviour { } fn poll(&mut self, cx: &mut Context<'_>) -> Poll>> { - match self.pending_query.poll_unpin(cx) { - Poll::Ready(Some(peers)) => { - self.honest_peers = peers.into_iter().collect(); - } - Poll::Ready(None) | Poll::Pending => {} + if let Poll::Ready(Some(Some(mut peers))) = self.pending_query.poll_unpin(cx) { + peers.truncate(self.max_honest_peers); + self.honest_peers = peers.into_iter().collect(); } if self.pending_query.is_terminated() && self.poll_interval.poll_tick(cx).is_ready() { let trust_provider = Arc::clone(&self.trust_provider); let freshness_hours = self.freshness_days.saturating_mul(24); - let fut: Fuse>> = async move { + let fut: Fuse>>> = async move { match trust_provider .peers_with_financially_relevant_swap(freshness_hours) .await { - Ok(peers) => peers, + Ok(peers) => Some(peers), Err(e) => { tracing::warn!(error = ?e, "Failed to query honest peers for connection-limit exemption"); - Vec::new() + None } } } @@ -437,6 +438,26 @@ mod tests { } } + struct OkThenErr { + peers: Vec, + calls: std::sync::atomic::AtomicUsize, + } + + #[async_trait::async_trait] + impl PeerTrust for OkThenErr { + async fn peers_with_financially_relevant_swap( + &self, + _freshness_hours: u64, + ) -> Result> { + let n = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + if n == 0 { + Ok(self.peers.clone()) + } else { + Err(anyhow::anyhow!("transient failure")) + } + } + } + fn addr() -> Multiaddr { "/ip4/127.0.0.1/tcp/1".parse().unwrap() } @@ -464,6 +485,7 @@ mod tests { limits, Arc::new(FixedTrust(vec![honest])), 24, + 500, Duration::from_secs(60), ); refresh_honest_peers(&mut behaviour).await; @@ -503,6 +525,7 @@ mod tests { limits, Arc::new(FixedTrust(vec![honest])), 24, + 500, Duration::from_secs(60), ); refresh_honest_peers(&mut behaviour).await; @@ -541,4 +564,50 @@ mod tests { "per-peer limit must apply even to honest peers" ); } + + #[tokio::test] + async fn honest_peer_set_is_capped() { + let peers: Vec = (0..5).map(|_| PeerId::random()).collect(); + let mut behaviour = Behaviour::new( + ConnectionLimits::default(), + Arc::new(FixedTrust(peers)), + 24, + 2, + Duration::from_secs(60), + ); + refresh_honest_peers(&mut behaviour).await; + + assert_eq!(behaviour.honest_peers.len(), 2); + } + + #[tokio::test] + async fn honest_peers_retained_when_refresh_query_fails() { + let honest = PeerId::random(); + let trust = Arc::new(OkThenErr { + peers: vec![honest], + calls: std::sync::atomic::AtomicUsize::new(0), + }); + let mut behaviour = Behaviour::new( + ConnectionLimits::default(), + trust, + 24, + 500, + Duration::from_millis(1), + ); + + refresh_honest_peers(&mut behaviour).await; + assert!(behaviour.honest_peers.contains(&honest)); + + let waker = futures::task::noop_waker(); + let mut cx = std::task::Context::from_waker(&waker); + for _ in 0..100 { + let _ = behaviour.poll(&mut cx); + tokio::time::sleep(Duration::from_millis(1)).await; + } + + assert!( + behaviour.honest_peers.contains(&honest), + "honest peers must survive a transient refresh failure" + ); + } } diff --git a/swap/src/network/swarm.rs b/swap/src/network/swarm.rs index b71bab26bb..f769e10006 100644 --- a/swap/src/network/swarm.rs +++ b/swap/src/network/swarm.rs @@ -45,6 +45,7 @@ pub fn asb( wormhole_swap_freshness_hours: u64, trust_provider: Arc, connection_limit_exemption_freshness_days: u64, + connection_limit_exemption_max_peers: u64, metrics_registry: Option<&mut Registry>, ) -> Result<( Swarm>, @@ -107,6 +108,7 @@ where connection_limits, trust_provider, connection_limit_exemption_freshness_days, + connection_limit_exemption_max_peers, // Passing None disables the wormhole behaviour entirely. if wormhole_enabled { wormhole_channels diff --git a/swap/src/network/wormhole/alice/mod.rs b/swap/src/network/wormhole/alice/mod.rs index 08d5e8892d..ca6b063c4b 100644 --- a/swap/src/network/wormhole/alice/mod.rs +++ b/swap/src/network/wormhole/alice/mod.rs @@ -104,7 +104,7 @@ pub struct Behaviour { /// Pending trust provider query result. /// Stored as an OptionFuture so we can poll it directly without Option dance. /// Wrapped with `Fuse` so we can use `is_terminated()` without extra state. - pending_query: OptionFuture>>>, + pending_query: OptionFuture>>>>, } impl Behaviour { @@ -381,28 +381,25 @@ impl NetworkBehaviour for Behaviour { } // Check if a pending trust provider query has completed - match self.pending_query.poll_unpin(cx) { - Poll::Ready(Some(peers)) => { - for peer_id in peers { - self.spawn_service_for_peer(peer_id); - } + if let Poll::Ready(Some(Some(peers))) = self.pending_query.poll_unpin(cx) { + for peer_id in peers { + self.spawn_service_for_peer(peer_id); } - Poll::Ready(None) | Poll::Pending => {} } // Check if it's time to poll the trust provider if self.pending_query.is_terminated() && self.poll_interval.poll_tick(cx).is_ready() { let trust_provider = Arc::clone(&self.trust_provider); let freshness_hours = self.swap_freshness_hours; - let fut: Fuse>> = async move { + let fut: Fuse>>> = async move { match trust_provider .peers_with_financially_relevant_swap(freshness_hours) .await { - Ok(peers) => peers, + Ok(peers) => Some(peers), Err(e) => { tracing::warn!(error = ?e, "Failed to query peers"); - Vec::new() + None } } } diff --git a/swap/tests/harness/mod.rs b/swap/tests/harness/mod.rs index 6ce72b2d83..392c712b37 100644 --- a/swap/tests/harness/mod.rs +++ b/swap/tests/harness/mod.rs @@ -434,6 +434,7 @@ async fn start_alice( 168, db.clone(), 30, + 500, None, ) .unwrap();