@@ -21,8 +21,8 @@ use tokio::sync::oneshot;
2121
2222use crate :: connection:: ConnectionManager ;
2323use crate :: liquidity:: {
24- select_lsps_for_protocol, LspConfig , LspNode , LIQUIDITY_REQUEST_TIMEOUT_SECS ,
25- LSPS_DISCOVERY_WAIT_TIMEOUT_SECS ,
24+ select_lsps_for_protocol, LspConfig , LspNode , PendingRequest , PendingRequestGuard ,
25+ LIQUIDITY_REQUEST_TIMEOUT_SECS , LSPS_DISCOVERY_WAIT_TIMEOUT_SECS ,
2626} ;
2727use crate :: logger:: { log_error, log_info, LdkLogger , Logger } ;
2828use crate :: runtime:: Runtime ;
@@ -35,11 +35,11 @@ where
3535{
3636 pub ( crate ) lsp_nodes : Arc < RwLock < Vec < LspNode > > > ,
3737 pub ( crate ) pending_opening_params_requests :
38- Mutex < HashMap < LSPSRequestId , oneshot :: Sender < LSPS1OpeningParamsResponse > > > ,
38+ Mutex < HashMap < LSPSRequestId , PendingRequest < LSPS1OpeningParamsResponse > > > ,
3939 pub ( crate ) pending_create_order_requests :
40- Mutex < HashMap < LSPSRequestId , oneshot :: Sender < LSPS1OrderStatus > > > ,
40+ Mutex < HashMap < LSPSRequestId , PendingRequest < LSPS1OrderStatus > > > ,
4141 pub ( crate ) pending_check_order_status_requests :
42- Mutex < HashMap < LSPSRequestId , oneshot :: Sender < LSPS1OrderStatus > > > ,
42+ Mutex < HashMap < LSPSRequestId , PendingRequest < LSPS1OrderStatus > > > ,
4343 pub ( crate ) discovery_done_rx : tokio:: sync:: watch:: Receiver < bool > ,
4444 pub ( crate ) liquidity_manager : Arc < LiquidityManager > ,
4545 pub ( crate ) logger : L ,
@@ -61,12 +61,17 @@ where
6161 } ) ?;
6262
6363 let ( request_sender, request_receiver) = oneshot:: channel ( ) ;
64- {
64+ let _pending_request = {
6565 let mut pending_opening_params_requests_lock =
6666 self . pending_opening_params_requests . lock ( ) . expect ( "lock" ) ;
6767 let request_id = client_handler. request_supported_options ( lsps1_node. node_id ) ;
68- pending_opening_params_requests_lock. insert ( request_id, request_sender) ;
69- }
68+ PendingRequestGuard :: insert (
69+ & self . pending_opening_params_requests ,
70+ & mut pending_opening_params_requests_lock,
71+ request_id,
72+ request_sender,
73+ )
74+ } ;
7075
7176 tokio:: time:: timeout ( Duration :: from_secs ( LIQUIDITY_REQUEST_TIMEOUT_SECS ) , request_receiver)
7277 . await
@@ -146,16 +151,21 @@ where
146151
147152 let ( request_sender, request_receiver) = oneshot:: channel ( ) ;
148153 let request_id;
149- {
154+ let _pending_request = {
150155 let mut pending_create_order_requests_lock =
151156 self . pending_create_order_requests . lock ( ) . expect ( "lock" ) ;
152157 request_id = client_handler. create_order (
153158 & lsps1_node. node_id ,
154159 order_params. clone ( ) ,
155160 Some ( refund_address) ,
156161 ) ;
157- pending_create_order_requests_lock. insert ( request_id. clone ( ) , request_sender) ;
158- }
162+ PendingRequestGuard :: insert (
163+ & self . pending_create_order_requests ,
164+ & mut pending_create_order_requests_lock,
165+ request_id. clone ( ) ,
166+ request_sender,
167+ )
168+ } ;
159169
160170 let response = tokio:: time:: timeout (
161171 Duration :: from_secs ( LIQUIDITY_REQUEST_TIMEOUT_SECS ) ,
@@ -191,12 +201,17 @@ where
191201 } ) ?;
192202
193203 let ( request_sender, request_receiver) = oneshot:: channel ( ) ;
194- {
204+ let _pending_request = {
195205 let mut pending_check_order_status_requests_lock =
196206 self . pending_check_order_status_requests . lock ( ) . expect ( "lock" ) ;
197207 let request_id = client_handler. check_order_status ( & lsp_node_id, order_id) ;
198- pending_check_order_status_requests_lock. insert ( request_id, request_sender) ;
199- }
208+ PendingRequestGuard :: insert (
209+ & self . pending_check_order_status_requests ,
210+ & mut pending_check_order_status_requests_lock,
211+ request_id,
212+ request_sender,
213+ )
214+ } ;
200215
201216 let response = tokio:: time:: timeout (
202217 Duration :: from_secs ( LIQUIDITY_REQUEST_TIMEOUT_SECS ) ,
@@ -229,15 +244,15 @@ where
229244 . iter ( )
230245 . any ( |n| n. node_id == counterparty_node_id)
231246 {
232- if let Some ( sender ) = self
247+ if let Some ( request ) = self
233248 . pending_opening_params_requests
234249 . lock ( )
235250 . expect ( "lock" )
236251 . remove ( & request_id)
237252 {
238253 let response = LSPS1OpeningParamsResponse { supported_options } ;
239254
240- match sender. send ( response) {
255+ match request . sender . send ( response) {
241256 Ok ( ( ) ) => ( ) ,
242257 Err ( _) => {
243258 log_error ! (
@@ -279,7 +294,7 @@ where
279294 . iter ( )
280295 . any ( |n| n. node_id == counterparty_node_id)
281296 {
282- if let Some ( sender ) =
297+ if let Some ( request ) =
283298 self . pending_create_order_requests . lock ( ) . expect ( "lock" ) . remove ( & request_id)
284299 {
285300 let response = LSPS1OrderStatus {
@@ -290,7 +305,7 @@ where
290305 counterparty_node_id,
291306 } ;
292307
293- match sender. send ( response) {
308+ match request . sender . send ( response) {
294309 Ok ( ( ) ) => ( ) ,
295310 Err ( _) => {
296311 log_error ! (
@@ -329,7 +344,7 @@ where
329344 . iter ( )
330345 . any ( |n| n. node_id == counterparty_node_id)
331346 {
332- if let Some ( sender ) = self
347+ if let Some ( request ) = self
333348 . pending_check_order_status_requests
334349 . lock ( )
335350 . expect ( "lock" )
@@ -343,7 +358,7 @@ where
343358 counterparty_node_id,
344359 } ;
345360
346- match sender. send ( response) {
361+ match request . sender . send ( response) {
347362 Ok ( ( ) ) => ( ) ,
348363 Err ( _) => {
349364 log_error ! (
0 commit comments