Skip to content

Commit 8304c3f

Browse files
committed
Clean up abandoned liquidity requests
Remove pending LSPS request state when callers time out or are cancelled so unresponsive services cannot grow request maps. Co-Authored-By: HAL 9000
1 parent 08d3b5d commit 8304c3f

3 files changed

Lines changed: 198 additions & 57 deletions

File tree

src/liquidity/client/lsps1.rs

Lines changed: 35 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -21,8 +21,8 @@ use tokio::sync::oneshot;
2121

2222
use crate::connection::ConnectionManager;
2323
use 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
};
2727
use crate::logger::{log_error, log_info, LdkLogger, Logger};
2828
use 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!(

src/liquidity/client/lsps2.rs

Lines changed: 24 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -26,8 +26,8 @@ use tokio::task::JoinSet;
2626

2727
use crate::connection::ConnectionManager;
2828
use crate::liquidity::{
29-
select_all_lsps_for_protocol, select_lsps_for_protocol, LspConfig, LspNode,
30-
LIQUIDITY_REQUEST_TIMEOUT_SECS, LSPS_DISCOVERY_WAIT_TIMEOUT_SECS,
29+
select_all_lsps_for_protocol, select_lsps_for_protocol, LspConfig, LspNode, PendingRequest,
30+
PendingRequestGuard, LIQUIDITY_REQUEST_TIMEOUT_SECS, LSPS_DISCOVERY_WAIT_TIMEOUT_SECS,
3131
};
3232
use crate::logger::{log_debug, log_error, log_info, LdkLogger};
3333
use crate::payment::store::LSPS2Parameters;
@@ -41,9 +41,9 @@ where
4141
{
4242
pub(crate) lsp_nodes: Arc<RwLock<Vec<LspNode>>>,
4343
pub(crate) pending_lsps2_fee_requests:
44-
Mutex<HashMap<LSPSRequestId, oneshot::Sender<LSPS2FeeResponse>>>,
44+
Mutex<HashMap<LSPSRequestId, PendingRequest<LSPS2FeeResponse>>>,
4545
pub(crate) pending_buy_requests:
46-
Mutex<HashMap<LSPSRequestId, oneshot::Sender<LSPS2BuyResponse>>>,
46+
Mutex<HashMap<LSPSRequestId, PendingRequest<LSPS2BuyResponse>>>,
4747
pub(crate) channel_manager: Arc<ChannelManager>,
4848
pub(crate) keys_manager: Arc<KeysManager>,
4949
pub(crate) discovery_done_rx: tokio::sync::watch::Receiver<bool>,
@@ -262,13 +262,18 @@ where
262262
})?;
263263

264264
let (fee_request_sender, fee_request_receiver) = oneshot::channel();
265-
{
265+
let _pending_request = {
266266
let mut pending_fee_requests_lock =
267267
self.pending_lsps2_fee_requests.lock().expect("lock");
268268
let request_id =
269269
client_handler.request_opening_params(lsps2_node.node_id, lsps2_node.token.clone());
270-
pending_fee_requests_lock.insert(request_id, fee_request_sender);
271-
}
270+
PendingRequestGuard::insert(
271+
&self.pending_lsps2_fee_requests,
272+
&mut pending_fee_requests_lock,
273+
request_id,
274+
fee_request_sender,
275+
)
276+
};
272277

273278
tokio::time::timeout(
274279
Duration::from_secs(LIQUIDITY_REQUEST_TIMEOUT_SECS),
@@ -298,7 +303,7 @@ where
298303
})?;
299304

300305
let (buy_request_sender, buy_request_receiver) = oneshot::channel();
301-
{
306+
let _pending_request = {
302307
let mut pending_buy_requests_lock = self.pending_buy_requests.lock().expect("lock");
303308
let request_id = client_handler
304309
.select_opening_params(lsps2_node.node_id, amount_msat, opening_fee_params)
@@ -310,8 +315,13 @@ where
310315
);
311316
Error::LiquidityRequestFailed
312317
})?;
313-
pending_buy_requests_lock.insert(request_id, buy_request_sender);
314-
}
318+
PendingRequestGuard::insert(
319+
&self.pending_buy_requests,
320+
&mut pending_buy_requests_lock,
321+
request_id,
322+
buy_request_sender,
323+
)
324+
};
315325

316326
let buy_response = tokio::time::timeout(
317327
Duration::from_secs(LIQUIDITY_REQUEST_TIMEOUT_SECS),
@@ -428,12 +438,12 @@ where
428438
.iter()
429439
.any(|n| n.node_id == counterparty_node_id)
430440
{
431-
if let Some(sender) =
441+
if let Some(request) =
432442
self.pending_lsps2_fee_requests.lock().expect("lock").remove(&request_id)
433443
{
434444
let response = LSPS2FeeResponse { opening_fee_params_menu };
435445

436-
match sender.send(response) {
446+
match request.sender.send(response) {
437447
Ok(()) => (),
438448
Err(_) => {
439449
log_error!(
@@ -474,12 +484,12 @@ where
474484
.iter()
475485
.any(|n| n.node_id == counterparty_node_id)
476486
{
477-
if let Some(sender) =
487+
if let Some(request) =
478488
self.pending_buy_requests.lock().expect("lock").remove(&request_id)
479489
{
480490
let response = LSPS2BuyResponse { intercept_scid, cltv_expiry_delta };
481491

482-
match sender.send(response) {
492+
match request.sender.send(response) {
483493
Ok(()) => (),
484494
Err(_) => {
485495
log_error!(

0 commit comments

Comments
 (0)