Skip to content

Commit a831210

Browse files
committed
Add forwarded payment tracking and statistics
Routing nodes and LSPs need to track forwarded payments so they can account for earned fees and profitability over time. Store forwarded payment data and aggregate it into channel and channel-pair statistics. Detailed tracking retains individual forwarding events for a configured window before aggregating them by channel pair. Stats mode keeps only per-channel forwarding totals for lightweight nodes. AI-assisted-by: OpenAI Codex
1 parent f2e44fd commit a831210

11 files changed

Lines changed: 1307 additions & 36 deletions

File tree

bindings/ldk_node.udl

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,8 @@ typedef dictionary ElectrumSyncConfig;
1111

1212
typedef dictionary TorConfig;
1313

14+
typedef enum ForwardedPaymentTrackingMode;
15+
1416
typedef interface NodeEntropy;
1517

1618
typedef enum WordCount;
@@ -137,6 +139,19 @@ interface Node {
137139
void remove_payment([ByRef]PaymentId payment_id);
138140
BalanceDetails list_balances();
139141
sequence<PaymentDetails> list_payments();
142+
ForwardedPaymentDetails? forwarded_payment([ByRef]ForwardedPaymentId forwarded_payment_id);
143+
[Throws=NodeError]
144+
ForwardedPaymentDetailsPage list_forwarded_payments(PageToken? page_token);
145+
ForwardedPaymentTrackingMode forwarded_payment_tracking_mode();
146+
ChannelForwardingStats? channel_forwarding_stats([ByRef]ChannelId channel_id);
147+
[Throws=NodeError]
148+
ChannelForwardingStatsPage list_channel_forwarding_stats(PageToken? page_token);
149+
[Throws=NodeError]
150+
ChannelPairForwardingStatsPage list_channel_pair_forwarding_stats(PageToken? page_token);
151+
[Throws=NodeError]
152+
ChannelPairForwardingStatsPage list_channel_pair_forwarding_stats_in_range(u64 start_timestamp, u64 end_timestamp, PageToken? page_token);
153+
[Throws=NodeError]
154+
ChannelPairForwardingStatsPage list_channel_pair_forwarding_stats_for_pair(ChannelId prev_channel_id, ChannelId next_channel_id, PageToken? page_token);
140155
sequence<PeerDetails> list_peers();
141156
sequence<ChannelDetails> list_channels();
142157
NetworkGraph network_graph();
@@ -379,6 +394,9 @@ typedef string OfferId;
379394
[Custom]
380395
typedef string PaymentId;
381396

397+
[Custom]
398+
typedef string ForwardedPaymentId;
399+
382400
[Custom]
383401
typedef string PaymentHash;
384402

@@ -391,6 +409,12 @@ typedef string PaymentSecret;
391409
[Custom]
392410
typedef string ChannelId;
393411

412+
[Custom]
413+
typedef string ChannelPairStatsId;
414+
415+
[Custom]
416+
typedef string PageToken;
417+
394418
[Custom]
395419
typedef string UserChannelId;
396420

@@ -417,3 +441,15 @@ typedef enum Event;
417441
typedef interface HRNResolverConfig;
418442

419443
typedef dictionary HumanReadableNamesConfig;
444+
445+
typedef dictionary ForwardedPaymentDetails;
446+
447+
typedef dictionary ChannelForwardingStats;
448+
449+
typedef dictionary ChannelPairForwardingStats;
450+
451+
typedef dictionary ForwardedPaymentDetailsPage;
452+
453+
typedef dictionary ChannelForwardingStatsPage;
454+
455+
typedef dictionary ChannelPairForwardingStatsPage;

src/builder.rs

Lines changed: 95 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,13 @@ use crate::io::utils::{
6464
};
6565
use crate::io::vss_store::VssStoreBuilder;
6666
use crate::io::{
67-
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
67+
self, CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
68+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
69+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
70+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
71+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
72+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
73+
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
6874
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
6975
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
7076
};
@@ -77,7 +83,8 @@ use crate::peer_store::PeerStore;
7783
use crate::runtime::{Runtime, RuntimeSpawner};
7884
use crate::tx_broadcaster::TransactionBroadcaster;
7985
use crate::types::{
80-
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreRef, DynStoreWrapper,
86+
AsyncPersister, ChainMonitor, ChannelForwardingStatsStore, ChannelManager,
87+
ChannelPairForwardingStatsStore, DynStore, DynStoreRef, DynStoreWrapper, ForwardedPaymentStore,
8188
GossipSync, Graph, HRNResolver, KeysManager, MessageRouter, OnionMessenger, PaymentStore,
8289
PeerManager, PendingPaymentStore,
8390
};
@@ -1397,24 +1404,48 @@ fn build_with_store_internal(
13971404

13981405
let kv_store_ref = Arc::clone(&kv_store);
13991406
let logger_ref = Arc::clone(&logger);
1400-
let (payment_store_res, node_metris_res, pending_payment_store_res) =
1401-
runtime.block_on(async move {
1402-
tokio::join!(
1403-
read_all_objects(
1404-
&*kv_store_ref,
1405-
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1406-
PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1407-
Arc::clone(&logger_ref),
1408-
),
1409-
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
1410-
read_all_objects(
1411-
&*kv_store_ref,
1412-
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1413-
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1414-
Arc::clone(&logger_ref),
1415-
)
1407+
let (
1408+
payment_store_res,
1409+
forwarded_payment_store_res,
1410+
channel_forwarding_stats_res,
1411+
channel_pair_forwarding_stats_res,
1412+
node_metris_res,
1413+
pending_payment_store_res,
1414+
) = runtime.block_on(async move {
1415+
tokio::join!(
1416+
read_all_objects(
1417+
&*kv_store_ref,
1418+
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1419+
PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1420+
Arc::clone(&logger_ref),
1421+
),
1422+
read_all_objects(
1423+
&*kv_store_ref,
1424+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1425+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1426+
Arc::clone(&logger_ref),
1427+
),
1428+
read_all_objects(
1429+
&*kv_store_ref,
1430+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
1431+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
1432+
Arc::clone(&logger_ref),
1433+
),
1434+
read_all_objects(
1435+
&*kv_store_ref,
1436+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
1437+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
1438+
Arc::clone(&logger_ref),
1439+
),
1440+
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
1441+
read_all_objects(
1442+
&*kv_store_ref,
1443+
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1444+
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1445+
Arc::clone(&logger_ref),
14161446
)
1417-
});
1447+
)
1448+
});
14181449

14191450
// Initialize the status fields.
14201451
let node_metrics = match node_metris_res {
@@ -1443,6 +1474,48 @@ fn build_with_store_internal(
14431474
},
14441475
};
14451476

1477+
let forwarded_payment_store = match forwarded_payment_store_res {
1478+
Ok(forwarded_payments) => Arc::new(ForwardedPaymentStore::new(
1479+
forwarded_payments,
1480+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1481+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1482+
Arc::clone(&kv_store),
1483+
Arc::clone(&logger),
1484+
)),
1485+
Err(e) => {
1486+
log_error!(logger, "Failed to read forwarded payment data from store: {}", e);
1487+
return Err(BuildError::ReadFailed);
1488+
},
1489+
};
1490+
1491+
let channel_forwarding_stats_store = match channel_forwarding_stats_res {
1492+
Ok(stats) => Arc::new(ChannelForwardingStatsStore::new(
1493+
stats,
1494+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1495+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1496+
Arc::clone(&kv_store),
1497+
Arc::clone(&logger),
1498+
)),
1499+
Err(e) => {
1500+
log_error!(logger, "Failed to read channel forwarding stats from store: {}", e);
1501+
return Err(BuildError::ReadFailed);
1502+
},
1503+
};
1504+
1505+
let channel_pair_forwarding_stats_store = match channel_pair_forwarding_stats_res {
1506+
Ok(stats) => Arc::new(ChannelPairForwardingStatsStore::new(
1507+
stats,
1508+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1509+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1510+
Arc::clone(&kv_store),
1511+
Arc::clone(&logger),
1512+
)),
1513+
Err(e) => {
1514+
log_error!(logger, "Failed to read channel pair forwarding stats from store: {}", e);
1515+
return Err(BuildError::ReadFailed);
1516+
},
1517+
};
1518+
14461519
let (chain_source, chain_tip_opt) = match chain_data_source_config {
14471520
Some(ChainDataSourceConfig::Esplora { server_url, headers, sync_config }) => {
14481521
let sync_config = sync_config.unwrap_or(EsploraSyncConfig::default());
@@ -2245,6 +2318,9 @@ fn build_with_store_internal(
22452318
scorer,
22462319
peer_store,
22472320
payment_store,
2321+
forwarded_payment_store,
2322+
channel_forwarding_stats_store,
2323+
channel_pair_forwarding_stats_store,
22482324
lnurl_auth,
22492325
is_running,
22502326
node_metrics,

src/config.rs

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,25 @@ pub(crate) const HRN_RESOLUTION_TIMEOUT_SECS: u64 = 5;
111111
// The timeout after which we abort an LNURL-auth operation.
112112
pub(crate) const LNURL_AUTH_TIMEOUT_SECS: u64 = 15;
113113

114+
/// The mode used for tracking forwarded payments.
115+
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
116+
#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
117+
pub enum ForwardedPaymentTrackingMode {
118+
/// Store individual forwarded payments until they are aggregated into channel-pair buckets.
119+
Detailed {
120+
/// Number of minutes to retain individual forwarded payments before aggregation.
121+
retention_minutes: u64,
122+
},
123+
/// Track only per-channel aggregate statistics.
124+
Stats,
125+
}
126+
127+
impl Default for ForwardedPaymentTrackingMode {
128+
fn default() -> Self {
129+
Self::Stats
130+
}
131+
}
132+
114133
#[derive(Debug, Clone)]
115134
#[cfg_attr(feature = "uniffi", derive(uniffi::Record))]
116135
/// Represents the configuration of an [`Node`] instance.
@@ -130,9 +149,10 @@ pub(crate) const LNURL_AUTH_TIMEOUT_SECS: u64 = 15;
130149
/// | `route_parameters` | None |
131150
/// | `tor_config` | None |
132151
/// | `hrn_config` | HumanReadableNamesConfig::default() |
152+
/// | `forwarded_payment_tracking_mode` | Stats |
133153
///
134-
/// See [`AnchorChannelsConfig`] and [`RouteParametersConfig`] for more information regarding their
135-
/// respective default values.
154+
/// See [`AnchorChannelsConfig`], [`RouteParametersConfig`], and
155+
/// [`ForwardedPaymentTrackingMode`] for more information regarding their respective default values.
136156
///
137157
/// [`Node`]: crate::Node
138158
pub struct Config {
@@ -205,6 +225,8 @@ pub struct Config {
205225
///
206226
/// [BIP 353]: https://github.com/bitcoin/bips/blob/master/bip-0353.mediawiki
207227
pub hrn_config: HumanReadableNamesConfig,
228+
/// The mode used for tracking forwarded payments.
229+
pub forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode,
208230
}
209231

210232
impl Default for Config {
@@ -221,6 +243,7 @@ impl Default for Config {
221243
route_parameters: None,
222244
node_alias: None,
223245
hrn_config: HumanReadableNamesConfig::default(),
246+
forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode::default(),
224247
}
225248
}
226249
}

src/data_store.rs

Lines changed: 61 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,11 +9,11 @@ use std::collections::HashMap;
99
use std::ops::Deref;
1010
use std::sync::{Arc, Mutex};
1111

12-
use lightning::util::persist::KVStore;
12+
use lightning::util::persist::{KVStore, PageToken, PaginatedKVStore};
1313
use lightning::util::ser::{Readable, Writeable};
1414

1515
use crate::logger::{log_error, LdkLogger};
16-
use crate::types::DynStore;
16+
use crate::types::{DynStore, DynStoreRef};
1717
use crate::Error;
1818

1919
pub(crate) trait StorableObject: Clone + Readable + Writeable {
@@ -179,6 +179,65 @@ where
179179
self.objects.lock().expect("lock").values().filter(f).cloned().collect::<Vec<SO>>()
180180
}
181181

182+
pub(crate) async fn list_page(
183+
&self, page_token: Option<PageToken>,
184+
) -> Result<(Vec<SO>, Option<PageToken>), Error> {
185+
let response = PaginatedKVStore::list_paginated(
186+
&DynStoreRef(Arc::clone(&self.kv_store)),
187+
&self.primary_namespace,
188+
&self.secondary_namespace,
189+
page_token,
190+
)
191+
.await
192+
.map_err(|e| {
193+
log_error!(
194+
self.logger,
195+
"Listing object data under {}/{} failed due to: {}",
196+
&self.primary_namespace,
197+
&self.secondary_namespace,
198+
e
199+
);
200+
Error::PersistenceFailed
201+
})?;
202+
203+
let mut objects = Vec::with_capacity(response.keys.len());
204+
for key in response.keys {
205+
let data = KVStore::read(
206+
&DynStoreRef(Arc::clone(&self.kv_store)),
207+
&self.primary_namespace,
208+
&self.secondary_namespace,
209+
&key,
210+
)
211+
.await
212+
.map_err(|e| {
213+
log_error!(
214+
self.logger,
215+
"Reading object data for key {}/{}/{} failed due to: {}",
216+
&self.primary_namespace,
217+
&self.secondary_namespace,
218+
key,
219+
e
220+
);
221+
Error::PersistenceFailed
222+
})?;
223+
224+
let object = SO::read(&mut &data[..]).map_err(|e| {
225+
log_error!(
226+
self.logger,
227+
"Failed to deserialize object data for key {}/{}/{}: {}",
228+
&self.primary_namespace,
229+
&self.secondary_namespace,
230+
key,
231+
e
232+
);
233+
Error::PersistenceFailed
234+
})?;
235+
objects.push(object);
236+
}
237+
238+
Ok((objects, response.next_page_token))
239+
}
240+
182241
async fn persist(&self, object: &SO) -> Result<(), Error> {
183242
let (store_key, data) = Self::encode_object(object);
184243
self.persist_encoded(store_key, data).await

0 commit comments

Comments
 (0)