Skip to content

Commit 3ae7578

Browse files
committed
Limit pending requests and peers in LSPS1 service
Add per-peer and global rate limiting to `LSPS1ServiceHandler` to prevent resource exhaustion, mirroring the existing LSPS2 pattern. Introduce `MAX_PENDING_REQUESTS_PER_PEER` (10), `MAX_TOTAL_PENDING_REQUESTS` (1000), and `MAX_TOTAL_PEERS` (100000) constants and enforce them in `handle_create_order_request`. Rejected requests receive a `CreateOrderError` with `LSPS0_CLIENT_REJECTED_ERROR_CODE`. A `total_pending_requests` atomic counter tracks the global count, and a `verify_pending_request_counter` debug assertion ensures it stays in sync. Co-Authored-By: HAL 9000
1 parent b37cfaf commit 3ae7578

3 files changed

Lines changed: 325 additions & 14 deletions

File tree

lightning-liquidity/src/lsps1/peer_state.rs

Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -353,6 +353,20 @@ impl PeerState {
353353
self.pending_requests.remove(request_id).ok_or(PeerStateError::UnknownRequestId)
354354
}
355355

356+
pub(super) fn pending_request_count(&self) -> usize {
357+
self.pending_requests.len()
358+
}
359+
360+
pub(super) fn pending_requests_and_channels(&self) -> usize {
361+
let pending_requests = self.pending_requests.len();
362+
let pending_orders = self
363+
.outbound_channels_by_order_id
364+
.iter()
365+
.filter(|(_, v)| v.is_pending_payment())
366+
.count();
367+
pending_requests + pending_orders
368+
}
369+
356370
pub(super) fn has_active_requests(&self) -> bool {
357371
!self.outbound_channels_by_order_id.is_empty()
358372
}
@@ -370,8 +384,10 @@ impl PeerState {
370384
self.pending_requests.is_empty() && self.outbound_channels_by_order_id.is_empty()
371385
}
372386

373-
pub(super) fn prune_pending_requests(&mut self) {
374-
self.pending_requests.clear()
387+
pub(super) fn prune_pending_requests(&mut self) -> usize {
388+
let num_pruned = self.pending_requests.len();
389+
self.pending_requests.clear();
390+
num_pruned
375391
}
376392

377393
pub(super) fn prune_expired_request_state(&mut self) {
@@ -433,6 +449,10 @@ impl ChannelOrder {
433449
self.state.channel_info()
434450
}
435451

452+
fn is_pending_payment(&self) -> bool {
453+
matches!(self.state, ChannelOrderState::ExpectingPayment { .. })
454+
}
455+
436456
fn is_prunable(&self) -> bool {
437457
let all_payment_details_expired;
438458
#[cfg(feature = "time")]

lightning-liquidity/src/lsps1/service.rs

Lines changed: 100 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ use crate::message_queue::MessageQueue;
3434
use crate::events::EventQueue;
3535
use crate::lsps0::ser::{
3636
LSPSDateTime, LSPSProtocolMessageHandler, LSPSRequestId, LSPSResponseError,
37+
LSPS0_CLIENT_REJECTED_ERROR_CODE,
3738
};
3839
use crate::persist::{
3940
LIQUIDITY_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE, LSPS1_SERVICE_PERSISTENCE_SECONDARY_NAMESPACE,
@@ -62,6 +63,10 @@ pub struct LSPS1ServiceConfig {
6263
pub supported_options: LSPS1Options,
6364
}
6465

66+
const MAX_PENDING_REQUESTS_PER_PEER: usize = 10;
67+
const MAX_TOTAL_PENDING_REQUESTS: usize = 1000;
68+
const MAX_TOTAL_PEERS: usize = 100000;
69+
6570
/// The main object allowing to send and receive bLIP-51 / LSPS1 messages.
6671
pub struct LSPS1ServiceHandler<
6772
ES: EntropySource,
@@ -78,6 +83,7 @@ pub struct LSPS1ServiceHandler<
7883
pending_messages: Arc<MessageQueue>,
7984
pending_events: Arc<EventQueue<K>>,
8085
per_peer_state: RwLock<HashMap<PublicKey, Mutex<PeerState>>>,
86+
total_pending_requests: AtomicUsize,
8187
persistence_in_flight: AtomicUsize,
8288
time_provider: TP,
8389
config: LSPS1ServiceConfig,
@@ -102,6 +108,7 @@ where
102108
pending_messages,
103109
pending_events,
104110
per_peer_state: RwLock::new(per_peer_state),
111+
total_pending_requests: AtomicUsize::new(0),
105112
persistence_in_flight: AtomicUsize::new(0),
106113
time_provider,
107114
config,
@@ -133,7 +140,8 @@ where
133140
let mut peer_state_lock = inner_state_lock.lock().unwrap();
134141
// We clean up the peer state, but leave removing the peer entry to the prune logic in
135142
// `persist` which removes it from the store.
136-
peer_state_lock.prune_pending_requests();
143+
let num_pruned = peer_state_lock.prune_pending_requests();
144+
self.total_pending_requests.fetch_sub(num_pruned, Ordering::Relaxed);
137145
peer_state_lock.prune_expired_request_state();
138146
}
139147
}
@@ -309,17 +317,75 @@ where
309317
{
310318
let mut outer_state_lock = self.per_peer_state.write().unwrap();
311319

320+
if !outer_state_lock.contains_key(counterparty_node_id)
321+
&& outer_state_lock.len() >= MAX_TOTAL_PEERS
322+
{
323+
let response = LSPS1Response::CreateOrderError(LSPSResponseError {
324+
code: LSPS0_CLIENT_REJECTED_ERROR_CODE,
325+
message: "Reached maximum number of pending requests. Please try again later."
326+
.to_string(),
327+
data: None,
328+
});
329+
let msg = LSPS1Message::Response(request_id, response).into();
330+
message_queue_notifier.enqueue(counterparty_node_id, msg);
331+
return Err(LightningError {
332+
err: format!(
333+
"Dropping request from peer {} due to reaching maximally allowed number of total peers: {}",
334+
counterparty_node_id, MAX_TOTAL_PEERS
335+
),
336+
action: ErrorAction::IgnoreAndLog(Level::Debug),
337+
});
338+
}
339+
340+
if self.total_pending_requests.load(Ordering::Relaxed) >= MAX_TOTAL_PENDING_REQUESTS {
341+
let response = LSPS1Response::CreateOrderError(LSPSResponseError {
342+
code: LSPS0_CLIENT_REJECTED_ERROR_CODE,
343+
message: "Reached maximum number of pending requests. Please try again later."
344+
.to_string(),
345+
data: None,
346+
});
347+
let msg = LSPS1Message::Response(request_id, response).into();
348+
message_queue_notifier.enqueue(counterparty_node_id, msg);
349+
return Err(LightningError {
350+
err: format!(
351+
"Reached maximum number of total pending requests: {}",
352+
MAX_TOTAL_PENDING_REQUESTS
353+
),
354+
action: ErrorAction::IgnoreAndLog(Level::Debug),
355+
});
356+
}
357+
312358
let inner_state_lock = outer_state_lock
313359
.entry(*counterparty_node_id)
314360
.or_insert(Mutex::new(PeerState::default()));
315361
let mut peer_state_lock = inner_state_lock.lock().unwrap();
316362

363+
if peer_state_lock.pending_requests_and_channels() >= MAX_PENDING_REQUESTS_PER_PEER {
364+
let response = LSPS1Response::CreateOrderError(LSPSResponseError {
365+
code: LSPS0_CLIENT_REJECTED_ERROR_CODE,
366+
message: "Reached maximum number of pending requests. Please try again later."
367+
.to_string(),
368+
data: None,
369+
});
370+
let msg = LSPS1Message::Response(request_id, response).into();
371+
message_queue_notifier.enqueue(counterparty_node_id, msg);
372+
return Err(LightningError {
373+
err: format!(
374+
"Peer {} reached maximum number of pending requests: {}",
375+
counterparty_node_id, MAX_PENDING_REQUESTS_PER_PEER
376+
),
377+
action: ErrorAction::IgnoreAndLog(Level::Debug),
378+
});
379+
}
380+
317381
let request = LSPS1Request::CreateOrder(params.clone());
318382
peer_state_lock.register_request(request_id.clone(), request).map_err(|e| {
319383
let err = format!("Failed to handle request due to: {}", e);
320384
let action = ErrorAction::IgnoreAndLog(Level::Error);
321385
LightningError { err, action }
322386
})?;
387+
388+
self.total_pending_requests.fetch_add(1, Ordering::Relaxed);
323389
}
324390

325391
event_queue_notifier.enqueue(LSPS1ServiceEvent::RequestForPaymentDetails {
@@ -356,6 +422,7 @@ where
356422
let err = format!("Failed to send response due to: {}", e);
357423
APIError::APIMisuseError { err }
358424
})?;
425+
self.total_pending_requests.fetch_sub(1, Ordering::Relaxed);
359426

360427
match request {
361428
LSPS1Request::CreateOrder(params) => {
@@ -453,6 +520,7 @@ where
453520
let err = format!("Failed to send response due to: {}", e);
454521
APIError::APIMisuseError { err }
455522
})?;
523+
self.total_pending_requests.fetch_sub(1, Ordering::Relaxed);
456524

457525
let response = LSPS1Response::CreateOrderError(LSPSResponseError {
458526
code: LSPS1_CREATE_ORDER_REQUEST_UNRECOGNIZED_OR_STALE_TOKEN_ERROR_CODE,
@@ -490,6 +558,7 @@ where
490558
let err = format!("Failed to send response due to: {}", e);
491559
APIError::APIMisuseError { err }
492560
})?;
561+
self.total_pending_requests.fetch_sub(1, Ordering::Relaxed);
493562

494563
let response = LSPS1Response::CreateOrderError(LSPSResponseError {
495564
code: LSPS1_CREATE_ORDER_REQUEST_OPTION_MISMATCH_ERROR_CODE,
@@ -693,6 +762,21 @@ where
693762
let bytes = self.entropy_source.get_secure_random_bytes();
694763
LSPS1OrderId(utils::hex_str(&bytes[0..16]))
695764
}
765+
766+
#[cfg(debug_assertions)]
767+
fn verify_pending_request_counter(&self) {
768+
let mut num_requests = 0;
769+
let outer_state_lock = self.per_peer_state.read().unwrap();
770+
for (_, inner) in outer_state_lock.iter() {
771+
let inner_state_lock = inner.lock().unwrap();
772+
num_requests += inner_state_lock.pending_request_count();
773+
}
774+
debug_assert_eq!(
775+
num_requests,
776+
self.total_pending_requests.load(Ordering::Relaxed),
777+
"total_pending_requests counter out-of-sync! This should never happen!"
778+
);
779+
}
696780
}
697781

698782
impl<ES: EntropySource, CM: Deref + Clone, K: KVStore + Clone, TP: Deref + Clone>
@@ -708,16 +792,21 @@ where
708792
&self, message: Self::ProtocolMessage, counterparty_node_id: &PublicKey,
709793
) -> Result<(), LightningError> {
710794
match message {
711-
LSPS1Message::Request(request_id, request) => match request {
712-
LSPS1Request::GetInfo(_) => {
713-
self.handle_get_info_request(request_id, counterparty_node_id)
714-
},
715-
LSPS1Request::CreateOrder(params) => {
716-
self.handle_create_order_request(request_id, counterparty_node_id, params)
717-
},
718-
LSPS1Request::GetOrder(params) => {
719-
self.handle_get_order_request(request_id, counterparty_node_id, params)
720-
},
795+
LSPS1Message::Request(request_id, request) => {
796+
let res = match request {
797+
LSPS1Request::GetInfo(_) => {
798+
self.handle_get_info_request(request_id, counterparty_node_id)
799+
},
800+
LSPS1Request::CreateOrder(params) => {
801+
self.handle_create_order_request(request_id, counterparty_node_id, params)
802+
},
803+
LSPS1Request::GetOrder(params) => {
804+
self.handle_get_order_request(request_id, counterparty_node_id, params)
805+
},
806+
};
807+
#[cfg(debug_assertions)]
808+
self.verify_pending_request_counter();
809+
res
721810
},
722811
_ => {
723812
debug_assert!(

0 commit comments

Comments
 (0)