@@ -65,7 +65,13 @@ use crate::io::utils::{
6565} ;
6666use crate :: io:: vss_store:: VssStoreBuilder ;
6767use 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::{
8288use crate :: runtime:: { Runtime , RuntimeSpawner } ;
8389use crate :: tx_broadcaster:: TransactionBroadcaster ;
8490use 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 } )
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,
0 commit comments