Skip to content

Commit dc4fa49

Browse files
authored
Merge pull request #815 from randomlogin/add-probing-service
Add probing service
2 parents df8fe95 + f5fbf42 commit dc4fa49

10 files changed

Lines changed: 1564 additions & 19 deletions

File tree

bindings/ldk_node.udl

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@ typedef dictionary TorConfig;
1313

1414
typedef interface NodeEntropy;
1515

16+
typedef interface ProbingConfig;
17+
1618
typedef enum WordCount;
1719

1820
[Remote]
@@ -32,6 +34,18 @@ interface LogWriter {
3234
void log(LogRecord record);
3335
};
3436

37+
interface ProbingConfigBuilder {
38+
[Name=high_degree]
39+
constructor(u64 top_node_count);
40+
[Name=random_walk]
41+
constructor(u64 max_hops);
42+
void set_interval(u64 secs);
43+
void set_max_locked_msat(u64 max_msat);
44+
void set_diversity_penalty_msat(u64 penalty_msat);
45+
void set_cooldown(u64 secs);
46+
ProbingConfig build();
47+
};
48+
3549
interface Builder {
3650
constructor();
3751
[Name=from_config]
@@ -59,6 +73,7 @@ interface Builder {
5973
void set_node_alias(string node_alias);
6074
[Throws=BuildError]
6175
void set_async_payments_role(AsyncPaymentsRole? role);
76+
void set_probing_config(ProbingConfig config);
6277
[Throws=BuildError]
6378
Node build(NodeEntropy node_entropy);
6479
[Throws=BuildError]

src/builder.rs

Lines changed: 91 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@ use crate::config::{
5050
default_user_config, may_announce_channel, AnnounceError, AsyncPaymentsRole,
5151
BitcoindRestClientConfig, Config, ElectrumSyncConfig, EsploraSyncConfig, HRNResolverConfig,
5252
TorConfig, DEFAULT_ESPLORA_SERVER_URL, DEFAULT_LOG_FILENAME, DEFAULT_LOG_LEVEL,
53+
DEFAULT_MAX_PROBE_AMOUNT_MSAT, DEFAULT_MIN_PROBE_AMOUNT_MSAT,
5354
};
5455
use crate::connection::ConnectionManager;
5556
use crate::entropy::NodeEntropy;
@@ -74,6 +75,10 @@ use crate::logger::{log_error, LdkLogger, LogLevel, LogWriter, Logger};
7475
use crate::message_handler::NodeCustomMessageHandler;
7576
use crate::payment::asynchronous::om_mailbox::OnionMessageMailbox;
7677
use crate::peer_store::PeerStore;
78+
use crate::probing::{
79+
HighDegreeStrategy, Prober, ProbingConfig, ProbingStrategy, ProbingStrategyKind,
80+
RandomWalkStrategy,
81+
};
7782
use crate::runtime::{Runtime, RuntimeSpawner};
7883
use crate::tx_broadcaster::TransactionBroadcaster;
7984
use crate::types::{
@@ -306,6 +311,7 @@ pub struct NodeBuilder {
306311
async_payments_role: Option<AsyncPaymentsRole>,
307312
runtime_handle: Option<tokio::runtime::Handle>,
308313
pathfinding_scores_sync_config: Option<PathfindingScoresSyncConfig>,
314+
probing_config: Option<ProbingConfig>,
309315
}
310316

311317
impl NodeBuilder {
@@ -323,6 +329,7 @@ impl NodeBuilder {
323329
let log_writer_config = None;
324330
let runtime_handle = None;
325331
let pathfinding_scores_sync_config = None;
332+
let probing_config = None;
326333
Self {
327334
config,
328335
chain_data_source_config,
@@ -332,6 +339,7 @@ impl NodeBuilder {
332339
runtime_handle,
333340
async_payments_role: None,
334341
pathfinding_scores_sync_config,
342+
probing_config,
335343
}
336344
}
337345

@@ -629,6 +637,31 @@ impl NodeBuilder {
629637
Ok(self)
630638
}
631639

640+
/// Sets background probing config.
641+
///
642+
/// Use [`ProbingConfigBuilder`] to build the configuration:
643+
/// ```no_run
644+
/// # #[cfg(not(feature = "uniffi"))]
645+
/// # {
646+
/// use std::time::Duration;
647+
/// use ldk_node::Builder;
648+
/// use ldk_node::probing::ProbingConfigBuilder;
649+
///
650+
/// let mut builder = Builder::new();
651+
/// builder.set_probing_config(
652+
/// ProbingConfigBuilder::high_degree(100)
653+
/// .interval(Duration::from_secs(30))
654+
/// .build()
655+
/// );
656+
/// # }
657+
/// ```
658+
///
659+
/// [`ProbingConfigBuilder`]: crate::probing::ProbingConfigBuilder
660+
pub fn set_probing_config(&mut self, config: ProbingConfig) -> &mut Self {
661+
self.probing_config = Some(config);
662+
self
663+
}
664+
632665
/// Builds a [`Node`] instance with a [`SqliteStore`] backend and according to the options
633666
/// previously configured.
634667
pub fn build(&self, node_entropy: NodeEntropy) -> Result<Node, BuildError> {
@@ -868,6 +901,7 @@ impl NodeBuilder {
868901
self.gossip_source_config.as_ref(),
869902
self.liquidity_source_config.as_ref(),
870903
self.pathfinding_scores_sync_config.as_ref(),
904+
self.probing_config.as_ref(),
871905
self.async_payments_role,
872906
seed_bytes,
873907
runtime,
@@ -1166,6 +1200,15 @@ impl ArcedNodeBuilder {
11661200
self.inner.write().expect("lock").set_async_payments_role(role).map(|_| ())
11671201
}
11681202

1203+
/// Configures background probing.
1204+
///
1205+
/// Use [`ProbingConfigBuilder`] to build the configuration.
1206+
///
1207+
/// [`ProbingConfigBuilder`]: crate::probing::ProbingConfigBuilder
1208+
pub fn set_probing_config(&self, config: Arc<ProbingConfig>) {
1209+
self.inner.write().expect("lock").set_probing_config((*config).clone());
1210+
}
1211+
11691212
/// Builds a [`Node`] instance with a [`SqliteStore`] backend and according to the options
11701213
/// previously configured.
11711214
pub fn build(&self, node_entropy: Arc<NodeEntropy>) -> Result<Arc<Node>, BuildError> {
@@ -1361,8 +1404,8 @@ fn build_with_store_internal(
13611404
gossip_source_config: Option<&GossipSourceConfig>,
13621405
liquidity_source_config: Option<&LiquiditySourceConfig>,
13631406
pathfinding_scores_sync_config: Option<&PathfindingScoresSyncConfig>,
1364-
async_payments_role: Option<AsyncPaymentsRole>, seed_bytes: [u8; 64], runtime: Arc<Runtime>,
1365-
logger: Arc<Logger>, kv_store: Arc<DynStore>,
1407+
probing_config: Option<&ProbingConfig>, async_payments_role: Option<AsyncPaymentsRole>,
1408+
seed_bytes: [u8; 64], runtime: Arc<Runtime>, logger: Arc<Logger>, kv_store: Arc<DynStore>,
13661409
) -> Result<Node, BuildError> {
13671410
optionally_install_rustls_cryptoprovider();
13681411

@@ -2219,6 +2262,51 @@ fn build_with_store_internal(
22192262
_leak_checker.0.push(Arc::downgrade(&wallet) as Weak<dyn Any + Send + Sync>);
22202263
}
22212264

2265+
let prober = probing_config.map(|probing_cfg| {
2266+
let strategy: Arc<dyn ProbingStrategy> = match &probing_cfg.kind {
2267+
ProbingStrategyKind::HighDegree { top_node_count } => {
2268+
// Dedicated router for probing so the diversity penalty doesn't interfere
2269+
// with real payments; shares the scorer so probe results still train it.
2270+
let mut probing_fee_params = ProbabilisticScoringFeeParameters::default();
2271+
if let Some(penalty) = probing_cfg.diversity_penalty_msat {
2272+
probing_fee_params.probing_diversity_penalty_msat = penalty;
2273+
}
2274+
let probing_router = Arc::new(DefaultRouter::new(
2275+
Arc::clone(&network_graph),
2276+
Arc::clone(&logger),
2277+
Arc::clone(&keys_manager),
2278+
Arc::clone(&scorer),
2279+
probing_fee_params,
2280+
));
2281+
Arc::new(HighDegreeStrategy::new(
2282+
Arc::clone(&network_graph),
2283+
Arc::clone(&channel_manager),
2284+
probing_router,
2285+
*top_node_count,
2286+
DEFAULT_MIN_PROBE_AMOUNT_MSAT,
2287+
DEFAULT_MAX_PROBE_AMOUNT_MSAT,
2288+
probing_cfg.cooldown,
2289+
config.probing_liquidity_limit_multiplier,
2290+
))
2291+
},
2292+
ProbingStrategyKind::RandomWalk { max_hops } => Arc::new(RandomWalkStrategy::new(
2293+
Arc::clone(&network_graph),
2294+
Arc::clone(&channel_manager),
2295+
*max_hops,
2296+
DEFAULT_MIN_PROBE_AMOUNT_MSAT,
2297+
DEFAULT_MAX_PROBE_AMOUNT_MSAT,
2298+
)),
2299+
ProbingStrategyKind::Custom(s) => Arc::clone(s),
2300+
};
2301+
Arc::new(Prober {
2302+
channel_manager: Arc::clone(&channel_manager),
2303+
logger: Arc::clone(&logger),
2304+
strategy,
2305+
interval: probing_cfg.interval,
2306+
max_locked_msat: probing_cfg.max_locked_msat,
2307+
})
2308+
});
2309+
22222310
Ok(Node {
22232311
runtime,
22242312
stop_sender,
@@ -2252,6 +2340,7 @@ fn build_with_store_internal(
22522340
om_mailbox,
22532341
async_payments_role,
22542342
hrn_resolver,
2343+
prober,
22552344
#[cfg(cycle_tests)]
22562345
_leak_checker,
22572346
})

src/config.rs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,12 @@ const DEFAULT_BDK_WALLET_SYNC_INTERVAL_SECS: u64 = 80;
2828
const DEFAULT_LDK_WALLET_SYNC_INTERVAL_SECS: u64 = 30;
2929
const DEFAULT_FEE_RATE_CACHE_UPDATE_INTERVAL_SECS: u64 = 60 * 10;
3030
const DEFAULT_PROBING_LIQUIDITY_LIMIT_MULTIPLIER: u64 = 3;
31+
pub(crate) const DEFAULT_PROBING_INTERVAL_SECS: u64 = 10;
32+
pub(crate) const MIN_PROBING_INTERVAL: Duration = Duration::from_millis(100);
33+
pub(crate) const DEFAULT_PROBED_NODE_COOLDOWN_SECS: u64 = 60 * 60; // 1 hour
34+
pub(crate) const DEFAULT_MAX_PROBE_LOCKED_MSAT: u64 = 100_000_000; // 100k sats
35+
pub(crate) const DEFAULT_MIN_PROBE_AMOUNT_MSAT: u64 = 1_000_000; // 1k sats
36+
pub(crate) const DEFAULT_MAX_PROBE_AMOUNT_MSAT: u64 = 10_000_000; // 10k sats
3137
const DEFAULT_ANCHOR_PER_CHANNEL_RESERVE_SATS: u64 = 25_000;
3238

3339
// The default timeout after which we abort a wallet syncing operation.

src/event.rs

Lines changed: 21 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ use crate::payment::store::{
5252
PaymentDetails, PaymentDetailsUpdate, PaymentDirection, PaymentKind, PaymentStatus,
5353
};
5454
use crate::payment::PaymentMetadata;
55+
use crate::probing::Prober;
5556
use crate::runtime::Runtime;
5657
use crate::types::{
5758
CustomTlvRecord, DynStore, KeysManager, OnionMessenger, PaymentStore, Sweeper, Wallet,
@@ -536,12 +537,13 @@ where
536537
payment_store: Arc<PaymentStore>,
537538
peer_store: Arc<PeerStore<L>>,
538539
keys_manager: Arc<KeysManager>,
539-
runtime: Arc<Runtime>,
540-
logger: L,
541-
config: Arc<Config>,
542540
static_invoice_store: Option<StaticInvoiceStore>,
543541
onion_messenger: Arc<OnionMessenger>,
544542
om_mailbox: Option<Arc<OnionMessageMailbox>>,
543+
prober: Option<Arc<Prober>>,
544+
runtime: Arc<Runtime>,
545+
logger: L,
546+
config: Arc<Config>,
545547
}
546548

547549
impl<L: Deref + Clone + Sync + Send + 'static> EventHandler<L>
@@ -556,8 +558,8 @@ where
556558
liquidity_source: Arc<LiquiditySource<Arc<Logger>>>, payment_store: Arc<PaymentStore>,
557559
peer_store: Arc<PeerStore<L>>, keys_manager: Arc<KeysManager>,
558560
static_invoice_store: Option<StaticInvoiceStore>, onion_messenger: Arc<OnionMessenger>,
559-
om_mailbox: Option<Arc<OnionMessageMailbox>>, runtime: Arc<Runtime>, logger: L,
560-
config: Arc<Config>,
561+
om_mailbox: Option<Arc<OnionMessageMailbox>>, prober: Option<Arc<Prober>>,
562+
runtime: Arc<Runtime>, logger: L, config: Arc<Config>,
561563
) -> Self {
562564
Self {
563565
event_queue,
@@ -571,12 +573,13 @@ where
571573
payment_store,
572574
peer_store,
573575
keys_manager,
574-
logger,
575-
runtime,
576-
config,
577576
static_invoice_store,
578577
onion_messenger,
579578
om_mailbox,
579+
prober,
580+
runtime,
581+
logger,
582+
config,
580583
}
581584
}
582585

@@ -1208,8 +1211,16 @@ where
12081211

12091212
LdkEvent::PaymentPathSuccessful { .. } => {},
12101213
LdkEvent::PaymentPathFailed { .. } => {},
1211-
LdkEvent::ProbeSuccessful { .. } => {},
1212-
LdkEvent::ProbeFailed { .. } => {},
1214+
LdkEvent::ProbeSuccessful { path, payment_id, .. } => {
1215+
if let Some(prober) = &self.prober {
1216+
prober.handle_background_probe_successful(&path, payment_id);
1217+
}
1218+
},
1219+
LdkEvent::ProbeFailed { path, payment_id, .. } => {
1220+
if let Some(prober) = &self.prober {
1221+
prober.handle_background_probe_failed(&path, payment_id);
1222+
}
1223+
},
12131224
LdkEvent::HTLCHandlingFailed { failure_type, .. } => {
12141225
self.liquidity_source
12151226
.lsps2_service()

src/ffi/types.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,7 @@ pub use crate::entropy::{generate_entropy_mnemonic, NodeEntropy, WordCount};
149149
use crate::error::Error;
150150
pub use crate::liquidity::LSPS1OrderStatus;
151151
pub use crate::logger::{LogLevel, LogRecord, LogWriter};
152+
pub use crate::probing::ProbingConfig;
152153
use crate::{hex_utils, SocketAddress, UserChannelId};
153154

154155
uniffi::custom_type!(PublicKey, String, {

src/lib.rs

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -101,10 +101,12 @@ pub mod logger;
101101
mod message_handler;
102102
pub mod payment;
103103
mod peer_store;
104+
pub mod probing;
104105
mod runtime;
105106
mod scoring;
106107
mod tx_broadcaster;
107108
mod types;
109+
mod util;
108110
mod wallet;
109111

110112
use std::default::Default;
@@ -172,6 +174,9 @@ use payment::{
172174
UnifiedPayment,
173175
};
174176
use peer_store::{PeerInfo, PeerStore};
177+
#[cfg(feature = "uniffi")]
178+
pub use probing::ArcedProbingConfigBuilder as ProbingConfigBuilder;
179+
use probing::{run_prober, Prober};
175180
use runtime::Runtime;
176181
pub use tokio;
177182
use types::{
@@ -250,6 +255,7 @@ pub struct Node {
250255
om_mailbox: Option<Arc<OnionMessageMailbox>>,
251256
async_payments_role: Option<AsyncPaymentsRole>,
252257
hrn_resolver: HRNResolver,
258+
prober: Option<Arc<Prober>>,
253259
#[cfg(cycle_tests)]
254260
_leak_checker: LeakChecker,
255261
}
@@ -610,11 +616,19 @@ impl Node {
610616
static_invoice_store,
611617
Arc::clone(&self.onion_messenger),
612618
self.om_mailbox.clone(),
619+
self.prober.clone(),
613620
Arc::clone(&self.runtime),
614621
Arc::clone(&self.logger),
615622
Arc::clone(&self.config),
616623
));
617624

625+
if let Some(prober) = self.prober.clone() {
626+
let stop_rx = self.stop_sender.subscribe();
627+
self.runtime.spawn_cancellable_background_task(async move {
628+
run_prober(prober, stop_rx).await;
629+
});
630+
}
631+
618632
// Setup background processing
619633
let background_persister = Arc::clone(&self.kv_store);
620634
let background_event_handler = Arc::clone(&event_handler);
@@ -1145,6 +1159,11 @@ impl Node {
11451159
))
11461160
}
11471161

1162+
/// Returns a reference to the [`Prober`], or `None` if no probing strategy is configured.
1163+
pub fn prober(&self) -> Option<&Prober> {
1164+
self.prober.as_deref()
1165+
}
1166+
11481167
/// Retrieve a list of known channels.
11491168
pub fn list_channels(&self) -> Vec<ChannelDetails> {
11501169
self.channel_manager

0 commit comments

Comments
 (0)