Skip to content

Commit f031041

Browse files
committed
Add forwarded payment tracking
Store only unambiguous single-HTLC forwarding events. Aggregate events into channel and channel-pair statistics. Use fixed one-hour buckets for detailed records. Expose bucket identifiers as opaque strings. AI-assisted-by: OpenAI Codex
1 parent abbed3e commit f031041

12 files changed

Lines changed: 2210 additions & 31 deletions

File tree

bindings/ldk_node.udl

Lines changed: 20 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 interface ProbingConfig;
@@ -113,6 +115,7 @@ interface Node {
113115
OnchainPayment onchain_payment();
114116
UnifiedPayment unified_payment();
115117
Liquidity liquidity();
118+
Forwarding forwarding();
116119
[Throws=NodeError]
117120
void lnurl_auth(string lnurl);
118121
[Throws=NodeError]
@@ -184,6 +187,8 @@ typedef interface UnifiedPayment;
184187

185188
typedef interface Liquidity;
186189

190+
typedef interface Forwarding;
191+
187192
[Error]
188193
enum NodeError {
189194
"AlreadyRunning",
@@ -407,6 +412,9 @@ typedef string PaymentSecret;
407412
[Custom]
408413
typedef string ChannelId;
409414

415+
[Custom]
416+
typedef string PageToken;
417+
410418
[Custom]
411419
typedef string UserChannelId;
412420

@@ -433,3 +441,15 @@ typedef enum Event;
433441
typedef interface HRNResolverConfig;
434442

435443
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: 76 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,13 @@ use crate::io::utils::{
6565
};
6666
use crate::io::vss_store::VssStoreBuilder;
6767
use crate::io::{
68-
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
68+
self, CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
69+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
70+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
71+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
72+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
73+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
74+
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
6975
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
7076
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
7177
};
@@ -82,7 +88,8 @@ use crate::probing::{
8288
use crate::runtime::{Runtime, RuntimeSpawner};
8389
use crate::tx_broadcaster::TransactionBroadcaster;
8490
use crate::types::{
85-
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreRef, DynStoreWrapper,
91+
AsyncPersister, ChainMonitor, ChannelForwardingStatsStore, ChannelManager,
92+
ChannelPairForwardingStatsStore, DynStore, DynStoreRef, DynStoreWrapper, ForwardedPaymentStore,
8693
GossipSync, Graph, HRNResolver, KeysManager, MessageRouter, OnionMessenger, PaymentStore,
8794
PeerManager, PendingPaymentStore,
8895
};
@@ -1439,24 +1446,34 @@ fn build_with_store_internal(
14391446

14401447
let kv_store_ref = Arc::clone(&kv_store);
14411448
let logger_ref = Arc::clone(&logger);
1442-
let (payment_store_res, node_metris_res, pending_payment_store_res) =
1443-
runtime.block_on(async move {
1444-
tokio::join!(
1445-
read_all_objects(
1446-
&*kv_store_ref,
1447-
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1448-
PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1449-
Arc::clone(&logger_ref),
1450-
),
1451-
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
1452-
read_all_objects(
1453-
&*kv_store_ref,
1454-
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1455-
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1456-
Arc::clone(&logger_ref),
1457-
)
1449+
let (
1450+
payment_store_res,
1451+
channel_forwarding_stats_res,
1452+
node_metris_res,
1453+
pending_payment_store_res,
1454+
) = runtime.block_on(async move {
1455+
tokio::join!(
1456+
read_all_objects(
1457+
&*kv_store_ref,
1458+
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1459+
PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1460+
Arc::clone(&logger_ref),
1461+
),
1462+
read_all_objects(
1463+
&*kv_store_ref,
1464+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
1465+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
1466+
Arc::clone(&logger_ref),
1467+
),
1468+
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
1469+
read_all_objects(
1470+
&*kv_store_ref,
1471+
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1472+
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1473+
Arc::clone(&logger_ref),
14581474
)
1459-
});
1475+
)
1476+
});
14601477

14611478
// Initialize the status fields.
14621479
let node_metrics = match node_metris_res {
@@ -1485,6 +1502,34 @@ fn build_with_store_internal(
14851502
},
14861503
};
14871504

1505+
let forwarded_payment_store = Arc::new(ForwardedPaymentStore::new(
1506+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1507+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1508+
Arc::clone(&kv_store),
1509+
Arc::clone(&logger),
1510+
));
1511+
1512+
let channel_forwarding_stats_store = match channel_forwarding_stats_res {
1513+
Ok(stats) => Arc::new(ChannelForwardingStatsStore::new(
1514+
stats,
1515+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1516+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1517+
Arc::clone(&kv_store),
1518+
Arc::clone(&logger),
1519+
)),
1520+
Err(e) => {
1521+
log_error!(logger, "Failed to read channel forwarding stats from store: {}", e);
1522+
return Err(BuildError::ReadFailed);
1523+
},
1524+
};
1525+
1526+
let channel_pair_forwarding_stats_store = Arc::new(ChannelPairForwardingStatsStore::new(
1527+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1528+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1529+
Arc::clone(&kv_store),
1530+
Arc::clone(&logger),
1531+
));
1532+
14881533
let (chain_source, chain_tip_opt) = match chain_data_source_config {
14891534
Some(ChainDataSourceConfig::Esplora { server_url, headers, sync_config }) => {
14901535
let sync_config = sync_config.unwrap_or(EsploraSyncConfig::default());
@@ -2306,6 +2351,14 @@ fn build_with_store_internal(
23062351
_leak_checker.0.push(Arc::downgrade(&wallet) as Weak<dyn Any + Send + Sync>);
23072352
}
23082353

2354+
let forwarded_payment_aggregation_retention_secs = match config.forwarded_payment_tracking_mode
2355+
{
2356+
crate::config::ForwardedPaymentTrackingMode::Detailed => {
2357+
Some(crate::payment::forwarding_store::FORWARDED_PAYMENT_AGGREGATION_BUCKET_SIZE_SECS)
2358+
},
2359+
crate::config::ForwardedPaymentTrackingMode::Stats => Some(0),
2360+
};
2361+
23092362
Ok(Node {
23102363
runtime,
23112364
stop_sender,
@@ -2333,6 +2386,10 @@ fn build_with_store_internal(
23332386
scorer,
23342387
peer_store,
23352388
payment_store,
2389+
forwarded_payment_store,
2390+
channel_forwarding_stats_store,
2391+
channel_pair_forwarding_stats_store,
2392+
forwarded_payment_aggregation_retention_secs,
23362393
lnurl_auth,
23372394
is_running,
23382395
node_metrics,

src/config.rs

Lines changed: 30 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,30 @@ pub(crate) const LIQUIDITY_DISCOVERY_RETRY_INITIAL_DELAY: Duration = Duration::f
136136
// thereafter until every configured LSP has been discovered.
137137
pub(crate) const LIQUIDITY_DISCOVERY_RETRY_MAX_DELAY: Duration = Duration::from_secs(60 * 60);
138138

139+
/// The mode used for tracking forwarded payments.
140+
///
141+
/// In either mode, a forward is tracked only when it has exactly one incoming HTLC and one outgoing
142+
/// HTLC, and LDK reports both the outbound amount and total fee.
143+
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
144+
#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
145+
pub enum ForwardedPaymentTrackingMode {
146+
/// Track eligible new forwarded payments only as per-channel aggregate statistics.
147+
///
148+
/// Any detailed records left by a previous configuration are aggregated and removed after their
149+
/// current one-hour bucket closes.
150+
Stats,
151+
/// Store eligible individual forwarded payments for the current and previous one-hour buckets.
152+
///
153+
/// Payments from older buckets are aggregated into channel-pair statistics and removed.
154+
Detailed,
155+
}
156+
157+
impl Default for ForwardedPaymentTrackingMode {
158+
fn default() -> Self {
159+
Self::Stats
160+
}
161+
}
162+
139163
#[derive(Debug, Clone)]
140164
#[cfg_attr(feature = "uniffi", derive(uniffi::Record))]
141165
/// Represents the configuration of an [`Node`] instance.
@@ -155,9 +179,10 @@ pub(crate) const LIQUIDITY_DISCOVERY_RETRY_MAX_DELAY: Duration = Duration::from_
155179
/// | `route_parameters` | None |
156180
/// | `tor_config` | None |
157181
/// | `hrn_config` | HumanReadableNamesConfig::default() |
182+
/// | `forwarded_payment_tracking_mode` | Stats |
158183
///
159-
/// See [`AnchorChannelsConfig`] and [`RouteParametersConfig`] for more information regarding their
160-
/// respective default values.
184+
/// See [`AnchorChannelsConfig`], [`RouteParametersConfig`], and
185+
/// [`ForwardedPaymentTrackingMode`] for more information regarding their respective default values.
161186
///
162187
/// [`Node`]: crate::Node
163188
pub struct Config {
@@ -219,6 +244,8 @@ pub struct Config {
219244
///
220245
/// [BIP 353]: https://github.com/bitcoin/bips/blob/master/bip-0353.mediawiki
221246
pub hrn_config: HumanReadableNamesConfig,
247+
/// The mode used for tracking forwarded payments.
248+
pub forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode,
222249
}
223250

224251
impl Default for Config {
@@ -235,6 +262,7 @@ impl Default for Config {
235262
route_parameters: None,
236263
node_alias: None,
237264
hrn_config: HumanReadableNamesConfig::default(),
265+
forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode::default(),
238266
}
239267
}
240268
}

0 commit comments

Comments
 (0)