Consistently use `wire::Message` for encoding network messages
What changed, and why it matters
This commit is a code cleanup in the Lightning Dev Kit's peer message handling. It changes internal functions so they consistently use a single `wire::Message` wrapper type when preparing messages to send, instead of accepting many different message types directly. It also switches from borrowing messages to moving them. There is no obvious security bug being fixed; it is primarily a maintainability and type-safety improvement.
No immediate security action required. Treat as a normal code-quality refactor. Reviewers may want to verify that the move semantics do not introduce unintended ownership issues with messages that were previously shared, and that the `#[cfg(simple_close)]` changes compile correctly with the feature disabled.
Security signals we found
Refactoring only: no new parsing, no new network input handling, no cryptographic changes
Move from generic `M: Type + Writeable` to concrete `wire::Message<T>` improves type discipline
Added `#[cfg(simple_close)]` guards to prevent referencing uncompiled message variants
No bounds checks, serialization length checks, or encryption logic were altered in behavior
Evidence from the diff
The patch refactors PeerManager::enqueue_message, PeerChannelEncryptor::encrypt_message, and related call sites to take wire::Message<T> rather than a generic M: Type + Writeable. It also converts references to owned values (move semantics) throughout the outbound message path. The change touches peer_handler.rs, peer_channel_encryptor.rs, channelmanager.rs, functional_test_utils.rs, and msgs.rs. A side change adds #[cfg(simple_close)] guards to SendClosingComplete and SendClosingSig variants and their match arms, likely to fix compilation when the simple_close feature is disabled. No vulnerability, memory-safety issue, or cryptographic flaw is visible in the diff.
Changed components
lightning/src/ln/peer_handler.rslightning/src/ln/peer_channel_encryptor.rslightning/src/ln/channelmanager.rslightning/src/ln/functional_test_utils.rslightning/src/ln/msgs.rsInspect captured patch +226 / −155
diff --git a/lightning/src/ln/channelmanager.rs b/lightning/src/ln/channelmanager.rs
index 399c51b..d52d853 100644
--- a/lightning/src/ln/channelmanager.rs
+++ b/lightning/src/ln/channelmanager.rs
@@ -13872,7 +13872,9 @@ where
&MessageSendEvent::UpdateHTLCs { .. } => false,
&MessageSendEvent::SendRevokeAndACK { .. } => false,
&MessageSendEvent::SendClosingSigned { .. } => false,
+ #[cfg(simple_close)]
&MessageSendEvent::SendClosingComplete { .. } => false,
+ #[cfg(simple_close)]
&MessageSendEvent::SendClosingSig { .. } => false,
&MessageSendEvent::SendShutdown { .. } => false,
&MessageSendEvent::SendChannelReestablish { .. } => false,
diff --git a/lightning/src/ln/functional_test_utils.rs b/lightning/src/ln/functional_test_utils.rs
index e31630a..e4e9c58 100644
--- a/lightning/src/ln/functional_test_utils.rs
+++ b/lightning/src/ln/functional_test_utils.rs
@@ -1124,7 +1124,9 @@ pub fn remove_first_msg_event_to_node(
MessageSendEvent::UpdateHTLCs { node_id, .. } => node_id == msg_node_id,
MessageSendEvent::SendRevokeAndACK { node_id, .. } => node_id == msg_node_id,
MessageSendEvent::SendClosingSigned { node_id, .. } => node_id == msg_node_id,
+ #[cfg(simple_close)]
MessageSendEvent::SendClosingComplete { node_id, .. } => node_id == msg_node_id,
+ #[cfg(simple_close)]
MessageSendEvent::SendClosingSig { node_id, .. } => node_id == msg_node_id,
MessageSendEvent::SendShutdown { node_id, .. } => node_id == msg_node_id,
MessageSendEvent::SendChannelReestablish { node_id, .. } => node_id == msg_node_id,
diff --git a/lightning/src/ln/msgs.rs b/lightning/src/ln/msgs.rs
index 8e230fa..0484ebe 100644
--- a/lightning/src/ln/msgs.rs
+++ b/lightning/src/ln/msgs.rs
@@ -1857,6 +1857,7 @@ pub enum MessageSendEvent {
msg: ClosingSigned,
},
/// Used to indicate that a `closing_complete` message should be sent to the peer with the given `node_id`.
+ #[cfg(simple_close)]
SendClosingComplete {
/// The node_id of the node which should receive this message
node_id: PublicKey,
@@ -1864,6 +1865,7 @@ pub enum MessageSendEvent {
msg: ClosingComplete,
},
/// Used to indicate that a `closing_sig` message should be sent to the peer with the given `node_id`.
+ #[cfg(simple_close)]
SendClosingSig {
/// The node_id of the node which should receive this message
node_id: PublicKey,
diff --git a/lightning/src/ln/peer_channel_encryptor.rs b/lightning/src/ln/peer_channel_encryptor.rs
index 09b970a..1d34d9a 100644
--- a/lightning/src/ln/peer_channel_encryptor.rs
+++ b/lightning/src/ln/peer_channel_encryptor.rs
@@ -565,12 +565,12 @@ impl PeerChannelEncryptor {
/// Encrypts the given message, returning the encrypted version.
/// panics if the length of `message`, once encoded, is greater than 65535 or if the Noise
/// handshake has not finished.
- pub fn encrypt_message<M: wire::Type>(&mut self, message: &M) -> Vec<u8> {
+ pub fn encrypt_message<T: wire::Type>(&mut self, message: wire::Message<T>) -> Vec<u8> {
// Allocate a buffer with 2KB, fitting most common messages. Reserve the first 16+2 bytes
// for the 2-byte message type prefix and its MAC.
let mut res = VecWriter(Vec::with_capacity(MSG_BUF_ALLOC_SIZE));
res.0.resize(16 + 2, 0);
- wire::write(message, &mut res).expect("In-memory messages must never fail to serialize");
+ wire::write(&message, &mut res).expect("In-memory messages must never fail to serialize");
self.encrypt_message_with_header_0s(&mut res.0);
res.0
diff --git a/lightning/src/ln/peer_handler.rs b/lightning/src/ln/peer_handler.rs
index c3b490e..8a6c6a7 100644
--- a/lightning/src/ln/peer_handler.rs
+++ b/lightning/src/ln/peer_handler.rs
@@ -29,7 +29,7 @@ use crate::ln::peer_channel_encryptor::{
};
use crate::ln::types::ChannelId;
use crate::ln::wire;
-use crate::ln::wire::{Encode, Type};
+use crate::ln::wire::{Encode, Message, Type};
use crate::onion_message::async_payments::{
AsyncPaymentsMessageHandler, HeldHtlcAvailable, OfferPaths, OfferPathsRequest, ReleaseHeldHtlc,
ServeStaticInvoice, StaticInvoicePersisted,
@@ -53,12 +53,14 @@ use crate::util::ser::{VecWriter, Writeable, Writer};
#[allow(unused_imports)]
use crate::prelude::*;
+use super::wire::CustomMessageReader;
use crate::io;
use crate::sync::{FairRwLock, Mutex, MutexGuard};
use core::convert::Infallible;
use core::ops::Deref;
use core::sync::atomic::{AtomicBool, AtomicI32, AtomicU32, Ordering};
use core::{cmp, fmt, hash, mem};
+
#[cfg(not(c_bindings))]
use {
crate::chain::chainmonitor::ChainMonitor,
@@ -1121,7 +1123,7 @@ pub struct PeerManager<
}
enum LogicalMessage<T: core::fmt::Debug + wire::Type + wire::TestEq> {
- FromWire(wire::Message<T>),
+ FromWire(Message<T>),
CommitmentSignedBatch(ChannelId, Vec<msgs::CommitmentSigned>),
}
@@ -1572,7 +1574,8 @@ where
if let Some(next_onion_message) =
handler.next_onion_message_for_peer(peer_node_id)
{
- self.enqueue_message(peer, &next_onion_message);
+ let msg = Message::OnionMessage(next_onion_message);
+ self.enqueue_message(peer, msg);
}
}
}
@@ -1590,16 +1593,20 @@ where
if let Some((announce, update_a_option, update_b_option)) =
self.message_handler.route_handler.get_next_channel_announcement(c)
{
- self.enqueue_message(peer, &announce);
+ peer.sync_status = InitSyncTracker::ChannelsSyncing(
+ announce.contents.short_channel_id + 1,
+ );
+ let msg = Message::ChannelAnnouncement(announce);
+ self.enqueue_message(peer, msg);
+
if let Some(update_a) = update_a_option {
- self.enqueue_message(peer, &update_a);
+ let msg = Message::ChannelUpdate(update_a);
+ self.enqueue_message(peer, msg);
}
if let Some(update_b) = update_b_option {
- self.enqueue_message(peer, &update_b);
+ let msg = Message::ChannelUpdate(update_b);
+ self.enqueue_message(peer, msg);
}
- peer.sync_status = InitSyncTracker::ChannelsSyncing(
- announce.contents.short_channel_id + 1,
- );
} else {
peer.sync_status =
InitSyncTracker::ChannelsSyncing(0xffff_ffff_ffff_ffff);
@@ -1608,8 +1615,9 @@ where
InitSyncTracker::ChannelsSyncing(c) if c == 0xffff_ffff_ffff_ffff => {
let handler = &self.message_handler.route_handler;
if let Some(msg) = handler.get_next_node_announcement(None) {
- self.enqueue_message(peer, &msg);
peer.sync_status = InitSyncTracker::NodesSyncing(msg.contents.node_id);
+ let msg = Message::NodeAnnouncement(msg);
+ self.enqueue_message(peer, msg);
} else {
peer.sync_status = InitSyncTracker::NoSyncRequested;
}
@@ -1618,8 +1626,9 @@ where
InitSyncTracker::NodesSyncing(sync_node_id) => {
let handler = &self.message_handler.route_handler;
if let Some(msg) = handler.get_next_node_announcement(Some(&sync_node_id)) {
- self.enqueue_message(peer, &msg);
peer.sync_status = InitSyncTracker::NodesSyncing(msg.contents.node_id);
+ let msg = Message::NodeAnnouncement(msg);
+ self.enqueue_message(peer, msg);
} else {
peer.sync_status = InitSyncTracker::NoSyncRequested;
}
@@ -1727,7 +1736,10 @@ where
}
/// Append a message to a peer's pending outbound/write buffer
- fn enqueue_message<M: wire::Type>(&self, peer: &mut Peer, message: &M) {
+ fn enqueue_message(
+ &self, peer: &mut Peer,
+ message: Message<<CMH::Target as CustomMessageReader>::CustomMessage>,
+ ) {
let their_node_id = peer.their_node_id.map(|p| p.0);
if their_node_id.is_some() {
let logger = WithContext::from(&self.logger, their_node_id, None, None);
@@ -1792,12 +1804,14 @@ where
},
msgs::ErrorAction::SendErrorMessage { msg } => {
log_debug!(logger, "Error handling message{}; sending error message with: {}", OptionalFromDebugger(&peer_node_id), e.err);
- self.enqueue_message($peer, &msg);
+ let msg = Message::Error(msg);
+ self.enqueue_message($peer, msg);
continue;
},
msgs::ErrorAction::SendWarningMessage { msg, log_level } => {
log_given_level!(logger, log_level, "Error handling message{}; sending warning message with: {}", OptionalFromDebugger(&peer_node_id), e.err);
- self.enqueue_message($peer, &msg);
+ let msg = Message::Warning(msg);
+ self.enqueue_message($peer, msg);
continue;
},
}
@@ -1892,7 +1906,8 @@ where
peer.their_socket_address.clone(),
),
};
- self.enqueue_message(peer, &resp);
+ let msg = Message::Init(resp);
+ self.enqueue_message(peer, msg);
},
NextNoiseStep::ActThree => {
let res = peer
@@ -1912,7 +1927,8 @@ where
peer.their_socket_address.clone(),
),
};
- self.enqueue_message(peer, &resp);
+ let msg = Message::Init(resp);
+ self.enqueue_message(peer, msg);
},
NextNoiseStep::NoiseComplete => {
if peer.pending_read_is_header {
@@ -1972,8 +1988,11 @@ where
let channel_id = ChannelId::new_zero();
let data = "Unsupported message compression: zlib"
.to_owned();
- let msg = msgs::WarningMessage { channel_id, data };
- self.enqueue_message(peer, &msg);
+ let msg = Message::Warning(msgs::WarningMessage {
+ channel_id,
+ data,
+ });
+ self.enqueue_message(peer, msg);
continue;
},
(_, Some(ty)) if is_gossip_msg(ty) => {
@@ -1983,8 +2002,11 @@ where
"Unreadable/bogus gossip message of type {}",
ty
);
- let msg = msgs::WarningMessage { channel_id, data };
- self.enqueue_message(peer, &msg);
+ let msg = Message::Warning(msgs::WarningMessage {
+ channel_id,
+ data,
+ });
+ self.enqueue_message(peer, msg);
continue;
},
(msgs::DecodeError::UnknownRequiredFeature, _) => {
@@ -2060,9 +2082,7 @@ where
/// Returns the message back if it needs to be broadcasted to all other peers.
fn handle_message(
&self, peer_mutex: &Mutex<Peer>, peer_lock: MutexGuard<Peer>,
- message: wire::Message<
- <<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage,
- >,
+ message: Message<<<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage>,
) -> Result<Option<BroadcastGossipMessage>, MessageHandlingError> {
let their_node_id = peer_lock
.their_node_id
@@ -2103,9 +2123,7 @@ where
// allow it to be subsequently processed by `do_handle_message_without_peer_lock`.
fn do_handle_message_holding_peer_lock<'a>(
&self, mut peer_lock: MutexGuard<Peer>,
- message: wire::Message<
- <<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage,
- >,
+ message: Message<<<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage>,
their_node_id: PublicKey, logger: &WithContext<'a, L>,
) -> Result<
Option<
@@ -2116,7 +2134,7 @@ where
peer_lock.received_message_since_timer_tick = true;
// Need an Init as first message
- if let wire::Message::Init(msg) = message {
+ if let Message::Init(msg) = message {
// Check if we have any compatible chains if the `networks` field is specified.
if let Some(networks) = &msg.networks {
let chan_handler = &self.message_handler.chan_handler;
@@ -2225,7 +2243,7 @@ where
// During splicing, commitment_signed messages need to be collected into a single batch
// before they are handled.
- if let wire::Message::StartBatch(msg) = message {
+ if let Message::StartBatch(msg) = message {
if peer_lock.message_batch.is_some() {
let error = format!(
"Peer {} sent start_batch for channel {} before previous batch completed",
@@ -2296,7 +2314,7 @@ where
return Ok(None);
}
- if let wire::Message::CommitmentSigned(msg) = message {
+ if let Message::CommitmentSigned(msg) = message {
if let Some(message_batch) = &mut peer_lock.message_batch {
let MessageBatchImpl::CommitmentSigned(ref mut messages) =
&mut message_batch.messages;
@@ -2325,7 +2343,7 @@ where
return Ok(None);
}
} else {
- return Ok(Some(LogicalMessage::FromWire(wire::Message::CommitmentSigned(msg))));
+ return Ok(Some(LogicalMessage::FromWire(Message::CommitmentSigned(msg))));
}
} else if let Some(message_batch) = &peer_lock.message_batch {
match message_batch.messages {
@@ -2341,7 +2359,7 @@ where
return Err(PeerHandleError {}.into());
}
- if let wire::Message::GossipTimestampFilter(_msg) = message {
+ if let Message::GossipTimestampFilter(_msg) = message {
// When supporting gossip messages, start initial gossip sync only after we receive
// a GossipTimestampFilter
if peer_lock.their_features.as_ref().unwrap().supports_gossip_queries()
@@ -2373,7 +2391,7 @@ where
return Ok(None);
}
- if let wire::Message::ChannelAnnouncement(ref _msg) = message {
+ if let Message::ChannelAnnouncement(ref _msg) = message {
peer_lock.received_channel_announce_since_backlogged = true;
}
@@ -2385,9 +2403,7 @@ where
// Returns the message back if it needs to be broadcasted to all other peers.
fn do_handle_message_without_peer_lock<'a>(
&self, peer_mutex: &Mutex<Peer>,
- message: wire::Message<
- <<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage,
- >,
+ message: Message<<<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage>,
their_node_id: PublicKey, logger: &WithContext<'a, L>,
) -> Result<Option<BroadcastGossipMessage>, MessageHandlingError> {
if is_gossip_msg(message.type_id()) {
@@ -2400,13 +2416,13 @@ where
match message {
// Setup and Control messages:
- wire::Message::Init(_) => {
+ Message::Init(_) => {
// Handled above
},
- wire::Message::GossipTimestampFilter(_) => {
+ Message::GossipTimestampFilter(_) => {
// Handled above
},
- wire::Message::Error(msg) => {
+ Message::Error(msg) => {
log_debug!(
logger,
"Got Err message from {}: {}",
@@ -2418,149 +2434,150 @@ where
return Err(PeerHandleError {}.into());
}
},
- wire::Message::Warning(msg) => {
+ Message::Warning(msg) => {
log_debug!(logger, "Got warning message: {}", PrintableString(&msg.data));
},
- wire::Message::Ping(msg) => {
+ Message::Ping(msg) => {
if msg.ponglen < 65532 {
let resp = msgs::Pong { byteslen: msg.ponglen };
- self.enqueue_message(&mut *peer_mutex.lock().unwrap(), &resp);
+ let msg = Message::Pong(resp);
+ self.enqueue_message(&mut *peer_mutex.lock().unwrap(), msg);
}
},
- wire::Message::Pong(_msg) => {
+ Message::Pong(_msg) => {
let mut peer_lock = peer_mutex.lock().unwrap();
peer_lock.awaiting_pong_timer_tick_intervals = 0;
peer_lock.msgs_sent_since_pong = 0;
},
// Channel messages:
- wire::Message::StartBatch(_msg) => {
+ Message::StartBatch(_msg) => {
debug_assert!(false);
},
- wire::Message::OpenChannel(msg) => {
+ Message::OpenChannel(msg) => {
self.message_handler.chan_handler.handle_open_channel(their_node_id, &msg);
},
- wire::Message::OpenChannelV2(_msg) => {
+ Message::OpenChannelV2(_msg) => {
self.message_handler.chan_handler.handle_open_channel_v2(their_node_id, &_msg);
},
- wire::Message::AcceptChannel(msg) => {
+ Message::AcceptChannel(msg) => {
self.message_handler.chan_handler.handle_accept_channel(their_node_id, &msg);
},
- wire::Message::AcceptChannelV2(msg) => {
+ Message::AcceptChannelV2(msg) => {
self.message_handler.chan_handler.handle_accept_channel_v2(their_node_id, &msg);
},
- wire::Message::FundingCreated(msg) => {
+ Message::FundingCreated(msg) => {
self.message_handler.chan_handler.handle_funding_created(their_node_id, &msg);
},
- wire::Message::FundingSigned(msg) => {
+ Message::FundingSigned(msg) => {
self.message_handler.chan_handler.handle_funding_signed(their_node_id, &msg);
},
- wire::Message::ChannelReady(msg) => {
+ Message::ChannelReady(msg) => {
self.message_handler.chan_handler.handle_channel_ready(their_node_id, &msg);
},
- wire::Message::PeerStorage(msg) => {
+ Message::PeerStorage(msg) => {
self.message_handler.chan_handler.handle_peer_storage(their_node_id, msg);
},
- wire::Message::PeerStorageRetrieval(msg) => {
+ Message::PeerStorageRetrieval(msg) => {
self.message_handler.chan_handler.handle_peer_storage_retrieval(their_node_id, msg);
},
// Quiescence messages:
- wire::Message::Stfu(msg) => {
+ Message::Stfu(msg) => {
self.message_handler.chan_handler.handle_stfu(their_node_id, &msg);
},
// Splicing messages:
- wire::Message::SpliceInit(msg) => {
+ Message::SpliceInit(msg) => {
self.message_handler.chan_handler.handle_splice_init(their_node_id, &msg);
},
- wire::Message::SpliceAck(msg) => {
+ Message::SpliceAck(msg) => {
self.message_handler.chan_handler.handle_splice_ack(their_node_id, &msg);
},
- wire::Message::SpliceLocked(msg) => {
+ Message::SpliceLocked(msg) => {
self.message_handler.chan_handler.handle_splice_locked(their_node_id, &msg);
},
// Interactive transaction construction messages:
- wire::Message::TxAddInput(msg) => {
+ Message::TxAddInput(msg) => {
self.message_handler.chan_handler.handle_tx_add_input(their_node_id, &msg);
},
- wire::Message::TxAddOutput(msg) => {
+ Message::TxAddOutput(msg) => {
self.message_handler.chan_handler.handle_tx_add_output(their_node_id, &msg);
},
- wire::Message::TxRemoveInput(msg) => {
+ Message::TxRemoveInput(msg) => {
self.message_handler.chan_handler.handle_tx_remove_input(their_node_id, &msg);
},
- wire::Message::TxRemoveOutput(msg) => {
+ Message::TxRemoveOutput(msg) => {
self.message_handler.chan_handler.handle_tx_remove_output(their_node_id, &msg);
},
- wire::Message::TxComplete(msg) => {
+ Message::TxComplete(msg) => {
self.message_handler.chan_handler.handle_tx_complete(their_node_id, &msg);
},
- wire::Message::TxSignatures(msg) => {
+ Message::TxSignatures(msg) => {
self.message_handler.chan_handler.handle_tx_signatures(their_node_id, &msg);
},
- wire::Message::TxInitRbf(msg) => {
+ Message::TxInitRbf(msg) => {
self.message_handler.chan_handler.handle_tx_init_rbf(their_node_id, &msg);
},
- wire::Message::TxAckRbf(msg) => {
+ Message::TxAckRbf(msg) => {
self.message_handler.chan_handler.handle_tx_ack_rbf(their_node_id, &msg);
},
- wire::Message::TxAbort(msg) => {
+ Message::TxAbort(msg) => {
self.message_handler.chan_handler.handle_tx_abort(their_node_id, &msg);
},
- wire::Message::Shutdown(msg) => {
+ Message::Shutdown(msg) => {
self.message_handler.chan_handler.handle_shutdown(their_node_id, &msg);
},
- wire::Message::ClosingSigned(msg) => {
+ Message::ClosingSigned(msg) => {
self.message_handler.chan_handler.handle_closing_signed(their_node_id, &msg);
},
#[cfg(simple_close)]
- wire::Message::ClosingComplete(msg) => {
+ Message::ClosingComplete(msg) => {
self.message_handler.chan_handler.handle_closing_complete(their_node_id, msg);
},
#[cfg(simple_close)]
- wire::Message::ClosingSig(msg) => {
+ Message::ClosingSig(msg) => {
self.message_handler.chan_handler.handle_closing_sig(their_node_id, msg);
},
// Commitment messages:
- wire::Message::UpdateAddHTLC(msg) => {
+ Message::UpdateAddHTLC(msg) => {
self.message_handler.chan_handler.handle_update_add_htlc(their_node_id, &msg);
},
- wire::Message::UpdateFulfillHTLC(msg) => {
+ Message::UpdateFulfillHTLC(msg) => {
self.message_handler.chan_handler.handle_update_fulfill_htlc(their_node_id, msg);
},
- wire::Message::UpdateFailHTLC(msg) => {
+ Message::UpdateFailHTLC(msg) => {
self.message_handler.chan_handler.handle_update_fail_htlc(their_node_id, &msg);
},
- wire::Message::UpdateFailMalformedHTLC(msg) => {
+ Message::UpdateFailMalformedHTLC(msg) => {
let chan_handler = &self.message_handler.chan_handler;
chan_handler.handle_update_fail_malformed_htlc(their_node_id, &msg);
},
- wire::Message::CommitmentSigned(msg) => {
+ Message::CommitmentSigned(msg) => {
self.message_handler.chan_handler.handle_commitment_signed(their_node_id, &msg);
},
- wire::Message::RevokeAndACK(msg) => {
+ Message::RevokeAndACK(msg) => {
self.message_handler.chan_handler.handle_revoke_and_ack(their_node_id, &msg);
},
- wire::Message::UpdateFee(msg) => {
+ Message::UpdateFee(msg) => {
self.message_handler.chan_handler.handle_update_fee(their_node_id, &msg);
},
- wire::Message::ChannelReestablish(msg) => {
+ Message::ChannelReestablish(msg) => {
self.message_handler.chan_handler.handle_channel_reestablish(their_node_id, &msg);
},
// Routing messages:
- wire::Message::AnnouncementSignatures(msg) => {
+ Message::AnnouncementSignatures(msg) => {
let chan_handler = &self.message_handler.chan_handler;
chan_handler.handle_announcement_signatures(their_node_id, &msg);
},
- wire::Message::ChannelAnnouncement(msg) => {
+ Message::ChannelAnnouncement(msg) => {
let route_handler = &self.message_handler.route_handler;
if route_handler
.handle_channel_announcement(Some(their_node_id), &msg)
@@ -2570,7 +2587,7 @@ where
}
self.update_gossip_backlogged();
},
- wire::Message::NodeAnnouncement(msg) => {
+ Message::NodeAnnouncement(msg) => {
let route_handler = &self.message_handler.route_handler;
if route_handler
.handle_node_announcement(Some(their_node_id), &msg)
@@ -2580,7 +2597,7 @@ where
}
self.update_gossip_backlogged();
},
- wire::Message::ChannelUpdate(msg) => {
+ Message::ChannelUpdate(msg) => {
let chan_handler = &self.message_handler.chan_handler;
chan_handler.handle_channel_update(their_node_id, &msg);
@@ -2594,31 +2611,31 @@ where
}
self.update_gossip_backlogged();
},
- wire::Message::QueryShortChannelIds(msg) => {
+ Message::QueryShortChannelIds(msg) => {
let route_handler = &self.message_handler.route_handler;
route_handler.handle_query_short_channel_ids(their_node_id, msg)?;
},
- wire::Message::ReplyShortChannelIdsEnd(msg) => {
+ Message::ReplyShortChannelIdsEnd(msg) => {
let route_handler = &self.message_handler.route_handler;
route_handler.handle_reply_short_channel_ids_end(their_node_id, msg)?;
},
- wire::Message::QueryChannelRange(msg) => {
+ Message::QueryChannelRange(msg) => {
let route_handler = &self.message_handler.route_handler;
route_handler.handle_query_channel_range(their_node_id, msg)?;
},
- wire::Message::ReplyChannelRange(msg) => {
+ Message::ReplyChannelRange(msg) => {
let route_handler = &self.message_handler.route_handler;
route_handler.handle_reply_channel_range(their_node_id, msg)?;
},
// Onion message:
- wire::Message::OnionMessage(msg) => {
+ Message::OnionMessage(msg) => {
let onion_message_handler = &self.message_handler.onion_message_handler;
onion_message_handler.handle_onion_message(their_node_id, &msg);
},
// Unknown messages:
- wire::Message::Unknown(type_id) if message.is_even() => {
+ Message::Unknown(type_id) if message.is_even() => {
log_debug!(
logger,
"Received unknown even message of type {}, disconnecting peer!",
@@ -2626,10 +2643,10 @@ where
);
return Err(PeerHandleError {}.into());
},
- wire::Message::Unknown(type_id) => {
+ Message::Unknown(type_id) => {
log_trace!(logger, "Received unknown odd message of type {}, ignoring", type_id);
},
- wire::Message::Custom(custom) => {
+ Message::Custom(custom) => {
let custom_message_handler = &self.message_handler.custom_message_handler;
custom_message_handler.handle_custom_message(custom, their_node_id)?;
},
@@ -2858,68 +2875,77 @@ where
// robustly gossip broadcast events even if a peer's message buffer is full.
let mut handle_event = |event, from_chan_handler| {
match event {
- MessageSendEvent::SendPeerStorage { ref node_id, ref msg } => {
+ MessageSendEvent::SendPeerStorage { ref node_id, msg } => {
log_debug!(
WithContext::from(&self.logger, Some(*node_id), None, None),
"Handling SendPeerStorage event in peer_handler for {}",
node_id,
);
+ let msg = Message::PeerStorage(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendPeerStorageRetrieval { ref node_id, ref msg } => {
+ MessageSendEvent::SendPeerStorageRetrieval { ref node_id, msg } => {
log_debug!(
WithContext::from(&self.logger, Some(*node_id), None, None),
"Handling SendPeerStorageRetrieval event in peer_handler for {}",
node_id,
);
+ let msg = Message::PeerStorageRetrieval(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendAcceptChannel { ref node_id, ref msg } => {
+ MessageSendEvent::SendAcceptChannel { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.common_fields.temporary_channel_id), None), "Handling SendAcceptChannel event in peer_handler for node {} for channel {}",
node_id,
&msg.common_fields.temporary_channel_id);
+ let msg = Message::AcceptChannel(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendAcceptChannelV2 { ref node_id, ref msg } => {
+ MessageSendEvent::SendAcceptChannelV2 { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.common_fields.temporary_channel_id), None), "Handling SendAcceptChannelV2 event in peer_handler for node {} for channel {}",
node_id,
&msg.common_fields.temporary_channel_id);
+ let msg = Message::AcceptChannelV2(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendOpenChannel { ref node_id, ref msg } => {
+ MessageSendEvent::SendOpenChannel { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.common_fields.temporary_channel_id), None), "Handling SendOpenChannel event in peer_handler for node {} for channel {}",
node_id,
&msg.common_fields.temporary_channel_id);
+ let msg = Message::OpenChannel(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendOpenChannelV2 { ref node_id, ref msg } => {
+ MessageSendEvent::SendOpenChannelV2 { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.common_fields.temporary_channel_id), None), "Handling SendOpenChannelV2 event in peer_handler for node {} for channel {}",
node_id,
&msg.common_fields.temporary_channel_id);
+ let msg = Message::OpenChannelV2(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendFundingCreated { ref node_id, ref msg } => {
+ MessageSendEvent::SendFundingCreated { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.temporary_channel_id), None), "Handling SendFundingCreated event in peer_handler for node {} for channel {} (which becomes {})",
node_id,
&msg.temporary_channel_id,
ChannelId::v1_from_funding_txid(msg.funding_txid.as_byte_array(), msg.funding_output_index));
// TODO: If the peer is gone we should generate a DiscardFunding event
// indicating to the wallet that they should just throw away this funding transaction
+ let msg = Message::FundingCreated(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendFundingSigned { ref node_id, ref msg } => {
+ MessageSendEvent::SendFundingSigned { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendFundingSigned event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::FundingSigned(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendChannelReady { ref node_id, ref msg } => {
+ MessageSendEvent::SendChannelReady { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendChannelReady event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::ChannelReady(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendStfu { ref node_id, ref msg } => {
+ MessageSendEvent::SendStfu { ref node_id, msg } => {
let logger = WithContext::from(
&self.logger,
Some(*node_id),
@@ -2929,9 +2955,10 @@ where
log_debug!(logger, "Handling SendStfu event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::Stfu(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendSpliceInit { ref node_id, ref msg } => {
+ MessageSendEvent::SendSpliceInit { ref node_id, msg } => {
let logger = WithContext::from(
&self.logger,
Some(*node_id),
@@ -2941,9 +2968,10 @@ where
log_debug!(logger, "Handling SendSpliceInit event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::SpliceInit(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendSpliceAck { ref node_id, ref msg } => {
+ MessageSendEvent::SendSpliceAck { ref node_id, msg } => {
let logger = WithContext::from(
&self.logger,
Some(*node_id),
@@ -2953,9 +2981,10 @@ where
log_debug!(logger, "Handling SendSpliceAck event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::SpliceAck(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendSpliceLocked { ref node_id, ref msg } => {
+ MessageSendEvent::SendSpliceLocked { ref node_id, msg } => {
let logger = WithContext::from(
&self.logger,
Some(*node_id),
@@ -2965,66 +2994,77 @@ where
log_debug!(logger, "Handling SendSpliceLocked event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::SpliceLocked(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxAddInput { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxAddInput { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxAddInput event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxAddInput(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxAddOutput { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxAddOutput { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxAddOutput event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxAddOutput(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxRemoveInput { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxRemoveInput { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxRemoveInput event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxRemoveInput(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxRemoveOutput { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxRemoveOutput { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxRemoveOutput event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxRemoveOutput(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxComplete { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxComplete { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxComplete event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxComplete(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxSignatures { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxSignatures { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxSignatures event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxSignatures(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxInitRbf { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxInitRbf { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxInitRbf event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxInitRbf(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxAckRbf { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxAckRbf { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxAckRbf event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxAckRbf(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendTxAbort { ref node_id, ref msg } => {
+ MessageSendEvent::SendTxAbort { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendTxAbort event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::TxAbort(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendAnnouncementSignatures { ref node_id, ref msg } => {
+ MessageSendEvent::SendAnnouncementSignatures { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendAnnouncementSignatures event in peer_handler for node {} for channel {})",
node_id,
&msg.channel_id);
+ let msg = Message::AnnouncementSignatures(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
MessageSendEvent::UpdateHTLCs {
@@ -3032,12 +3072,12 @@ where
ref channel_id,
updates:
msgs::CommitmentUpdate {
- ref update_add_htlcs,
- ref update_fulfill_htlcs,
- ref update_fail_htlcs,
- ref update_fail_malformed_htlcs,
- ref update_fee,
- ref commitment_signed,
+ update_add_htlcs,
+ update_fulfill_htlcs,
+ update_fail_htlcs,
+ update_fail_malformed_htlcs,
+ update_fee,
+ commitment_signed,
},
} => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(*channel_id), None), "Handling UpdateHTLCs event in peer_handler for node {} with {} adds, {} fulfills, {} fails, {} commits for channel {}",
@@ -3049,18 +3089,23 @@ where
channel_id);
let mut peer = get_peer_for_forwarding!(node_id)?;
for msg in update_fulfill_htlcs {
+ let msg = Message::UpdateFulfillHTLC(msg);
self.enqueue_message(&mut *peer, msg);
}
for msg in update_fail_htlcs {
+ let msg = Message::UpdateFailHTLC(msg);
self.enqueue_message(&mut *peer, msg);
}
for msg in update_fail_malformed_htlcs {
+ let msg = Message::UpdateFailMalformedHTLC(msg);
self.enqueue_message(&mut *peer, msg);
}
for msg in update_add_htlcs {
+ let msg = Message::UpdateAddHTLC(msg);
self.enqueue_message(&mut *peer, msg);
}
- if let &Some(ref msg) = update_fee {
+ if let Some(msg) = update_fee {
+ let msg = Message::UpdateFee(msg);
self.enqueue_message(&mut *peer, msg);
}
if commitment_signed.len() > 1 {
@@ -3069,37 +3114,45 @@ where
batch_size: commitment_signed.len() as u16,
message_type: Some(msgs::CommitmentSigned::TYPE),
};
- self.enqueue_message(&mut *peer, &msg);
+ let msg = Message::StartBatch(msg);
+ self.enqueue_message(&mut *peer, msg);
}
for msg in commitment_signed {
+ let msg = Message::CommitmentSigned(msg);
self.enqueue_message(&mut *peer, msg);
}
},
- MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
+ MessageSendEvent::SendRevokeAndACK { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendRevokeAndACK event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::RevokeAndACK(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendClosingSigned { ref node_id, ref msg } => {
+ MessageSendEvent::SendClosingSigned { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendClosingSigned event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::ClosingSigned(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendClosingComplete { ref node_id, ref msg } => {
+ #[cfg(simple_close)]
+ MessageSendEvent::SendClosingComplete { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendClosingComplete event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::ClosingComplete(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendClosingSig { ref node_id, ref msg } => {
+ #[cfg(simple_close)]
+ MessageSendEvent::SendClosingSig { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendClosingSig event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::ClosingSig(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendShutdown { ref node_id, ref msg } => {
+ MessageSendEvent::SendShutdown { ref node_id, msg } => {
log_debug!(
WithContext::from(
&self.logger,
@@ -3109,23 +3162,27 @@ where
),
"Handling Shutdown event in peer_handler",
);
+ let msg = Message::Shutdown(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
+ MessageSendEvent::SendChannelReestablish { ref node_id, msg } => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), Some(msg.channel_id), None), "Handling SendChannelReestablish event in peer_handler for node {} for channel {}",
node_id,
&msg.channel_id);
+ let msg = Message::ChannelReestablish(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
MessageSendEvent::SendChannelAnnouncement {
ref node_id,
- ref msg,
- ref update_msg,
+ msg,
+ update_msg,
} => {
log_debug!(WithContext::from(&self.logger, Some(*node_id), None, None), "Handling SendChannelAnnouncement event in peer_handler for node {} for short channel id {}",
node_id,
msg.contents.short_channel_id);
+ let msg = Message::ChannelAnnouncement(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
+ let update_msg = Message::ChannelUpdate(update_msg);
self.enqueue_message(
&mut *get_peer_for_forwarding!(node_id)?,
update_msg,
@@ -3216,12 +3273,13 @@ where
_ => {},
}
},
- MessageSendEvent::SendChannelUpdate { ref node_id, ref msg } => {
+ MessageSendEvent::SendChannelUpdate { ref node_id, msg } => {
log_trace!(
WithContext::from(&self.logger, Some(*node_id), None, None),
"Handling SendChannelUpdate event in peer_handler for channel {}",
msg.contents.short_channel_id
);
+ let msg = Message::ChannelUpdate(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
MessageSendEvent::HandleError { node_id, action } => {
@@ -3239,7 +3297,7 @@ where
// about to disconnect the peer and do it after we finish
// processing most messages.
let msg = msg.map(|msg| {
- wire::Message::<<<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage>::Error(msg)
+ Message::<<<CMH as Deref>::Target as wire::CustomMessageReader>::CustomMessage>::Error(msg)
});
peers_to_disconnect.insert(node_id, msg);
},
@@ -3250,7 +3308,7 @@ where
// about to disconnect the peer and do it after we finish
// processing most messages.
peers_to_disconnect
- .insert(node_id, Some(wire::Message::Warning(msg)));
+ .insert(node_id, Some(Message::Warning(msg)));
},
msgs::ErrorAction::IgnoreAndLog(level) => {
log_given_level!(
@@ -3266,22 +3324,21 @@ where
"Received a HandleError event to be ignored",
);
},
- msgs::ErrorAction::SendErrorMessage { ref msg } => {
+ msgs::ErrorAction::SendErrorMessage { msg } => {
log_trace!(logger, "Handling SendErrorMessage HandleError event in peer_handler with message {}",
msg.data);
+ let msg = Message::Error(msg);
self.enqueue_message(
&mut *get_peer_for_forwarding!(&node_id)?,
msg,
);
},
- msgs::ErrorAction::SendWarningMessage {
- ref msg,
- ref log_level,
- } => {
+ msgs::ErrorAction::SendWarningMessage { msg, ref log_level } => {
log_given_level!(logger, *log_level, "Handling SendWarningMessage HandleError event in peer_handler with message {}",
msg.data);
+ let msg = Message::Warning(msg);
self.enqueue_message(
&mut *get_peer_for_forwarding!(&node_id)?,
msg,
@@ -3289,33 +3346,37 @@ where
},
}
},
- MessageSendEvent::SendChannelRangeQuery { ref node_id, ref msg } => {
+ MessageSendEvent::SendChannelRangeQuery { ref node_id, msg } => {
log_gossip!(WithContext::from(&self.logger, Some(*node_id), None, None), "Handling SendChannelRangeQuery event in peer_handler with first_blocknum={}, number_of_blocks={}",
msg.first_blocknum,
msg.number_of_blocks);
+ let msg = Message::QueryChannelRange(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendShortIdsQuery { ref node_id, ref msg } => {
+ MessageSendEvent::SendShortIdsQuery { ref node_id, msg } => {
log_gossip!(WithContext::from(&self.logger, Some(*node_id), None, None), "Handling SendShortIdsQuery event in peer_handler with num_scids={}",
msg.short_channel_ids.len());
+ let msg = Message::QueryShortChannelIds(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendReplyChannelRange { ref node_id, ref msg } => {
+ MessageSendEvent::SendReplyChannelRange { ref node_id, msg } => {
log_gossip!(WithContext::from(&self.logger, Some(*node_id), None, None), "Handling SendReplyChannelRange event in peer_handler with num_scids={} first_blocknum={} number_of_blocks={}, sync_complete={}",
msg.short_channel_ids.len(),
msg.first_blocknum,
msg.number_of_blocks,
msg.sync_complete);
+ let msg = Message::ReplyChannelRange(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
- MessageSendEvent::SendGossipTimestampFilter { ref node_id, ref msg } => {
+ MessageSendEvent::SendGossipTimestampFilter { ref node_id, msg } => {
log_gossip!(WithContext::from(&self.logger, Some(*node_id), None, None), "Handling SendGossipTimestampFilter event in peer_handler with first_timestamp={}, timestamp_range={}",
msg.first_timestamp,
msg.timestamp_range);
+ let msg = Message::GossipTimestampFilter(msg);
self.enqueue_message(&mut *get_peer_for_forwarding!(node_id)?, msg);
},
}
@@ -3351,7 +3412,8 @@ where
} else {
continue;
};
- self.enqueue_message(&mut peer, &msg);
+ let msg = Message::Custom(msg);
+ self.enqueue_message(&mut peer, msg);
}
for (descriptor, peer_mutex) in peers.iter() {
@@ -3381,7 +3443,7 @@ where
if let Some(peer_mutex) = peers.remove(&descriptor) {
let mut peer = peer_mutex.lock().unwrap();
if let Some(msg) = msg {
- self.enqueue_message(&mut *peer, &msg);
+ self.enqueue_message(&mut *peer, msg);
// This isn't guaranteed to work, but if there is enough free
// room in the send buffer, put the error message there...
self.do_attempt_write_data(&mut descriptor, &mut *peer, false);
@@ -3506,7 +3568,9 @@ where
if peer.awaiting_pong_timer_tick_intervals == 0 {
peer.awaiting_pong_timer_tick_intervals = -1;
let ping = msgs::Ping { ponglen: 0, byteslen: 64 };
- self.enqueue_message(peer, &ping);
+ let msg: Message<<CMH::Target as CustomMessageReader>::CustomMessage> =
+ Message::Ping(ping);
+ self.enqueue_message(peer, msg);
}
}
@@ -3577,7 +3641,8 @@ where
peer.awaiting_pong_timer_tick_intervals = 1;
let ping = msgs::Ping { ponglen: 0, byteslen: 64 };
- self.enqueue_message(&mut *peer, &ping);
+ let msg = Message::Ping(ping);
+ self.enqueue_message(&mut *peer, msg);
break;
}
self.do_attempt_write_data(
@@ -4226,7 +4291,7 @@ mod tests {
.push(MessageSendEvent::SendShutdown { node_id: their_id, msg: msg.clone() });
peers[0].message_handler.chan_handler = &a_chan_handler;
- b_chan_handler.expect_receive_msg(wire::Message::Shutdown(msg));
+ b_chan_handler.expect_receive_msg(Message::Shutdown(msg));
peers[1].message_handler.chan_handler = &b_chan_handler;
peers[0].process_events();
@@ -4261,7 +4326,8 @@ mod tests {
peers[0].read_event(&mut fd_dup, &act_three).unwrap();
let not_init_msg = msgs::Ping { ponglen: 4, byteslen: 0 };
- let msg_bytes = dup_encryptor.encrypt_message(¬_init_msg);
+ let msg: Message<()> = Message::Ping(not_init_msg);
+ let msg_bytes = dup_encryptor.encrypt_message(msg);
assert!(peers[0].read_event(&mut fd_dup, &msg_bytes).is_err());
}
@@ -4639,13 +4705,12 @@ mod tests {
{
let peers = peer_a.peers.read().unwrap();
let mut peer_b = peers.get(&fd_a).unwrap().lock().unwrap();
- peer_a.enqueue_message(
- &mut peer_b,
- &msgs::WarningMessage {
- channel_id: ChannelId([0; 32]),
- data: "no disconnect plz".to_string(),
- },
- );
+ let warning = msgs::WarningMessage {
+ channel_id: ChannelId([0; 32]),
+ data: "no disconnect plz".to_string(),
+ };
+ let msg = Message::Warning(warning);
+ peer_a.enqueue_message(&mut peer_b, msg);
}
peer_a.process_events();
let msg = fd_a.outbound_data.lock().unwrap().split_off(0);
Why this scored 19/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.