Skip to content

Commit d49b681

Browse files
committed
Move peer persistence onto async KV storage
Persist peer store updates through async KVStore operations. The synchronous node APIs keep bridging at their runtime boundary while async event handling awaits peer persistence directly. Co-Authored-By: HAL 9000
1 parent e67b7a1 commit d49b681

4 files changed

Lines changed: 82 additions & 60 deletions

File tree

src/event.rs

Lines changed: 32 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -1524,36 +1524,38 @@ where
15241524
},
15251525
};
15261526

1527-
let network_graph = self.network_graph.read_only();
1528-
let channels =
1529-
self.channel_manager.list_channels_with_counterparty(&counterparty_node_id);
1530-
if let Some(pending_channel) =
1531-
channels.into_iter().find(|c| c.channel_id == channel_id)
1532-
{
1533-
if !pending_channel.is_outbound
1534-
&& self.peer_store.get_peer(&counterparty_node_id).is_none()
1535-
{
1536-
if let Some(address) = network_graph
1537-
.nodes()
1538-
.get(&NodeId::from_pubkey(&counterparty_node_id))
1539-
.and_then(|node_info| node_info.announcement_info.as_ref())
1540-
.and_then(|ann_info| ann_info.addresses().first())
1541-
{
1542-
let peer = PeerInfo {
1543-
node_id: counterparty_node_id,
1544-
address: address.clone(),
1545-
};
1546-
1547-
self.peer_store.add_peer(peer).unwrap_or_else(|e| {
1548-
log_error!(
1549-
self.logger,
1550-
"Failed to add peer {} to peer store: {}",
1551-
counterparty_node_id,
1552-
e
1553-
);
1554-
});
1555-
}
1556-
}
1527+
let peer_to_store = {
1528+
let network_graph = self.network_graph.read_only();
1529+
let channels =
1530+
self.channel_manager.list_channels_with_counterparty(&counterparty_node_id);
1531+
channels
1532+
.into_iter()
1533+
.find(|c| c.channel_id == channel_id)
1534+
.filter(|pending_channel| {
1535+
!pending_channel.is_outbound
1536+
&& self.peer_store.get_peer(&counterparty_node_id).is_none()
1537+
})
1538+
.and_then(|_| {
1539+
network_graph
1540+
.nodes()
1541+
.get(&NodeId::from_pubkey(&counterparty_node_id))
1542+
.and_then(|node_info| node_info.announcement_info.as_ref())
1543+
.and_then(|ann_info| ann_info.addresses().first())
1544+
.map(|address| PeerInfo {
1545+
node_id: counterparty_node_id,
1546+
address: address.clone(),
1547+
})
1548+
})
1549+
};
1550+
if let Some(peer) = peer_to_store {
1551+
self.peer_store.add_peer(peer).await.unwrap_or_else(|e| {
1552+
log_error!(
1553+
self.logger,
1554+
"Failed to add peer {} to peer store: {}",
1555+
counterparty_node_id,
1556+
e
1557+
);
1558+
});
15571559
}
15581560
},
15591561
LdkEvent::ChannelReady {

src/lib.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1121,7 +1121,7 @@ impl Node {
11211121
log_info!(self.logger, "Connected to peer {}@{}. ", peer_info.node_id, peer_info.address);
11221122

11231123
if persist {
1124-
self.peer_store.add_peer(peer_info)?;
1124+
self.runtime.block_on(self.peer_store.add_peer(peer_info))?;
11251125
}
11261126

11271127
Ok(())
@@ -1138,7 +1138,7 @@ impl Node {
11381138

11391139
log_info!(self.logger, "Disconnecting peer {}..", counterparty_node_id);
11401140

1141-
match self.peer_store.remove_peer(&counterparty_node_id) {
1141+
match self.runtime.block_on(self.peer_store.remove_peer(&counterparty_node_id)) {
11421142
Ok(()) => {},
11431143
Err(e) => {
11441144
log_error!(self.logger, "Failed to remove peer {}: {}", counterparty_node_id, e)
@@ -1255,7 +1255,7 @@ impl Node {
12551255
zero_reserve_string,
12561256
peer_info.node_id
12571257
);
1258-
self.peer_store.add_peer(peer_info)?;
1258+
self.runtime.block_on(self.peer_store.add_peer(peer_info))?;
12591259
Ok(UserChannelId(user_channel_id))
12601260
},
12611261
Err(e) => {
@@ -1861,7 +1861,7 @@ impl Node {
18611861

18621862
// Check if this was the last open channel, if so, forget the peer.
18631863
if open_channels.len() == 1 {
1864-
self.peer_store.remove_peer(&counterparty_node_id)?;
1864+
self.runtime.block_on(self.peer_store.remove_peer(&counterparty_node_id))?;
18651865
}
18661866
}
18671867

src/payment/bolt11.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -242,7 +242,7 @@ impl Bolt11Payment {
242242
self.payment_store.insert(payment)?;
243243

244244
// Persist LSP peer to make sure we reconnect on restart.
245-
self.peer_store.add_peer(peer_info)?;
245+
self.runtime.block_on(self.peer_store.add_peer(peer_info))?;
246246

247247
Ok(invoice)
248248
}

src/peer_store.rs

Lines changed: 45 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ use std::sync::{Arc, RwLock};
1111

1212
use bitcoin::secp256k1::PublicKey;
1313
use lightning::impl_writeable_tlv_based;
14-
use lightning::util::persist::KVStoreSync;
14+
use lightning::util::persist::KVStore;
1515
use lightning::util::ser::{Readable, ReadableArgs, Writeable, Writer};
1616

1717
use crate::io::{
@@ -27,6 +27,7 @@ where
2727
L::Target: LdkLogger,
2828
{
2929
peers: RwLock<HashMap<PublicKey, PeerInfo>>,
30+
mutation_lock: tokio::sync::Mutex<()>,
3031
kv_store: Arc<DynStore>,
3132
logger: L,
3233
}
@@ -37,44 +38,60 @@ where
3738
{
3839
pub(crate) fn new(kv_store: Arc<DynStore>, logger: L) -> Self {
3940
let peers = RwLock::new(HashMap::new());
40-
Self { peers, kv_store, logger }
41+
let mutation_lock = tokio::sync::Mutex::new(());
42+
Self { peers, mutation_lock, kv_store, logger }
4143
}
4244

43-
pub(crate) fn add_peer(&self, peer_info: PeerInfo) -> Result<(), Error> {
44-
let mut locked_peers = self.peers.write().expect("lock");
45-
46-
if locked_peers.contains_key(&peer_info.node_id) {
47-
return Ok(());
48-
}
49-
50-
locked_peers.insert(peer_info.node_id, peer_info);
51-
self.persist_peers(&*locked_peers)
45+
pub(crate) async fn add_peer(&self, peer_info: PeerInfo) -> Result<(), Error> {
46+
let _guard = self.mutation_lock.lock().await;
47+
let data = {
48+
let mut locked_peers = self.peers.write().expect("lock");
49+
if locked_peers.contains_key(&peer_info.node_id) {
50+
return Ok(());
51+
}
52+
locked_peers.insert(peer_info.node_id, peer_info);
53+
PeerStoreSerWrapper(&locked_peers).encode()
54+
};
55+
self.persist_peers(data).await
5256
}
5357

54-
pub(crate) fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
55-
let mut locked_peers = self.peers.write().expect("lock");
56-
57-
locked_peers.remove(node_id);
58-
self.persist_peers(&*locked_peers)
58+
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
59+
let _guard = self.mutation_lock.lock().await;
60+
let data = {
61+
let mut locked_peers = self.peers.write().expect("lock");
62+
locked_peers.remove(node_id);
63+
PeerStoreSerWrapper(&locked_peers).encode()
64+
};
65+
self.persist_peers(data).await
5966
}
6067

68+
/// Returns the current in-memory peer set.
69+
///
70+
/// The async mutation lock serializes `add_peer` and `remove_peer`, but this synchronous
71+
/// reader cannot wait on it. Until peer-store reads are async, callers may observe peer
72+
/// changes that are still being persisted.
6173
pub(crate) fn list_peers(&self) -> Vec<PeerInfo> {
6274
self.peers.read().expect("lock").values().cloned().collect()
6375
}
6476

77+
/// Returns the current in-memory peer info for `node_id`.
78+
///
79+
/// The async mutation lock serializes `add_peer` and `remove_peer`, but this synchronous
80+
/// reader cannot wait on it. Until peer-store reads are async, callers may observe peer
81+
/// changes that are still being persisted.
6582
pub(crate) fn get_peer(&self, node_id: &PublicKey) -> Option<PeerInfo> {
6683
self.peers.read().expect("lock").get(node_id).cloned()
6784
}
6885

69-
fn persist_peers(&self, locked_peers: &HashMap<PublicKey, PeerInfo>) -> Result<(), Error> {
70-
let data = PeerStoreSerWrapper(&*locked_peers).encode();
71-
KVStoreSync::write(
86+
async fn persist_peers(&self, data: Vec<u8>) -> Result<(), Error> {
87+
KVStore::write(
7288
&*self.kv_store,
7389
PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
7490
PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
7591
PEER_INFO_PERSISTENCE_KEY,
7692
data,
7793
)
94+
.await
7895
.map_err(|e| {
7996
log_error!(
8097
self.logger,
@@ -101,7 +118,8 @@ where
101118
let (kv_store, logger) = args;
102119
let read_peers: PeerStoreDeserWrapper = Readable::read(reader)?;
103120
let peers: RwLock<HashMap<PublicKey, PeerInfo>> = RwLock::new(read_peers.0);
104-
Ok(Self { peers, kv_store, logger })
121+
let mutation_lock = tokio::sync::Mutex::new(());
122+
Ok(Self { peers, mutation_lock, kv_store, logger })
105123
}
106124
}
107125

@@ -158,8 +176,8 @@ mod tests {
158176
use crate::io::test_utils::InMemoryStore;
159177
use crate::types::DynStoreWrapper;
160178

161-
#[test]
162-
fn peer_info_persistence() {
179+
#[tokio::test]
180+
async fn peer_info_persistence() {
163181
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
164182
let logger = Arc::new(TestLogger::new());
165183
let peer_store = PeerStore::new(Arc::clone(&store), Arc::clone(&logger));
@@ -170,22 +188,24 @@ mod tests {
170188
.unwrap();
171189
let address = SocketAddress::from_str("127.0.0.1:9738").unwrap();
172190
let expected_peer_info = PeerInfo { node_id, address };
173-
assert!(KVStoreSync::read(
191+
assert!(KVStore::read(
174192
&*store,
175193
PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
176194
PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
177195
PEER_INFO_PERSISTENCE_KEY,
178196
)
197+
.await
179198
.is_err());
180-
peer_store.add_peer(expected_peer_info.clone()).unwrap();
199+
peer_store.add_peer(expected_peer_info.clone()).await.unwrap();
181200

182201
// Check we can read back what we persisted.
183-
let persisted_bytes = KVStoreSync::read(
202+
let persisted_bytes = KVStore::read(
184203
&*store,
185204
PEER_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
186205
PEER_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
187206
PEER_INFO_PERSISTENCE_KEY,
188207
)
208+
.await
189209
.unwrap();
190210
let deser_peer_store =
191211
PeerStore::read(&mut &persisted_bytes[..], (Arc::clone(&store), logger)).unwrap();

0 commit comments

Comments
 (0)