Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions src/connection.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,10 @@ where
self.do_connect_peer(node_id, addr).await
}

pub(crate) fn disconnect_peer(&self, node_id: PublicKey) {
self.peer_manager.disconnect_by_node_id(node_id);
}

pub(crate) async fn do_connect_peer(
&self, node_id: PublicKey, addr: SocketAddress,
) -> Result<(), Error> {
Expand Down
100 changes: 84 additions & 16 deletions src/event.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,7 +34,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_liquidity::lsps2::utils::compute_opening_fee;
use lightning_types::payment::{PaymentHash, PaymentPreimage};

use crate::config::{may_announce_channel, Config};
use crate::config::{may_announce_channel, Config, PEER_RECONNECTION_INTERVAL};
use crate::connection::ConnectionManager;
use crate::data_store::DataStoreUpdateResult;
use crate::fee_estimator::ConfirmationTarget;
Expand DownExpand Up@@ -583,6 +583,68 @@ where
}
}

fn remove_peer_after_reconnect(&self, peer_info: PeerInfo, closed_channel_id: ChannelId) {
let channel_manager = Arc::clone(&self.channel_manager);
let connection_manager = Arc::clone(&self.connection_manager);
let peer_store = Arc::clone(&self.peer_store);
let logger = self.logger.clone();
self.runtime.spawn_cancellable_background_task(async move {
let has_other_channels = || {
channel_manager
.list_channels_with_counterparty(&peer_info.node_id)
.iter()
.any(|c| c.channel_id != closed_channel_id)
};

if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

// Ensure a connected peer cannot be mistaken for a completed recovery reconnect.
// With no other channels left, reconnecting once gives `channel_reestablish` a chance
// to retransmit the force-close error before we stop persisting the peer.
connection_manager.disconnect_peer(peer_info.node_id);

loop {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels() {
return;
}

match connection_manager
.connect_peer_if_necessary(peer_info.node_id, peer_info.address.clone())
.await
{
Ok(()) => {
if peer_store.get_peer(&peer_info.node_id).is_none() || has_other_channels()
{
return;
}
if let Err(e) = peer_store.remove_peer(&peer_info.node_id).await {
log_error!(
logger,
"Failed to remove peer {} from peer store: {}",
peer_info.node_id,
e
);
} else {
return;
}
},
Err(e) => {
log_debug!(
logger,
"Failed to reconnect peer {} before removing from peer store: {}",
peer_info.node_id,
e
);
},
}

tokio::time::sleep(PEER_RECONNECTION_INTERVAL).await;
}
});
}

async fn fail_claimable_payment(
&self, payment_id: PaymentId, payment_hash: &PaymentHash,
) -> Result<(), ReplayEvent> {
Expand DownExpand Up@@ -1627,25 +1689,24 @@ where
let counterparty_node_id = counterparty_node_id
.expect("counterparty_node_id is always set since LDK 0.0.117");

// Drop the peer once its last channel with us has reached a terminal state
// that reconnection cannot recover. Every closure reason is terminal except
// `HolderForceClosed`: when *we* force-close, we keep reconnecting so that
// `channel_reestablish` can drive recovery (see `Node::close_channel_internal`).
// Drop the peer once its last channel with us has reached a terminal state.
// For `HolderForceClosed`, retain it through one recovery reconnect so that
// `channel_reestablish` can retransmit the force-close error before cleanup.
// This also cleans up peers persisted for a channel that closed before funding
// (e.g. `CounterpartyCoopClosedUnfundedChannel`), which would otherwise be
// retried forever.
// We exclude `channel_id` from the count because LDK emits `ChannelClosed`
// before removing it from its internal list.
let dont_reconnect = !matches!(reason, ClosureReason::HolderForceClosed { .. });

if dont_reconnect {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

if !has_other_channels {
let has_other_channels = self
.channel_manager
.list_channels_with_counterparty(&counterparty_node_id)
.iter()
.any(|c| c.channel_id != channel_id);

let peer_to_reconnect = if !has_other_channels {
if matches!(reason, ClosureReason::HolderForceClosed { .. }) {
self.peer_store.get_peer(&counterparty_node_id)
} else {
if let Err(e) = self.peer_store.remove_peer(&counterparty_node_id).await {
log_error!(
self.logger,
Expand All@@ -1655,8 +1716,11 @@ where
);
return Err(ReplayEvent());
}
None
}
}
} else {
None
};

let event = Event::ChannelClosed {
channel_id,
Expand All@@ -1672,6 +1736,10 @@ where
return Err(ReplayEvent());
},
};

if let Some(peer_info) = peer_to_reconnect {
self.remove_peer_after_reconnect(peer_info, channel_id);
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
Expand Down
10 changes: 4 additions & 6 deletions src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2030,12 +2030,10 @@ impl Node {
}

// Peer store cleanup is handled centrally in the `ChannelClosed` event handler,
// which drops the peer once its last channel reaches a terminal state that
// reconnection cannot recover. We intentionally do nothing here so that a
// force-closed peer is retained, letting the background reconnection task keep
// firing and drive the `channel_reestablish` recovery flow. This is especially
// important against LND peers, which don't always handle force-closure error
// messages correctly.
// which retains a force-closed peer through one recovery reconnect before
// dropping it. This lets `channel_reestablish` drive the recovery flow, which is
// especially important against LND peers that don't always handle force-closure
// error messages correctly.
}

Ok(())
Expand Down
70 changes: 66 additions & 4 deletions src/peer_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,11 +58,14 @@ where
pub(crate) async fn remove_peer(&self, node_id: &PublicKey) -> Result<(), Error> {
let _guard = self.mutation_lock.lock().await;
let data = {
let mut locked_peers = self.peers.write().expect("lock");
locked_peers.remove(node_id);
PeerStoreSerWrapper(&locked_peers).encode()
let locked_peers = self.peers.read().expect("lock");
let mut updated_peers = locked_peers.clone();
updated_peers.remove(node_id);
PeerStoreSerWrapper(&updated_peers).encode()
};
self.persist_peers(data).await
self.persist_peers(data).await?;
self.peers.write().expect("lock").remove(node_id);
Ok(())
}

/// Returns the current in-memory peer set.
Expand DownExpand Up@@ -170,12 +173,52 @@ mod tests {
use std::str::FromStr;
use std::sync::Arc;

use bitcoin::io;
use lightning::util::persist::{PageToken, PaginatedKVStore, PaginatedListResponse};
use lightning::util::test_utils::TestLogger;

use super::*;
use crate::io::test_utils::InMemoryStore;
use crate::types::DynStoreWrapper;

struct FailingStore;

impl KVStore for FailingStore {
fn read(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str,
) -> impl std::future::Future<Output = Result<Vec<u8>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "read failed")) }
}

fn write(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _buf: Vec<u8>,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "write failed")) }
}

fn remove(
&self, _primary_namespace: &str, _secondary_namespace: &str, _key: &str, _lazy: bool,
) -> impl std::future::Future<Output = Result<(), io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "remove failed")) }
}

fn list(
&self, _primary_namespace: &str, _secondary_namespace: &str,
) -> impl std::future::Future<Output = Result<Vec<String>, io::Error>> + 'static + Send {
async { Err(io::Error::new(io::ErrorKind::Other, "list failed")) }
}
}

impl PaginatedKVStore for FailingStore {
fn list_paginated(
&self, _primary_namespace: &str, _secondary_namespace: &str,
_page_token: Option<PageToken>,
) -> impl std::future::Future<Output = Result<PaginatedListResponse, io::Error>> + 'static + Send
{
async { Err(io::Error::new(io::ErrorKind::Other, "list_paginated failed")) }
}
}

#[tokio::test]
async fn peer_info_persistence() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
Expand DownExpand Up@@ -215,4 +258,23 @@ mod tests {
assert_eq!(peers[0], expected_peer_info);
assert_eq!(deser_peer_store.get_peer(&node_id), Some(expected_peer_info));
}

#[tokio::test]
async fn remove_peer_does_not_mutate_memory_if_persist_fails() {
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(FailingStore));
let logger = Arc::new(TestLogger::new());
let node_id = PublicKey::from_str(
"0276607124ebe6a6c9338517b6f485825b27c2dcc0b9fc2aa6a4c0df91194e5993",
)
.unwrap();
let peer_info =
PeerInfo { node_id, address: SocketAddress::from_str("127.0.0.1:9738").unwrap() };
let mut peers = HashMap::new();
peers.insert(node_id, peer_info.clone());
let persisted_bytes = PeerStoreSerWrapper(&peers).encode();
let peer_store = PeerStore::read(&mut &persisted_bytes[..], (store, logger)).unwrap();

assert_eq!(Err(Error::PersistenceFailed), peer_store.remove_peer(&node_id).await);
assert_eq!(Some(peer_info), peer_store.get_peer(&node_id));
}
}
7 changes: 4 additions & 3 deletions tests/common/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -1702,10 +1702,11 @@ pub(crate) async fn do_channel_full_cycle<E: ElectrumApi>(
}

if force_close {
// Peer retained after local force-close to allow channel_reestablish recovery.
// The recovery reconnect completed while the force-close settled, so the peer no longer
// needs to remain persisted.
assert!(
node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should remain persisted in node_a peer store after locally-initiated force-close"
!node_a.list_peers().iter().any(|p| p.node_id == node_b.node_id() && p.is_persisted),
"node_b should be removed from node_a peer store after the recovery reconnect"
);
assert_any_node_has_onchain_tx_type(
&[("node_a", &node_a), ("node_b", &node_b)],
Expand Down
Loading