@@ -34,6 +34,7 @@ use crate::message_queue::MessageQueue;
3434use crate :: events:: EventQueue ;
3535use crate :: lsps0:: ser:: {
3636 LSPSDateTime , LSPSProtocolMessageHandler , LSPSRequestId , LSPSResponseError ,
37+ LSPS0_CLIENT_REJECTED_ERROR_CODE ,
3738} ;
3839use 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.
6671pub 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
698782impl < 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