diff --git a/CHANGELOG.md b/CHANGELOG.md index 7ad0f0eba69..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) @@ -59,6 +60,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/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 080e1f97366..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 @@ -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; @@ -298,12 +299,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(); @@ -314,6 +315,16 @@ public void processMessage(final Capability capability, final Message message) { } else { maybeResponseData = ethMessages.dispatch(ethMessage, capability); } + } catch (final FramingException e) { + LOG.atDebug() + .setMessage("Disconnecting peer {} due to decompression failure for message code {}") + .addArgument(ethPeer::getLoggableId) + .addArgument(code) + .setCause(e) + .log(); + + ethPeer.disconnect( + DisconnectMessage.DisconnectReason.BREACH_OF_PROTOCOL_MALFORMED_MESSAGE_RECEIVED); } catch (final RLPException e) { LOG.atDebug() .setMessage( 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 0b42c35a578..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 @@ -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; @@ -94,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) { @@ -103,33 +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; } - // This will handle responses - ethPeers.dispatchMessage(ethPeer, ethMessage, getSupportedProtocol()); - - // This will handle requests 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, 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.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 d66e62104af..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 @@ -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,16 @@ private void handleMessage( promise.complete(r); peer.recordUsefulResponse(); }); + } catch (final FramingException e) { + // Peer sent us data that failed to decompress - disconnect + 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) { // Peer sent us malformed data - disconnect LOG.debug( 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..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 @@ -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,30 @@ 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( - () -> - transactionsMessageProcessor.processNewPooledTransactionHashesMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + final NewPooledTransactionHashesMessage transactionsMessage; + try { + transactionsMessage = + NewPooledTransactionHashesMessage.readFrom(rawMessage, capability); + } 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); + 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 b6bbf442dab..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 @@ -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,29 @@ 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( - () -> - transactionsMessageProcessor.processTransactionsMessage( - message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive)); + () -> { + if (message.getPeer().isDisconnected()) { + return; + } + final TransactionsMessage transactionsMessage; + try { + transactionsMessage = TransactionsMessage.readFrom(rawMessage); + } 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); + return; + } + transactionsMessageProcessor.processTransactionsMessage( + message.getPeer(), transactionsMessage, startedAt, txMsgKeepAlive); + }); } } 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/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/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/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..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 @@ -27,7 +27,7 @@ protected AbstractMessageData(final Bytes 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..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 @@ -14,19 +14,66 @@ */ 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 volatile 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. */ + public RawMessage(final int code, final byte[] compressedData) { + super(Bytes.EMPTY); + this.code = code; + this.compressedData = compressedData; + this.decompressed = false; + this.decompressedData = null; } @Override public int getCode() { return code; } + + @Override + public Bytes getData() { + if (!decompressed) { + 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() { + final byte[] compressed = compressedData; + if (compressed != null) { + return compressor.uncompressedLength(compressed); + } + return decompressedData.size(); + } } 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()); 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(); + } + } +}