From 21a3ffdb1d8589d9178919596554c2e02b6c07df Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Thu, 12 Mar 2026 16:20:34 +1000 Subject: [PATCH 1/6] Defer Snappy decompression of P2P messages until worker thread processing RawMessage now supports a compressed constructor that defers Snappy decompression until getData() is first called on the worker thread. Messages stay in their compressed form while queued in the tx worker pool. Additionally, worker threads now skip processing for messages from already-disconnected peers, and decompression/deserialization failures disconnect the peer with BREACH_OF_PROTOCOL. Co-Authored-By: Claude Opus 4.6 Signed-off-by: stefan.pingel@consensys.net --- ...PooledTransactionHashesMessageHandler.java | 30 ++++++++++++++--- .../TransactionsMessageHandler.java | 29 +++++++++++++--- .../rlpx/connections/netty/ApiHandler.java | 20 ++++++----- .../ethereum/p2p/rlpx/framing/Framer.java | 33 +++++++++---------- .../p2p/rlpx/wire/AbstractMessageData.java | 4 +-- .../ethereum/p2p/rlpx/wire/RawMessage.java | 30 +++++++++++++++++ 6 files changed, 109 insertions(+), 37 deletions(-) diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java index f2e9378fd4b..37f80851ace 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java @@ -22,12 +22,19 @@ import org.hyperledger.besu.ethereum.eth.manager.EthScheduler; import org.hyperledger.besu.ethereum.eth.messages.NewPooledTransactionHashesMessage; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Capability; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; import java.time.Duration; import java.time.Instant; import java.util.concurrent.atomic.AtomicBoolean; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + class NewPooledTransactionHashesMessageHandler implements EthMessages.MessageCallback { + private static final Logger LOG = + LoggerFactory.getLogger(NewPooledTransactionHashesMessageHandler.class); private final NewPooledTransactionHashesMessageProcessor transactionsMessageProcessor; private final EthScheduler scheduler; @@ -47,13 +54,28 @@ public NewPooledTransactionHashesMessageHandler( public void exec(final EthMessage message) { if (isEnabled.get()) { final Capability capability = message.getPeer().getConnection().capability(EthProtocol.NAME); - final NewPooledTransactionHashesMessage transactionsMessage = - NewPooledTransactionHashesMessage.readFrom(message.getData(), capability); + final MessageData rawMessage = message.getData(); final Instant startedAt = now(); scheduler.scheduleTxWorkerTask( - () -> + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + try { + final NewPooledTransactionHashesMessage transactionsMessage = + NewPooledTransactionHashesMessage.readFrom(rawMessage, capability); transactionsMessageProcessor.processNewPooledTransactionHashesMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + } catch (final Exception e) { + LOG.debug( + "Malformed pooled transaction hashes message received (BREACH_OF_PROTOCOL), disconnecting: {}", + message.getPeer(), + e); + message + .getPeer() + .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } + }); } } diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java index b6bbf442dab..930d99ff55b 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java @@ -20,12 +20,18 @@ import org.hyperledger.besu.ethereum.eth.manager.EthMessages; import org.hyperledger.besu.ethereum.eth.manager.EthScheduler; import org.hyperledger.besu.ethereum.eth.messages.TransactionsMessage; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; import java.time.Duration; import java.time.Instant; import java.util.concurrent.atomic.AtomicBoolean; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + class TransactionsMessageHandler implements EthMessages.MessageCallback { + private static final Logger LOG = LoggerFactory.getLogger(TransactionsMessageHandler.class); private final TransactionsMessageProcessor transactionsMessageProcessor; private final EthScheduler scheduler; @@ -44,13 +50,28 @@ public TransactionsMessageHandler( @Override public void exec(final EthMessage message) { if (isEnabled.get()) { - final TransactionsMessage transactionsMessage = - TransactionsMessage.readFrom(message.getData()); + final MessageData rawMessage = message.getData(); final Instant startedAt = now(); scheduler.scheduleTxWorkerTask( - () -> + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + try { + final TransactionsMessage transactionsMessage = + TransactionsMessage.readFrom(rawMessage); transactionsMessageProcessor.processTransactionsMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + } catch (final Exception e) { + LOG.debug( + "Malformed transactions message received (BREACH_OF_PROTOCOL), disconnecting: {}", + message.getPeer(), + e); + message + .getPeer() + .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } + }); } } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java index d0a689a0521..acedf401a26 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java @@ -99,15 +99,17 @@ protected void channelRead0(final ChannelHandlerContext ctx, final MessageData o } return; } - LOG.atTrace() - .addMarker(P2P_MESSAGE_MARKER) - .setMessage("Received {} from {} via protocol {}") - .addArgument(message) - .addArgument(connection.getPeerInfo()) - .addArgument(demultiplexed.getCapability()) - .addKeyValue("rawData", message.getData()) - .addKeyValue("decodedData", message::toStringDecoded) - .log(); + if (LOG.isTraceEnabled()) { + LOG.atTrace() + .addMarker(P2P_MESSAGE_MARKER) + .setMessage("Received {} from {} via protocol {}") + .addArgument(message) + .addArgument(connection.getPeerInfo()) + .addArgument(demultiplexed.getCapability()) + .addKeyValue("rawData", message.getData()) + .addKeyValue("decodedData", message::toStringDecoded) + .log(); + } connectionEventDispatcher.dispatchMessage(demultiplexed.getCapability(), connection, message); } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java index 1302b3765d4..c2ba9571944 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java @@ -282,38 +282,35 @@ private MessageData processFrame(final ByteBuf f, final int frameSize) { final int id = idbv.isZero() || idbv.size() == 0 ? 0 : idbv.get(0); // Write message data to ByteBuf, decompressing as necessary - final Bytes data; if (compressionEnabled) { final byte[] compressedMessageData = Arrays.copyOfRange(frameData, 1, frameData.length - pad); final int uncompressedLength = compressor.uncompressedLength(compressedMessageData); if (uncompressedLength >= LENGTH_MAX_MESSAGE_FRAME) { throw error("Message size %s in excess of maximum length.", uncompressedLength); } - Bytes _data; - try { - final byte[] decompressedMessageData = compressor.decompress(compressedMessageData); - _data = Bytes.wrap(decompressedMessageData); - compressionSuccessful = true; - } catch (final FramingException fe) { - if (compressionSuccessful) { - throw fe; - } else { - // OpenEthereum/Parity does not implement EIP-706 - // If failing on the first packet downgrade to uncompressed + + if (!compressionSuccessful) { + // First compressed message: decompress eagerly to validate and handle + // the OpenEthereum/Parity fallback (non-Snappy peer detection via EIP-706) + try { + final byte[] decompressedMessageData = compressor.decompress(compressedMessageData); + compressionSuccessful = true; + return new RawMessage(id, Bytes.wrap(decompressedMessageData)); + } catch (final FramingException fe) { compressionEnabled = false; LOG.debug("Snappy decompression failed: downgrading to uncompressed"); final int messageLength = frameSize - LENGTH_MESSAGE_ID; - _data = Bytes.wrap(frameData, 1, messageLength); + return new RawMessage(id, Bytes.wrap(frameData, 1, messageLength)); } + } else { + // Subsequent messages: store compressed, decompress lazily on getData() + return new RawMessage(id, compressedMessageData); } - data = _data; } else { - // Move data to a ByteBuf final int messageLength = frameSize - LENGTH_MESSAGE_ID; - data = Bytes.wrap(frameData, 1, messageLength); + final Bytes data = Bytes.wrap(frameData, 1, messageLength); + return new RawMessage(id, data); } - - return new RawMessage(id, data); } private void validateMac(final byte[] candidateMac, final byte[] expectedMac) { diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java index c6c613b90b0..d8818ae9ae1 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java @@ -20,14 +20,14 @@ public abstract class AbstractMessageData implements MessageData { - protected final Bytes data; + protected Bytes data; protected AbstractMessageData(final Bytes data) { this.data = data; } @Override - public final int getSize() { + public int getSize() { return data.size(); } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java index b4b0bc20c59..e6d8bdace92 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java @@ -14,19 +14,49 @@ */ package org.hyperledger.besu.ethereum.p2p.rlpx.wire; +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.SnappyCompressor; + import org.apache.tuweni.bytes.Bytes; public final class RawMessage extends AbstractMessageData { + private static final SnappyCompressor compressor = new SnappyCompressor(); + private final int code; + private final byte[] compressedData; + /** Constructor for uncompressed messages. */ public RawMessage(final int code, final Bytes data) { super(data); this.code = code; + this.compressedData = null; + } + + /** Constructor for compressed messages — decompression is deferred until getData() is called. */ + public RawMessage(final int code, final byte[] compressedData) { + super(Bytes.EMPTY); + this.code = code; + this.compressedData = compressedData; } @Override public int getCode() { return code; } + + @Override + public Bytes getData() { + if (compressedData != null && data == Bytes.EMPTY) { + data = Bytes.wrap(compressor.decompress(compressedData)); + } + return data; + } + + @Override + public int getSize() { + if (compressedData == null) { + return data.size(); + } + return compressor.uncompressedLength(compressedData); + } } From a4962198590d720eab40277a814b391b494f3e36 Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Thu, 12 Mar 2026 16:20:34 +1000 Subject: [PATCH 2/6] Defer Snappy decompression of P2P messages until worker thread processing RawMessage now supports a compressed constructor that defers Snappy decompression until getData() is first called on the worker thread. Messages stay in their compressed form while queued in the tx worker pool. Additionally, worker threads now skip processing for messages from already-disconnected peers, and decompression/deserialization failures disconnect the peer with BREACH_OF_PROTOCOL. Co-Authored-By: Claude Opus 4.6 Signed-off-by: stefan.pingel@consensys.net --- ...PooledTransactionHashesMessageHandler.java | 30 ++++++++++++++--- .../TransactionsMessageHandler.java | 29 +++++++++++++--- .../rlpx/connections/netty/ApiHandler.java | 20 ++++++----- .../ethereum/p2p/rlpx/framing/Framer.java | 33 +++++++++---------- .../p2p/rlpx/wire/AbstractMessageData.java | 4 +-- .../ethereum/p2p/rlpx/wire/RawMessage.java | 31 +++++++++++++++++ 6 files changed, 110 insertions(+), 37 deletions(-) diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java index f2e9378fd4b..37f80851ace 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java @@ -22,12 +22,19 @@ import org.hyperledger.besu.ethereum.eth.manager.EthScheduler; import org.hyperledger.besu.ethereum.eth.messages.NewPooledTransactionHashesMessage; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Capability; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; import java.time.Duration; import java.time.Instant; import java.util.concurrent.atomic.AtomicBoolean; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + class NewPooledTransactionHashesMessageHandler implements EthMessages.MessageCallback { + private static final Logger LOG = + LoggerFactory.getLogger(NewPooledTransactionHashesMessageHandler.class); private final NewPooledTransactionHashesMessageProcessor transactionsMessageProcessor; private final EthScheduler scheduler; @@ -47,13 +54,28 @@ public NewPooledTransactionHashesMessageHandler( public void exec(final EthMessage message) { if (isEnabled.get()) { final Capability capability = message.getPeer().getConnection().capability(EthProtocol.NAME); - final NewPooledTransactionHashesMessage transactionsMessage = - NewPooledTransactionHashesMessage.readFrom(message.getData(), capability); + final MessageData rawMessage = message.getData(); final Instant startedAt = now(); scheduler.scheduleTxWorkerTask( - () -> + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + try { + final NewPooledTransactionHashesMessage transactionsMessage = + NewPooledTransactionHashesMessage.readFrom(rawMessage, capability); transactionsMessageProcessor.processNewPooledTransactionHashesMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + } catch (final Exception e) { + LOG.debug( + "Malformed pooled transaction hashes message received (BREACH_OF_PROTOCOL), disconnecting: {}", + message.getPeer(), + e); + message + .getPeer() + .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } + }); } } diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java index b6bbf442dab..930d99ff55b 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java @@ -20,12 +20,18 @@ import org.hyperledger.besu.ethereum.eth.manager.EthMessages; import org.hyperledger.besu.ethereum.eth.manager.EthScheduler; import org.hyperledger.besu.ethereum.eth.messages.TransactionsMessage; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; import java.time.Duration; import java.time.Instant; import java.util.concurrent.atomic.AtomicBoolean; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + class TransactionsMessageHandler implements EthMessages.MessageCallback { + private static final Logger LOG = LoggerFactory.getLogger(TransactionsMessageHandler.class); private final TransactionsMessageProcessor transactionsMessageProcessor; private final EthScheduler scheduler; @@ -44,13 +50,28 @@ public TransactionsMessageHandler( @Override public void exec(final EthMessage message) { if (isEnabled.get()) { - final TransactionsMessage transactionsMessage = - TransactionsMessage.readFrom(message.getData()); + final MessageData rawMessage = message.getData(); final Instant startedAt = now(); scheduler.scheduleTxWorkerTask( - () -> + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + try { + final TransactionsMessage transactionsMessage = + TransactionsMessage.readFrom(rawMessage); transactionsMessageProcessor.processTransactionsMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + } catch (final Exception e) { + LOG.debug( + "Malformed transactions message received (BREACH_OF_PROTOCOL), disconnecting: {}", + message.getPeer(), + e); + message + .getPeer() + .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } + }); } } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java index d0a689a0521..acedf401a26 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/ApiHandler.java @@ -99,15 +99,17 @@ protected void channelRead0(final ChannelHandlerContext ctx, final MessageData o } return; } - LOG.atTrace() - .addMarker(P2P_MESSAGE_MARKER) - .setMessage("Received {} from {} via protocol {}") - .addArgument(message) - .addArgument(connection.getPeerInfo()) - .addArgument(demultiplexed.getCapability()) - .addKeyValue("rawData", message.getData()) - .addKeyValue("decodedData", message::toStringDecoded) - .log(); + if (LOG.isTraceEnabled()) { + LOG.atTrace() + .addMarker(P2P_MESSAGE_MARKER) + .setMessage("Received {} from {} via protocol {}") + .addArgument(message) + .addArgument(connection.getPeerInfo()) + .addArgument(demultiplexed.getCapability()) + .addKeyValue("rawData", message.getData()) + .addKeyValue("decodedData", message::toStringDecoded) + .log(); + } connectionEventDispatcher.dispatchMessage(demultiplexed.getCapability(), connection, message); } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java index 1302b3765d4..c2ba9571944 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/Framer.java @@ -282,38 +282,35 @@ private MessageData processFrame(final ByteBuf f, final int frameSize) { final int id = idbv.isZero() || idbv.size() == 0 ? 0 : idbv.get(0); // Write message data to ByteBuf, decompressing as necessary - final Bytes data; if (compressionEnabled) { final byte[] compressedMessageData = Arrays.copyOfRange(frameData, 1, frameData.length - pad); final int uncompressedLength = compressor.uncompressedLength(compressedMessageData); if (uncompressedLength >= LENGTH_MAX_MESSAGE_FRAME) { throw error("Message size %s in excess of maximum length.", uncompressedLength); } - Bytes _data; - try { - final byte[] decompressedMessageData = compressor.decompress(compressedMessageData); - _data = Bytes.wrap(decompressedMessageData); - compressionSuccessful = true; - } catch (final FramingException fe) { - if (compressionSuccessful) { - throw fe; - } else { - // OpenEthereum/Parity does not implement EIP-706 - // If failing on the first packet downgrade to uncompressed + + if (!compressionSuccessful) { + // First compressed message: decompress eagerly to validate and handle + // the OpenEthereum/Parity fallback (non-Snappy peer detection via EIP-706) + try { + final byte[] decompressedMessageData = compressor.decompress(compressedMessageData); + compressionSuccessful = true; + return new RawMessage(id, Bytes.wrap(decompressedMessageData)); + } catch (final FramingException fe) { compressionEnabled = false; LOG.debug("Snappy decompression failed: downgrading to uncompressed"); final int messageLength = frameSize - LENGTH_MESSAGE_ID; - _data = Bytes.wrap(frameData, 1, messageLength); + return new RawMessage(id, Bytes.wrap(frameData, 1, messageLength)); } + } else { + // Subsequent messages: store compressed, decompress lazily on getData() + return new RawMessage(id, compressedMessageData); } - data = _data; } else { - // Move data to a ByteBuf final int messageLength = frameSize - LENGTH_MESSAGE_ID; - data = Bytes.wrap(frameData, 1, messageLength); + final Bytes data = Bytes.wrap(frameData, 1, messageLength); + return new RawMessage(id, data); } - - return new RawMessage(id, data); } private void validateMac(final byte[] candidateMac, final byte[] expectedMac) { diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java index c6c613b90b0..d8818ae9ae1 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java @@ -20,14 +20,14 @@ public abstract class AbstractMessageData implements MessageData { - protected final Bytes data; + protected Bytes data; protected AbstractMessageData(final Bytes data) { this.data = data; } @Override - public final int getSize() { + public int getSize() { return data.size(); } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java index b4b0bc20c59..3a0cdc9ad1e 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java @@ -14,19 +14,50 @@ */ package org.hyperledger.besu.ethereum.p2p.rlpx.wire; +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.SnappyCompressor; + import org.apache.tuweni.bytes.Bytes; public final class RawMessage extends AbstractMessageData { + private static final SnappyCompressor compressor = new SnappyCompressor(); + private final int code; + private byte[] compressedData; + /** Constructor for uncompressed messages. */ public RawMessage(final int code, final Bytes data) { super(data); this.code = code; + this.compressedData = null; + } + + /** Constructor for compressed messages — decompression is deferred until getData() is called. */ + public RawMessage(final int code, final byte[] compressedData) { + super(Bytes.EMPTY); + this.code = code; + this.compressedData = compressedData; } @Override public int getCode() { return code; } + + @Override + public Bytes getData() { + if (compressedData != null && data == Bytes.EMPTY) { + data = Bytes.wrap(compressor.decompress(compressedData)); + compressedData = null; + } + return data; + } + + @Override + public int getSize() { + if (compressedData == null) { + return data.size(); + } + return compressor.uncompressedLength(compressedData); + } } From 3f2e28e90d8c0972f5fb4e3b37d409d5b7bc790a Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Tue, 17 Mar 2026 14:31:42 +1000 Subject: [PATCH 3/6] fix: address PR review feedback for deferred Snappy decompression - Restore `data` to final in AbstractMessageData to preserve immutability contract for all subclasses - Use private volatile fields and a `decompressed` flag in RawMessage instead of mutating parent state or using fragile reference equality - Narrow try/catch in TransactionsMessageHandler and NewPooledTransactionHashesMessageHandler to only cover the decode step, avoiding false-positive peer disconnects from processing exceptions Co-Authored-By: Claude Opus 4.6 Signed-off-by: stefan.pingel@consensys.net --- ...PooledTransactionHashesMessageHandler.java | 8 +++++--- .../TransactionsMessageHandler.java | 9 +++++---- .../p2p/rlpx/wire/AbstractMessageData.java | 2 +- .../ethereum/p2p/rlpx/wire/RawMessage.java | 20 ++++++++++++------- 4 files changed, 24 insertions(+), 15 deletions(-) diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java index 37f80851ace..a71e8d3008e 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandler.java @@ -61,11 +61,10 @@ public void exec(final EthMessage message) { if (message.getPeer().isDisconnected()) { return; } + final NewPooledTransactionHashesMessage transactionsMessage; try { - final NewPooledTransactionHashesMessage transactionsMessage = + transactionsMessage = NewPooledTransactionHashesMessage.readFrom(rawMessage, capability); - transactionsMessageProcessor.processNewPooledTransactionHashesMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); } catch (final Exception e) { LOG.debug( "Malformed pooled transaction hashes message received (BREACH_OF_PROTOCOL), disconnecting: {}", @@ -74,7 +73,10 @@ public void exec(final EthMessage message) { message .getPeer() .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + return; } + transactionsMessageProcessor.processNewPooledTransactionHashesMessage( + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); }); } } diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java index 930d99ff55b..c1bed5975b2 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/transactions/TransactionsMessageHandler.java @@ -57,11 +57,9 @@ public void exec(final EthMessage message) { if (message.getPeer().isDisconnected()) { return; } + final TransactionsMessage transactionsMessage; try { - final TransactionsMessage transactionsMessage = - TransactionsMessage.readFrom(rawMessage); - transactionsMessageProcessor.processTransactionsMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + transactionsMessage = TransactionsMessage.readFrom(rawMessage); } catch (final Exception e) { LOG.debug( "Malformed transactions message received (BREACH_OF_PROTOCOL), disconnecting: {}", @@ -70,7 +68,10 @@ public void exec(final EthMessage message) { message .getPeer() .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + return; } + transactionsMessageProcessor.processTransactionsMessage( + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); }); } } diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java index d8818ae9ae1..04ba4e0eb22 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/AbstractMessageData.java @@ -20,7 +20,7 @@ public abstract class AbstractMessageData implements MessageData { - protected Bytes data; + protected final Bytes data; protected AbstractMessageData(final Bytes data) { this.data = data; diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java index 3a0cdc9ad1e..bf6d03decad 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java @@ -23,13 +23,17 @@ public final class RawMessage extends AbstractMessageData { private static final SnappyCompressor compressor = new SnappyCompressor(); private final int code; - private byte[] compressedData; + private final byte[] compressedData; + private volatile boolean decompressed; + private volatile Bytes decompressedData; /** Constructor for uncompressed messages. */ public RawMessage(final int code, final Bytes data) { super(data); this.code = code; this.compressedData = null; + this.decompressed = true; + this.decompressedData = data; } /** Constructor for compressed messages — decompression is deferred until getData() is called. */ @@ -37,6 +41,8 @@ public RawMessage(final int code, final byte[] compressedData) { super(Bytes.EMPTY); this.code = code; this.compressedData = compressedData; + this.decompressed = false; + this.decompressedData = null; } @Override @@ -46,17 +52,17 @@ public int getCode() { @Override public Bytes getData() { - if (compressedData != null && data == Bytes.EMPTY) { - data = Bytes.wrap(compressor.decompress(compressedData)); - compressedData = null; + if (!decompressed) { + decompressedData = Bytes.wrap(compressor.decompress(compressedData)); + decompressed = true; } - return data; + return decompressedData; } @Override public int getSize() { - if (compressedData == null) { - return data.size(); + if (decompressed) { + return decompressedData.size(); } return compressor.uncompressedLength(compressedData); } From 8bc3321a5e2694afe3295dfc5a8197d4d41ef7f3 Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Thu, 19 Mar 2026 14:32:56 +1000 Subject: [PATCH 4/6] fix some things based on comments Signed-off-by: stefan.pingel@consensys.net --- .../eth/manager/EthProtocolManager.java | 20 ++- .../eth/manager/snap/SnapProtocolManager.java | 16 +- .../manager/task/AbstractPeerRequestTask.java | 9 ++ .../p2p/rlpx/connections/netty/DeFramer.java | 6 +- .../ethereum/p2p/rlpx/wire/RawMessage.java | 22 ++- .../p2p/rlpx/wire/RawMessageTest.java | 142 ++++++++++++++++++ 6 files changed, 200 insertions(+), 15 deletions(-) create mode 100644 ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessageTest.java diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java index 86143925ca2..76349a06d01 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java @@ -37,6 +37,7 @@ import org.hyperledger.besu.ethereum.p2p.network.ProtocolManager; import org.hyperledger.besu.ethereum.p2p.rlpx.connections.PeerConnection; import org.hyperledger.besu.ethereum.p2p.rlpx.connections.PeerConnection.PeerNotConnected; +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.FramingException; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Capability; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Message; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; @@ -297,12 +298,12 @@ public void processMessage(final Capability capability, final Message message) { return; } - // This will handle responses - ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); - - // This will handle requests Optional maybeResponseData = Optional.empty(); try { + // This will handle responses + ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); + + // This will handle requests if (EthProtocol.requestIdCompatible(code)) { final Map.Entry requestIdAndEthMessage = ethMessage.getData().unwrapMessageData(); @@ -313,6 +314,17 @@ public void processMessage(final Capability capability, final Message message) { } else { maybeResponseData = ethMessages.dispatch(ethMessage, capability); } + } catch (final FramingException e) { + LOG.atDebug() + .setMessage( + "Failed to decompress message with code {} (BREACH_OF_PROTOCOL), disconnecting: {}, {}") + .addArgument(code) + .addArgument(ethPeer::toString) + .addArgument(e::toString) + .log(); + + ethPeer.disconnect( + DisconnectMessage.DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); } catch (final RLPException e) { LOG.atDebug() .setMessage("Received malformed message {} (BREACH_OF_PROTOCOL), disconnecting: {}, {}") diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java index cd805983460..a31dfcaad0a 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java @@ -24,6 +24,7 @@ import org.hyperledger.besu.ethereum.eth.sync.snapsync.SnapSyncConfiguration; import org.hyperledger.besu.ethereum.p2p.network.ProtocolManager; import org.hyperledger.besu.ethereum.p2p.rlpx.connections.PeerConnection; +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.FramingException; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.AbstractSnapMessageData; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Capability; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.Message; @@ -113,18 +114,25 @@ public void processMessage(final Capability cap, final Message message) { return; } - // This will handle responses - ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); - - // This will handle requests Optional maybeResponseData = Optional.empty(); try { + // This will handle responses + ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); + + // This will handle requests final Map.Entry requestIdAndEthMessage = ethMessage.getData().unwrapMessageData(); maybeResponseData = snapMessages .dispatch(new EthMessage(ethPeer, requestIdAndEthMessage.getValue()), cap) .map(responseData -> responseData.wrapMessageData(requestIdAndEthMessage.getKey())); + } catch (final FramingException e) { + LOG.debug( + "Failed to decompress message with code {} (BREACH_OF_PROTOCOL), disconnecting: {}", + code, + ethPeer, + e); + ethPeer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); } catch (final RLPException e) { LOG.debug( "Received malformed message {} , disconnecting: {}", messageData.getData(), ethPeer, e); diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java index d66e62104af..ed38aa166c2 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java @@ -20,6 +20,7 @@ import org.hyperledger.besu.ethereum.eth.manager.PendingPeerRequest; import org.hyperledger.besu.ethereum.eth.manager.RequestManager; import org.hyperledger.besu.ethereum.eth.manager.exceptions.PeerBreachedProtocolException; +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.FramingException; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; import org.hyperledger.besu.ethereum.rlp.RLPException; @@ -111,6 +112,14 @@ private void handleMessage( promise.complete(r); peer.recordUsefulResponse(); }); + } catch (final FramingException e) { + // Peer sent us data that failed to decompress - disconnect + LOG.debug( + "Disconnecting with BREACH_OF_PROTOCOL due to decompression failure: {}", + peer.getLoggableId(), + e); + peer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + promise.completeExceptionally(new PeerBreachedProtocolException()); } catch (final RLPException e) { // Peer sent us malformed data - disconnect LOG.debug( diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/DeFramer.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/DeFramer.java index 42601f7134d..ab01bef7075 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/DeFramer.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/connections/netty/DeFramer.java @@ -31,6 +31,7 @@ import org.hyperledger.besu.ethereum.p2p.rlpx.wire.CapabilityMultiplexer; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.MessageData; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.PeerInfo; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.RawMessage; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.SubProtocol; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage; import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.HelloMessage; @@ -54,6 +55,7 @@ import io.netty.handler.codec.ByteToMessageDecoder; import io.netty.handler.codec.DecoderException; import io.netty.handler.timeout.IdleStateHandler; +import org.apache.tuweni.bytes.Bytes; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -227,7 +229,9 @@ protected void decode(final ChannelHandlerContext ctx, final ByteBuf in, final L "Message received before HELLO's exchanged (BREACH_OF_PROTOCOL), disconnecting. Peer: {}, Code: {}, Data: {}", expectedPeer.map(Peer::getEnodeURLString).orElse("unknown"), message.getCode(), - message.getData().toString()); + message instanceof RawMessage raw && raw.getCompressedData() != null + ? "snappy compressed data: " + Bytes.wrap(raw.getCompressedData()) + : message.getData().toString()); ctx.writeAndFlush( new OutboundMessage( null, diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java index bf6d03decad..0d9d1d2bf87 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java @@ -23,7 +23,7 @@ public final class RawMessage extends AbstractMessageData { private static final SnappyCompressor compressor = new SnappyCompressor(); private final int code; - private final byte[] compressedData; + private byte[] compressedData; private volatile boolean decompressed; private volatile Bytes decompressedData; @@ -53,17 +53,27 @@ public int getCode() { @Override public Bytes getData() { if (!decompressed) { - decompressedData = Bytes.wrap(compressor.decompress(compressedData)); - decompressed = true; + synchronized (this) { + if (!decompressed) { + decompressedData = Bytes.wrap(compressor.decompress(compressedData)); + decompressed = true; + compressedData = null; + } + } } return decompressedData; } + public byte[] getCompressedData() { + return compressedData; + } + @Override public int getSize() { - if (decompressed) { - return decompressedData.size(); + final byte[] compressed = compressedData; + if (compressed != null) { + return compressor.uncompressedLength(compressed); } - return compressor.uncompressedLength(compressedData); + return decompressedData.size(); } } diff --git a/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessageTest.java b/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessageTest.java new file mode 100644 index 00000000000..94c52fcf495 --- /dev/null +++ b/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessageTest.java @@ -0,0 +1,142 @@ +/* + * Copyright contributors to Hyperledger Besu. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + * + * SPDX-License-Identifier: Apache-2.0 + */ +package org.hyperledger.besu.ethereum.p2p.rlpx.wire; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.assertj.core.api.Assertions.assertThat; + +import org.hyperledger.besu.ethereum.p2p.rlpx.framing.SnappyCompressor; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; + +import org.apache.tuweni.bytes.Bytes; +import org.junit.jupiter.api.Test; + +class RawMessageTest { + + private static final SnappyCompressor compressor = new SnappyCompressor(); + private static final int CODE = 0x10; + private static final byte[] ORIGINAL_DATA = "hello world, this is a test message".getBytes(UTF_8); + private static final byte[] COMPRESSED_DATA = compressor.compress(ORIGINAL_DATA); + + @Test + void uncompressedConstructorReturnsDataImmediately() { + final Bytes data = Bytes.wrap(ORIGINAL_DATA); + final RawMessage message = new RawMessage(CODE, data); + + assertThat(message.getCode()).isEqualTo(CODE); + assertThat(message.getData()).isEqualTo(data); + assertThat(message.getCompressedData()).isNull(); + } + + @Test + void uncompressedConstructorGetSizeIsCorrect() { + final Bytes data = Bytes.wrap(ORIGINAL_DATA); + final RawMessage message = new RawMessage(CODE, data); + + assertThat(message.getSize()).isEqualTo(ORIGINAL_DATA.length); + } + + @Test + void compressedConstructorDefersDecompression() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + assertThat(message.getCompressedData()).isNotNull(); + } + + @Test + void compressedConstructorGetDataDecompressesCorrectly() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + final Bytes result = message.getData(); + + assertThat(result).isEqualTo(Bytes.wrap(ORIGINAL_DATA)); + } + + @Test + void compressedConstructorGetSizeBeforeDecompression() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + assertThat(message.getSize()).isEqualTo(ORIGINAL_DATA.length); + // compressed data should still be present (getSize does not trigger decompression) + assertThat(message.getCompressedData()).isNotNull(); + } + + @Test + void compressedConstructorGetSizeAfterDecompression() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + message.getData(); // trigger decompression + assertThat(message.getSize()).isEqualTo(ORIGINAL_DATA.length); + } + + @Test + void compressedDataIsNulledAfterDecompression() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + message.getData(); + assertThat(message.getCompressedData()).isNull(); + } + + @Test + void getDataIsIdempotent() { + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + + final Bytes first = message.getData(); + final Bytes second = message.getData(); + + assertThat(first).isEqualTo(Bytes.wrap(ORIGINAL_DATA)); + assertThat(first).isSameAs(second); + } + + @Test + void concurrentGetDataDecompressesOnlyOnce() throws Exception { + final int threadCount = 16; + final RawMessage message = new RawMessage(CODE, COMPRESSED_DATA.clone()); + final CyclicBarrier barrier = new CyclicBarrier(threadCount); + final ExecutorService executor = Executors.newFixedThreadPool(threadCount); + + try { + final List> futures = new ArrayList<>(); + for (int i = 0; i < threadCount; i++) { + futures.add( + executor.submit( + () -> { + barrier.await(); + return message.getData(); + })); + } + + Bytes firstResult = null; + for (final Future future : futures) { + final Bytes result = future.get(); + assertThat(result).isEqualTo(Bytes.wrap(ORIGINAL_DATA)); + if (firstResult == null) { + firstResult = result; + } else { + // all threads should get the exact same Bytes instance + assertThat(result).isSameAs(firstResult); + } + } + } finally { + executor.shutdown(); + } + } +} From 4ca652ebbb3b69191669ef2f84d1f50267a6cec4 Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Fri, 20 Mar 2026 16:59:36 +1000 Subject: [PATCH 5/6] review comments Signed-off-by: stefan.pingel@consensys.net --- CHANGELOG.md | 1 + .../ethereum/p2p/rlpx/framing/FramerTest.java | 57 +++++++++++++++++++ 2 files changed, 58 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7ad0f0eba69..cd5b000c967 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -59,6 +59,7 @@ are provided with different values, using input as per the execution-apis spec i - Fix addMod case with 256bit moduluses [#10001](https://github.com/besu-eth/besu/pull/10001) - Performance improvements on MOD variant instructions while converting from byte[] to longs [#9976](https://github.com/besu-eth/besu/pull/9976) - Implement DIV and SDIV with long limbs [#9923](https://github.com/besu-eth/besu/pull/9923) +- Defer Snappy decompression of inbound P2P messages from the Netty I/O thread to the worker thread, reducing memory held in the transaction worker queue to compressed size [#10048](https://github.com/besu-eth/besu/pull/10048) ## 26.2.0 diff --git a/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/FramerTest.java b/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/FramerTest.java index 403a0d48cc9..182fabe9a16 100644 --- a/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/FramerTest.java +++ b/ethereum/p2p/src/test/java/org/hyperledger/besu/ethereum/p2p/rlpx/framing/FramerTest.java @@ -271,6 +271,63 @@ public void compressionWorks() { assertThat(receivingFramer.isCompressionSuccessful()).isTrue(); } + @Test + public void firstMessageDecompressesEagerlySubsequentMessagesDeferred() { + final HandshakeSecrets sendSecrets = + new HandshakeSecrets( + Bytes.fromHexString( + "0x75b3ee95adff0c529a05efd7612aa1dbe5057eb9facdde0dfc837ad143da1d43") + .toArray(), + Bytes.fromHexString( + "0x030dfd1566f4800c4842c177f7d476b64ae2b99a2aa0ab5600aa2f41a8710575") + .toArray(), + Bytes.fromHexString( + "0xc9d3385b1588a5969cba312f8c29bedb4cb9d56ec0cf825436addc1ec644f1d6") + .toArray()); + final HandshakeSecrets recvSecrets = + new HandshakeSecrets( + Bytes.fromHexString( + "0x75b3ee95adff0c529a05efd7612aa1dbe5057eb9facdde0dfc837ad143da1d43") + .toArray(), + Bytes.fromHexString( + "0x030dfd1566f4800c4842c177f7d476b64ae2b99a2aa0ab5600aa2f41a8710575") + .toArray(), + Bytes.fromHexString( + "0xc9d3385b1588a5969cba312f8c29bedb4cb9d56ec0cf825436addc1ec644f1d6") + .toArray()); + final Framer sendingFramer = new Framer(sendSecrets); + final Framer receivingFramer = new Framer(recvSecrets); + sendingFramer.enableCompression(); + receivingFramer.enableCompression(); + + final Bytes payload1 = DisconnectMessage.create(DisconnectReason.TIMEOUT).getData(); + final Bytes payload2 = DisconnectMessage.create(DisconnectReason.BREACH_OF_PROTOCOL).getData(); + + // Frame two messages + final ByteBuf out1 = Unpooled.buffer(); + sendingFramer.frame(new RawMessage(0x01, payload1), out1); + final ByteBuf out2 = Unpooled.buffer(); + sendingFramer.frame(new RawMessage(0x01, payload2), out2); + + // First message: eagerly decompressed (compressionSuccessful was false) + final MessageData msg1 = receivingFramer.deframe(out1); + assertThat(msg1).isNotNull(); + final RawMessage raw1 = (RawMessage) msg1; + assertThat(raw1.getCompressedData()).isNull(); + assertThat(raw1.getData()).isEqualTo(payload1); + assertThat(receivingFramer.isCompressionSuccessful()).isTrue(); + + // Second message: deferred decompression (compressionSuccessful was true) + final MessageData msg2 = receivingFramer.deframe(out2); + assertThat(msg2).isNotNull(); + final RawMessage raw2 = (RawMessage) msg2; + assertThat(raw2.getCompressedData()).isNotNull(); + // getData() triggers lazy decompression + assertThat(raw2.getData()).isEqualTo(payload2); + // Compressed data released after decompression + assertThat(raw2.getCompressedData()).isNull(); + } + private HandshakeSecrets secretsFrom(final JsonNode td, final boolean swap) { final byte[] aes = decodeHexDump(td.get("aes_secret").asText()); final byte[] mac = decodeHexDump(td.get("mac_secret").asText()); From ca61602b0c9b764f844335b9985a228f367725f4 Mon Sep 17 00:00:00 2001 From: "stefan.pingel@consensys.net" Date: Mon, 23 Mar 2026 15:40:27 +1000 Subject: [PATCH 6/6] Address PR review comments for deferred Snappy decompression - Make compressedData volatile in RawMessage for safe access from getCompressedData() - Standardize disconnect debug log messages across EthProtocolManager, SnapProtocolManager, and AbstractPeerRequestTask - Move AbstractSnapMessageData.create() inside try-catch in SnapProtocolManager so FramingException from decompression is properly caught and peer disconnected - Add disconnect-on-decompression-failure tests for EthProtocolManager, SnapProtocolManager, and NewPooledTransactionHashesMessageHandler - Add changelog entry Co-Authored-By: Claude Opus 4.6 Signed-off-by: stefan.pingel@consensys.net --- CHANGELOG.md | 1 + .../eth/manager/EthProtocolManager.java | 7 +- .../eth/manager/snap/SnapProtocolManager.java | 34 +++---- .../manager/task/AbstractPeerRequestTask.java | 10 ++- .../eth/manager/EthProtocolManagerTest.java | 22 +++++ .../manager/snap/SnapProtocolManagerTest.java | 90 +++++++++++++++++++ ...edTransactionHashesMessageHandlerTest.java | 65 ++++++++++++++ .../ethereum/p2p/rlpx/wire/RawMessage.java | 2 +- 8 files changed, 205 insertions(+), 26 deletions(-) create mode 100644 ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManagerTest.java create mode 100644 ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandlerTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index cd5b000c967..676f53cfbde 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,6 +36,7 @@ are provided with different values, using input as per the execution-apis spec i - Wait for peers before starting chain download. Prevents an OutOfMemory (OOM) error when the node has zero peers [#9979](https://github.com/hyperledger/besu/pull/9979) ### Additions and Improvements +- Defer Snappy decompression of inbound P2P messages from the Netty I/O thread to the worker thread, reducing memory held in the transaction worker queue to compressed size [#10048](https://github.com/besu-eth/besu/pull/10048) - Add IPv6 dual-stack support for DiscV5 peer discovery (enabled via `--Xv5-discovery-enabled`): new `--p2p-host-ipv6`, `--p2p-interface-ipv6`, and `--p2p-port-ipv6` CLI options enable a second UDP discovery socket; `--p2p-ipv6-outbound-enabled` controls whether IPv6 is preferred for outbound connections when a peer advertises both address families [#9763](https://github.com/hyperledger/besu/pull/9763); RLPx now also binds a second TCP socket on the IPv6 interface so IPv6-only peers can establish connections [#9873](https://github.com/hyperledger/besu/pull/9873) - `--net-restrict` now supports IPv6 CIDR notation (e.g. `fd00::/64`) in addition to IPv4, enabling subnet-based peer filtering in IPv6 and dual-stack deployments [#10028](https://github.com/besu-eth/besu/pull/10028) - Stop EngineQosTimer as part of shutdown [#9903](https://github.com/hyperledger/besu/pull/9903) diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java index 5bdc4179967..85d2c6df497 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManager.java @@ -317,11 +317,10 @@ public void processMessage(final Capability capability, final Message message) { } } catch (final FramingException e) { LOG.atDebug() - .setMessage( - "Failed to decompress message with code {} (BREACH_OF_PROTOCOL), disconnecting: {}, {}") + .setMessage("Disconnecting peer {} due to decompression failure for message code {}") + .addArgument(ethPeer::getLoggableId) .addArgument(code) - .addArgument(ethPeer::toString) - .addArgument(e::toString) + .setCause(e) .log(); ethPeer.disconnect( diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java index b3f9f4c0997..f95ac54f15f 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManager.java @@ -95,8 +95,7 @@ public void awaitStop() throws InterruptedException {} */ @Override public void processMessage(final Capability cap, final Message message) { - final MessageData messageData = AbstractSnapMessageData.create(message); - final int code = messageData.getCode(); + final int code = message.getData().getCode(); LOG.trace("Process snap message {}, {}", cap, code); final EthPeer ethPeer = ethPeers.peer(message.getConnection()); if (ethPeer == null) { @@ -104,40 +103,41 @@ public void processMessage(final Capability cap, final Message message) { "Ignoring message received from unknown peer connection: {}", message.getConnection()); return; } - final EthMessage ethMessage = new EthMessage(ethPeer, messageData); + + final EthMessage ethMessage = new EthMessage(ethPeer, message.getData()); if (!ethPeer.validateReceivedMessage(ethMessage, getSupportedProtocol())) { - LOG.debug( - "Unsolicited message {} received from, disconnecting: {}", - ethMessage.getData().getCode(), - ethPeer); + LOG.debug("Unsolicited message {} received from, disconnecting: {}", code, ethPeer); ethPeer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_UNSOLICITED_MESSAGE_RECEIVED); return; } Optional maybeResponseData = Optional.empty(); try { + final MessageData messageData = AbstractSnapMessageData.create(message); + final EthMessage decodedEthMessage = new EthMessage(ethPeer, messageData); + // This will handle responses - ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); + ethPeers.dispatchMessage(ethPeer, decodedEthMessage, getSupportedProtocol()); // This will handle requests final Map.Entry requestIdAndEthMessage = - ethMessage.getData().unwrapMessageData(); + decodedEthMessage.getData().unwrapMessageData(); maybeResponseData = snapMessages .dispatch(new EthMessage(ethPeer, requestIdAndEthMessage.getValue()), cap) .map(responseData -> responseData.wrapMessageData(requestIdAndEthMessage.getKey())); } catch (final FramingException e) { - LOG.debug( - "Failed to decompress message with code {} (BREACH_OF_PROTOCOL), disconnecting: {}", - code, - ethPeer, - e); + LOG.atDebug() + .setMessage("Disconnecting peer {} due to decompression failure for message code {}") + .addArgument(ethPeer::getLoggableId) + .addArgument(code) + .setCause(e) + .log(); ethPeer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); } catch (final RLPException e) { LOG.debug( - "Received malformed message code={} data={} (BREACH_OF_PROTOCOL), disconnecting: {}", - messageData.getCode(), - messageData.getData(), + "Received malformed message code={} (BREACH_OF_PROTOCOL), disconnecting: {}", + code, ethPeer, e); ethPeer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); diff --git a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java index ed38aa166c2..3bb2220f111 100644 --- a/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java +++ b/ethereum/eth/src/main/java/org/hyperledger/besu/ethereum/eth/manager/task/AbstractPeerRequestTask.java @@ -114,10 +114,12 @@ private void handleMessage( }); } catch (final FramingException e) { // Peer sent us data that failed to decompress - disconnect - LOG.debug( - "Disconnecting with BREACH_OF_PROTOCOL due to decompression failure: {}", - peer.getLoggableId(), - e); + LOG.atDebug() + .setMessage("Disconnecting peer {} due to decompression failure for message code {}") + .addArgument(peer::getLoggableId) + .addArgument(message.getCode()) + .setCause(e) + .log(); peer.disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); promise.completeExceptionally(new PeerBreachedProtocolException()); } catch (final RLPException e) { diff --git a/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManagerTest.java b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManagerTest.java index 20525bb55f2..fd9890668f9 100644 --- a/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManagerTest.java +++ b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/EthProtocolManagerTest.java @@ -133,6 +133,28 @@ public static void setup() { assertThat(blockchainSetupUtil.getMaxBlockNumber()).isGreaterThanOrEqualTo(20L); } + @Test + public void disconnectOnDecompressionFailure() { + try (final EthProtocolManager ethManager = + EthProtocolManagerTestBuilder.builder() + .setProtocolSchedule(protocolSchedule) + .setBlockchain(blockchain) + .setEthScheduler(new DeterministicEthScheduler(() -> false)) + .setWorldStateArchive(protocolContext.getWorldStateArchive()) + .setTransactionPool(transactionPool) + .setEthereumWireProtocolConfiguration(EthProtocolConfiguration.DEFAULT) + .build()) { + // Create a RawMessage with invalid compressed data that will throw FramingException + final MessageData messageData = + new RawMessage(EthProtocolMessages.GET_BLOCK_HEADERS, new byte[] {0x01, 0x02, 0x03}); + final MockPeerConnection peer = setupPeer(ethManager, (cap, msg, conn) -> {}); + ethManager.processMessage(EthProtocol.ETH68, new DefaultMessage(peer, messageData)); + assertThat(peer.isDisconnected()).isTrue(); + assertThat(peer.getDisconnectReason()) + .contains(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } + } + @Test public void handleMalformedRequestIdMessage() { try (final EthProtocolManager ethManager = diff --git a/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManagerTest.java b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManagerTest.java new file mode 100644 index 00000000000..8901ba4fd04 --- /dev/null +++ b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/manager/snap/SnapProtocolManagerTest.java @@ -0,0 +1,90 @@ +/* + * Copyright contributors to Hyperledger Besu. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + * + * SPDX-License-Identifier: Apache-2.0 + */ +package org.hyperledger.besu.ethereum.eth.manager.snap; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.when; + +import org.hyperledger.besu.ethereum.ProtocolContext; +import org.hyperledger.besu.ethereum.core.Synchronizer; +import org.hyperledger.besu.ethereum.eth.SnapProtocol; +import org.hyperledger.besu.ethereum.eth.manager.EthMessages; +import org.hyperledger.besu.ethereum.eth.manager.EthPeer; +import org.hyperledger.besu.ethereum.eth.manager.EthPeers; +import org.hyperledger.besu.ethereum.eth.manager.MockPeerConnection; +import org.hyperledger.besu.ethereum.eth.sync.snapsync.SnapSyncConfiguration; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.DefaultMessage; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.RawMessage; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; +import org.hyperledger.besu.ethereum.worldstate.WorldStateStorageCoordinator; + +import java.util.Collections; +import java.util.HashSet; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.junit.jupiter.MockitoSettings; +import org.mockito.quality.Strictness; + +@ExtendWith(MockitoExtension.class) +@MockitoSettings(strictness = Strictness.LENIENT) +class SnapProtocolManagerTest { + + @Mock private WorldStateStorageCoordinator worldStateStorageCoordinator; + @Mock private SnapSyncConfiguration snapConfig; + @Mock private EthPeers ethPeers; + @Mock private EthMessages snapMessages; + @Mock private ProtocolContext protocolContext; + @Mock private Synchronizer synchronizer; + @Mock private EthPeer ethPeer; + + private SnapProtocolManager snapProtocolManager; + + @BeforeEach + void setUp() { + when(snapConfig.isSnapServerEnabled()).thenReturn(false); + snapProtocolManager = + new SnapProtocolManager( + worldStateStorageCoordinator, + snapConfig, + ethPeers, + snapMessages, + protocolContext, + synchronizer); + } + + @Test + void disconnectsPeerOnDecompressionFailure() { + final MockPeerConnection peerConnection = + new MockPeerConnection( + new HashSet<>(Collections.singletonList(SnapProtocol.SNAP1)), (cap, msg, conn) -> {}); + when(ethPeers.peer(peerConnection)).thenReturn(ethPeer); + when(ethPeer.validateReceivedMessage(any(), any())).thenReturn(true); + + // Create a RawMessage with invalid compressed data that will throw FramingException + final RawMessage badMessage = new RawMessage(0x00, new byte[] {0x01, 0x02, 0x03}); + snapProtocolManager.processMessage( + SnapProtocol.SNAP1, new DefaultMessage(peerConnection, badMessage)); + + assertThat(peerConnection.isDisconnected()).isFalse(); + // ethPeer (mock) receives the disconnect call + org.mockito.Mockito.verify(ethPeer) + .disconnect(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); + } +} diff --git a/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandlerTest.java b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandlerTest.java new file mode 100644 index 00000000000..ab24cbc1443 --- /dev/null +++ b/ethereum/eth/src/test/java/org/hyperledger/besu/ethereum/eth/transactions/NewPooledTransactionHashesMessageHandlerTest.java @@ -0,0 +1,65 @@ +/* + * Copyright contributors to Hyperledger Besu. + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + * + * SPDX-License-Identifier: Apache-2.0 + */ +package org.hyperledger.besu.ethereum.eth.transactions; + +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import org.hyperledger.besu.ethereum.eth.EthProtocol; +import org.hyperledger.besu.ethereum.eth.manager.EthMessage; +import org.hyperledger.besu.ethereum.eth.manager.EthPeer; +import org.hyperledger.besu.ethereum.p2p.rlpx.connections.PeerConnection; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.RawMessage; +import org.hyperledger.besu.ethereum.p2p.rlpx.wire.messages.DisconnectMessage.DisconnectReason; +import org.hyperledger.besu.testutil.DeterministicEthScheduler; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +class NewPooledTransactionHashesMessageHandlerTest { + + @Mock private NewPooledTransactionHashesMessageProcessor processor; + @Mock private EthPeer peer; + @Mock private PeerConnection peerConnection; + + private NewPooledTransactionHashesMessageHandler handler; + + @BeforeEach + void setUp() { + final DeterministicEthScheduler scheduler = new DeterministicEthScheduler(); + handler = new NewPooledTransactionHashesMessageHandler(scheduler, processor, 300); + handler.setEnabled(); + when(peer.getConnection()).thenReturn(peerConnection); + when(peerConnection.capability(EthProtocol.NAME)).thenReturn(EthProtocol.ETH68); + } + + @Test + void disconnectsPeerOnDecompressionFailure() { + // Create a RawMessage with invalid compressed data that will throw FramingException + final RawMessage badMessage = new RawMessage(0x08, new byte[] {0x01, 0x02, 0x03}); + final EthMessage ethMessage = new EthMessage(peer, badMessage); + + handler.exec(ethMessage); + + verify(peer).disconnect(eq(DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED)); + verifyNoInteractions(processor); + } +} diff --git a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java index 0d9d1d2bf87..7d6b3600229 100644 --- a/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java +++ b/ethereum/p2p/src/main/java/org/hyperledger/besu/ethereum/p2p/rlpx/wire/RawMessage.java @@ -23,7 +23,7 @@ public final class RawMessage extends AbstractMessageData { private static final SnappyCompressor compressor = new SnappyCompressor(); private final int code; - private byte[] compressedData; + private volatile byte[] compressedData; private volatile boolean decompressed; private volatile Bytes decompressedData;