p2p: add consensus_encoding to v1 network message
What changed, and why it matters
This commit adds a new way to encode and decode Bitcoin peer-to-peer (P2P) network messages in the rust-bitcoin library. It is a feature/refactoring change that introduces a new 'consensus_encoding' implementation for v1 network messages, replacing the old serialization approach in an example program. There is no indication this fixes a security bug or introduces a vulnerability.
No security action required. Review as normal code-quality/refactoring change.
Security signals we found
No security-relevant signals detected in the diff or commit message.
The decoder enforces a maximum payload length via MAX_MSG_SIZE, which is a standard defensive check, not a fix for a known issue.
The change is framed as adding consensus_encoding support, not as a security patch.
Evidence from the diff
The commit implements new encoding::Encodable and encoding::Decodable traits for RawNetworkMessage in p2p/src/message.rs. It adds NetworkMessageEncoder, RawNetworkMessageEncoder, NetworkMessageDecoder, RawNetworkMessageDecoder, and a RawNetworkMessageDecodeError type. The decoder performs two-phase parsing: fixed-size header first, then variable-length payload, with a check that payload_len does not exceed MAX_MSG_SIZE. The example p2p/examples/handshake.rs is updated to use the new encoding::encode_to_writer and encoding::decode_from_read functions. The payload decoding currently delegates back to the legacy bitcoin::consensus::encode::Decodable implementation for each command type, with TODO comments noting future migration.
Changed components
p2p/src/message.rsp2p/examples/handshake.rsInspect captured patch +520 / −5
diff --git a/p2p/examples/handshake.rs b/p2p/examples/handshake.rs
index 77aadf5a..690ffc06 100644
--- a/p2p/examples/handshake.rs
+++ b/p2p/examples/handshake.rs
@@ -1,9 +1,8 @@
-use std::io::{BufReader, Write};
+use std::io::BufReader;
use std::net::{IpAddr, Ipv4Addr, Shutdown, SocketAddr, TcpStream};
use std::time::{SystemTime, UNIX_EPOCH};
use std::{env, process};
-use bitcoin::consensus::{encode, Decodable};
use bitcoin_p2p_messages::message_network::{ClientSoftwareVersion, UserAgent, UserAgentVersion};
use bitcoin_p2p_messages::{
self, address, message, message_network, Magic, ProtocolVersion, ServiceFlags,
@@ -36,7 +35,7 @@ fn main() {
if let Ok(mut stream) = TcpStream::connect(address) {
// Send the message
- let _ = stream.write_all(encode::serialize(&first_message).as_slice());
+ encoding::encode_to_writer(&first_message, &mut stream).unwrap();
println!("Sent version message");
// Setup StreamReader
@@ -44,7 +43,9 @@ fn main() {
let mut stream_reader = BufReader::new(read_stream);
loop {
// Loop and retrieve new messages
- let reply = message::RawNetworkMessage::consensus_decode(&mut stream_reader).unwrap();
+ let reply =
+ encoding::decode_from_read::<message::RawNetworkMessage, _>(&mut stream_reader)
+ .unwrap();
match reply.payload() {
message::NetworkMessage::Version(_) => {
println!("Received version message: {:?}", reply.payload());
@@ -54,7 +55,7 @@ fn main() {
message::NetworkMessage::Verack,
);
- let _ = stream.write_all(encode::serialize(&second_message).as_slice());
+ encoding::encode_to_writer(&second_message, &mut stream).unwrap();
println!("Sent verack message");
}
message::NetworkMessage::Verack => {
diff --git a/p2p/src/message.rs b/p2p/src/message.rs
index fca41136..e87cc88b 100644
--- a/p2p/src/message.rs
+++ b/p2p/src/message.rs
@@ -680,6 +680,520 @@ impl Encodable for RawNetworkMessage {
}
}
+struct NetworkMessageEncoder {
+ buffer: Vec<u8>,
+ exhausted: bool,
+}
+
+impl NetworkMessageEncoder {
+ fn new(msg: &NetworkMessage) -> Self {
+ let mut buffer = Vec::new();
+ // TODO: delegate to internal encoders once migrated to consensus_encoding.
+ bitcoin::consensus::encode::Encodable::consensus_encode(msg, &mut buffer)
+ .expect("encoding to vec cannot fail");
+ Self { buffer, exhausted: false }
+ }
+}
+
+impl encoding::Encoder for NetworkMessageEncoder {
+ fn current_chunk(&self) -> &[u8] {
+ if self.exhausted {
+ &[]
+ } else {
+ &self.buffer
+ }
+ }
+
+ fn advance(&mut self) -> bool {
+ self.exhausted = true;
+ false
+ }
+}
+
+encoding::encoder_newtype! {
+ /// Encoder for [`RawNetworkMessage`].
+ pub struct RawNetworkMessageEncoder(
+ encoding::Encoder2<
+ encoding::Encoder4<
+ encoding::ArrayEncoder<4>,
+ encoding::ArrayEncoder<12>,
+ encoding::ArrayEncoder<4>,
+ encoding::ArrayEncoder<4>,
+ >,
+ NetworkMessageEncoder,
+ >
+ );
+}
+
+impl encoding::Encodable for RawNetworkMessage {
+ type Encoder<'e> = RawNetworkMessageEncoder;
+
+ fn encoder(&self) -> Self::Encoder<'_> {
+ RawNetworkMessageEncoder(encoding::Encoder2::new(
+ encoding::Encoder4::new(
+ encoding::ArrayEncoder::without_length_prefix(self.magic.to_bytes()),
+ self.command().encoder(),
+ encoding::ArrayEncoder::without_length_prefix(self.payload_len.to_le_bytes()),
+ encoding::ArrayEncoder::without_length_prefix(self.checksum),
+ ),
+ NetworkMessageEncoder::new(&self.payload),
+ ))
+ }
+}
+
+struct NetworkMessageDecoder {
+ command: CommandString,
+ payload_len: usize,
+ buffer: Vec<u8>,
+}
+
+impl NetworkMessageDecoder {
+ fn new(command: CommandString, payload_len: usize) -> Self {
+ Self { command, payload_len, buffer: Vec::new() }
+ }
+}
+
+impl encoding::Decoder for NetworkMessageDecoder {
+ type Output = NetworkMessage;
+ type Error = RawNetworkMessageDecodeError;
+
+ fn push_bytes(&mut self, bytes: &mut &[u8]) -> Result<bool, Self::Error> {
+ let remaining = self.payload_len - self.buffer.len();
+ let copy_len = bytes.len().min(remaining);
+
+ self.buffer.extend_from_slice(&bytes[..copy_len]);
+ *bytes = &bytes[copy_len..];
+
+ Ok(self.buffer.len() < self.payload_len)
+ }
+
+ fn end(self) -> Result<Self::Output, Self::Error> {
+ let payload_bytes = self.buffer;
+
+ // Validate payload length matches actual data.
+ if payload_bytes.len() != self.payload_len {
+ return Err(RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload));
+ }
+
+ // TODO: delegate to internal decoders once migrated to consensus_encoding.
+ let mut mem_d = payload_bytes.as_slice();
+ let payload = match self.command.as_ref() {
+ "version" => NetworkMessage::Version(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "verack" => NetworkMessage::Verack,
+ "addr" => NetworkMessage::Addr(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "inv" => NetworkMessage::Inv(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getdata" => NetworkMessage::GetData(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "notfound" => NetworkMessage::NotFound(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getblocks" => NetworkMessage::GetBlocks(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getheaders" => NetworkMessage::GetHeaders(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "mempool" => NetworkMessage::MemPool,
+ "block" => NetworkMessage::Block(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "headers" => NetworkMessage::Headers(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "sendheaders" => NetworkMessage::SendHeaders,
+ "getaddr" => NetworkMessage::GetAddr,
+ "ping" => NetworkMessage::Ping(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "pong" => NetworkMessage::Pong(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "merkleblock" => NetworkMessage::MerkleBlock(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "filterload" => NetworkMessage::FilterLoad(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "filteradd" => NetworkMessage::FilterAdd(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "filterclear" => NetworkMessage::FilterClear,
+ "getcfilters" => NetworkMessage::GetCFilters(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "cfilter" => NetworkMessage::CFilter(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getcfheaders" => NetworkMessage::GetCFHeaders(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "cfheaders" => NetworkMessage::CFHeaders(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getcfcheckpt" => NetworkMessage::GetCFCheckpt(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "cfcheckpt" => NetworkMessage::CFCheckpt(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "sendcmpct" => NetworkMessage::SendCmpct(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "cmpctblock" => NetworkMessage::CmpctBlock(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "getblocktxn" => NetworkMessage::GetBlockTxn(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "blocktxn" => NetworkMessage::BlockTxn(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "tx" => NetworkMessage::Tx(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "alert" => NetworkMessage::Alert(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "reject" => NetworkMessage::Reject(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "feefilter" => NetworkMessage::FeeFilter(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "wtxidrelay" => NetworkMessage::WtxidRelay,
+ "addrv2" => NetworkMessage::AddrV2(
+ bitcoin::consensus::encode::Decodable::consensus_decode_from_finite_reader(
+ &mut mem_d,
+ )
+ .map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Payload)
+ })?,
+ ),
+ "sendaddrv2" => NetworkMessage::SendAddrV2,
+ _ => NetworkMessage::Unknown { command: self.command, payload: payload_bytes },
+ };
+
+ Ok(payload)
+ }
+
+ fn read_limit(&self) -> usize { self.payload_len - self.buffer.len() }
+}
+
+enum DecoderState {
+ ReadingHeader {
+ header_decoder: encoding::Decoder4<
+ encoding::ArrayDecoder<4>,
+ CommandStringDecoder,
+ encoding::ArrayDecoder<4>,
+ encoding::ArrayDecoder<4>,
+ >,
+ },
+ ReadingPayload {
+ magic_bytes: [u8; 4],
+ payload_len_bytes: [u8; 4],
+ checksum: [u8; 4],
+ payload_decoder: NetworkMessageDecoder,
+ },
+}
+
+/// Decoder for [`RawNetworkMessage`].
+///
+/// This decoder implements a two-phase decoding process for Bitcoin V1 P2P messages.
+/// It first decodes the fixed-sized header. It then uses the payload length information
+/// to decode the dynamically sized network message.
+pub struct RawNetworkMessageDecoder {
+ state: DecoderState,
+}
+
+impl encoding::Decoder for RawNetworkMessageDecoder {
+ type Output = RawNetworkMessage;
+ type Error = RawNetworkMessageDecodeError;
+
+ fn push_bytes(&mut self, bytes: &mut &[u8]) -> Result<bool, Self::Error> {
+ match &mut self.state {
+ DecoderState::ReadingHeader { header_decoder } => {
+ let need_more = header_decoder.push_bytes(bytes).map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Header)
+ })?;
+
+ if !need_more {
+ // Header complete, extract values and transition to payload state.
+ let old_state = core::mem::replace(
+ &mut self.state,
+ DecoderState::ReadingHeader {
+ header_decoder: encoding::Decoder4::new(
+ encoding::ArrayDecoder::new(),
+ CommandStringDecoder { inner: encoding::ArrayDecoder::new() },
+ encoding::ArrayDecoder::new(),
+ encoding::ArrayDecoder::new(),
+ ),
+ },
+ );
+
+ let header_decoder = match old_state {
+ DecoderState::ReadingHeader { header_decoder } => header_decoder,
+ _ => unreachable!("we are in ReadingHeader state"),
+ };
+
+ let (magic_bytes, command, payload_len_bytes, checksum) =
+ header_decoder.end().map_err(|_| {
+ RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Header)
+ })?;
+
+ let payload_len = u32::from_le_bytes(payload_len_bytes) as usize;
+ if payload_len > MAX_MSG_SIZE {
+ return Err(RawNetworkMessageDecodeError(
+ RawNetworkMessageDecodeErrorInner::PayloadTooLarge,
+ ));
+ }
+
+ let payload_decoder = NetworkMessageDecoder::new(command, payload_len);
+ self.state = DecoderState::ReadingPayload {
+ magic_bytes,
+ payload_len_bytes,
+ checksum,
+ payload_decoder,
+ };
+
+ // Continue with any remaining bytes.
+ return self.push_bytes(bytes);
+ }
+
+ Ok(need_more)
+ }
+ DecoderState::ReadingPayload { payload_decoder, .. } =>
+ payload_decoder.push_bytes(bytes),
+ }
+ }
+
+ fn end(self) -> Result<Self::Output, Self::Error> {
+ match self.state {
+ DecoderState::ReadingHeader { .. } =>
+ Err(RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner::Header)),
+ DecoderState::ReadingPayload {
+ magic_bytes,
+ payload_len_bytes,
+ checksum,
+ payload_decoder,
+ ..
+ } => {
+ let payload = payload_decoder.end()?;
+
+ Ok(RawNetworkMessage {
+ magic: Magic::from_bytes(magic_bytes),
+ payload,
+ payload_len: u32::from_le_bytes(payload_len_bytes),
+ checksum,
+ })
+ }
+ }
+ }
+
+ fn read_limit(&self) -> usize {
+ match &self.state {
+ DecoderState::ReadingHeader { header_decoder } => header_decoder.read_limit(),
+ DecoderState::ReadingPayload { payload_decoder, .. } => payload_decoder.read_limit(),
+ }
+ }
+}
+
+impl encoding::Decodable for RawNetworkMessage {
+ type Decoder = RawNetworkMessageDecoder;
+
+ fn decoder() -> Self::Decoder {
+ RawNetworkMessageDecoder {
+ state: DecoderState::ReadingHeader {
+ header_decoder: encoding::Decoder4::new(
+ encoding::ArrayDecoder::new(),
+ CommandStringDecoder { inner: encoding::ArrayDecoder::new() },
+ encoding::ArrayDecoder::new(),
+ encoding::ArrayDecoder::new(),
+ ),
+ },
+ }
+ }
+}
+
+/// Error decoding a raw network message.
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct RawNetworkMessageDecodeError(RawNetworkMessageDecodeErrorInner);
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+enum RawNetworkMessageDecodeErrorInner {
+ /// Error decoding the message header.
+ Header,
+ /// Payload length exceeds maximum allowed message size.
+ PayloadTooLarge,
+ /// Error decoding the message payload.
+ Payload,
+}
+
+impl fmt::Display for RawNetworkMessageDecodeError {
+ fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
+ match self.0 {
+ RawNetworkMessageDecodeErrorInner::Header => {
+ write!(f, "error decoding message header")
+ }
+ RawNetworkMessageDecodeErrorInner::PayloadTooLarge => {
+ write!(f, "payload length exceeds maximum allowed message size")
+ }
+ RawNetworkMessageDecodeErrorInner::Payload => {
+ write!(f, "error decoding message payload")
+ }
+ }
+ }
+}
+
+#[cfg(feature = "std")]
+impl std::error::Error for RawNetworkMessageDecodeError {}
+
impl Encodable for V2NetworkMessage {
fn consensus_encode<W: Write + ?Sized>(&self, writer: &mut W) -> Result<usize, io::Error> {
// A subset of message types are optimized to only use one byte to encode the command.
Why this scored 15/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.