Read persisted event queue state in `LiquidityManager::new`
What changed, and why it matters
This commit fixes a bug where the LiquidityManager would start with an empty event queue after a restart, even if important events had been saved to disk before shutdown. Now it reads the saved queue when it starts up, so no events are lost. This is a reliability/correctness fix rather than a typical security vulnerability.
Review whether any other persisted state is also not being loaded on `LiquidityManager` startup. Consider adding tests that simulate a restart and verify the event queue is restored. No urgent security patch is indicated by the diff alone.
Security signals we found
State loss on restart: previously persisted events were ignored on startup
Fixes consistency between on-disk state and in-memory queue
No input validation changes or cryptographic changes
No privilege escalation, injection, or memory-safety issues evident
Evidence from the diff
The change adds read_event_queue in persist.rs and calls it from LiquidityManager::new in manager.rs. Previously EventQueue::new always created an empty VecDeque; now it accepts a pre-populated queue. The manager awaits the persisted queue and passes it to EventQueue::new, ensuring events survive process restarts. Tests and internal call sites were updated to pass an empty VecDeque::new().
Changed components
lightning-liquidity/src/events/event_queue.rslightning-liquidity/src/events/mod.rslightning-liquidity/src/lsps0/client.rslightning-liquidity/src/lsps5/client.rslightning-liquidity/src/manager.rslightning-liquidity/src/persist.rsInspect captured patch +49 / −8
diff --git a/lightning-liquidity/src/events/event_queue.rs b/lightning-liquidity/src/events/event_queue.rs
index 1c8282f..9904fc0 100644
--- a/lightning-liquidity/src/events/event_queue.rs
+++ b/lightning-liquidity/src/events/event_queue.rs
@@ -39,8 +39,8 @@ impl<K: Deref + Clone> EventQueue<K>
where
K::Target: KVStore,
{
- pub fn new(kv_store: K) -> Self {
- let queue = Arc::new(Mutex::new(VecDeque::new()));
+ pub fn new(queue: VecDeque<LiquidityEvent>, kv_store: K) -> Self {
+ let queue = Arc::new(Mutex::new(queue));
let waker = Arc::new(Mutex::new(None));
Self {
queue,
@@ -266,7 +266,7 @@ mod tests {
use std::time::Duration;
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(kv_store));
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
assert_eq!(event_queue.next_event(), None);
let secp_ctx = Secp256k1::new();
diff --git a/lightning-liquidity/src/events/mod.rs b/lightning-liquidity/src/events/mod.rs
index 82e480a..c39b8b9 100644
--- a/lightning-liquidity/src/events/mod.rs
+++ b/lightning-liquidity/src/events/mod.rs
@@ -17,8 +17,8 @@
mod event_queue;
-pub(crate) use event_queue::EventQueue;
pub use event_queue::MAX_EVENT_QUEUE_SIZE;
+pub(crate) use event_queue::{EventQueue, EventQueueDeserWrapper};
use crate::lsps0;
use crate::lsps1;
diff --git a/lightning-liquidity/src/lsps0/client.rs b/lightning-liquidity/src/lsps0/client.rs
index 2efadd4..56dfb24 100644
--- a/lightning-liquidity/src/lsps0/client.rs
+++ b/lightning-liquidity/src/lsps0/client.rs
@@ -117,6 +117,7 @@ where
#[cfg(test)]
mod tests {
+ use alloc::collections::VecDeque;
use alloc::string::ToString;
use alloc::sync::Arc;
@@ -133,7 +134,7 @@ mod tests {
let pending_messages = Arc::new(MessageQueue::new());
let entropy_source = Arc::new(TestEntropy {});
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(kv_store));
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
let lsps0_handler = Arc::new(LSPS0ClientHandler::new(
entropy_source,
diff --git a/lightning-liquidity/src/lsps5/client.rs b/lightning-liquidity/src/lsps5/client.rs
index e464e88..1045c1b 100644
--- a/lightning-liquidity/src/lsps5/client.rs
+++ b/lightning-liquidity/src/lsps5/client.rs
@@ -475,7 +475,7 @@ mod tests {
let message_queue = Arc::new(MessageQueue::new());
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(kv_store));
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
let client = LSPS5ClientHandler::new(
test_entropy_source,
Arc::clone(&message_queue),
diff --git a/lightning-liquidity/src/manager.rs b/lightning-liquidity/src/manager.rs
index 6a899c6..27e82f5 100644
--- a/lightning-liquidity/src/manager.rs
+++ b/lightning-liquidity/src/manager.rs
@@ -24,7 +24,9 @@ use crate::lsps5::client::{LSPS5ClientConfig, LSPS5ClientHandler};
use crate::lsps5::msgs::LSPS5Message;
use crate::lsps5::service::{LSPS5ServiceConfig, LSPS5ServiceHandler};
use crate::message_queue::MessageQueue;
-use crate::persist::{read_lsps2_service_peer_states, read_lsps5_service_peer_states};
+use crate::persist::{
+ read_event_queue, read_lsps2_service_peer_states, read_lsps5_service_peer_states,
+};
use crate::lsps1::client::{LSPS1ClientConfig, LSPS1ClientHandler};
use crate::lsps1::msgs::LSPS1Message;
@@ -383,7 +385,8 @@ where
client_config: Option<LiquidityClientConfig>, time_provider: TP,
) -> Result<Self, lightning::io::Error> {
let pending_messages = Arc::new(MessageQueue::new());
- let pending_events = Arc::new(EventQueue::new(kv_store.clone()));
+ let persisted_queue = read_event_queue(kv_store.clone()).await?.unwrap_or_default();
+ let pending_events = Arc::new(EventQueue::new(persisted_queue, kv_store.clone()));
let ignored_peers = RwLock::new(new_hash_set());
let mut supported_protocols = Vec::new();
diff --git a/lightning-liquidity/src/persist.rs b/lightning-liquidity/src/persist.rs
index 5b8a63e..ec0d5a6 100644
--- a/lightning-liquidity/src/persist.rs
+++ b/lightning-liquidity/src/persist.rs
@@ -9,6 +9,7 @@
//! Types and utils for persistence.
+use crate::events::{EventQueueDeserWrapper, LiquidityEvent};
use crate::lsps2::service::PeerState as LSPS2ServicePeerState;
use crate::lsps5::service::PeerState as LSPS5ServicePeerState;
use crate::prelude::{new_hash_map, HashMap};
@@ -20,6 +21,8 @@ use lightning::util::ser::Readable;
use bitcoin::secp256k1::PublicKey;
+use alloc::collections::VecDeque;
+
use core::ops::Deref;
use core::str::FromStr;
@@ -48,6 +51,40 @@ pub const LSPS2_SERVICE_PERSISTENCE_SECONDARY_NAMESPACE: &str = "lsps2_service";
/// [`LSPS5ServiceHandler`]: crate::lsps5::service::LSPS5ServiceHandler
pub const LSPS5_SERVICE_PERSISTENCE_SECONDARY_NAMESPACE: &str = "lsps5_service";
+pub(crate) async fn read_event_queue<K: Deref>(
+ kv_store: K,
+) -> Result<Option<VecDeque<LiquidityEvent>>, lightning::io::Error>
+where
+ K::Target: KVStore,
+{
+ let read_fut = kv_store.read(
+ LIQUIDITY_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE,
+ LIQUIDITY_MANAGER_EVENT_QUEUE_PERSISTENCE_SECONDARY_NAMESPACE,
+ LIQUIDITY_MANAGER_EVENT_QUEUE_PERSISTENCE_KEY,
+ );
+
+ let mut reader = match read_fut.await {
+ Ok(r) => Cursor::new(r),
+ Err(e) => {
+ if e.kind() == lightning::io::ErrorKind::NotFound {
+ // Key wasn't found, no error but first time running.
+ return Ok(None);
+ } else {
+ return Err(e);
+ }
+ },
+ };
+
+ let queue: EventQueueDeserWrapper = Readable::read(&mut reader).map_err(|_| {
+ lightning::io::Error::new(
+ lightning::io::ErrorKind::InvalidData,
+ "Failed to deserialize liquidity event queue",
+ )
+ })?;
+
+ Ok(Some(queue.0))
+}
+
pub(crate) async fn read_lsps2_service_peer_states<K: Deref>(
kv_store: K,
) -> Result<HashMap<PublicKey, Mutex<LSPS2ServicePeerState>>, lightning::io::Error>
Why this scored 31/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.