Let `LiquidityManager` take a `KVStore` and add `Sync` wrapper
What changed, and why it matters
This commit is a routine internal refactor for the rust-lightning project. It adds a key-value store parameter to the LiquidityManager component and introduces a synchronous wrapper (LiquidityManagerSync) so existing code can keep using the manager while the underlying implementation becomes asynchronous. There is no indication of a security fix or vulnerability being addressed.
No security action required; treat as normal dependency/API refactor.
Security signals we found
No strong security signals were identified.
Evidence from the diff
The change threads a KVStore through LiquidityManager, adds a KVStoreSync-wrapping LiquidityManagerSync facade, and updates background-processor and test code to use the new sync wrapper. It is purely architectural plumbing for future async persistence work. No security-sensitive logic is modified; no bounds checks, cryptographic operations, or network parsing behavior are changed beyond type signatures and delegation.
Changed components
lightning-liquidity/src/manager.rslightning-liquidity/src/lib.rslightning-background-processor/src/lib.rsfuzz/src/lsps_message.rslightning-liquidity/tests/common/mod.rsInspect captured patch +535 / −43
diff --git a/fuzz/src/lsps_message.rs b/fuzz/src/lsps_message.rs
index 299b9f0..2b5ee68 100644
--- a/fuzz/src/lsps_message.rs
+++ b/fuzz/src/lsps_message.rs
@@ -21,7 +21,7 @@ use lightning::util::test_utils::{
};
use lightning_liquidity::lsps0::ser::LSPS_MESSAGE_TYPE_ID;
-use lightning_liquidity::LiquidityManager;
+use lightning_liquidity::LiquidityManagerSync;
use core::time::Duration;
@@ -77,12 +77,13 @@ pub fn do_test(data: &[u8]) {
genesis_block.header.time,
));
- let liquidity_manager = Arc::new(LiquidityManager::new(
+ let liquidity_manager = Arc::new(LiquidityManagerSync::new(
Arc::clone(&keys_manager),
Arc::clone(&keys_manager),
Arc::clone(&manager),
None::<Arc<dyn Filter + Send + Sync>>,
None,
+ kv_store,
None,
None,
));
diff --git a/lightning-background-processor/Cargo.toml b/lightning-background-processor/Cargo.toml
index 47d3211..415676f 100644
--- a/lightning-background-processor/Cargo.toml
+++ b/lightning-background-processor/Cargo.toml
@@ -31,6 +31,7 @@ possiblyrandom = { version = "0.2", path = "../possiblyrandom", default-features
tokio = { version = "1.35", features = [ "macros", "rt", "rt-multi-thread", "sync", "time" ] }
lightning = { version = "0.2.0", path = "../lightning", features = ["_test_utils"] }
lightning-invoice = { version = "0.34.0", path = "../lightning-invoice" }
+lightning-liquidity = { version = "0.2.0", path = "../lightning-liquidity", default-features = false, features = ["_test_utils"] }
lightning-persister = { version = "0.2.0", path = "../lightning-persister" }
[lints]
diff --git a/lightning-background-processor/src/lib.rs b/lightning-background-processor/src/lib.rs
index 95adc65..f10a3e2 100644
--- a/lightning-background-processor/src/lib.rs
+++ b/lightning-background-processor/src/lib.rs
@@ -69,6 +69,8 @@ use lightning::util::wakers::Sleeper;
use lightning_rapid_gossip_sync::RapidGossipSync;
use lightning_liquidity::ALiquidityManager;
+#[cfg(feature = "std")]
+use lightning_liquidity::ALiquidityManagerSync;
use core::ops::Deref;
use core::time::Duration;
@@ -424,6 +426,31 @@ pub const NO_LIQUIDITY_MANAGER: Option<
CM = &DynChannelManager,
Filter = dyn chain::Filter,
C = &dyn chain::Filter,
+ KVStore = dyn lightning::util::persist::KVStore,
+ K = &dyn lightning::util::persist::KVStore,
+ TimeProvider = dyn lightning_liquidity::utils::time::TimeProvider,
+ TP = &dyn lightning_liquidity::utils::time::TimeProvider,
+ > + Send
+ + Sync,
+ >,
+> = None;
+
+/// When initializing a background processor without a liquidity manager, this can be used to avoid
+/// specifying a concrete `LiquidityManagerSync` type.
+#[cfg(all(not(c_bindings), feature = "std"))]
+pub const NO_LIQUIDITY_MANAGER_SYNC: Option<
+ Arc<
+ dyn ALiquidityManagerSync<
+ EntropySource = dyn EntropySource,
+ ES = &dyn EntropySource,
+ NodeSigner = dyn lightning::sign::NodeSigner,
+ NS = &dyn lightning::sign::NodeSigner,
+ AChannelManager = DynChannelManager,
+ CM = &DynChannelManager,
+ Filter = dyn chain::Filter,
+ C = &dyn chain::Filter,
+ KVStoreSync = dyn lightning::util::persist::KVStoreSync,
+ KS = &dyn lightning::util::persist::KVStoreSync,
TimeProvider = dyn lightning_liquidity::utils::time::TimeProvider,
TP = &dyn lightning_liquidity::utils::time::TimeProvider,
> + Send
@@ -731,7 +758,7 @@ use futures_util::{dummy_waker, Joiner, OptionalSelector, Selector, SelectorOutp
/// # type P2PGossipSync<UL> = lightning::routing::gossip::P2PGossipSync<Arc<NetworkGraph>, Arc<UL>, Arc<Logger>>;
/// # type ChannelManager<B, F, FE> = lightning::ln::channelmanager::SimpleArcChannelManager<ChainMonitor<B, F, FE>, B, FE, Logger>;
/// # type OnionMessenger<B, F, FE> = lightning::onion_message::messenger::OnionMessenger<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<Logger>, Arc<ChannelManager<B, F, FE>>, Arc<lightning::onion_message::messenger::DefaultMessageRouter<Arc<NetworkGraph>, Arc<Logger>, Arc<lightning::sign::KeysManager>>>, Arc<ChannelManager<B, F, FE>>, lightning::ln::peer_handler::IgnoringMessageHandler, lightning::ln::peer_handler::IgnoringMessageHandler, lightning::ln::peer_handler::IgnoringMessageHandler>;
-/// # type LiquidityManager<B, F, FE> = lightning_liquidity::LiquidityManager<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<ChannelManager<B, F, FE>>, Arc<F>, Arc<DefaultTimeProvider>>;
+/// # type LiquidityManager<B, F, FE> = lightning_liquidity::LiquidityManager<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<ChannelManager<B, F, FE>>, Arc<F>, Arc<Store>, Arc<DefaultTimeProvider>>;
/// # type Scorer = RwLock<lightning::routing::scoring::ProbabilisticScorer<Arc<NetworkGraph>, Arc<Logger>>>;
/// # type PeerManager<B, F, FE, UL> = lightning::ln::peer_handler::SimpleArcPeerManager<SocketDescriptor, ChainMonitor<B, F, FE>, B, FE, Arc<UL>, Logger, F, StoreSync>;
/// # type OutputSweeper<B, D, FE, F, O> = lightning::util::sweep::OutputSweeper<Arc<B>, Arc<D>, Arc<FE>, Arc<F>, Arc<Store>, Arc<Logger>, Arc<O>>;
@@ -1450,7 +1477,7 @@ impl BackgroundProcessor {
CM::Target: AChannelManager,
OM::Target: AOnionMessenger,
PM::Target: APeerManager,
- LM::Target: ALiquidityManager,
+ LM::Target: ALiquidityManagerSync,
D::Target: ChangeDestinationSourceSync,
O::Target: 'static + OutputSpender,
K::Target: 'static + KVStoreSync,
@@ -1793,7 +1820,7 @@ mod tests {
use lightning::util::test_utils;
use lightning::{get_event, get_event_msg};
use lightning_liquidity::utils::time::DefaultTimeProvider;
- use lightning_liquidity::LiquidityManager;
+ use lightning_liquidity::{ALiquidityManagerSync, LiquidityManagerSync};
use lightning_persister::fs_store::FilesystemStore;
use lightning_rapid_gossip_sync::RapidGossipSync;
use std::collections::VecDeque;
@@ -1890,11 +1917,12 @@ mod tests {
IgnoringMessageHandler,
>;
- type LM = LiquidityManager<
+ type LM = LiquidityManagerSync<
Arc<KeysManager>,
Arc<KeysManager>,
Arc<ChannelManager>,
Arc<dyn Filter + Sync + Send>,
+ Arc<Persister>,
Arc<DefaultTimeProvider>,
>;
@@ -2342,12 +2370,13 @@ mod tests {
Arc::clone(&logger),
Arc::clone(&keys_manager),
));
- let liquidity_manager = Arc::new(LiquidityManager::new(
+ let liquidity_manager = Arc::new(LiquidityManagerSync::new(
Arc::clone(&keys_manager),
Arc::clone(&keys_manager),
Arc::clone(&manager),
None,
None,
+ Arc::clone(&kv_store),
None,
None,
));
@@ -2727,7 +2756,7 @@ mod tests {
Some(Arc::clone(&nodes[0].messenger)),
nodes[0].rapid_gossip_sync(),
Arc::clone(&nodes[0].peer_manager),
- Some(Arc::clone(&nodes[0].liquidity_manager)),
+ Some(nodes[0].liquidity_manager.get_lm_async()),
Some(nodes[0].sweeper.sweeper_async()),
Arc::clone(&nodes[0].logger),
Some(Arc::clone(&nodes[0].scorer)),
@@ -3236,7 +3265,7 @@ mod tests {
Some(Arc::clone(&nodes[0].messenger)),
nodes[0].rapid_gossip_sync(),
Arc::clone(&nodes[0].peer_manager),
- Some(Arc::clone(&nodes[0].liquidity_manager)),
+ Some(nodes[0].liquidity_manager.get_lm_async()),
Some(nodes[0].sweeper.sweeper_async()),
Arc::clone(&nodes[0].logger),
Some(Arc::clone(&nodes[0].scorer)),
@@ -3451,7 +3480,7 @@ mod tests {
Some(Arc::clone(&nodes[0].messenger)),
nodes[0].no_gossip_sync(),
Arc::clone(&nodes[0].peer_manager),
- Some(Arc::clone(&nodes[0].liquidity_manager)),
+ Some(nodes[0].liquidity_manager.get_lm_async()),
Some(nodes[0].sweeper.sweeper_async()),
Arc::clone(&nodes[0].logger),
Some(Arc::clone(&nodes[0].scorer)),
@@ -3500,7 +3529,7 @@ mod tests {
crate::NO_ONION_MESSENGER,
nodes[0].no_gossip_sync(),
Arc::clone(&nodes[0].peer_manager),
- crate::NO_LIQUIDITY_MANAGER,
+ crate::NO_LIQUIDITY_MANAGER_SYNC,
Some(Arc::clone(&nodes[0].sweeper)),
Arc::clone(&nodes[0].logger),
Some(Arc::clone(&nodes[0].scorer)),
diff --git a/lightning-liquidity/Cargo.toml b/lightning-liquidity/Cargo.toml
index f301e4f..ff270a7 100644
--- a/lightning-liquidity/Cargo.toml
+++ b/lightning-liquidity/Cargo.toml
@@ -18,6 +18,7 @@ default = ["std", "time"]
std = ["lightning/std"]
time = ["std"]
backtrace = ["dep:backtrace"]
+_test_utils = []
[dependencies]
lightning = { version = "0.2.0", path = "../lightning", default-features = false }
diff --git a/lightning-liquidity/src/lib.rs b/lightning-liquidity/src/lib.rs
index 6c26b21..e8875e1 100644
--- a/lightning-liquidity/src/lib.rs
+++ b/lightning-liquidity/src/lib.rs
@@ -73,5 +73,6 @@ mod tests;
pub mod utils;
pub use manager::{
- ALiquidityManager, LiquidityClientConfig, LiquidityManager, LiquidityServiceConfig,
+ ALiquidityManager, ALiquidityManagerSync, LiquidityClientConfig, LiquidityManager,
+ LiquidityManagerSync, LiquidityServiceConfig,
};
diff --git a/lightning-liquidity/src/manager.rs b/lightning-liquidity/src/manager.rs
index 3d49679..a64efa4 100644
--- a/lightning-liquidity/src/manager.rs
+++ b/lightning-liquidity/src/manager.rs
@@ -45,6 +45,7 @@ use lightning::ln::peer_handler::CustomMessageHandler;
use lightning::ln::wire::CustomMessageReader;
use lightning::sign::{EntropySource, NodeSigner};
use lightning::util::logger::Level;
+use lightning::util::persist::{KVStore, KVStoreSync, KVStoreSyncWrapper};
use lightning::util::ser::{LengthLimitedRead, LengthReadable};
use lightning::util::wakers::Future;
@@ -108,12 +109,17 @@ pub trait ALiquidityManager {
type Filter: Filter + ?Sized;
/// A type that may be dereferenced to [`Self::Filter`].
type C: Deref<Target = Self::Filter> + Clone;
+ /// A type implementing [`KVStore`].
+ type KVStore: KVStore + ?Sized;
+ /// A type that may be dereferenced to [`Self::KVStore`].
+ type K: Deref<Target = Self::KVStore> + Clone;
/// A type implementing [`TimeProvider`].
type TimeProvider: TimeProvider + ?Sized;
/// A type that may be dereferenced to [`Self::TimeProvider`].
type TP: Deref<Target = Self::TimeProvider> + Clone;
/// Returns a reference to the actual [`LiquidityManager`] object.
- fn get_lm(&self) -> &LiquidityManager<Self::ES, Self::NS, Self::CM, Self::C, Self::TP>;
+ fn get_lm(&self)
+ -> &LiquidityManager<Self::ES, Self::NS, Self::CM, Self::C, Self::K, Self::TP>;
}
impl<
@@ -121,13 +127,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > ALiquidityManager for LiquidityManager<ES, NS, CM, C, TP>
+ > ALiquidityManager for LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
type EntropySource = ES::Target;
@@ -138,9 +146,109 @@ where
type CM = CM;
type Filter = C::Target;
type C = C;
+ type KVStore = K::Target;
+ type K = K;
type TimeProvider = TP::Target;
type TP = TP;
- fn get_lm(&self) -> &LiquidityManager<ES, NS, CM, C, TP> {
+ fn get_lm(&self) -> &LiquidityManager<ES, NS, CM, C, K, TP> {
+ self
+ }
+}
+
+/// A trivial trait which describes any [`LiquidityManagerSync`].
+///
+/// This is not exported to bindings users as general cover traits aren't useful in other
+/// languages.
+pub trait ALiquidityManagerSync {
+ /// A type implementing [`EntropySource`]
+ type EntropySource: EntropySource + ?Sized;
+ /// A type that may be dereferenced to [`Self::EntropySource`].
+ type ES: Deref<Target = Self::EntropySource> + Clone;
+ /// A type implementing [`NodeSigner`]
+ type NodeSigner: NodeSigner + ?Sized;
+ /// A type that may be dereferenced to [`Self::NodeSigner`].
+ type NS: Deref<Target = Self::NodeSigner> + Clone;
+ /// A type implementing [`AChannelManager`]
+ type AChannelManager: AChannelManager + ?Sized;
+ /// A type that may be dereferenced to [`Self::AChannelManager`].
+ type CM: Deref<Target = Self::AChannelManager> + Clone;
+ /// A type implementing [`Filter`].
+ type Filter: Filter + ?Sized;
+ /// A type that may be dereferenced to [`Self::Filter`].
+ type C: Deref<Target = Self::Filter> + Clone;
+ /// A type implementing [`KVStoreSync`].
+ type KVStoreSync: KVStoreSync + ?Sized;
+ /// A type that may be dereferenced to [`Self::KVStoreSync`].
+ type KS: Deref<Target = Self::KVStoreSync> + Clone;
+ /// A type implementing [`TimeProvider`].
+ type TimeProvider: TimeProvider + ?Sized;
+ /// A type that may be dereferenced to [`Self::TimeProvider`].
+ type TP: Deref<Target = Self::TimeProvider> + Clone;
+ /// Returns the inner async [`LiquidityManager`] for testing purposes.
+ #[cfg(any(test, feature = "_test_utils"))]
+ fn get_lm_async(
+ &self,
+ ) -> Arc<
+ LiquidityManager<
+ Self::ES,
+ Self::NS,
+ Self::CM,
+ Self::C,
+ Arc<KVStoreSyncWrapper<Self::KS>>,
+ Self::TP,
+ >,
+ >;
+ /// Returns a reference to the actual [`LiquidityManager`] object.
+ fn get_lm(
+ &self,
+ ) -> &LiquidityManagerSync<Self::ES, Self::NS, Self::CM, Self::C, Self::KS, Self::TP>;
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > ALiquidityManagerSync for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ type EntropySource = ES::Target;
+ type ES = ES;
+ type NodeSigner = NS::Target;
+ type NS = NS;
+ type AChannelManager = CM::Target;
+ type CM = CM;
+ type Filter = C::Target;
+ type C = C;
+ type KVStoreSync = KS::Target;
+ type KS = KS;
+ type TimeProvider = TP::Target;
+ type TP = TP;
+ /// Returns the inner async [`LiquidityManager`] for testing purposes.
+ #[cfg(any(test, feature = "_test_utils"))]
+ fn get_lm_async(
+ &self,
+ ) -> Arc<
+ LiquidityManager<
+ Self::ES,
+ Self::NS,
+ Self::CM,
+ Self::C,
+ Arc<KVStoreSyncWrapper<Self::KS>>,
+ Self::TP,
+ >,
+ > {
+ Arc::clone(&self.inner)
+ }
+ fn get_lm(&self) -> &LiquidityManagerSync<ES, NS, CM, C, KS, TP> {
self
}
}
@@ -169,12 +277,14 @@ pub struct LiquidityManager<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
> where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
pending_messages: Arc<MessageQueue>,
@@ -195,21 +305,29 @@ pub struct LiquidityManager<
_client_config: Option<LiquidityClientConfig>,
best_block: RwLock<Option<BestBlock>>,
_chain_source: Option<C>,
+ kv_store: K,
}
#[cfg(feature = "time")]
-impl<ES: Deref + Clone, NS: Deref + Clone, CM: Deref + Clone, C: Deref + Clone>
- LiquidityManager<ES, NS, CM, C, Arc<DefaultTimeProvider>>
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ K: Deref + Clone,
+ > LiquidityManager<ES, NS, CM, C, K, Arc<DefaultTimeProvider>>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
{
/// Constructor for the [`LiquidityManager`] using the default system clock
pub fn new(
entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, service_config: Option<LiquidityServiceConfig>,
+ chain_params: Option<ChainParameters>, kv_store: K,
+ service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>,
) -> Self {
let time_provider = Arc::new(DefaultTimeProvider);
@@ -219,6 +337,7 @@ where
channel_manager,
chain_source,
chain_params,
+ kv_store,
service_config,
client_config,
time_provider,
@@ -231,13 +350,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > LiquidityManager<ES, NS, CM, C, TP>
+ > LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
/// Constructor for the [`LiquidityManager`] with a custom time provider.
@@ -248,7 +369,8 @@ where
/// [`LiquidityClientConfig`] and [`LiquidityServiceConfig`].
pub fn new_with_custom_time_provider(
entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, service_config: Option<LiquidityServiceConfig>,
+ chain_params: Option<ChainParameters>, kv_store: K,
+ service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>, time_provider: TP,
) -> Self {
let pending_messages = Arc::new(MessageQueue::new());
@@ -373,6 +495,7 @@ where
_client_config: client_config,
best_block: RwLock::new(chain_params.map(|chain_params| chain_params.best_block)),
_chain_source: chain_source,
+ kv_store,
}
}
@@ -388,7 +511,7 @@ where
/// Returns a reference to the LSPS1 client-side handler.
///
- /// The returned hendler allows to initiate the LSPS1 client-side flow, i.e., allows to request
+ /// The returned handler allows to initiate the LSPS1 client-side flow, i.e., allows to request
/// channels from the configured LSP.
pub fn lsps1_client_handler(&self) -> Option<&LSPS1ClientHandler<ES>> {
self.lsps1_client_handler.as_ref()
@@ -402,7 +525,7 @@ where
/// Returns a reference to the LSPS2 client-side handler.
///
- /// The returned hendler allows to initiate the LSPS2 client-side flow. That is, it allows to
+ /// The returned handler allows to initiate the LSPS2 client-side flow. That is, it allows to
/// retrieve all necessary data to create 'just-in-time' invoices that, when paid, will have
/// the configured LSP open a 'just-in-time' channel.
pub fn lsps2_client_handler(&self) -> Option<&LSPS2ClientHandler<ES>> {
@@ -411,21 +534,21 @@ where
/// Returns a reference to the LSPS2 server-side handler.
///
- /// The returned hendler allows to initiate the LSPS2 service-side flow.
+ /// The returned handler allows to initiate the LSPS2 service-side flow.
pub fn lsps2_service_handler(&self) -> Option<&LSPS2ServiceHandler<CM>> {
self.lsps2_service_handler.as_ref()
}
/// Returns a reference to the LSPS5 client-side handler.
///
- /// The returned hendler allows to initiate the LSPS5 client-side flow. That is, it allows to
+ /// The returned handler allows to initiate the LSPS5 client-side flow. That is, it allows to
pub fn lsps5_client_handler(&self) -> Option<&LSPS5ClientHandler<ES>> {
self.lsps5_client_handler.as_ref()
}
/// Returns a reference to the LSPS5 server-side handler.
///
- /// The returned hendler allows to initiate the LSPS5 service-side flow.
+ /// The returned handler allows to initiate the LSPS5 service-side flow.
pub fn lsps5_service_handler(&self) -> Option<&LSPS5ServiceHandler<CM, NS, TP>> {
self.lsps5_service_handler.as_ref()
}
@@ -441,15 +564,10 @@ where
/// Blocks the current thread until next event is ready and returns it.
///
- /// Typically you would spawn a thread or task that calls this in a loop.
- ///
- /// **Note**: Users must handle events as soon as possible to avoid an increased event queue
- /// memory footprint. We will start dropping any generated events after
- /// [`MAX_EVENT_QUEUE_SIZE`] has been reached.
- ///
- /// [`MAX_EVENT_QUEUE_SIZE`]: crate::events::MAX_EVENT_QUEUE_SIZE
+ /// Only available via the [`LiquidityManagerSync`] interface to avoid having users
+ /// accidentally blocking their async contexts.
#[cfg(feature = "std")]
- pub fn wait_next_event(&self) -> LiquidityEvent {
+ pub(crate) fn wait_next_event(&self) -> LiquidityEvent {
self.pending_events.wait_next_event()
}
@@ -608,13 +726,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageReader for LiquidityManager<ES, NS, CM, C, TP>
+ > CustomMessageReader for LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
type CustomMessage = RawLSPSMessage;
@@ -636,13 +756,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageHandler for LiquidityManager<ES, NS, CM, C, TP>
+ > CustomMessageHandler for LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
fn handle_custom_message(
@@ -766,13 +888,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > Listen for LiquidityManager<ES, NS, CM, C, TP>
+ > Listen for LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
fn filtered_block_connected(
@@ -808,13 +932,15 @@ impl<
NS: Deref + Clone,
CM: Deref + Clone,
C: Deref + Clone,
+ K: Deref + Clone,
TP: Deref + Clone,
- > Confirm for LiquidityManager<ES, NS, CM, C, TP>
+ > Confirm for LiquidityManager<ES, NS, CM, C, K, TP>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
+ K::Target: KVStore,
TP::Target: TimeProvider,
{
fn transactions_confirmed(
@@ -845,3 +971,330 @@ where
Vec::new()
}
}
+
+/// A synchroneous wrapper around [`LiquidityManager`] to be used in contexts where async is not
+/// available.
+pub struct LiquidityManagerSync<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+> where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ inner: Arc<LiquidityManager<ES, NS, CM, C, Arc<KVStoreSyncWrapper<KS>>, TP>>,
+}
+
+#[cfg(feature = "time")]
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ > LiquidityManagerSync<ES, NS, CM, C, KS, Arc<DefaultTimeProvider>>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ KS::Target: KVStoreSync,
+ C::Target: Filter,
+{
+ /// Constructor for the [`LiquidityManagerSync`] using the default system clock
+ ///
+ /// Wraps [`LiquidityManager::new`].
+ pub fn new(
+ entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
+ chain_params: Option<ChainParameters>, kv_store_sync: KS,
+ service_config: Option<LiquidityServiceConfig>,
+ client_config: Option<LiquidityClientConfig>,
+ ) -> Self {
+ let kv_store = Arc::new(KVStoreSyncWrapper(kv_store_sync));
+ let inner = Arc::new(LiquidityManager::new(
+ entropy_source,
+ node_signer,
+ channel_manager,
+ chain_source,
+ chain_params,
+ kv_store,
+ service_config,
+ client_config,
+ ));
+ Self { inner }
+ }
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ /// Constructor for the [`LiquidityManagerSync`] with a custom time provider.
+ ///
+ /// Wraps [`LiquidityManager::new_with_custom_time_provider`].
+ pub fn new_with_custom_time_provider(
+ entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
+ chain_params: Option<ChainParameters>, kv_store_sync: KS,
+ service_config: Option<LiquidityServiceConfig>,
+ client_config: Option<LiquidityClientConfig>, time_provider: TP,
+ ) -> Self {
+ let kv_store = Arc::new(KVStoreSyncWrapper(kv_store_sync));
+ let inner = Arc::new(LiquidityManager::new_with_custom_time_provider(
+ entropy_source,
+ node_signer,
+ channel_manager,
+ chain_source,
+ chain_params,
+ kv_store,
+ service_config,
+ client_config,
+ time_provider,
+ ));
+ Self { inner }
+ }
+
+ /// Returns a reference to the LSPS0 client-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps0_client_handler`].
+ pub fn lsps0_client_handler(&self) -> &LSPS0ClientHandler<ES> {
+ self.inner.lsps0_client_handler()
+ }
+
+ /// Returns a reference to the LSPS0 server-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps0_service_handler`].
+ pub fn lsps0_service_handler(&self) -> Option<&LSPS0ServiceHandler> {
+ self.inner.lsps0_service_handler()
+ }
+
+ /// Returns a reference to the LSPS1 client-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps1_client_handler`].
+ pub fn lsps1_client_handler(&self) -> Option<&LSPS1ClientHandler<ES>> {
+ self.inner.lsps1_client_handler()
+ }
+
+ /// Returns a reference to the LSPS1 server-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps1_service_handler`].
+ #[cfg(lsps1_service)]
+ pub fn lsps1_service_handler(&self) -> Option<&LSPS1ServiceHandler<ES, CM, C>> {
+ self.inner.lsps1_service_handler()
+ }
+
+ /// Returns a reference to the LSPS2 client-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps2_client_handler`].
+ pub fn lsps2_client_handler(&self) -> Option<&LSPS2ClientHandler<ES>> {
+ self.inner.lsps2_client_handler()
+ }
+
+ /// Returns a reference to the LSPS2 server-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps2_service_handler`].
+ pub fn lsps2_service_handler(&self) -> Option<&LSPS2ServiceHandler<CM>> {
+ self.inner.lsps2_service_handler()
+ }
+
+ /// Returns a reference to the LSPS5 client-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps5_client_handler`].
+ pub fn lsps5_client_handler(&self) -> Option<&LSPS5ClientHandler<ES>> {
+ self.inner.lsps5_client_handler()
+ }
+
+ /// Returns a reference to the LSPS5 server-side handler.
+ ///
+ /// Wraps [`LiquidityManager::lsps5_service_handler`].
+ pub fn lsps5_service_handler(&self) -> Option<&LSPS5ServiceHandler<CM, NS, TP>> {
+ self.inner.lsps5_service_handler()
+ }
+
+ /// Returns a [`Future`] that will complete when the next batch of pending messages is ready to
+ /// be processed.
+ ///
+ /// Wraps [`LiquidityManager::get_pending_msgs_future`].
+ pub fn get_pending_msgs_future(&self) -> Future {
+ self.inner.get_pending_msgs_future()
+ }
+
+ /// Blocks the current thread until next event is ready and returns it.
+ ///
+ /// Typically you would spawn a thread or task that calls this in a loop.
+ ///
+ /// **Note**: Users must handle events as soon as possible to avoid an increased event queue
+ /// memory footprint. We will start dropping any generated events after
+ /// [`MAX_EVENT_QUEUE_SIZE`] has been reached.
+ ///
+ /// [`MAX_EVENT_QUEUE_SIZE`]: crate::events::MAX_EVENT_QUEUE_SIZE
+ #[cfg(feature = "std")]
+ pub fn wait_next_event(&self) -> LiquidityEvent {
+ self.inner.wait_next_event()
+ }
+
+ /// Returns `Some` if an event is ready.
+ ///
+ /// Wraps [`LiquidityManager::next_event`].
+ pub fn next_event(&self) -> Option<LiquidityEvent> {
+ self.inner.next_event()
+ }
+
+ /// Returns and clears all events without blocking.
+ ///
+ /// Wraps [`LiquidityManager::get_and_clear_pending_events`].
+ pub fn get_and_clear_pending_events(&self) -> Vec<LiquidityEvent> {
+ self.inner.get_and_clear_pending_events()
+ }
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > CustomMessageReader for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ type CustomMessage = RawLSPSMessage;
+
+ fn read<RD: LengthLimitedRead>(
+ &self, message_type: u16, buffer: &mut RD,
+ ) -> Result<Option<Self::CustomMessage>, lightning::ln::msgs::DecodeError> {
+ self.inner.read(message_type, buffer)
+ }
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > CustomMessageHandler for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ fn handle_custom_message(
+ &self, msg: Self::CustomMessage, sender_node_id: PublicKey,
+ ) -> Result<(), lightning::ln::msgs::LightningError> {
+ self.inner.handle_custom_message(msg, sender_node_id)
+ }
+
+ fn get_and_clear_pending_msg(&self) -> Vec<(PublicKey, Self::CustomMessage)> {
+ self.inner.get_and_clear_pending_msg()
+ }
+
+ fn provided_node_features(&self) -> NodeFeatures {
+ self.inner.provided_node_features()
+ }
+
+ fn provided_init_features(&self, their_node_id: PublicKey) -> InitFeatures {
+ self.inner.provided_init_features(their_node_id)
+ }
+
+ fn peer_disconnected(&self, counterparty_node_id: bitcoin::secp256k1::PublicKey) {
+ self.inner.peer_disconnected(counterparty_node_id)
+ }
+ fn peer_connected(
+ &self, counterparty_node_id: bitcoin::secp256k1::PublicKey,
+ init_msg: &lightning::ln::msgs::Init, inbound: bool,
+ ) -> Result<(), ()> {
+ self.inner.peer_connected(counterparty_node_id, init_msg, inbound)
+ }
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > Listen for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ fn filtered_block_connected(
+ &self, header: &bitcoin::block::Header, txdata: &chain::transaction::TransactionData,
+ height: u32,
+ ) {
+ self.inner.filtered_block_connected(header, txdata, height)
+ }
+
+ fn blocks_disconnected(&self, fork_point: BestBlock) {
+ self.inner.blocks_disconnected(fork_point);
+ }
+}
+
+impl<
+ ES: Deref + Clone,
+ NS: Deref + Clone,
+ CM: Deref + Clone,
+ C: Deref + Clone,
+ KS: Deref + Clone,
+ TP: Deref + Clone,
+ > Confirm for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+where
+ ES::Target: EntropySource,
+ NS::Target: NodeSigner,
+ CM::Target: AChannelManager,
+ C::Target: Filter,
+ KS::Target: KVStoreSync,
+ TP::Target: TimeProvider,
+{
+ fn transactions_confirmed(
+ &self, header: &bitcoin::block::Header, txdata: &chain::transaction::TransactionData,
+ height: u32,
+ ) {
+ self.inner.transactions_confirmed(header, txdata, height)
+ }
+
+ fn transaction_unconfirmed(&self, txid: &bitcoin::Txid) {
+ self.inner.transaction_unconfirmed(txid)
+ }
+
+ fn best_block_updated(&self, header: &bitcoin::block::Header, height: u32) {
+ self.inner.best_block_updated(header, height)
+ }
+
+ fn get_relevant_txids(&self) -> Vec<(bitcoin::Txid, u32, Option<bitcoin::BlockHash>)> {
+ self.inner.get_relevant_txids()
+ }
+}
diff --git a/lightning-liquidity/tests/common/mod.rs b/lightning-liquidity/tests/common/mod.rs
index 013378f..c7ec116 100644
--- a/lightning-liquidity/tests/common/mod.rs
+++ b/lightning-liquidity/tests/common/mod.rs
@@ -1,12 +1,12 @@
#![cfg(test)]
use lightning_liquidity::utils::time::TimeProvider;
-use lightning_liquidity::{LiquidityClientConfig, LiquidityManager, LiquidityServiceConfig};
+use lightning_liquidity::{LiquidityClientConfig, LiquidityManagerSync, LiquidityServiceConfig};
use lightning::chain::{BestBlock, Filter};
use lightning::ln::channelmanager::ChainParameters;
use lightning::ln::functional_test_utils::{Node, TestChannelManager};
-use lightning::util::test_utils::TestKeysInterface;
+use lightning::util::test_utils::{TestKeysInterface, TestStore};
use bitcoin::Network;
@@ -27,23 +27,27 @@ pub(crate) fn create_service_and_client_nodes<'a, 'b, 'c>(
network: Network::Testnet,
best_block: BestBlock::from_network(Network::Testnet),
};
- let service_lm = LiquidityManager::new_with_custom_time_provider(
+ let service_kv_store = Arc::new(TestStore::new(false));
+ let service_lm = LiquidityManagerSync::new_with_custom_time_provider(
nodes[0].keys_manager,
nodes[0].keys_manager,
nodes[0].node,
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params.clone()),
+ service_kv_store,
Some(service_config),
None,
Arc::clone(&time_provider),
);
- let client_lm = LiquidityManager::new_with_custom_time_provider(
+ let client_kv_store = Arc::new(TestStore::new(false));
+ let client_lm = LiquidityManagerSync::new_with_custom_time_provider(
nodes[1].keys_manager,
nodes[1].keys_manager,
nodes[1].node,
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params),
+ client_kv_store,
None,
Some(client_config),
time_provider,
@@ -58,11 +62,12 @@ pub(crate) fn create_service_and_client_nodes<'a, 'b, 'c>(
pub(crate) struct LiquidityNode<'a, 'b, 'c> {
pub inner: Node<'a, 'b, 'c>,
- pub liquidity_manager: LiquidityManager<
+ pub liquidity_manager: LiquidityManagerSync<
&'c TestKeysInterface,
&'c TestKeysInterface,
&'a TestChannelManager<'b, 'c>,
Arc<dyn Filter + Send + Sync>,
+ Arc<TestStore>,
Arc<dyn TimeProvider + Send + Sync>,
>,
}
@@ -70,11 +75,12 @@ pub(crate) struct LiquidityNode<'a, 'b, 'c> {
impl<'a, 'b, 'c> LiquidityNode<'a, 'b, 'c> {
pub fn new(
node: Node<'a, 'b, 'c>,
- liquidity_manager: LiquidityManager<
+ liquidity_manager: LiquidityManagerSync<
&'c TestKeysInterface,
&'c TestKeysInterface,
&'a TestChannelManager<'b, 'c>,
Arc<dyn Filter + Send + Sync>,
+ Arc<TestStore>,
Arc<dyn TimeProvider + Send + Sync>,
>,
) -> Self {
Why this scored 15/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.