Skip to content

Commit 4a2b945

Browse files
committed
Add forwarded payment tracking
Store individual forwarding events and aggregate them into channel and channel-pair statistics for fee and profitability tracking. Use fixed one-hour buckets for detailed records and expose their identifiers as opaque strings. Infer channel-pair allocations for multi-HTLC forwards by matching incoming and outgoing amounts in FIFO order. AI-assisted-by: OpenAI Codex
1 parent 0d201d6 commit 4a2b945

12 files changed

Lines changed: 2638 additions & 33 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: 114 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,48 @@ 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+
forwarded_payment_store_res,
1452+
channel_forwarding_stats_res,
1453+
channel_pair_forwarding_stats_res,
1454+
node_metris_res,
1455+
pending_payment_store_res,
1456+
) = runtime.block_on(async move {
1457+
tokio::join!(
1458+
read_all_objects(
1459+
&*kv_store_ref,
1460+
PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1461+
PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1462+
Arc::clone(&logger_ref),
1463+
),
1464+
read_all_objects(
1465+
&*kv_store_ref,
1466+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1467+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1468+
Arc::clone(&logger_ref),
1469+
),
1470+
read_all_objects(
1471+
&*kv_store_ref,
1472+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
1473+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
1474+
Arc::clone(&logger_ref),
1475+
),
1476+
read_all_objects(
1477+
&*kv_store_ref,
1478+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE,
1479+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE,
1480+
Arc::clone(&logger_ref),
1481+
),
1482+
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
1483+
read_all_objects(
1484+
&*kv_store_ref,
1485+
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
1486+
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
1487+
Arc::clone(&logger_ref),
14581488
)
1459-
});
1489+
)
1490+
});
14601491

14611492
// Initialize the status fields.
14621493
let node_metrics = match node_metris_res {
@@ -1485,6 +1516,55 @@ fn build_with_store_internal(
14851516
},
14861517
};
14871518

1519+
let (forwarded_payment_store, has_stored_forwarded_payments) = match forwarded_payment_store_res
1520+
{
1521+
Ok(forwarded_payments) => {
1522+
let has_stored_forwarded_payments = !forwarded_payments.is_empty();
1523+
(
1524+
Arc::new(ForwardedPaymentStore::new(
1525+
forwarded_payments,
1526+
FORWARDED_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1527+
FORWARDED_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1528+
Arc::clone(&kv_store),
1529+
Arc::clone(&logger),
1530+
)),
1531+
has_stored_forwarded_payments,
1532+
)
1533+
},
1534+
Err(e) => {
1535+
log_error!(logger, "Failed to read forwarded payment data from store: {}", e);
1536+
return Err(BuildError::ReadFailed);
1537+
},
1538+
};
1539+
1540+
let channel_forwarding_stats_store = match channel_forwarding_stats_res {
1541+
Ok(stats) => Arc::new(ChannelForwardingStatsStore::new(
1542+
stats,
1543+
CHANNEL_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1544+
CHANNEL_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1545+
Arc::clone(&kv_store),
1546+
Arc::clone(&logger),
1547+
)),
1548+
Err(e) => {
1549+
log_error!(logger, "Failed to read channel forwarding stats from store: {}", e);
1550+
return Err(BuildError::ReadFailed);
1551+
},
1552+
};
1553+
1554+
let channel_pair_forwarding_stats_store = match channel_pair_forwarding_stats_res {
1555+
Ok(stats) => Arc::new(ChannelPairForwardingStatsStore::new(
1556+
stats,
1557+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
1558+
CHANNEL_PAIR_FORWARDING_STATS_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
1559+
Arc::clone(&kv_store),
1560+
Arc::clone(&logger),
1561+
)),
1562+
Err(e) => {
1563+
log_error!(logger, "Failed to read channel pair forwarding stats from store: {}", e);
1564+
return Err(BuildError::ReadFailed);
1565+
},
1566+
};
1567+
14881568
let (chain_source, chain_tip_opt) = match chain_data_source_config {
14891569
Some(ChainDataSourceConfig::Esplora { server_url, headers, sync_config }) => {
14901570
let sync_config = sync_config.unwrap_or(EsploraSyncConfig::default());
@@ -2306,6 +2386,17 @@ fn build_with_store_internal(
23062386
})
23072387
});
23082388

2389+
let forwarded_payment_aggregation_retention_secs = match config.forwarded_payment_tracking_mode
2390+
{
2391+
crate::config::ForwardedPaymentTrackingMode::Detailed => {
2392+
Some(crate::payment::store::FORWARDED_PAYMENT_AGGREGATION_BUCKET_SIZE_SECS)
2393+
},
2394+
crate::config::ForwardedPaymentTrackingMode::Stats if has_stored_forwarded_payments => {
2395+
Some(0)
2396+
},
2397+
_ => None,
2398+
};
2399+
23092400
Ok(Node {
23102401
runtime,
23112402
stop_sender,
@@ -2333,6 +2424,10 @@ fn build_with_store_internal(
23332424
scorer,
23342425
peer_store,
23352426
payment_store,
2427+
forwarded_payment_store,
2428+
channel_forwarding_stats_store,
2429+
channel_pair_forwarding_stats_store,
2430+
forwarded_payment_aggregation_retention_secs,
23362431
lnurl_auth,
23372432
is_running,
23382433
node_metrics,

src/config.rs

Lines changed: 27 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,27 @@ 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+
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
141+
#[cfg_attr(feature = "uniffi", derive(uniffi::Enum))]
142+
pub enum ForwardedPaymentTrackingMode {
143+
/// Track new forwarded payments only as per-channel aggregate statistics.
144+
///
145+
/// Any detailed records left by a previous configuration are aggregated and removed after their
146+
/// current one-hour bucket closes.
147+
Stats,
148+
/// Store individual forwarded payments for the current and previous one-hour buckets.
149+
///
150+
/// Payments from older buckets are aggregated into channel-pair statistics and removed.
151+
Detailed,
152+
}
153+
154+
impl Default for ForwardedPaymentTrackingMode {
155+
fn default() -> Self {
156+
Self::Stats
157+
}
158+
}
159+
139160
#[derive(Debug, Clone)]
140161
#[cfg_attr(feature = "uniffi", derive(uniffi::Record))]
141162
/// Represents the configuration of an [`Node`] instance.
@@ -155,9 +176,10 @@ pub(crate) const LIQUIDITY_DISCOVERY_RETRY_MAX_DELAY: Duration = Duration::from_
155176
/// | `route_parameters` | None |
156177
/// | `tor_config` | None |
157178
/// | `hrn_config` | HumanReadableNamesConfig::default() |
179+
/// | `forwarded_payment_tracking_mode` | Stats |
158180
///
159-
/// See [`AnchorChannelsConfig`] and [`RouteParametersConfig`] for more information regarding their
160-
/// respective default values.
181+
/// See [`AnchorChannelsConfig`], [`RouteParametersConfig`], and
182+
/// [`ForwardedPaymentTrackingMode`] for more information regarding their respective default values.
161183
///
162184
/// [`Node`]: crate::Node
163185
pub struct Config {
@@ -219,6 +241,8 @@ pub struct Config {
219241
///
220242
/// [BIP 353]: https://github.com/bitcoin/bips/blob/master/bip-0353.mediawiki
221243
pub hrn_config: HumanReadableNamesConfig,
244+
/// The mode used for tracking forwarded payments.
245+
pub forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode,
222246
}
223247

224248
impl Default for Config {
@@ -235,6 +259,7 @@ impl Default for Config {
235259
route_parameters: None,
236260
node_alias: None,
237261
hrn_config: HumanReadableNamesConfig::default(),
262+
forwarded_payment_tracking_mode: ForwardedPaymentTrackingMode::default(),
238263
}
239264
}
240265
}

0 commit comments

Comments
 (0)