Skip `EventQueue` persistence if unnecessary
What changed, and why it matters
This change is a performance and reliability optimization, not a security fix. It teaches the liquidity service's event queue to skip writing to disk when nothing has changed, and to wake the background persister only when a save is actually needed. Previously, the queue may have been persisted repeatedly even when unchanged. The patch also ensures that if a persistence write fails, the system marks the queue as needing to be saved again so it won't lose data.
Treat as a routine optimization commit. No security response required. If reviewing for reliability, verify that the new `needs_persist` flag is correctly set on all queue-mutating paths and that persistence failures are retried by the background persister loop.
Security signals we found
No security-relevant signals in commit message or diff
Change is framed as an optimization ('Skip ... persistence if unnecessary')
Adds retry-on-failure behavior for persistence errors
No input validation, serialization bounds, or trust-boundary changes observed
Evidence from the diff
The commit refactors EventQueue to track a needs_persist flag inside a new QueueState struct. It wires in a Notifier so that enqueue/dequeue operations wake the background persister only when the persistent state changes, and persist() returns early if no changes occurred. On persistence failure it resets needs_persist = true so a retry will happen. The change is defensive and reduces I/O; it does not alter authorization, cryptography, network parsing, or access control.
Changed components
lightning-liquidity/src/events/event_queue.rslightning-liquidity/src/manager.rslightning-liquidity/src/lsps0/client.rslightning-liquidity/src/lsps5/client.rsInspect captured patch +117 / −32
diff --git a/lightning-liquidity/src/events/event_queue.rs b/lightning-liquidity/src/events/event_queue.rs
index 9904fc0..dc68ee6 100644
--- a/lightning-liquidity/src/events/event_queue.rs
+++ b/lightning-liquidity/src/events/event_queue.rs
@@ -20,6 +20,7 @@ use lightning::util::persist::KVStore;
use lightning::util::ser::{
BigSize, CollectionLength, FixedLengthReader, Readable, Writeable, Writer,
};
+use lightning::util::wakers::Notifier;
/// The maximum queue size we allow before starting to drop events.
pub const MAX_EVENT_QUEUE_SIZE: usize = 1000;
@@ -28,50 +29,73 @@ pub(crate) struct EventQueue<K: Deref + Clone>
where
K::Target: KVStore,
{
- queue: Arc<Mutex<VecDeque<LiquidityEvent>>>,
+ state: Arc<Mutex<QueueState>>,
waker: Arc<Mutex<Option<Waker>>>,
#[cfg(feature = "std")]
condvar: Arc<crate::sync::Condvar>,
kv_store: K,
+ persist_notifier: Arc<Notifier>,
}
impl<K: Deref + Clone> EventQueue<K>
where
K::Target: KVStore,
{
- pub fn new(queue: VecDeque<LiquidityEvent>, kv_store: K) -> Self {
- let queue = Arc::new(Mutex::new(queue));
+ pub fn new(
+ queue: VecDeque<LiquidityEvent>, kv_store: K, persist_notifier: Arc<Notifier>,
+ ) -> Self {
+ let state = Arc::new(Mutex::new(QueueState { queue, needs_persist: false }));
let waker = Arc::new(Mutex::new(None));
Self {
- queue,
+ state,
waker,
#[cfg(feature = "std")]
condvar: Arc::new(crate::sync::Condvar::new()),
kv_store,
+ persist_notifier,
}
}
pub fn next_event(&self) -> Option<LiquidityEvent> {
- self.queue.lock().unwrap().pop_front()
+ let event_opt = {
+ let mut state_lock = self.state.lock().unwrap();
+ if state_lock.queue.is_empty() {
+ // Skip notifying below if nothing changed.
+ return None;
+ }
+
+ state_lock.needs_persist = true;
+ state_lock.queue.pop_front()
+ };
+
+ self.persist_notifier.notify();
+
+ event_opt
}
pub async fn next_event_async(&self) -> LiquidityEvent {
- EventFuture { event_queue: Arc::clone(&self.queue), waker: Arc::clone(&self.waker) }.await
+ EventFuture {
+ queue_state: Arc::clone(&self.state),
+ waker: Arc::clone(&self.waker),
+ persist_notifier: Arc::clone(&self.persist_notifier),
+ }
+ .await
}
#[cfg(feature = "std")]
pub fn wait_next_event(&self) -> LiquidityEvent {
- let mut queue = self
+ let mut state_lock = self
.condvar
- .wait_while(self.queue.lock().unwrap(), |queue: &mut VecDeque<LiquidityEvent>| {
- queue.is_empty()
+ .wait_while(self.state.lock().unwrap(), |state_lock: &mut QueueState| {
+ state_lock.queue.is_empty()
})
.unwrap();
- let event = queue.pop_front().expect("non-empty queue");
- let should_notify = !queue.is_empty();
+ let event = state_lock.queue.pop_front().expect("non-empty queue");
+ let should_notify = !state_lock.queue.is_empty();
+ state_lock.needs_persist = true;
- drop(queue);
+ drop(state_lock);
if should_notify {
if let Some(waker) = self.waker.lock().unwrap().take() {
@@ -81,11 +105,28 @@ where
self.condvar.notify_one();
}
+ self.persist_notifier.notify();
+
event
}
pub fn get_and_clear_pending_events(&self) -> Vec<LiquidityEvent> {
- self.queue.lock().unwrap().split_off(0).into()
+ let mut state_lock = self.state.lock().unwrap();
+
+ let needs_persist = !state_lock.queue.is_empty();
+ let events = state_lock.queue.split_off(0).into();
+
+ if needs_persist {
+ state_lock.needs_persist = true;
+ }
+
+ drop(state_lock);
+
+ if needs_persist {
+ self.persist_notifier.notify();
+ }
+
+ events
}
// Returns an [`EventQueueNotifierGuard`] that will notify about new event when dropped.
@@ -94,20 +135,38 @@ where
}
pub async fn persist(&self) -> Result<(), lightning::io::Error> {
- let queue = self.queue.lock().unwrap();
- let encoded = EventQueueSerWrapper(&queue).encode();
+ let fut = {
+ let mut state_lock = self.state.lock().unwrap();
+
+ if !state_lock.needs_persist {
+ return Ok(());
+ }
- self.kv_store
- .write(
+ state_lock.needs_persist = false;
+ let encoded = EventQueueSerWrapper(&state_lock.queue).encode();
+
+ self.kv_store.write(
LIQUIDITY_MANAGER_PERSISTENCE_PRIMARY_NAMESPACE,
LIQUIDITY_MANAGER_EVENT_QUEUE_PERSISTENCE_SECONDARY_NAMESPACE,
LIQUIDITY_MANAGER_EVENT_QUEUE_PERSISTENCE_KEY,
encoded,
)
- .await
+ };
+
+ fut.await.map_err(|e| {
+ self.state.lock().unwrap().needs_persist = true;
+ e
+ })?;
+
+ Ok(())
}
}
+struct QueueState {
+ queue: VecDeque<LiquidityEvent>,
+ needs_persist: bool,
+}
+
// A guard type that will notify about new events when dropped.
#[must_use]
pub(crate) struct EventQueueNotifierGuard<'a, K: Deref + Clone>(&'a EventQueue<K>)
@@ -119,9 +178,10 @@ where
K::Target: KVStore,
{
pub fn enqueue<E: Into<LiquidityEvent>>(&self, event: E) {
- let mut queue = self.0.queue.lock().unwrap();
- if queue.len() < MAX_EVENT_QUEUE_SIZE {
- queue.push_back(event.into());
+ let mut state_lock = self.0.state.lock().unwrap();
+ if state_lock.queue.len() < MAX_EVENT_QUEUE_SIZE {
+ state_lock.queue.push_back(event.into());
+ state_lock.needs_persist = true;
} else {
return;
}
@@ -133,7 +193,10 @@ where
K::Target: KVStore,
{
fn drop(&mut self) {
- let should_notify = !self.0.queue.lock().unwrap().is_empty();
+ let (should_notify, should_persist_notify) = {
+ let state_lock = self.0.state.lock().unwrap();
+ (!state_lock.queue.is_empty(), state_lock.needs_persist)
+ };
if should_notify {
if let Some(waker) = self.0.waker.lock().unwrap().take() {
@@ -143,12 +206,17 @@ where
#[cfg(feature = "std")]
self.0.condvar.notify_one();
}
+
+ if should_persist_notify {
+ self.0.persist_notifier.notify();
+ }
}
}
struct EventFuture {
- event_queue: Arc<Mutex<VecDeque<LiquidityEvent>>>,
+ queue_state: Arc<Mutex<QueueState>>,
waker: Arc<Mutex<Option<Waker>>>,
+ persist_notifier: Arc<Notifier>,
}
impl Future for EventFuture {
@@ -157,12 +225,22 @@ impl Future for EventFuture {
fn poll(
self: core::pin::Pin<&mut Self>, cx: &mut core::task::Context<'_>,
) -> core::task::Poll<Self::Output> {
- if let Some(event) = self.event_queue.lock().unwrap().pop_front() {
- Poll::Ready(event)
- } else {
- *self.waker.lock().unwrap() = Some(cx.waker().clone());
- Poll::Pending
+ let (res, should_persist_notify) = {
+ let mut state_lock = self.queue_state.lock().unwrap();
+ if let Some(event) = state_lock.queue.pop_front() {
+ state_lock.needs_persist = true;
+ (Poll::Ready(event), true)
+ } else {
+ *self.waker.lock().unwrap() = Some(cx.waker().clone());
+ (Poll::Pending, false)
+ }
+ };
+
+ if should_persist_notify {
+ self.persist_notifier.notify();
}
+
+ res
}
}
@@ -266,7 +344,8 @@ mod tests {
use std::time::Duration;
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
+ let persist_notifier = Arc::new(Notifier::new());
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store, persist_notifier));
assert_eq!(event_queue.next_event(), None);
let secp_ctx = Secp256k1::new();
diff --git a/lightning-liquidity/src/lsps0/client.rs b/lightning-liquidity/src/lsps0/client.rs
index 9c019aa..7f26d74 100644
--- a/lightning-liquidity/src/lsps0/client.rs
+++ b/lightning-liquidity/src/lsps0/client.rs
@@ -136,7 +136,8 @@ mod tests {
let pending_messages = Arc::new(MessageQueue::new(notifier));
let entropy_source = Arc::new(TestEntropy {});
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
+ let persist_notifier = Arc::new(Notifier::new());
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store, persist_notifier));
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 a643238..1c6f8b8 100644
--- a/lightning-liquidity/src/lsps5/client.rs
+++ b/lightning-liquidity/src/lsps5/client.rs
@@ -477,7 +477,8 @@ mod tests {
let message_queue = Arc::new(MessageQueue::new(notifier));
let kv_store = Arc::new(KVStoreSyncWrapper(Arc::new(TestStore::new(false))));
- let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store));
+ let persist_notifier = Arc::new(Notifier::new());
+ let event_queue = Arc::new(EventQueue::new(VecDeque::new(), kv_store, persist_notifier));
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 3ba1638..8ff274a 100644
--- a/lightning-liquidity/src/manager.rs
+++ b/lightning-liquidity/src/manager.rs
@@ -389,7 +389,11 @@ where
let pending_messages =
Arc::new(MessageQueue::new(Arc::clone(&pending_msgs_or_needs_persist_notifier)));
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 pending_events = Arc::new(EventQueue::new(
+ persisted_queue,
+ kv_store.clone(),
+ Arc::clone(&pending_msgs_or_needs_persist_notifier),
+ ));
let ignored_peers = RwLock::new(new_hash_set());
let mut supported_protocols = Vec::new();
Why this scored 17/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.