From 9a5d1b6595447150384567e3f6278d9ca33025d7 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 14 Sep 2026 18:12:55 +0800 Subject: [PATCH 1/7] fix the bug of update lastactivetime with invalid inv --- .../tron/core/net/P2pEventHandlerImpl.java | 1 - .../net/messagehandler/BlockMsgHandler.java | 73 ++-- .../messagehandler/InventoryMsgHandler.java | 7 - .../tron/core/net/peer/PeerConnection.java | 9 +- .../tron/core/net/service/adv/AdvService.java | 58 ++- .../service/fetchblock/FetchBlockService.java | 91 +++-- .../core/net/service/sync/SyncService.java | 12 +- .../tron/core/net/PeerBlockTestSupport.java | 44 +++ .../messagehandler/BlockContributionTest.java | 320 ++++++++++++++++ .../tron/core/net/peer/PeerBlockIdleTest.java | 42 +++ .../fetchblock/FetchBlockRetryTest.java | 357 ++++++++++++++++++ .../sync/SyncBlockContributionTest.java | 115 ++++++ 12 files changed, 1048 insertions(+), 81 deletions(-) create mode 100644 framework/src/test/java/org/tron/core/net/PeerBlockTestSupport.java create mode 100644 framework/src/test/java/org/tron/core/net/messagehandler/BlockContributionTest.java create mode 100644 framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java create mode 100644 framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java create mode 100644 framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java diff --git a/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java b/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java index 9dd950ae57b..510bc49dcf2 100644 --- a/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java +++ b/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java @@ -252,7 +252,6 @@ private void updateLastInteractiveTime(PeerConnection peer, TronMessage msg) { switch (type) { case SYNC_BLOCK_CHAIN: case BLOCK_CHAIN_INVENTORY: - case BLOCK: flag = true; break; case FETCH_INV_DATA: diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index 452209d575f..b30b9f2a2ae 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -9,6 +9,7 @@ import org.springframework.stereotype.Component; import org.tron.common.prometheus.MetricKeys; import org.tron.common.prometheus.Metrics; +import org.tron.common.utils.Sha256Hash; import org.tron.core.Constant; import org.tron.core.capsule.BlockCapsule; import org.tron.core.capsule.BlockCapsule.BlockId; @@ -77,6 +78,12 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep check(peer, blockMessage); } + if (blockCapsule.getNum() <= 0 + || blockCapsule.getParentHashStr().size() != Sha256Hash.LENGTH + || new BlockId(blockCapsule.getParentHash()).getNum() != blockCapsule.getNum() - 1) { + throw new P2pException(TypeEnum.BAD_BLOCK, "block number does not follow parent"); + } + blockMessage.sanitize(); if (peer.getSyncBlockRequested().containsKey(blockId)) { @@ -89,7 +96,12 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep if (peer.isRelayPeer()) { peer.getAdvInvSpread().put(item, now); } - Long time = peer.getAdvInvRequest().remove(item); + Long time = peer.getAdvInvRequest().get(item); + long interval = blockId.getNum() - tronNetDelegate.getHeadBlockId().getNum(); + BlockResult result = processBlock(peer, blockMessage.getBlockCapsule()); + if (result == BlockResult.ACCEPTED || result == BlockResult.SYNC_REQUIRED) { + peer.setBlockRcvTime(System.currentTimeMillis()); + } if (null != time) { MetricsUtil.histogramUpdateUnCheck(MetricsKey.NET_LATENCY_FETCH_BLOCK + peer.getInetAddress(), now - time); @@ -98,9 +110,6 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep } Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_RECEIVE_DELAY, (now - blockMessage.getBlockCapsule().getTimeStamp()) / Metrics.MILLISECONDS_PER_SECOND); - fetchBlockService.blockFetchSuccess(blockId); - long interval = blockId.getNum() - tronNetDelegate.getHeadBlockId().getNum(); - processBlock(peer, blockMessage.getBlockCapsule()); logger.info( "Receive block/interval {}/{} from {} fetch/delay {}/{}ms, " + "txs/process {}/{}ms, witness: {}", @@ -125,46 +134,60 @@ private void check(PeerConnection peer, BlockMessage msg) throws P2pException { } } - private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pException { + private BlockResult processBlock(PeerConnection peer, BlockCapsule block) throws P2pException { BlockId blockId = block.getBlockId(); - boolean flag = tronNetDelegate.validBlock(block); - if (!flag) { - logger.warn("Receive a bad block from {}, {}, {}", + boolean activeWitness = tronNetDelegate.validBlock(block); + // Keep the request until validation succeeds so disconnect can retry invalid data. + fetchBlockService.blockFetchSuccess(blockId); + peer.getAdvInvRequest().remove(new Item(blockId, InventoryType.BLOCK)); + if (!activeWitness) { + logger.warn("Receive a block from an inactive witness, peer {}, block {}, witness {}", peer.getInetSocketAddress(), blockId.getString(), Hex.toHexString(block.getWitnessAddress().toByteArray())); - return; + syncService.startSync(peer); + return BlockResult.STATE_FAILED; + } + + peer.setLastInteractiveTime(System.currentTimeMillis()); + long headNum = tronNetDelegate.getHeadBlockId().getNum(); + if (block.getNum() < headNum || tronNetDelegate.containBlock(blockId)) { + logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); + return BlockResult.IGNORED; } if (!tronNetDelegate.containBlock(block.getParentBlockId())) { logger.warn("Get unlink block {} from {}, head is {}", blockId.getString(), peer.getInetAddress(), tronNetDelegate.getHeadBlockId().getString()); syncService.startSync(peer); - return; - } - - long headNum = tronNetDelegate.getHeadBlockId().getNum(); - if (block.getNum() < headNum) { - logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); - return; + return BlockResult.SYNC_REQUIRED; } broadcast(new BlockMessage(block)); try { tronNetDelegate.processBlock(block, false); - peer.setBlockRcvTime(System.currentTimeMillis()); - witnessProductBlockService.validWitnessProductTwoBlock(block); - - Item item = new Item(blockId, InventoryType.BLOCK); - tronNetDelegate.getActivePeer().forEach(p -> { - if (p.getAdvInvReceive().getIfPresent(item) != null) { - p.setBlockBothHave(blockId); - } - }); } catch (Exception e) { logger.warn("Process adv block {} from peer {} failed. reason: {}", blockId, peer.getInetAddress(), e.getMessage()); + syncService.startSync(peer); + return BlockResult.STATE_FAILED; } + if (tronNetDelegate.isHitDown()) { + return BlockResult.IGNORED; + } + witnessProductBlockService.validWitnessProductTwoBlock(block); + + Item item = new Item(blockId, InventoryType.BLOCK); + tronNetDelegate.getActivePeer().forEach(p -> { + if (p.getAdvInvReceive().getIfPresent(item) != null) { + p.setBlockBothHave(blockId); + } + }); + return BlockResult.ACCEPTED; + } + + private enum BlockResult { + ACCEPTED, SYNC_REQUIRED, IGNORED, STATE_FAILED } private void broadcast(BlockMessage blockMessage) { diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java index f96f7f0b0ff..e83563ecfd7 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java @@ -6,7 +6,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.tron.common.utils.Sha256Hash; -import org.tron.core.capsule.BlockCapsule.BlockId; import org.tron.core.config.args.Args; import org.tron.core.exception.P2pException; import org.tron.core.exception.P2pException.TypeEnum; @@ -44,12 +43,6 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep Item item = new Item(id, type); peer.getAdvInvReceive().put(item, System.currentTimeMillis()); advService.addInv(item); - if (type.equals(InventoryType.BLOCK) && peer.getAdvInvSpread().getIfPresent(item) == null) { - long headNum = tronNetDelegate.getHeadBlockId().getNum(); - if (new BlockId(id).getNum() > headNum) { - peer.setLastInteractiveTime(System.currentTimeMillis()); - } - } } } diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index 7d7457cf2fc..25ad47ce73e 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -10,7 +10,6 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.util.Deque; -import java.util.HashSet; import java.util.List; import java.util.Locale; import java.util.Map; @@ -158,7 +157,7 @@ public class PeerConnection { private volatile Pair, Long> syncChainRequested = null; @Setter @Getter - private Set syncBlockInProcess = new HashSet<>(); + private Set syncBlockInProcess = ConcurrentHashMap.newKeySet(); @Setter @Getter private volatile boolean needSyncFromPeer = true; @@ -193,6 +192,12 @@ public boolean isIdle() { return advInvRequest.isEmpty() && isSyncIdle(); } + public boolean isBlockIdle() { + return advInvRequest.keySet().stream() + .noneMatch(item -> item.getType() == Protocol.Inventory.InventoryType.BLOCK) + && isSyncIdle() && syncBlockInProcess.isEmpty(); + } + public boolean isSyncIdle() { return syncBlockRequested.isEmpty() && syncChainRequested == null; } diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 2241ce8d1ce..0a6431a938c 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -8,11 +8,13 @@ import com.google.common.cache.CacheBuilder; import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; import java.util.Comparator; import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map.Entry; +import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -26,6 +28,7 @@ import org.tron.common.utils.Sha256Hash; import org.tron.common.utils.Time; import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.core.config.Parameter.NetConstants; import org.tron.core.config.args.Args; import org.tron.core.net.TronNetDelegate; import org.tron.core.net.message.adv.BlockMessage; @@ -251,13 +254,32 @@ public void fastForward(BlockMessage msg) { */ public void onDisconnect(PeerConnection peer) { + fetchBlockService.onDisconnect(peer); if (!peer.getAdvInvRequest().isEmpty()) { peer.getAdvInvRequest().keySet().forEach(item -> { - if (tronNetDelegate.getActivePeer().stream() - .anyMatch(p -> !p.equals(peer) && p.getAdvInvReceive().getIfPresent(item) != null)) { - invToFetch.put(item, System.currentTimeMillis()); - } else { - invToFetchCache.invalidate(item); + synchronized (this) { + Collection peers = tronNetDelegate.getActivePeer().stream() + .filter(p -> !p.equals(peer) && !p.isDisconnect()).collect(Collectors.toList()); + if (item.getType() == InventoryType.BLOCK + && (blockCache.getIfPresent(item) != null + || tronNetDelegate.containBlock(new BlockId(item.getHash())))) { + return; + } + if (item.getType() == InventoryType.BLOCK) { + Optional pending = peers.stream() + .filter(p -> p.getAdvInvRequest().containsKey(item)).findFirst(); + if (pending.isPresent()) { + fetchBlockService.fetchBlock(Collections.singletonList(item.getHash()), + pending.get()); + return; + } + } + if (peers.stream().anyMatch(p -> p.getAdvInvReceive().getIfPresent(item) != null)) { + invToFetch.put(item, System.currentTimeMillis()); + } else { + invToFetch.remove(item); + invToFetchCache.invalidate(item); + } } }); } @@ -269,7 +291,9 @@ public void onDisconnect(PeerConnection peer) { private void consumerInvToFetch() { Collection peers = tronNetDelegate.getActivePeer().stream() - .filter(peer -> peer.isIdle()) + .filter(peer -> !peer.isDisconnect()) + .collect(Collectors.toList()); + Collection trxPeers = peers.stream().filter(PeerConnection::isIdle) .collect(Collectors.toList()); InvSender invSender = new InvSender(); synchronized (this) { @@ -278,14 +302,30 @@ private void consumerInvToFetch() { } long now = System.currentTimeMillis(); invToFetch.forEach((item, time) -> { - if (time < now - TIMEOUT) { + long timeout = item.getType() == InventoryType.BLOCK ? NetConstants.ADV_TIME_OUT : TIMEOUT; + if (time < now - timeout) { logger.info("This obj is too late to fetch, type: {} hash: {}", item.getType(), item.getHash()); invToFetch.remove(item); invToFetchCache.invalidate(item); return; } - peers.stream().filter(peer -> { + if (item.getType() == InventoryType.BLOCK) { + if (blockCache.getIfPresent(item) != null + || tronNetDelegate.containBlock(new BlockId(item.getHash())) + || peers.stream().anyMatch(peer -> peer.getAdvInvRequest().containsKey(item))) { + invToFetch.remove(item); + return; + } + fetchBlockService.selectBlockPeer(peers, item, now).ifPresent(peer -> { + if (peer.checkAndPutAdvInvRequest(item, now)) { + invSender.add(item, peer); + invToFetch.remove(item); + } + }); + return; + } + trxPeers.stream().filter(peer -> { Long t = peer.getAdvInvReceive().getIfPresent(item); return t != null && now - t < TIMEOUT && invSender.getSize(peer) < MAX_TRX_FETCH_PER_PEER; }).sorted(Comparator.comparingInt(peer -> invSender.getSize(peer))) @@ -387,8 +427,8 @@ void sendFetch() { send.forEach((peer, ids) -> ids.forEach((key, value) -> { if (key.equals(InventoryType.BLOCK)) { value.sort(Comparator.comparingLong(value1 -> new BlockId(value1).getNum())); - peer.sendMessage(new FetchInvDataMessage(value, key)); fetchBlockService.fetchBlock(value, peer); + peer.sendMessage(new FetchInvDataMessage(value, key)); } else { peer.sendMessage(new FetchInvDataMessage(value, key)); } diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index bda2646abbc..7a6b74bf9cc 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -1,5 +1,6 @@ package org.tron.core.net.service.fetchblock; +import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.List; @@ -7,7 +8,6 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import lombok.Getter; -import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; @@ -16,6 +16,7 @@ import org.tron.common.utils.Sha256Hash; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.BlockCapsule; +import org.tron.core.config.Parameter.NetConstants; import org.tron.core.metrics.MetricsKey; import org.tron.core.metrics.MetricsUtil; import org.tron.core.net.TronNetDelegate; @@ -34,7 +35,7 @@ public class FetchBlockService { @Autowired private ChainBaseManager chainBaseManager; - private FetchBlockInfo fetchBlockInfo = null; + private volatile FetchBlockInfo fetchBlockInfo; private final long fetchTimeOut = CommonParameter.getInstance().fetchBlockTimeout; @@ -59,63 +60,86 @@ public void close() { ExecutorServiceManager.shutdownAndAwaitTermination(fetchBlockWorkerExecutor, esName); } - public void fetchBlock(List sha256HashList, PeerConnection peer) { + public synchronized void fetchBlock(List sha256HashList, PeerConnection peer) { if (sha256HashList.size() > 0) { logger.info("Begin fetch block {} from {}", new BlockCapsule.BlockId(sha256HashList.get(0)).getString(), peer.getInetAddress()); } - if (null != fetchBlockInfo) { + long headNum = chainBaseManager.getHeadBlockNum(); + if (fetchBlockInfo != null + && new BlockCapsule.BlockId(fetchBlockInfo.getHash()).getNum() > headNum) { return; } + fetchBlockInfo = null; sha256HashList.stream().filter(sha256Hash -> new BlockCapsule.BlockId(sha256Hash).getNum() - == chainBaseManager.getHeadBlockNum() + 1) + == headNum + 1) .findFirst().ifPresent(sha256Hash -> { - long now = System.currentTimeMillis(); - fetchBlockInfo = new FetchBlockInfo(sha256Hash, peer, now); - logger.info("Set fetchBlockInfo, block: {}, peer: {}, time: {}", sha256Hash, - peer.getInetAddress(), now); + Long requestTime = peer.getAdvInvRequest().get(new Item(sha256Hash, InventoryType.BLOCK)); + if (requestTime != null) { + fetchBlockInfo = new FetchBlockInfo(sha256Hash, peer, requestTime); + logger.info("Set fetchBlockInfo, block: {}, peer: {}, time: {}", sha256Hash, + peer.getInetAddress(), requestTime); + } }); } - public void blockFetchSuccess(Sha256Hash sha256Hash) { - logger.info("Fetch block success, {}", new BlockCapsule.BlockId(sha256Hash).getString()); + public synchronized void blockFetchSuccess(Sha256Hash sha256Hash) { FetchBlockInfo fetchBlockInfoTemp = this.fetchBlockInfo; if (null == fetchBlockInfoTemp || !fetchBlockInfoTemp.getHash().equals(sha256Hash)) { return; } + logger.info("Fetch block success, {}", new BlockCapsule.BlockId(sha256Hash).getString()); this.fetchBlockInfo = null; } - private void fetchBlockProcess(FetchBlockInfo fetchBlock) { - if (null == fetchBlock) { + public synchronized void onDisconnect(PeerConnection peer) { + if (fetchBlockInfo != null && fetchBlockInfo.getPeer().equals(peer)) { + fetchBlockInfo = null; + } + } + + public Optional selectBlockPeer(Collection peers, + Item item, long now) { + return peers.stream() + .filter(peer -> !peer.isDisconnect()) + .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) + .filter(PeerConnection::isBlockIdle) + .filter(peer -> { + Long received = peer.getAdvInvReceive().getIfPresent(item); + return received != null && received >= now - NetConstants.ADV_TIME_OUT; + }) + .min(Comparator.comparingDouble(this::getPeerTop75)); + } + + private synchronized void fetchBlockProcess(FetchBlockInfo fetchBlock) { + if (fetchBlock == null || fetchBlockInfo != fetchBlock) { + return; + } + if (new BlockCapsule.BlockId(fetchBlock.getHash()).getNum() + <= chainBaseManager.getHeadBlockNum()) { + fetchBlockInfo = null; + return; + } + long now = System.currentTimeMillis(); + if (now - fetchBlock.getTime() >= NetConstants.ADV_TIME_OUT) { + // PeerStatusCheck owns the final deadline and disconnects the responsible provider. return; } Item item = new Item(fetchBlock.getHash(), InventoryType.BLOCK); - Optional optionalPeerConnection = tronNetDelegate.getActivePeer().stream() - .filter(PeerConnection::isIdle) - .filter(filterPeer -> !filterPeer.equals(fetchBlock.getPeer())) - .filter(filterPeer -> filterPeer.getAdvInvReceive().getIfPresent(item) != null) - .filter(filterPeer -> getPeerTop75(filterPeer) - <= CommonParameter.getInstance().fetchBlockTimeout) - .min(Comparator.comparingDouble(this::getPeerTop75)); + Optional optionalPeerConnection = selectBlockPeer( + tronNetDelegate.getActivePeer(), item, now); if (optionalPeerConnection.isPresent()) { optionalPeerConnection.ifPresent(firstPeer -> { if (shouldFetchBlock(firstPeer, fetchBlock) - && firstPeer.checkAndPutAdvInvRequest(item, System.currentTimeMillis())) { + && firstPeer.checkAndPutAdvInvRequest(item, now)) { + this.fetchBlockInfo = new FetchBlockInfo(item.getHash(), firstPeer, now); firstPeer.sendMessage(new FetchInvDataMessage(Collections.singletonList(item.getHash()), item.getType())); - this.fetchBlockInfo = null; } }); - } else { - if (System.currentTimeMillis() - fetchBlock.getTime() >= fetchTimeOut) { - logger.info("Clear fetchBlockInfo due to fetch block {} timeout {}ms", - fetchBlock.getHash(), fetchTimeOut); - this.fetchBlockInfo = null; - } } } @@ -140,16 +164,13 @@ private double getPeerTop75(PeerConnection peerConnection) { private static class FetchBlockInfo { @Getter - @Setter - private PeerConnection peer; + private final PeerConnection peer; @Getter - @Setter - private Sha256Hash hash; + private final Sha256Hash hash; @Getter - @Setter - private long time; + private final long time; public FetchBlockInfo(Sha256Hash hash, PeerConnection peer, long time) { this.peer = peer; @@ -159,4 +180,4 @@ public FetchBlockInfo(Sha256Hash hash, PeerConnection peer, long time) { } -} \ No newline at end of file +} diff --git a/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java index bd656d9c41e..c873f135118 100644 --- a/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java +++ b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java @@ -336,8 +336,16 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection) BlockId blockId = block.getBlockId(); try { tronNetDelegate.validSignature(block); + boolean useful = block.getNum() >= tronNetDelegate.getHeadBlockId().getNum() + && !tronNetDelegate.containBlock(blockId); tronNetDelegate.processBlock(block, true); - peerConnection.setBlockRcvTime(System.currentTimeMillis()); + if (tronNetDelegate.isHitDown()) { + return; + } + peerConnection.setLastInteractiveTime(System.currentTimeMillis()); + if (useful) { + peerConnection.setBlockRcvTime(System.currentTimeMillis()); + } pbftDataSyncHandler.processPBFTCommitData(block); } catch (P2pException p2pException) { logger.error("Process sync block {} failed, type: {}", @@ -366,7 +374,7 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection) syncNext(peer); } } else { - peer.disconnect(ReasonCode.BAD_BLOCK); + peer.disconnect(ReasonCode.SYNC_FAIL); } } } diff --git a/framework/src/test/java/org/tron/core/net/PeerBlockTestSupport.java b/framework/src/test/java/org/tron/core/net/PeerBlockTestSupport.java new file mode 100644 index 00000000000..b754b1df98f --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/PeerBlockTestSupport.java @@ -0,0 +1,44 @@ +package org.tron.core.net; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; + +import com.google.protobuf.ByteString; +import java.net.InetSocketAddress; +import org.tron.common.utils.ReflectUtils; +import org.tron.common.utils.Sha256Hash; +import org.tron.core.capsule.BlockCapsule; +import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.core.net.peer.PeerConnection; +import org.tron.p2p.connection.Channel; + +public final class PeerBlockTestSupport { + + private PeerBlockTestSupport() { + } + + public static PeerConnection peer(int port) { + PeerConnection peer = spy(new PeerConnection()); + Channel channel = mock(Channel.class); + InetSocketAddress address = new InetSocketAddress("127.0.0.1", port); + when(channel.getInetSocketAddress()).thenReturn(address); + when(channel.getInetAddress()).thenReturn(address.getAddress()); + ReflectUtils.setFieldValue(peer, "channel", channel); + peer.setNeedSyncFromPeer(false); + peer.setNeedSyncFromUs(false); + peer.setLastInteractiveTime(1L); + doNothing().when(peer).sendMessage(any()); + doNothing().when(peer).disconnect(any()); + return peer; + } + + public static BlockCapsule block(long number) { + BlockCapsule block = new BlockCapsule(number, new BlockId(Sha256Hash.ZERO_HASH, number - 1), + System.currentTimeMillis() - 1_000, ByteString.copyFromUtf8("witness")); + block.setMerkleRoot(); + return block; + } +} diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/BlockContributionTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/BlockContributionTest.java new file mode 100644 index 00000000000..1ac40121214 --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/messagehandler/BlockContributionTest.java @@ -0,0 +1,320 @@ +package org.tron.core.net.messagehandler; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.google.protobuf.ByteString; +import java.lang.reflect.Method; +import java.util.Arrays; +import java.util.Collections; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.tron.common.crypto.ECKey; +import org.tron.common.utils.ReflectUtils; +import org.tron.core.capsule.BlockCapsule; +import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.core.db.Manager; +import org.tron.core.exception.P2pException; +import org.tron.core.exception.P2pException.TypeEnum; +import org.tron.core.net.P2pEventHandlerImpl; +import org.tron.core.net.PeerBlockTestSupport; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.message.TronMessage; +import org.tron.core.net.message.adv.BlockMessage; +import org.tron.core.net.peer.Item; +import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.service.adv.AdvService; +import org.tron.core.net.service.fetchblock.FetchBlockService; +import org.tron.core.net.service.sync.SyncService; +import org.tron.core.services.WitnessProductBlockService; +import org.tron.core.store.DynamicPropertiesStore; +import org.tron.core.store.WitnessScheduleStore; +import org.tron.protos.Protocol; +import org.tron.protos.Protocol.Inventory.InventoryType; + +public class BlockContributionTest { + + private BlockMsgHandler handler; + private TronNetDelegate delegate; + private FetchBlockService fetch; + private AdvService adv; + private SyncService sync; + private PeerConnection provider; + private BlockCapsule block; + private Item item; + private long requestTime; + + @Before + public void setUp() throws Exception { + handler = new BlockMsgHandler(); + delegate = mock(TronNetDelegate.class); + fetch = mock(FetchBlockService.class); + adv = mock(AdvService.class); + sync = mock(SyncService.class); + ReflectUtils.setFieldValue(handler, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(handler, "fetchBlockService", fetch); + ReflectUtils.setFieldValue(handler, "advService", adv); + ReflectUtils.setFieldValue(handler, "syncService", sync); + ReflectUtils.setFieldValue(handler, "witnessProductBlockService", + mock(WitnessProductBlockService.class)); + ReflectUtils.setFieldValue(handler, "fastForward", false); + provider = PeerBlockTestSupport.peer(18888); + block = PeerBlockTestSupport.block(85636071); + item = new Item(block.getBlockId(), InventoryType.BLOCK); + requestTime = System.currentTimeMillis() - 1_000; + provider.getAdvInvRequest().put(item, requestTime); + when(delegate.getHeadBlockId()).thenReturn(block.getParentBlockId()); + when(delegate.getActivePeer()).thenReturn(Collections.singletonList(provider)); + when(delegate.containBlock(block.getParentBlockId())).thenReturn(true); + when(delegate.validBlock(block)).thenReturn(true); + } + + @Test + public void testValidationFetchBroadcastAndExecutionOrder() throws Exception { + when(delegate.validBlock(block)).thenAnswer(call -> { + Assert.assertEquals(Long.valueOf(requestTime), provider.getAdvInvRequest().get(item)); + verify(fetch, never()).blockFetchSuccess(any()); + return true; + }); + doAnswer(call -> { + verify(fetch).blockFetchSuccess(block.getBlockId()); + Assert.assertFalse(provider.getAdvInvRequest().containsKey(item)); + Assert.assertTrue(provider.getLastInteractiveTime() > 1); + Assert.assertEquals(0, provider.getBlockRcvTime()); + verify(delegate, never()).processBlock(any(), eq(false)); + return null; + }).when(adv).broadcast(any(BlockMessage.class)); + doAnswer(call -> { + verify(adv).broadcast(any(BlockMessage.class)); + Assert.assertEquals(0, provider.getBlockRcvTime()); + return null; + }).when(delegate).processBlock(block, false); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertTrue(provider.getBlockRcvTime() > 0); + } + + @Test + public void testOnlyProviderGetsContribution() throws Exception { + PeerConnection advertiser = PeerBlockTestSupport.peer(18889); + advertiser.getAdvInvReceive().put(item, System.currentTimeMillis()); + when(delegate.getActivePeer()).thenReturn(Arrays.asList(provider, advertiser)); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertEquals(0, advertiser.getBlockRcvTime()); + } + + @Test + public void testBadMerkleRetainsRequestAndDoesNotCompleteFetch() throws Exception { + assertInvalid(TypeEnum.BLOCK_MERKLE_INVALID); + } + + @Test + public void testBadSignatureRetainsRequestAndDoesNotCompleteFetch() throws Exception { + assertInvalid(TypeEnum.BLOCK_SIGN_INVALID); + } + + @Test + public void testInactiveWitnessDoesNotImproveEitherTimestamp() throws Exception { + when(delegate.validBlock(block)).thenReturn(false); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertFalse(provider.getAdvInvRequest().containsKey(item)); + verify(sync).startSync(provider); + verify(adv, never()).broadcast(any()); + verify(provider, never()).disconnect(any()); + } + + @Test + public void testMissingParentCreditsValidatedProviderAndStartsSync() throws Exception { + when(delegate.containBlock(block.getParentBlockId())).thenReturn(false); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertTrue(provider.getBlockRcvTime() > 0); + verify(sync).startSync(provider); + verify(adv, never()).broadcast(any()); + verify(delegate, never()).processBlock(any(), eq(false)); + } + + @Test + public void testOldOrphanDoesNotGetContribution() throws Exception { + when(delegate.getHeadBlockId()).thenReturn(new BlockId(block.getBlockId(), block.getNum() + 1)); + when(delegate.containBlock(block.getParentBlockId())).thenReturn(false); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertEquals(0, provider.getBlockRcvTime()); + verify(sync, never()).startSync(any()); + verify(adv, never()).broadcast(any()); + } + + @Test + public void testAlreadyKnownBlockDoesNotGetContribution() throws Exception { + when(delegate.containBlock(block.getBlockId())).thenReturn(true); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertEquals(0, provider.getBlockRcvTime()); + verify(adv, never()).broadcast(any()); + } + + @Test + public void testExecutionFailureStartsSyncWithoutBlamingProvider() throws Exception { + doThrow(new P2pException(TypeEnum.BAD_BLOCK, "local execution state")) + .when(delegate).processBlock(block, false); + + handler.processMessage(provider, new BlockMessage(block)); + + verify(adv).broadcast(any(BlockMessage.class)); + verify(sync).startSync(provider); + verify(provider, never()).disconnect(any()); + Assert.assertTrue(provider.getLastInteractiveTime() > 1); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testShutdownDoesNotGetContribution() throws Exception { + when(delegate.isHitDown()).thenReturn(true); + + handler.processMessage(provider, new BlockMessage(block)); + + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testUnrequestedBlockCannotImproveTimestamps() throws Exception { + provider.getAdvInvRequest().clear(); + try { + handler.processMessage(provider, new BlockMessage(block)); + Assert.fail("Expected an unrequested block to be rejected"); + } catch (P2pException e) { + Assert.assertEquals(TypeEnum.BAD_MESSAGE, e.getType()); + } + verify(delegate, never()).validBlock(any()); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testInvalidParentHeightCannotCompleteFetch() throws Exception { + BlockCapsule invalid = new BlockCapsule(block.getInstance().toBuilder() + .setBlockHeader(block.getInstance().getBlockHeader().toBuilder() + .setRawData(block.getInstance().getBlockHeader().getRawData().toBuilder() + .setNumber(block.getNum() + 1))).build()); + Item invalidItem = new Item(invalid.getBlockId(), InventoryType.BLOCK); + provider.getAdvInvRequest().put(invalidItem, requestTime); + try { + handler.processMessage(provider, new BlockMessage(invalid)); + Assert.fail("Expected invalid parent height to be rejected"); + } catch (P2pException e) { + Assert.assertEquals(TypeEnum.BAD_BLOCK, e.getType()); + } + Assert.assertEquals(Long.valueOf(requestTime), provider.getAdvInvRequest().get(invalidItem)); + verify(fetch, never()).blockFetchSuccess(any()); + verify(adv, never()).broadcast(any()); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + } + + @Test + public void testDispatcherDoesNotUnconditionallyCreditBlocks() throws Exception { + Method update = P2pEventHandlerImpl.class.getDeclaredMethod("updateLastInteractiveTime", + PeerConnection.class, TronMessage.class); + update.setAccessible(true); + + update.invoke(new P2pEventHandlerImpl(), provider, new BlockMessage(block)); + + Assert.assertEquals(1, provider.getLastInteractiveTime()); + } + + private void assertInvalid(TypeEnum type) throws Exception { + when(delegate.validBlock(block)).thenThrow(new P2pException(type, "invalid data")); + try { + handler.processMessage(provider, new BlockMessage(block)); + Assert.fail("Expected invalid block to be rejected"); + } catch (P2pException e) { + Assert.assertEquals(type, e.getType()); + } + Assert.assertEquals(Long.valueOf(requestTime), provider.getAdvInvRequest().get(item)); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + verify(fetch, never()).blockFetchSuccess(any()); + verify(adv, never()).broadcast(any()); + } + + @Test(timeout = 20_000) + public void testRealSignedBlockPassesValidationBeforeContribution() throws Exception { + useRealValidation(); + handler.processMessage(provider, new BlockMessage(block)); + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertFalse(provider.getAdvInvRequest().containsKey(item)); + verify(fetch).blockFetchSuccess(block.getBlockId()); + } + + @Test(timeout = 20_000) + public void testRealBadSignatureRetainsRequest() throws Exception { + useRealValidation(); + block = new BlockCapsule(block.getInstance().toBuilder() + .setBlockHeader(block.getInstance().getBlockHeader().toBuilder() + .setWitnessSignature(ByteString.copyFrom(new byte[65]))).build()); + assertRealValidationFailure(TypeEnum.BLOCK_SIGN_INVALID); + } + + @Test(timeout = 20_000) + public void testRealTamperedBodyRetainsRequest() throws Exception { + useRealValidation(); + block = new BlockCapsule(block.getInstance().toBuilder() + .addTransactions(Protocol.Transaction.newBuilder() + .setRawData(Protocol.Transaction.raw.newBuilder() + .setData(ByteString.copyFromUtf8("tampered body")))).build()); + assertRealValidationFailure(TypeEnum.BLOCK_MERKLE_INVALID); + } + + private void useRealValidation() throws Exception { + ECKey key = new ECKey(); + block = new BlockCapsule(block.getNum(), block.getParentBlockId(), block.getTimeStamp(), + ByteString.copyFrom(key.getAddress())); + block.setMerkleRoot(); + block.sign(key.getPrivKeyBytes()); + item = new Item(block.getBlockId(), InventoryType.BLOCK); + provider.getAdvInvRequest().clear(); + provider.getAdvInvRequest().put(item, requestTime); + + TronNetDelegate validator = new TronNetDelegate(); + Manager manager = mock(Manager.class); + when(manager.getDynamicPropertiesStore()).thenReturn(mock(DynamicPropertiesStore.class)); + WitnessScheduleStore witnesses = mock(WitnessScheduleStore.class); + when(witnesses.getActiveWitnesses()) + .thenReturn(Collections.singletonList(block.getWitnessAddress())); + ReflectUtils.setFieldValue(validator, "dbManager", manager); + ReflectUtils.setFieldValue(validator, "witnessScheduleStore", witnesses); + when(delegate.validBlock(any())).thenAnswer(call -> validator.validBlock(call.getArgument(0))); + } + + private void assertRealValidationFailure(TypeEnum type) throws Exception { + Assert.assertEquals(item.getHash(), block.getBlockId()); + try { + handler.processMessage(provider, new BlockMessage(block)); + Assert.fail("Expected cryptographic validation to reject the block"); + } catch (P2pException e) { + Assert.assertEquals(type, e.getType()); + } + Assert.assertEquals(Long.valueOf(requestTime), provider.getAdvInvRequest().get(item)); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + verify(fetch, never()).blockFetchSuccess(any()); + } +} diff --git a/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java b/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java new file mode 100644 index 00000000000..3e7a7d59ed2 --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java @@ -0,0 +1,42 @@ +package org.tron.core.net.peer; + +import java.util.ArrayDeque; +import org.junit.Assert; +import org.junit.Test; +import org.tron.common.utils.Pair; +import org.tron.common.utils.Sha256Hash; +import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.protos.Protocol.Inventory.InventoryType; + +public class PeerBlockIdleTest { + + @Test + public void testTransactionRequestDoesNotBlockBlockFetch() { + PeerConnection peer = new PeerConnection(); + Item trx = new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX); + peer.getAdvInvRequest().put(trx, System.currentTimeMillis()); + Assert.assertFalse(peer.isIdle()); + Assert.assertTrue(peer.isBlockIdle()); + + Item block = new Item(new BlockId(), InventoryType.BLOCK); + peer.getAdvInvRequest().put(block, System.currentTimeMillis()); + Assert.assertFalse(peer.isBlockIdle()); + peer.getAdvInvRequest().remove(block); + Assert.assertTrue(peer.isBlockIdle()); + } + + @Test + public void testSyncRequestAndProcessingExcludeBlockProvider() { + PeerConnection peer = new PeerConnection(); + peer.getSyncBlockRequested().put(new BlockId(), System.currentTimeMillis()); + Assert.assertFalse(peer.isBlockIdle()); + peer.getSyncBlockRequested().clear(); + peer.setSyncChainRequested(new Pair<>(new ArrayDeque<>(), System.currentTimeMillis())); + Assert.assertFalse(peer.isBlockIdle()); + peer.setSyncChainRequested(null); + peer.getSyncBlockInProcess().add(new BlockId()); + Assert.assertFalse(peer.isBlockIdle()); + peer.getSyncBlockInProcess().clear(); + Assert.assertTrue(peer.isBlockIdle()); + } +} diff --git a/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java new file mode 100644 index 00000000000..2059c0c013c --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java @@ -0,0 +1,357 @@ +package org.tron.core.net.service.fetchblock; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.lang.reflect.Method; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.tron.common.parameter.CommonParameter; +import org.tron.common.utils.ReflectUtils; +import org.tron.common.utils.Sha256Hash; +import org.tron.core.ChainBaseManager; +import org.tron.core.capsule.BlockCapsule; +import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.core.config.Parameter.NetConstants; +import org.tron.core.exception.P2pException; +import org.tron.core.exception.P2pException.TypeEnum; +import org.tron.core.net.PeerBlockTestSupport; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.message.adv.BlockMessage; +import org.tron.core.net.message.adv.FetchInvDataMessage; +import org.tron.core.net.message.adv.InventoryMessage; +import org.tron.core.net.messagehandler.BlockMsgHandler; +import org.tron.core.net.messagehandler.InventoryMsgHandler; +import org.tron.core.net.peer.Item; +import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.peer.PeerStatusCheck; +import org.tron.core.net.service.adv.AdvService; +import org.tron.protos.Protocol.Inventory.InventoryType; +import org.tron.protos.Protocol.ReasonCode; + +public class FetchBlockRetryTest { + + private FetchBlockService fetch; + private AdvService adv; + private InventoryMsgHandler inventoryHandler; + private TronNetDelegate delegate; + private ChainBaseManager chain; + private PeerConnection first; + private PeerConnection second; + private List peers; + private BlockCapsule block; + private Item item; + private int previousTimeout; + + @Before + public void setUp() { + previousTimeout = CommonParameter.getInstance().fetchBlockTimeout; + CommonParameter.getInstance().fetchBlockTimeout = 1_000; + fetch = new FetchBlockService(); + adv = new AdvService(); + inventoryHandler = new InventoryMsgHandler(); + delegate = mock(TronNetDelegate.class); + chain = mock(ChainBaseManager.class); + first = PeerBlockTestSupport.peer(18888); + second = PeerBlockTestSupport.peer(18889); + peers = new ArrayList<>(Arrays.asList(first, second)); + block = PeerBlockTestSupport.block(85636071); + item = new Item(block.getBlockId(), InventoryType.BLOCK); + when(delegate.getActivePeer()).thenReturn(peers); + when(delegate.getHeadBlockId()).thenReturn(block.getParentBlockId()); + when(chain.getHeadBlockNum()).thenReturn(block.getNum() - 1); + ReflectUtils.setFieldValue(fetch, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(fetch, "chainBaseManager", chain); + ReflectUtils.setFieldValue(adv, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(adv, "fetchBlockService", fetch); + ReflectUtils.setFieldValue(inventoryHandler, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(inventoryHandler, "advService", adv); + } + + @After + public void tearDown() { + fetch.close(); + adv.close(); + CommonParameter.getInstance().fetchBlockTimeout = previousTimeout; + } + + @Test + public void testShortTimeoutWithoutAlternativeRetainsState() throws Exception { + beginAgedFetch(2_000); + Object state = state(); + + tick(state); + + Assert.assertSame(state, state()); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + verify(second, never()).sendMessage(any()); + } + + @Test + public void testLateInventoryBypassesCacheForProviderSwitch() throws Exception { + inventoryHandler.processMessage(first, inventory()); + Assert.assertNotNull(state()); + Assert.assertFalse(adv.addInv(item)); + ageState(2_000); + long originalRequest = first.getAdvInvRequest().get(item); + second.getAdvInvRequest().put(new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX), + System.currentTimeMillis()); + + inventoryHandler.processMessage(second, inventory()); + tick(state()); + + verify(second).sendMessage(any(FetchInvDataMessage.class)); + Assert.assertNotNull(state()); + Assert.assertSame(second, ReflectUtils.getFieldObject(state(), "peer")); + Assert.assertEquals(Long.valueOf(originalRequest), first.getAdvInvRequest().get(item)); + Assert.assertTrue(second.getAdvInvRequest().containsKey(item)); + Assert.assertEquals(1, first.getLastInteractiveTime()); + Assert.assertEquals(1, second.getLastInteractiveTime()); + Assert.assertEquals(0, second.getBlockRcvTime()); + } + + @Test + public void testFirstProviderCanHavePendingTransactions() throws Exception { + first.getAdvInvRequest().put(new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX), + System.currentTimeMillis()); + + inventoryHandler.processMessage(first, inventory()); + + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + verify(first).sendMessage(any(FetchInvDataMessage.class)); + } + + @Test + public void testSameSelectorExcludesBusySyncingAndDisconnectedPeers() { + first.getAdvInvReceive().put(item, System.currentTimeMillis()); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + first.getSyncBlockInProcess().add(block.getParentBlockId()); + Assert.assertSame(second, + fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).get()); + second.setNeedSyncFromPeer(true); + Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + second.setNeedSyncFromPeer(false); + when(second.getChannel().isDisconnect()).thenReturn(true); + Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + } + + @Test + public void testSelectorRejectsStaleInventory() { + first.getAdvInvReceive().put(item, + System.currentTimeMillis() - NetConstants.ADV_TIME_OUT - 1_000); + Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + } + + @Test + public void testOnlyOriginalProviderGetsFinalTimeout() throws Exception { + beginAgedFetch(2_000); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + first.getAdvInvRequest().put(item, + System.currentTimeMillis() - NetConstants.ADV_TIME_OUT - 1_000); + PeerStatusCheck status = new PeerStatusCheck(); + ReflectUtils.setFieldValue(status, "tronNetDelegate", delegate); + try { + status.statusCheck(); + } finally { + status.close(); + } + + verify(first).disconnect(ReasonCode.TIME_OUT); + verify(second, never()).disconnect(any()); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + Assert.assertTrue(second.getAdvInvRequest().containsKey(item)); + } + + @Test + public void testFinalDeadlineDoesNotDiscardResponsibility() throws Exception { + beginAgedFetch(NetConstants.ADV_TIME_OUT + 1_000); + Object state = state(); + + tick(state); + + Assert.assertSame(state, state()); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + verify(first, never()).disconnect(any()); + } + + @Test + public void testDisconnectInvalidatesCacheWhenNoAlternativeExists() throws Exception { + inventoryHandler.processMessage(first, inventory()); + disconnect(first); + Assert.assertNull(state()); + + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + Assert.assertTrue(adv.addInv(item)); + verify(second).sendMessage(any(FetchInvDataMessage.class)); + } + + @Test + public void testInvalidBlockDisconnectRetriesFromAlternative() throws Exception { + inventoryHandler.processMessage(first, inventory()); + inventoryHandler.processMessage(second, inventory()); + BlockMsgHandler handler = new BlockMsgHandler(); + ReflectUtils.setFieldValue(handler, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(handler, "fetchBlockService", fetch); + ReflectUtils.setFieldValue(handler, "fastForward", false); + when(delegate.validBlock(block)) + .thenThrow(new P2pException(TypeEnum.BLOCK_MERKLE_INVALID, "bad")); + try { + handler.processMessage(first, new BlockMessage(block)); + Assert.fail("Expected a bad Merkle root to be rejected"); + } catch (P2pException e) { + Assert.assertEquals(TypeEnum.BLOCK_MERKLE_INVALID, e.getType()); + } + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + + disconnect(first); + + Assert.assertTrue(second.getAdvInvRequest().containsKey(item)); + verify(second).sendMessage(any(FetchInvDataMessage.class)); + Assert.assertNotNull(state()); + } + + @Test + public void testOriginalDisconnectDoesNotDuplicateOutstandingBackupRequest() throws Exception { + beginAgedFetch(2_000); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + Object backupState = state(); + + disconnect(first); + + Assert.assertSame(backupState, state()); + verify(second, times(1)).sendMessage(any(FetchInvDataMessage.class)); + } + + @Test + public void testBackupDisconnectRestoresTrackingOfOriginalRequest() throws Exception { + beginAgedFetch(2_000); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + + disconnect(second); + + Assert.assertNotNull(state()); + Assert.assertSame(first, ReflectUtils.getFieldObject(state(), "peer")); + verify(first, never()).sendMessage(any()); + } + + @Test + public void testStaleWorkerCannotOverwriteNewFetch() throws Exception { + beginAgedFetch(2_000); + Object oldState = state(); + fetch.blockFetchSuccess(item.getHash()); + Item next = new Item(PeerBlockTestSupport.block(block.getNum() + 1).getBlockId(), + InventoryType.BLOCK); + when(chain.getHeadBlockNum()).thenReturn(block.getNum()); + second.getAdvInvRequest().put(next, System.currentTimeMillis()); + fetch.fetchBlock(Collections.singletonList(next.getHash()), second); + Object newState = state(); + + tick(oldState); + fetch.blockFetchSuccess(item.getHash()); + + Assert.assertSame(newState, state()); + verify(second, never()).sendMessage(any()); + } + + @Test + public void testImmediateFirstResponseDoesNotResurrectCompletedFetch() throws Exception { + doAnswer(call -> { + Assert.assertNotNull(state()); + fetch.blockFetchSuccess(item.getHash()); + first.getAdvInvRequest().remove(item); + return null; + }).when(first).sendMessage(any(FetchInvDataMessage.class)); + + inventoryHandler.processMessage(first, inventory()); + + Assert.assertNull(state()); + } + + @Test + public void testImmediateBackupResponseDoesNotResurrectCompletedFetch() throws Exception { + beginAgedFetch(2_000); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + doAnswer(call -> { + fetch.blockFetchSuccess(item.getHash()); + second.getAdvInvRequest().remove(item); + return null; + }).when(second).sendMessage(any(FetchInvDataMessage.class)); + + tick(state()); + + Assert.assertNull(state()); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + } + + private InventoryMessage inventory() { + return new InventoryMessage(Collections.singletonList(item.getHash()), InventoryType.BLOCK); + } + + @Test + public void testNewHeadReplacesObsoleteFetchWithoutLosingOriginalRequest() { + beginAgedFetch(2_000); + Item next = new Item(PeerBlockTestSupport.block(block.getNum() + 1).getBlockId(), + InventoryType.BLOCK); + when(chain.getHeadBlockNum()).thenReturn(block.getNum()); + second.getAdvInvRequest().put(next, System.currentTimeMillis()); + + fetch.fetchBlock(Collections.singletonList(next.getHash()), second); + + Assert.assertEquals(next.getHash(), ReflectUtils.getFieldObject(state(), "hash")); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + } + + @Test + public void testWorkerStopsObsoleteFetchAfterHeadAdvances() throws Exception { + beginAgedFetch(2_000); + when(chain.getHeadBlockNum()).thenReturn(block.getNum()); + + tick(state()); + + Assert.assertNull(state()); + Assert.assertTrue(first.getAdvInvRequest().containsKey(item)); + } + + private void beginAgedFetch(long age) { + first.getAdvInvRequest().put(item, System.currentTimeMillis() - age); + fetch.fetchBlock(Collections.singletonList(item.getHash()), first); + Assert.assertNotNull(state()); + } + + private void ageState(long age) { + long time = System.currentTimeMillis() - age; + ReflectUtils.setFieldValue(state(), "time", time); + first.getAdvInvRequest().put(item, time); + } + + private Object state() { + return ReflectUtils.getFieldObject(fetch, "fetchBlockInfo"); + } + + private void tick(Object state) throws Exception { + Class stateClass = Class.forName(FetchBlockService.class.getName() + "$FetchBlockInfo"); + Method method = FetchBlockService.class.getDeclaredMethod("fetchBlockProcess", stateClass); + method.setAccessible(true); + method.invoke(fetch, state); + } + + private void disconnect(PeerConnection peer) { + when(peer.getChannel().isDisconnect()).thenReturn(true); + peers.remove(peer); + adv.onDisconnect(peer); + } +} diff --git a/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java new file mode 100644 index 00000000000..ec3c4b9cbaa --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java @@ -0,0 +1,115 @@ +package org.tron.core.net.service.sync; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.lang.reflect.Method; +import java.util.Arrays; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.tron.common.utils.ReflectUtils; +import org.tron.core.capsule.BlockCapsule; +import org.tron.core.exception.P2pException; +import org.tron.core.exception.P2pException.TypeEnum; +import org.tron.core.net.PeerBlockTestSupport; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.messagehandler.PbftDataSyncHandler; +import org.tron.core.net.peer.PeerConnection; +import org.tron.protos.Protocol.ReasonCode; + +public class SyncBlockContributionTest { + + private SyncService sync; + private TronNetDelegate delegate; + private PeerConnection provider; + private PeerConnection other; + private BlockCapsule block; + + @Before + public void setUp() { + sync = new SyncService(); + delegate = mock(TronNetDelegate.class); + provider = PeerBlockTestSupport.peer(18888); + other = PeerBlockTestSupport.peer(18889); + block = PeerBlockTestSupport.block(10); + provider.getSyncBlockToFetch().add(block.getBlockId()); + other.getSyncBlockToFetch().add(block.getBlockId()); + when(delegate.getActivePeer()).thenReturn(Arrays.asList(provider, other)); + when(delegate.getHeadBlockId()).thenReturn(block.getParentBlockId()); + ReflectUtils.setFieldValue(sync, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(sync, "pbftDataSyncHandler", mock(PbftDataSyncHandler.class)); + } + + @After + public void tearDown() { + sync.close(); + } + + @Test + public void testOnlyValidatedProviderGetsTimestamps() throws Exception { + process(); + Assert.assertTrue(provider.getLastInteractiveTime() > 1); + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertEquals(1, other.getLastInteractiveTime()); + Assert.assertEquals(0, other.getBlockRcvTime()); + } + + @Test + public void testInvalidSignatureOnlyBlamesProvider() throws Exception { + doThrow(new P2pException(TypeEnum.BLOCK_SIGN_INVALID, "bad signature")) + .when(delegate).validSignature(block); + process(); + verify(provider).disconnect(ReasonCode.BAD_BLOCK); + verify(other, never()).disconnect(any()); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testStateFailureDoesNotBanPeersOrImproveTimestamps() throws Exception { + doThrow(new P2pException(TypeEnum.BAD_BLOCK, "state failure")) + .when(delegate).processBlock(block, true); + process(); + verify(provider).disconnect(ReasonCode.SYNC_FAIL); + verify(other).disconnect(ReasonCode.SYNC_FAIL); + verify(provider, never()).disconnect(ReasonCode.BAD_BLOCK); + verify(other, never()).disconnect(ReasonCode.BAD_BLOCK); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testShutdownDoesNotImproveTimestamps() throws Exception { + when(delegate.isHitDown()).thenReturn(true); + process(); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testOldSyncBlockDoesNotGetContribution() throws Exception { + when(delegate.getHeadBlockId()).thenReturn(PeerBlockTestSupport.block(11).getBlockId()); + process(); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + @Test + public void testKnownSyncBlockDoesNotGetContribution() throws Exception { + when(delegate.containBlock(block.getBlockId())).thenReturn(true); + process(); + Assert.assertEquals(0, provider.getBlockRcvTime()); + } + + private void process() throws Exception { + Method method = SyncService.class.getDeclaredMethod("processSyncBlock", + BlockCapsule.class, PeerConnection.class); + method.setAccessible(true); + method.invoke(sync, block, provider); + } +} From 1afae715e75706b6380b345a6cc870bcb22cf5dc Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 14 Sep 2026 22:22:54 +0800 Subject: [PATCH 2/7] add comment for 3 idle methods in PeerConnection --- .../tron/core/net/peer/PeerConnection.java | 19 ++++++++++++++++++- .../service/fetchblock/FetchBlockService.java | 2 +- .../tron/core/net/peer/PeerBlockIdleTest.java | 14 +++++++------- 3 files changed, 26 insertions(+), 9 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index 25ad47ce73e..4e41a40e410 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -188,16 +188,33 @@ public void setBlockBothHave(BlockId blockId) { this.blockBothHaveUpdateTime = System.currentTimeMillis(); } + /** + * Returns whether there are no outstanding inventory or sync requests. + * + *

Sync blocks already being processed and sync direction flags are not checked. + */ public boolean isIdle() { return advInvRequest.isEmpty() && isSyncIdle(); } - public boolean isBlockIdle() { + /** + * Returns whether there are no outstanding block or sync requests and no sync blocks + * being processed. Outstanding transaction inventory requests do not make this check fail. + * + *

Callers must check connection state and sync direction flags separately before fetching. + */ + public boolean isBlockFetchIdle() { return advInvRequest.keySet().stream() .noneMatch(item -> item.getType() == Protocol.Inventory.InventoryType.BLOCK) && isSyncIdle() && syncBlockInProcess.isEmpty(); } + /** + * Returns whether there are no outstanding sync block or chain-summary requests. + * + *

This does not mean synchronization is complete: received blocks may still be processing, + * and sync direction flags are not checked. + */ public boolean isSyncIdle() { return syncBlockRequested.isEmpty() && syncChainRequested == null; } diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index 7a6b74bf9cc..8ff9b779405 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -105,7 +105,7 @@ public Optional selectBlockPeer(Collection peers return peers.stream() .filter(peer -> !peer.isDisconnect()) .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) - .filter(PeerConnection::isBlockIdle) + .filter(PeerConnection::isBlockFetchIdle) .filter(peer -> { Long received = peer.getAdvInvReceive().getIfPresent(item); return received != null && received >= now - NetConstants.ADV_TIME_OUT; diff --git a/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java b/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java index 3e7a7d59ed2..69d4aef4fe7 100644 --- a/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java +++ b/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java @@ -16,27 +16,27 @@ public void testTransactionRequestDoesNotBlockBlockFetch() { Item trx = new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX); peer.getAdvInvRequest().put(trx, System.currentTimeMillis()); Assert.assertFalse(peer.isIdle()); - Assert.assertTrue(peer.isBlockIdle()); + Assert.assertTrue(peer.isBlockFetchIdle()); Item block = new Item(new BlockId(), InventoryType.BLOCK); peer.getAdvInvRequest().put(block, System.currentTimeMillis()); - Assert.assertFalse(peer.isBlockIdle()); + Assert.assertFalse(peer.isBlockFetchIdle()); peer.getAdvInvRequest().remove(block); - Assert.assertTrue(peer.isBlockIdle()); + Assert.assertTrue(peer.isBlockFetchIdle()); } @Test public void testSyncRequestAndProcessingExcludeBlockProvider() { PeerConnection peer = new PeerConnection(); peer.getSyncBlockRequested().put(new BlockId(), System.currentTimeMillis()); - Assert.assertFalse(peer.isBlockIdle()); + Assert.assertFalse(peer.isBlockFetchIdle()); peer.getSyncBlockRequested().clear(); peer.setSyncChainRequested(new Pair<>(new ArrayDeque<>(), System.currentTimeMillis())); - Assert.assertFalse(peer.isBlockIdle()); + Assert.assertFalse(peer.isBlockFetchIdle()); peer.setSyncChainRequested(null); peer.getSyncBlockInProcess().add(new BlockId()); - Assert.assertFalse(peer.isBlockIdle()); + Assert.assertFalse(peer.isBlockFetchIdle()); peer.getSyncBlockInProcess().clear(); - Assert.assertTrue(peer.isBlockIdle()); + Assert.assertTrue(peer.isBlockFetchIdle()); } } From 1f4eb95cc4709369a30fa3f1ee5807cda14ec515 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 28 Sep 2026 17:32:28 +0800 Subject: [PATCH 3/7] fetch first block uses invSender's size, retry block uses Top75 latency --- .../tron/core/net/service/adv/AdvService.java | 33 +-- .../service/fetchblock/FetchBlockService.java | 24 +-- .../fetchblock/FetchBlockRetryTest.java | 203 +++++++++++++++++- 3 files changed, 227 insertions(+), 33 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 0a6431a938c..98d2e39f42c 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -301,6 +301,7 @@ private void consumerInvToFetch() { return; } long now = System.currentTimeMillis(); + List blocks = new ArrayList<>(); invToFetch.forEach((item, time) -> { long timeout = item.getType() == InventoryType.BLOCK ? NetConstants.ADV_TIME_OUT : TIMEOUT; if (time < now - timeout) { @@ -311,18 +312,7 @@ private void consumerInvToFetch() { return; } if (item.getType() == InventoryType.BLOCK) { - if (blockCache.getIfPresent(item) != null - || tronNetDelegate.containBlock(new BlockId(item.getHash())) - || peers.stream().anyMatch(peer -> peer.getAdvInvRequest().containsKey(item))) { - invToFetch.remove(item); - return; - } - fetchBlockService.selectBlockPeer(peers, item, now).ifPresent(peer -> { - if (peer.checkAndPutAdvInvRequest(item, now)) { - invSender.add(item, peer); - invToFetch.remove(item); - } - }); + blocks.add(item); return; } trxPeers.stream().filter(peer -> { @@ -336,6 +326,25 @@ private void consumerInvToFetch() { invToFetch.remove(item); }); }); + + // Reserve peers for earlier blocks before later blocks can make them busy. + blocks.sort(Comparator.comparingLong(item -> new BlockId(item.getHash()).getNum())); + blocks.forEach(item -> { + if (blockCache.getIfPresent(item) != null + || tronNetDelegate.containBlock(new BlockId(item.getHash())) + || peers.stream().anyMatch(peer -> peer.getAdvInvRequest().containsKey(item))) { + invToFetch.remove(item); + return; + } + peers.stream().filter(peer -> fetchBlockService.canFetchBlock(peer, item, now)) + .min(Comparator.comparingInt(invSender::getSize)) + .ifPresent(peer -> { + if (peer.checkAndPutAdvInvRequest(item, now)) { + invSender.add(item, peer); + invToFetch.remove(item); + } + }); + }); } invSender.sendFetch(); diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index 8ff9b779405..8122c150015 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -1,6 +1,5 @@ package org.tron.core.net.service.fetchblock; -import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.List; @@ -100,17 +99,13 @@ public synchronized void onDisconnect(PeerConnection peer) { } } - public Optional selectBlockPeer(Collection peers, - Item item, long now) { - return peers.stream() - .filter(peer -> !peer.isDisconnect()) - .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) - .filter(PeerConnection::isBlockFetchIdle) - .filter(peer -> { - Long received = peer.getAdvInvReceive().getIfPresent(item); - return received != null && received >= now - NetConstants.ADV_TIME_OUT; - }) - .min(Comparator.comparingDouble(this::getPeerTop75)); + public boolean canFetchBlock(PeerConnection peer, Item item, long now) { + if (peer.isDisconnect() || peer.isNeedSyncFromPeer() || peer.isNeedSyncFromUs() + || !peer.isBlockFetchIdle()) { + return false; + } + Long received = peer.getAdvInvReceive().getIfPresent(item); + return received != null && received >= now - NetConstants.ADV_TIME_OUT; } private synchronized void fetchBlockProcess(FetchBlockInfo fetchBlock) { @@ -128,8 +123,9 @@ private synchronized void fetchBlockProcess(FetchBlockInfo fetchBlock) { return; } Item item = new Item(fetchBlock.getHash(), InventoryType.BLOCK); - Optional optionalPeerConnection = selectBlockPeer( - tronNetDelegate.getActivePeer(), item, now); + Optional optionalPeerConnection = tronNetDelegate.getActivePeer().stream() + .filter(peer -> canFetchBlock(peer, item, now)) + .min(Comparator.comparingDouble(this::getPeerTop75)); if (optionalPeerConnection.isPresent()) { optionalPeerConnection.ifPresent(firstPeer -> { diff --git a/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java index 2059c0c013c..a07d849720c 100644 --- a/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java +++ b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java @@ -1,22 +1,29 @@ package org.tron.core.net.service.fetchblock; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockStatic; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import com.codahale.metrics.MetricRegistry; import java.lang.reflect.Method; +import java.net.InetSocketAddress; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; +import java.util.Map; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.MockedStatic; import org.tron.common.parameter.CommonParameter; import org.tron.common.utils.ReflectUtils; import org.tron.common.utils.Sha256Hash; @@ -26,6 +33,8 @@ import org.tron.core.config.Parameter.NetConstants; import org.tron.core.exception.P2pException; import org.tron.core.exception.P2pException.TypeEnum; +import org.tron.core.metrics.MetricsKey; +import org.tron.core.metrics.MetricsUtil; import org.tron.core.net.PeerBlockTestSupport; import org.tron.core.net.TronNetDelegate; import org.tron.core.net.message.adv.BlockMessage; @@ -133,24 +142,138 @@ public void testFirstProviderCanHavePendingTransactions() throws Exception { } @Test - public void testSameSelectorExcludesBusySyncingAndDisconnectedPeers() { + public void testInitialFetchPrefersFewerBatchRequestsOverLatency() { + Item transaction = new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX); + long now = System.currentTimeMillis(); + first.getAdvInvReceive().put(transaction, now); + first.getAdvInvReceive().put(item, now); + second.getAdvInvReceive().put(item, now); + Assert.assertTrue(adv.addInv(transaction)); + + try (MockedStatic metrics = mockBlockLatencies(first, 10L, second, 500L)) { + Assert.assertTrue(adv.addInv(item)); + } + + assertBlockRequests(second, item); + Assert.assertTrue(first.getAdvInvRequest().containsKey(transaction)); + Assert.assertFalse(first.getAdvInvRequest().containsKey(item)); + Assert.assertTrue(second.getAdvInvRequest().containsKey(item)); + Assert.assertSame(second, ReflectUtils.getFieldObject(state(), "peer")); + } + + @Test + public void testInitialFetchDoesNotPreferPeerWithoutLatencySamples() { + long now = System.currentTimeMillis(); + first.getAdvInvReceive().put(item, now); + second.getAdvInvReceive().put(item, now); + + try (MockedStatic metrics = mockBlockLatencies(first, 500L, second, null)) { + Assert.assertTrue(adv.addInv(item)); + } + + assertBlockRequests(first, item); + verify(second, never()).sendMessage(any()); + Assert.assertSame(first, ReflectUtils.getFieldObject(state(), "peer")); + } + + @Test + public void testBackupFetchStillPrefersLowerLatency() throws Exception { + PeerConnection faster = PeerBlockTestSupport.peer(18890); + peers.add(faster); + beginAgedFetch(2_000); + Long originalRequestTime = first.getAdvInvRequest().get(item); + long now = System.currentTimeMillis(); + second.getAdvInvReceive().put(item, now); + faster.getAdvInvReceive().put(item, now); + Item transaction = new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX); + faster.getAdvInvRequest().put(transaction, now); + + try (MockedStatic metrics = mockBlockLatencies(second, 500L, faster, 10L)) { + tick(state()); + } + + assertBlockRequests(faster, item); + verify(second, never()).sendMessage(any()); + Assert.assertTrue(faster.getAdvInvRequest().containsKey(transaction)); + Assert.assertTrue(faster.getAdvInvRequest().containsKey(item)); + Assert.assertEquals(originalRequestTime, first.getAdvInvRequest().get(item)); + Assert.assertSame(faster, ReflectUtils.getFieldObject(state(), "peer")); + } + + @Test + public void testQueuedBlocksAreFetchedInHeightOrder() throws Exception { + List blocks = queueBlocksInReverseHeightOrder(); + Item transaction = new Item(Sha256Hash.ZERO_HASH, InventoryType.TRX); + first.getAdvInvRequest().put(transaction, System.currentTimeMillis()); + + consumeInventory(); + + assertBlockRequests(first, blocks.get(0)); + Assert.assertTrue(first.getAdvInvRequest().containsKey(blocks.get(0))); + Assert.assertTrue(first.getAdvInvRequest().containsKey(transaction)); + Assert.assertFalse(first.getAdvInvRequest().containsKey(blocks.get(1))); + Assert.assertFalse(pendingInventory().containsKey(blocks.get(0))); + Assert.assertTrue(pendingInventory().containsKey(blocks.get(1))); + Assert.assertEquals(blocks.get(0).getHash(), ReflectUtils.getFieldObject(state(), "hash")); + } + + @Test + public void testQueuedNextBlockFetchedAfterHeadAdvances() throws Exception { + List blocks = queueBlocksInReverseHeightOrder(); + consumeInventory(); + Assert.assertTrue(first.getAdvInvRequest().containsKey(blocks.get(0))); + Assert.assertFalse(first.getAdvInvRequest().containsKey(blocks.get(1))); + + fetch.blockFetchSuccess(blocks.get(0).getHash()); + first.getAdvInvRequest().remove(blocks.get(0)); + when(chain.getHeadBlockNum()).thenReturn(block.getNum()); + when(delegate.getHeadBlockId()).thenReturn(new BlockId(blocks.get(0).getHash())); + + consumeInventory(); + + assertBlockRequests(first, blocks.get(0), blocks.get(1)); + Assert.assertFalse(first.getAdvInvRequest().containsKey(blocks.get(0))); + Assert.assertTrue(first.getAdvInvRequest().containsKey(blocks.get(1))); + Assert.assertTrue(pendingInventory().isEmpty()); + Assert.assertEquals(blocks.get(1).getHash(), ReflectUtils.getFieldObject(state(), "hash")); + } + + @Test + public void testOnlyHigherBlockCanStillBeFetched() throws Exception { + Item higher = new Item(new BlockId(Sha256Hash.ZERO_HASH, block.getNum() + 1), + InventoryType.BLOCK); + + inventoryHandler.processMessage(first, + new InventoryMessage(Collections.singletonList(higher.getHash()), InventoryType.BLOCK)); + + assertBlockRequests(first, higher); + Assert.assertTrue(first.getAdvInvRequest().containsKey(higher)); + Assert.assertTrue(pendingInventory().isEmpty()); + } + + @Test + public void testBlockEligibilityExcludesBusySyncingAndDisconnectedPeers() { first.getAdvInvReceive().put(item, System.currentTimeMillis()); second.getAdvInvReceive().put(item, System.currentTimeMillis()); first.getSyncBlockInProcess().add(block.getParentBlockId()); - Assert.assertSame(second, - fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).get()); + Assert.assertFalse(fetch.canFetchBlock(first, item, System.currentTimeMillis())); + Assert.assertTrue(fetch.canFetchBlock(second, item, System.currentTimeMillis())); second.setNeedSyncFromPeer(true); - Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + Assert.assertFalse(fetch.canFetchBlock(second, item, System.currentTimeMillis())); second.setNeedSyncFromPeer(false); + second.setNeedSyncFromUs(true); + Assert.assertFalse(fetch.canFetchBlock(second, item, System.currentTimeMillis())); + second.setNeedSyncFromUs(false); when(second.getChannel().isDisconnect()).thenReturn(true); - Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + Assert.assertFalse(fetch.canFetchBlock(second, item, System.currentTimeMillis())); } @Test - public void testSelectorRejectsStaleInventory() { + public void testBlockEligibilityRequiresRecentInventory() { + Assert.assertFalse(fetch.canFetchBlock(first, item, System.currentTimeMillis())); first.getAdvInvReceive().put(item, System.currentTimeMillis() - NetConstants.ADV_TIME_OUT - 1_000); - Assert.assertFalse(fetch.selectBlockPeer(peers, item, System.currentTimeMillis()).isPresent()); + Assert.assertFalse(fetch.canFetchBlock(first, item, System.currentTimeMillis())); } @Test @@ -332,6 +455,72 @@ private void beginAgedFetch(long age) { Assert.assertNotNull(state()); } + private List queueBlocksInReverseHeightOrder() throws Exception { + peers.remove(second); + byte[] hash = new byte[Sha256Hash.LENGTH]; + hash[hash.length - 1] = 1; + Item lower = new Item(new BlockId(hash, block.getNum()), InventoryType.BLOCK); + Item higher = new Item(new BlockId(Sha256Hash.ZERO_HASH, block.getNum() + 1), + InventoryType.BLOCK); + Item busy = new Item(block.getParentBlockId(), InventoryType.BLOCK); + first.getAdvInvRequest().put(busy, System.currentTimeMillis()); + + inventoryHandler.processMessage(first, + new InventoryMessage(Arrays.asList(lower.getHash(), higher.getHash()), + InventoryType.BLOCK)); + + verify(first, never()).sendMessage(any(FetchInvDataMessage.class)); + Assert.assertEquals(2, pendingInventory().size()); + // Assert the fixture order so the test cannot pass merely because the map visits lower first. + Assert.assertEquals(Arrays.asList(higher, lower), + new ArrayList<>(pendingInventory().keySet())); + first.getAdvInvRequest().remove(busy); + return Arrays.asList(lower, higher); + } + + private Map pendingInventory() { + return (Map) ReflectUtils.getFieldObject(adv, "invToFetch"); + } + + private void consumeInventory() throws Exception { + Method method = AdvService.class.getDeclaredMethod("consumerInvToFetch"); + method.setAccessible(true); + method.invoke(adv); + } + + private void assertBlockRequests(PeerConnection peer, Item... blocks) { + ArgumentCaptor requests = + ArgumentCaptor.forClass(FetchInvDataMessage.class); + verify(peer, times(blocks.length)).sendMessage(requests.capture()); + for (int i = 0; i < blocks.length; i++) { + Assert.assertEquals(InventoryType.BLOCK, requests.getAllValues().get(i).getInventoryType()); + Assert.assertEquals(Collections.singletonList(blocks[i].getHash()), + requests.getAllValues().get(i).getHashList()); + } + } + + private MockedStatic mockBlockLatencies(PeerConnection left, Long leftLatency, + PeerConnection right, Long rightLatency) { + MetricRegistry registry = new MetricRegistry(); + PeerConnection[] candidates = {left, right}; + Long[] latencies = {leftLatency, rightLatency}; + for (int i = 0; i < candidates.length; i++) { + PeerConnection peer = candidates[i]; + InetSocketAddress address = new InetSocketAddress("127.0.0." + (i + 2), + peer.getInetSocketAddress().getPort()); + when(peer.getChannel().getInetAddress()).thenReturn(address.getAddress()); + when(peer.getChannel().getInetSocketAddress()).thenReturn(address); + if (latencies[i] != null) { + registry.histogram(MetricsKey.NET_LATENCY_FETCH_BLOCK + peer.getInetAddress()) + .update(latencies[i]); + } + } + MockedStatic metrics = mockStatic(MetricsUtil.class); + metrics.when(() -> MetricsUtil.getHistogram(anyString())) + .thenAnswer(call -> registry.histogram(call.getArgument(0))); + return metrics; + } + private void ageState(long age) { long time = System.currentTimeMillis() - age; ReflectUtils.setFieldValue(state(), "time", time); From 29312f39fab90814c2913ed65cc7db605787b473 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 28 Sep 2026 17:33:08 +0800 Subject: [PATCH 4/7] delete unused code --- .../tron/core/net/service/adv/AdvService.java | 21 ------------------- 1 file changed, 21 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 98d2e39f42c..90f1fb41e8b 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -232,27 +232,6 @@ public void broadcast(Message msg) { } } - /* - public void fastForward(BlockMessage msg) { - Item item = new Item(msg.getBlockId(), InventoryType.BLOCK); - List peers = tronNetDelegate.getActivePeer().stream() - .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) - .filter(peer -> peer.getAdvInvReceive().getIfPresent(item) == null - && peer.getAdvInvSpread().getIfPresent(item) == null) - .collect(Collectors.toList()); - - if (!fastForward) { - peers = peers.stream().filter(peer -> peer.isFastForwardPeer()).collect(Collectors.toList()); - } - - peers.forEach(peer -> { - peer.fastSend(msg); - peer.getAdvInvSpread().put(item, System.currentTimeMillis()); - peer.setFastForwardBlock(msg.getBlockId()); - }); - } - */ - public void onDisconnect(PeerConnection peer) { fetchBlockService.onDisconnect(peer); if (!peer.getAdvInvRequest().isEmpty()) { From fbc11c6a04fcb6ff5338447d0b555ec5d19c3654 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 28 Sep 2026 18:21:17 +0800 Subject: [PATCH 5/7] add comment for AdvService --- .../tron/core/net/service/adv/AdvService.java | 46 ++++++++++++++----- 1 file changed, 34 insertions(+), 12 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 90f1fb41e8b..1771d9d136e 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -44,6 +44,7 @@ @Slf4j(topic = "net") @Component public class AdvService { + private final int MAX_INV_TO_FETCH_CACHE_SIZE = 100_000; private final int MAX_TRX_CACHE_SIZE = 50_000; private final int MAX_BLOCK_CACHE_SIZE = 10; @@ -162,8 +163,8 @@ public Message getMessage(Item item) { public int fastBroadcastTransaction(TransactionMessage msg) { List peers = tronNetDelegate.getActivePeer().stream() - .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) - .collect(Collectors.toList()); + .filter(peer -> !peer.isNeedSyncFromPeer() && !peer.isNeedSyncFromUs()) + .collect(Collectors.toList()); if (peers.size() == 0) { logger.warn("Broadcast transaction {} failed, no connection", msg.getMessageId()); @@ -179,9 +180,9 @@ public int fastBroadcastTransaction(TransactionMessage msg) { InventoryMessage inventoryMessage = new InventoryMessage(list, InventoryType.TRX); int peersCount = 0; - for (PeerConnection peer: peers) { + for (PeerConnection peer : peers) { if (peer.getAdvInvReceive().getIfPresent(item) == null - && peer.getAdvInvSpread().getIfPresent(item) == null) { + && peer.getAdvInvSpread().getIfPresent(item) == null) { peersCount++; peer.getAdvInvSpread().put(item, Time.getCurrentMillis()); peer.sendMessage(inventoryMessage); @@ -233,12 +234,14 @@ public void broadcast(Message msg) { } public void onDisconnect(PeerConnection peer) { + // Release this peer's retry tracking before looking for another outstanding request. fetchBlockService.onDisconnect(peer); if (!peer.getAdvInvRequest().isEmpty()) { peer.getAdvInvRequest().keySet().forEach(item -> { synchronized (this) { Collection peers = tronNetDelegate.getActivePeer().stream() .filter(p -> !p.equals(peer) && !p.isDisconnect()).collect(Collectors.toList()); + // Another provider may have already delivered the block. if (item.getType() == InventoryType.BLOCK && (blockCache.getIfPresent(item) != null || tronNetDelegate.containBlock(new BlockId(item.getHash())))) { @@ -248,14 +251,17 @@ public void onDisconnect(PeerConnection peer) { Optional pending = peers.stream() .filter(p -> p.getAdvInvRequest().containsKey(item)).findFirst(); if (pending.isPresent()) { + // Reuse the pending request and its original timestamp instead of sending it again. fetchBlockService.fetchBlock(Collections.singletonList(item.getHash()), pending.get()); return; } } if (peers.stream().anyMatch(p -> p.getAdvInvReceive().getIfPresent(item) != null)) { + // Let normal scheduling select a replacement from the remaining advertisers. invToFetch.put(item, System.currentTimeMillis()); } else { + // Allow future announcements to enqueue the item again when no source remains. invToFetch.remove(item); invToFetchCache.invalidate(item); } @@ -263,18 +269,21 @@ public void onDisconnect(PeerConnection peer) { }); } - if (invToFetch.size() > 0) { + if (!invToFetch.isEmpty()) { consumerInvToFetch(); } } private void consumerInvToFetch() { + // Snapshot connected peers; block fetching can still use peers with pending TRX requests. Collection peers = tronNetDelegate.getActivePeer().stream() .filter(peer -> !peer.isDisconnect()) .collect(Collectors.toList()); Collection trxPeers = peers.stream().filter(PeerConnection::isIdle) .collect(Collectors.toList()); + // Batch sizes count requests assigned in this pass, not all outstanding requests on a peer. InvSender invSender = new InvSender(); + // Coordinate queue updates with inventory admission, other consumers, and disconnect recovery. synchronized (this) { if (invToFetch.isEmpty() || peers.isEmpty()) { return; @@ -284,21 +293,28 @@ private void consumerInvToFetch() { invToFetch.forEach((item, time) -> { long timeout = item.getType() == InventoryType.BLOCK ? NetConstants.ADV_TIME_OUT : TIMEOUT; if (time < now - timeout) { + // Release the deduplication entry so a later announcement can enqueue this item again. logger.info("This obj is too late to fetch, type: {} hash: {}", item.getType(), - item.getHash()); + item.getHash()); invToFetch.remove(item); invToFetchCache.invalidate(item); return; } if (item.getType() == InventoryType.BLOCK) { + // Defer block allocation until all queued blocks can be ordered by height. blocks.add(item); return; } - trxPeers.stream().filter(peer -> { - Long t = peer.getAdvInvReceive().getIfPresent(item); - return t != null && now - t < TIMEOUT && invSender.getSize(peer) < MAX_TRX_FETCH_PER_PEER; - }).sorted(Comparator.comparingInt(peer -> invSender.getSize(peer))) - .findFirst().ifPresent(peer -> { + // Use recent advertisers, enforce the per-peer TRX batch limit, and prefer smaller batches. + trxPeers.stream() + .filter(peer -> { + Long t = peer.getAdvInvReceive().getIfPresent(item); + return t != null && now - t < TIMEOUT + && invSender.getSize(peer) < MAX_TRX_FETCH_PER_PEER; + }) + .min(Comparator.comparingInt(invSender::getSize)) + .ifPresent(peer -> { + // Reserve the request before sending; an existing request must not be sent again. if (peer.checkAndPutAdvInvRequest(item, now)) { invSender.add(item, peer); } @@ -309,13 +325,17 @@ private void consumerInvToFetch() { // Reserve peers for earlier blocks before later blocks can make them busy. blocks.sort(Comparator.comparingLong(item -> new BlockId(item.getHash()).getNum())); blocks.forEach(item -> { + // A cached, stored, or already requested block no longer needs an initial fetch. if (blockCache.getIfPresent(item) != null || tronNetDelegate.containBlock(new BlockId(item.getHash())) || peers.stream().anyMatch(peer -> peer.getAdvInvRequest().containsKey(item))) { invToFetch.remove(item); return; } - peers.stream().filter(peer -> fetchBlockService.canFetchBlock(peer, item, now)) + // Recheck eligibility after earlier reservations. Initial fetches prefer smaller batches; + // latency ranking is reserved for backup fetches in FetchBlockService. + peers.stream() + .filter(peer -> fetchBlockService.canFetchBlock(peer, item, now)) .min(Comparator.comparingInt(invSender::getSize)) .ifPresent(peer -> { if (peer.checkAndPutAdvInvRequest(item, now)) { @@ -324,8 +344,10 @@ private void consumerInvToFetch() { } }); }); + // Items without an eligible provider remain queued for a later pass. } + // Send outside the scheduling lock; sendFetch() registers block retry state before sending. invSender.sendFetch(); } From 018b8c4210d3fd56a2863c17d95f6ee2f9f77c64 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Mon, 28 Sep 2026 18:37:03 +0800 Subject: [PATCH 6/7] format AdvService and FetchBlockService --- .../org/tron/core/net/service/adv/AdvService.java | 4 ++-- .../net/service/fetchblock/FetchBlockService.java | 12 +++++------- 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 1771d9d136e..4f37b7bf407 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -377,7 +377,7 @@ private synchronized void consumerInvToSpread() { invSender.sendInv(); } - class InvSender { + private class InvSender { private HashMap>> send = new HashMap<>(); @@ -433,7 +433,7 @@ public void sendInv() { })); } - void sendFetch() { + private void sendFetch() { send.forEach((peer, ids) -> ids.forEach((key, value) -> { if (key.equals(InventoryType.BLOCK)) { value.sort(Comparator.comparingLong(value1 -> new BlockId(value1).getNum())); diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index 8122c150015..a2b0fcdaf10 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -60,7 +60,7 @@ public void close() { } public synchronized void fetchBlock(List sha256HashList, PeerConnection peer) { - if (sha256HashList.size() > 0) { + if (!sha256HashList.isEmpty()) { logger.info("Begin fetch block {} from {}", new BlockCapsule.BlockId(sha256HashList.get(0)).getString(), peer.getInetAddress()); @@ -71,9 +71,10 @@ public synchronized void fetchBlock(List sha256HashList, PeerConnect return; } fetchBlockInfo = null; - sha256HashList.stream().filter(sha256Hash -> new BlockCapsule.BlockId(sha256Hash).getNum() - == headNum + 1) - .findFirst().ifPresent(sha256Hash -> { + sha256HashList.stream() + .filter(sha256Hash -> new BlockCapsule.BlockId(sha256Hash).getNum() == headNum + 1) + .findFirst() + .ifPresent(sha256Hash -> { Long requestTime = peer.getAdvInvRequest().get(new Item(sha256Hash, InventoryType.BLOCK)); if (requestTime != null) { fetchBlockInfo = new FetchBlockInfo(sha256Hash, peer, requestTime); @@ -83,7 +84,6 @@ public synchronized void fetchBlock(List sha256HashList, PeerConnect }); } - public synchronized void blockFetchSuccess(Sha256Hash sha256Hash) { FetchBlockInfo fetchBlockInfoTemp = this.fetchBlockInfo; if (null == fetchBlockInfoTemp || !fetchBlockInfoTemp.getHash().equals(sha256Hash)) { @@ -173,7 +173,5 @@ public FetchBlockInfo(Sha256Hash hash, PeerConnection peer, long time) { this.hash = hash; this.time = time; } - } - } From fa3b81823d6684a6828c24886ddb160b9b0a177f Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Tue, 29 Sep 2026 12:36:36 +0800 Subject: [PATCH 7/7] fix the amplification of fetchblock; set last active time when inventory is confirmed; revert to BAD_BLOCK when processSyncBlock --- .../tron/core/net/P2pEventHandlerImpl.java | 2 +- .../net/messagehandler/BlockMsgHandler.java | 3 +- .../messagehandler/InventoryMsgHandler.java | 2 +- .../tron/core/net/peer/PeerConnection.java | 14 +- .../tron/core/net/service/adv/AdvService.java | 29 ++ .../service/fetchblock/FetchBlockService.java | 15 +- .../core/net/service/sync/SyncService.java | 9 +- .../BlockInventoryActivityTest.java | 408 ++++++++++++++++++ .../core/net/peer/PeerConnectionTest.java | 13 + .../fetchblock/FetchBlockRetryTest.java | 162 ++++++- .../sync/SyncBlockContributionTest.java | 37 +- 11 files changed, 676 insertions(+), 18 deletions(-) create mode 100644 framework/src/test/java/org/tron/core/net/messagehandler/BlockInventoryActivityTest.java diff --git a/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java b/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java index 510bc49dcf2..ae6c4486921 100644 --- a/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java +++ b/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java @@ -261,7 +261,7 @@ private void updateLastInteractiveTime(PeerConnection peer, TronMessage msg) { break; } if (flag) { - peer.setLastInteractiveTime(System.currentTimeMillis()); + peer.updateLastInteractiveTime(System.currentTimeMillis()); } } diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index b30b9f2a2ae..d2df7304dc1 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -148,7 +148,7 @@ private BlockResult processBlock(PeerConnection peer, BlockCapsule block) throws return BlockResult.STATE_FAILED; } - peer.setLastInteractiveTime(System.currentTimeMillis()); + peer.updateLastInteractiveTime(System.currentTimeMillis()); long headNum = tronNetDelegate.getHeadBlockId().getNum(); if (block.getNum() < headNum || tronNetDelegate.containBlock(blockId)) { logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); @@ -175,6 +175,7 @@ private BlockResult processBlock(PeerConnection peer, BlockCapsule block) throws if (tronNetDelegate.isHitDown()) { return BlockResult.IGNORED; } + advService.confirmBlockInventory(blockId); witnessProductBlockService.validWitnessProductTwoBlock(block); Item item = new Item(blockId, InventoryType.BLOCK); diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java index e83563ecfd7..30f1009346b 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/InventoryMsgHandler.java @@ -41,7 +41,7 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep for (Sha256Hash id : inventoryMessage.getHashList()) { Item item = new Item(id, type); - peer.getAdvInvReceive().put(item, System.currentTimeMillis()); + advService.recordInventory(peer, item, System.currentTimeMillis()); advService.addInv(item); } } diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index 4e41a40e410..1c5cb0290b9 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -121,6 +121,11 @@ public class PeerConnection { private Cache advInvReceive = CacheBuilder.newBuilder().maximumSize(invCacheSize) .expireAfterWrite(1, TimeUnit.HOURS).recordStats().build(); + // Eligible block INV receipt times, pending successful block processing before activity updates. + @Getter + private Cache advBlockInvReceive = CacheBuilder.newBuilder().maximumSize(100) + .expireAfterWrite(1, TimeUnit.MINUTES).build(); + @Setter @Getter private Cache advInvSpread = CacheBuilder.newBuilder().maximumSize(invCacheSize) @@ -174,7 +179,7 @@ public void setChannel(Channel channel) { this.isRelayPeer = true; } this.nodeStatistics = TronStatsManager.getNodeStatistics(channel.getInetAddress()); - lastInteractiveTime = System.currentTimeMillis(); + updateLastInteractiveTime(System.currentTimeMillis()); p2pRateLimiter.register(SYNC_BLOCK_CHAIN.asByte(), Args.getInstance().getRateLimiterSyncBlockChain()); p2pRateLimiter.register(FETCH_INV_DATA.asByte(), @@ -188,6 +193,12 @@ public void setBlockBothHave(BlockId blockId) { this.blockBothHaveUpdateTime = System.currentTimeMillis(); } + public synchronized void updateLastInteractiveTime(long time) { + if (time > lastInteractiveTime) { + lastInteractiveTime = time; + } + } + /** * Returns whether there are no outstanding inventory or sync requests. * @@ -247,6 +258,7 @@ public void onDisconnect() { syncService.onDisconnect(this); advService.onDisconnect(this); advInvReceive.invalidateAll(); + advBlockInvReceive.invalidateAll(); advInvSpread.invalidateAll(); advInvRequest.clear(); syncBlockIdCache.invalidateAll(); diff --git a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 4f37b7bf407..f5133306b38 100644 --- a/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java +++ b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java @@ -117,6 +117,35 @@ public synchronized void addInvToCache(Item item) { invToFetch.remove(item); } + public void recordInventory(PeerConnection peer, Item item, long receivedAt) { + if (item.getType() != InventoryType.BLOCK) { + peer.getAdvInvReceive().put(item, receivedAt); + return; + } + synchronized (this) { + // Capture eligibility before fetching; validation may advance the head immediately. + if (peer.getAdvInvSpread().getIfPresent(item) == null + && new BlockId(item.getHash()).getNum() > tronNetDelegate.getHeadBlockId().getNum()) { + peer.getAdvBlockInvReceive().put(item, receivedAt); + } + peer.getAdvInvReceive().put(item, receivedAt); + } + } + + /** + * Confirms eligible announcements only after the block has been successfully processed. + */ + public synchronized void confirmBlockInventory(BlockId blockId) { + Item item = new Item(blockId, InventoryType.BLOCK); + // Share the registration lock so confirmation cannot miss an eligible announcement. + tronNetDelegate.getActivePeer().forEach(peer -> { + Long receivedAt = peer.getAdvBlockInvReceive().asMap().remove(item); + if (receivedAt != null) { + peer.updateLastInteractiveTime(receivedAt); + } + }); + } + public boolean addInv(Item item) { if (fastForward && item.getType().equals(InventoryType.TRX)) { return false; diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index a2b0fcdaf10..953270fb7d5 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -1,5 +1,6 @@ package org.tron.core.net.service.fetchblock; +import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.List; @@ -40,6 +41,8 @@ public class FetchBlockService { private static final double BLOCK_FETCH_LEFT_TIME_PERCENT = 0.5; + private static final int MAX_IN_FLIGHT_REQUESTS_PER_BLOCK = 2; + private final String esName = "fetch-block"; private final ScheduledExecutorService fetchBlockWorkerExecutor = @@ -123,7 +126,15 @@ private synchronized void fetchBlockProcess(FetchBlockInfo fetchBlock) { return; } Item item = new Item(fetchBlock.getHash(), InventoryType.BLOCK); - Optional optionalPeerConnection = tronNetDelegate.getActivePeer().stream() + Collection peers = tronNetDelegate.getActivePeer(); + long pendingRequests = peers.stream() + .filter(peer -> peer.getAdvInvRequest().containsKey(item)) + .count(); + // Keep retry tracking and timeout responsibility while the original and backup are pending. + if (pendingRequests >= MAX_IN_FLIGHT_REQUESTS_PER_BLOCK) { + return; + } + Optional optionalPeerConnection = peers.stream() .filter(peer -> canFetchBlock(peer, item, now)) .min(Comparator.comparingDouble(this::getPeerTop75)); @@ -143,7 +154,7 @@ private boolean shouldFetchBlock(PeerConnection newPeer, FetchBlockInfo fetchBlo double newPeerTop75 = getPeerTop75(newPeer); double oldPeerTop75 = getPeerTop75(fetchBlock.getPeer()); long oldPeerSpendTime = System.currentTimeMillis() - fetchBlock.getTime(); - if (oldPeerTop75 > fetchTimeOut || oldPeerSpendTime >= fetchTimeOut) { + if (oldPeerSpendTime >= fetchTimeOut) { return true; } diff --git a/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java index c873f135118..573a1b0c6a5 100644 --- a/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java +++ b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java @@ -32,6 +32,7 @@ import org.tron.core.net.messagehandler.PbftDataSyncHandler; import org.tron.core.net.peer.PeerConnection; import org.tron.core.net.peer.TronState; +import org.tron.core.net.service.adv.AdvService; import org.tron.protos.Protocol.Inventory.InventoryType; import org.tron.protos.Protocol.ReasonCode; @@ -45,6 +46,9 @@ public class SyncService { @Autowired private PbftDataSyncHandler pbftDataSyncHandler; + @Autowired + private AdvService advService; + private Map blockWaitToProcess = new ConcurrentHashMap<>(); private Map blockJustReceived = new ConcurrentHashMap<>(); @@ -342,7 +346,8 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection) if (tronNetDelegate.isHitDown()) { return; } - peerConnection.setLastInteractiveTime(System.currentTimeMillis()); + advService.confirmBlockInventory(blockId); + peerConnection.updateLastInteractiveTime(System.currentTimeMillis()); if (useful) { peerConnection.setBlockRcvTime(System.currentTimeMillis()); } @@ -374,7 +379,7 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection) syncNext(peer); } } else { - peer.disconnect(ReasonCode.SYNC_FAIL); + peer.disconnect(ReasonCode.BAD_BLOCK); } } } diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/BlockInventoryActivityTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/BlockInventoryActivityTest.java new file mode 100644 index 00000000000..464f37bca06 --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/messagehandler/BlockInventoryActivityTest.java @@ -0,0 +1,408 @@ +package org.tron.core.net.messagehandler; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; + +import com.google.common.base.Ticker; +import java.util.Arrays; +import java.util.Collections; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.mockito.MockedStatic; +import org.tron.common.utils.ReflectUtils; +import org.tron.common.utils.Sha256Hash; +import org.tron.core.capsule.BlockCapsule; +import org.tron.core.capsule.BlockCapsule.BlockId; +import org.tron.core.exception.P2pException; +import org.tron.core.exception.P2pException.TypeEnum; +import org.tron.core.net.PeerBlockTestSupport; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.message.adv.BlockMessage; +import org.tron.core.net.message.adv.InventoryMessage; +import org.tron.core.net.peer.Item; +import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.service.adv.AdvService; +import org.tron.core.net.service.fetchblock.FetchBlockService; +import org.tron.core.net.service.sync.SyncService; +import org.tron.core.services.WitnessProductBlockService; +import org.tron.protos.Protocol.Inventory.InventoryType; + +public class BlockInventoryActivityTest { + + private AdvService adv; + private TronNetDelegate delegate; + private InventoryMsgHandler inventory; + private BlockMsgHandler blocks; + private PeerConnection provider; + private PeerConnection advertiser; + private BlockCapsule block; + private Item item; + private AtomicReference head; + private AtomicBoolean accepted; + + @Before + public void setUp() throws Exception { + adv = spy(new AdvService()); + delegate = mock(TronNetDelegate.class); + inventory = new InventoryMsgHandler(); + blocks = new BlockMsgHandler(); + provider = PeerBlockTestSupport.peer(18888); + advertiser = PeerBlockTestSupport.peer(18889); + block = PeerBlockTestSupport.block(101); + item = new Item(block.getBlockId(), InventoryType.BLOCK); + head = new AtomicReference<>(block.getParentBlockId()); + accepted = new AtomicBoolean(); + provider.getAdvInvRequest().put(item, System.currentTimeMillis()); + when(delegate.getActivePeer()).thenReturn(Arrays.asList(provider, advertiser)); + when(delegate.getHeadBlockId()).thenAnswer(call -> head.get()); + when(delegate.containBlock(block.getParentBlockId())).thenReturn(true); + when(delegate.containBlock(block.getBlockId())).thenAnswer(call -> accepted.get()); + when(delegate.validBlock(block)).thenReturn(true); + doAnswer(call -> { + acceptBlock(); + return null; + }).when(delegate).processBlock(eq(block), anyBoolean()); + ReflectUtils.setFieldValue(adv, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(adv, "fastForward", false); + // Exercise inventory admission independently of provider selection and background workers. + doReturn(false).when(adv).addInv(any()); + ReflectUtils.setFieldValue(inventory, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(inventory, "advService", adv); + ReflectUtils.setFieldValue(inventory, "transactionsMsgHandler", + mock(TransactionsMsgHandler.class)); + ReflectUtils.setFieldValue(blocks, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(blocks, "advService", adv); + ReflectUtils.setFieldValue(blocks, "fetchBlockService", mock(FetchBlockService.class)); + ReflectUtils.setFieldValue(blocks, "syncService", mock(SyncService.class)); + ReflectUtils.setFieldValue(blocks, "witnessProductBlockService", + mock(WitnessProductBlockService.class)); + ReflectUtils.setFieldValue(blocks, "fastForward", false); + } + + @After + public void tearDown() { + adv.close(); + } + + @Test + public void testAnnouncementWaitsForSuccessfulBlockProcessing() throws Exception { + record(100); + doAnswer(call -> { + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertNotNull(adv.getMessage(item)); + acceptBlock(); + return null; + }).when(delegate).processBlock(block, false); + + processBlock(); + + Assert.assertEquals(100, advertiser.getLastInteractiveTime()); + Assert.assertEquals(0, advertiser.getBlockRcvTime()); + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } + + @Test + public void testInventoryIsRegisteredBeforeFetchCanComplete() throws Exception { + AtomicLong receivedAt = new AtomicLong(); + doAnswer(call -> { + Long pending = advertiser.getAdvBlockInvReceive().getIfPresent(item); + Assert.assertNotNull(pending); + receivedAt.set(pending); + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + processBlock(); + return false; + }).when(adv).addInv(item); + + inventory.processMessage(advertiser, inventory(item)); + + Assert.assertEquals(block.getBlockId(), head.get()); + Assert.assertEquals(receivedAt.get(), advertiser.getLastInteractiveTime()); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } + + @Test + public void testEqualAndLowerInventoryCannotRefreshActivity() throws Exception { + head.set(block.getBlockId()); + inventory.processMessage(advertiser, inventory(item)); + Item lower = new Item(block.getParentBlockId(), InventoryType.BLOCK); + inventory.processMessage(advertiser, inventory(lower)); + + adv.confirmBlockInventory(block.getBlockId()); + adv.confirmBlockInventory(block.getParentBlockId()); + + Assert.assertEquals(0, advertiser.getAdvBlockInvReceive().size()); + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + } + + @Test + public void testPreviouslySpreadInventoryCannotRefreshActivity() throws Exception { + advertiser.getAdvInvSpread().put(item, System.currentTimeMillis()); + inventory.processMessage(advertiser, inventory(item)); + + processBlock(); + + Assert.assertEquals(0, advertiser.getAdvBlockInvReceive().size()); + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + } + + @Test + public void testTransactionInventoryCannotRefreshBlockActivity() throws Exception { + inventory.processMessage(advertiser, inventory(new Item(item.getHash(), InventoryType.TRX))); + + processBlock(); + + Assert.assertEquals(0, advertiser.getAdvBlockInvReceive().size()); + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + } + + @Test + public void testSpreadingAfterReceiptDoesNotChangeEligibility() throws Exception { + record(100); + advertiser.getAdvInvSpread().put(item, 200L); + + processBlock(); + + Assert.assertEquals(100, advertiser.getLastInteractiveTime()); + } + + @Test + public void testLateDuplicateCannotReplaceEligibleReceiptTime() throws Exception { + record(100); + doAnswer(call -> { + acceptBlock(); + adv.recordInventory(advertiser, item, 200); + Assert.assertEquals(Long.valueOf(200), advertiser.getAdvInvReceive().getIfPresent(item)); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvBlockInvReceive().getIfPresent(item)); + return null; + }).when(delegate).processBlock(block, false); + + processBlock(); + adv.recordInventory(advertiser, item, 300); + adv.confirmBlockInventory(block.getBlockId()); + + Assert.assertEquals(100, advertiser.getLastInteractiveTime()); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } + + @Test + public void testDifferentHashAtSameHeightCannotConfirmInventory() throws Exception { + record(100); + BlockId different = new BlockId(Sha256Hash.ZERO_HASH, block.getNum()); + Assert.assertNotEquals(block.getBlockId(), different); + + adv.confirmBlockInventory(different); + + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvBlockInvReceive().getIfPresent(item)); + processBlock(); + Assert.assertEquals(100, advertiser.getLastInteractiveTime()); + } + + @Test + public void testBadSignatureDoesNotConfirmInventory() throws Exception { + assertInvalidBlockDoesNotConfirm(TypeEnum.BLOCK_SIGN_INVALID); + } + + @Test + public void testBadMerkleDoesNotConfirmInventory() throws Exception { + assertInvalidBlockDoesNotConfirm(TypeEnum.BLOCK_MERKLE_INVALID); + } + + @Test + public void testInactiveWitnessDoesNotConfirmInventory() throws Exception { + record(100); + when(delegate.validBlock(block)).thenReturn(false); + + processBlock(); + + assertUnconfirmed(); + } + + @Test + public void testExecutionFailureDoesNotConfirmCachedInventory() throws Exception { + record(100); + doThrow(new P2pException(TypeEnum.BAD_BLOCK, "execution failed")) + .when(delegate).processBlock(block, false); + + processBlock(); + + Assert.assertNotNull(adv.getMessage(item)); + assertUnconfirmed(); + } + + @Test + public void testShutdownDoesNotConfirmInventory() throws Exception { + record(100); + when(delegate.isHitDown()).thenReturn(true); + + processBlock(); + + assertUnconfirmed(); + } + + @Test + public void testConfirmationDoesNotOverwriteNewerInteraction() throws Exception { + record(100); + advertiser.updateLastInteractiveTime(200); + + processBlock(); + + Assert.assertEquals(200, advertiser.getLastInteractiveTime()); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } + + @Test + public void testPendingInventoryEvictionDoesNotAffectBlockFetch() throws Exception { + record(100); + Long requestTime = provider.getAdvInvRequest().get(item); + for (int i = 1; i <= 200; i++) { + Item another = new Item(new BlockId(block.getBlockId(), block.getNum() + i), + InventoryType.BLOCK); + adv.recordInventory(advertiser, another, 200); + } + Assert.assertTrue(advertiser.getAdvBlockInvReceive().size() > 0); + Assert.assertTrue(advertiser.getAdvBlockInvReceive().size() <= 100); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvInvReceive().getIfPresent(item)); + Assert.assertEquals(requestTime, provider.getAdvInvRequest().get(item)); + + processBlock(); + + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertEquals(0, advertiser.getBlockRcvTime()); + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertFalse(provider.getAdvInvRequest().containsKey(item)); + } + + @Test + public void testPendingInventoryExpiresAfterOneMinuteWithoutExtendingOnRead() throws Exception { + AtomicLong nanos = new AtomicLong(); + Ticker ticker = new Ticker() { + @Override + public long read() { + return nanos.get(); + } + }; + try (MockedStatic tickers = mockStatic(Ticker.class)) { + tickers.when(Ticker::systemTicker).thenReturn(ticker); + // Construct the production cache with a controlled clock, retaining its actual policy. + advertiser = PeerBlockTestSupport.peer(18890); + when(delegate.getActivePeer()).thenReturn(Arrays.asList(provider, advertiser)); + record(100); + Long requestTime = provider.getAdvInvRequest().get(item); + nanos.set(TimeUnit.MINUTES.toNanos(1) - 1); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvBlockInvReceive().getIfPresent(item)); + nanos.incrementAndGet(); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvInvReceive().getIfPresent(item)); + Assert.assertEquals(requestTime, provider.getAdvInvRequest().get(item)); + + processBlock(); + + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertEquals(0, advertiser.getBlockRcvTime()); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + Assert.assertTrue(provider.getBlockRcvTime() > 0); + Assert.assertFalse(provider.getAdvInvRequest().containsKey(item)); + } + } + + @Test(timeout = 10_000) + public void testConcurrentConfirmationCannotMissEligibleInventory() throws Exception { + CountDownLatch headRead = new CountDownLatch(1); + CountDownLatch finishRegistration = new CountDownLatch(1); + CountDownLatch confirmationStarted = new CountDownLatch(1); + ExecutorService executor = Executors.newFixedThreadPool(2); + when(delegate.getHeadBlockId()).thenAnswer(call -> { + BlockId snapshot = head.get(); + headRead.countDown(); + Assert.assertTrue(finishRegistration.await(5, TimeUnit.SECONDS)); + return snapshot; + }); + try { + Future registration = executor.submit(() -> adv.recordInventory(advertiser, item, 100)); + Assert.assertTrue(headRead.await(5, TimeUnit.SECONDS)); + Assert.assertNull(advertiser.getAdvInvReceive().getIfPresent(item)); + acceptBlock(); + Future confirmation = executor.submit(() -> { + confirmationStarted.countDown(); + adv.confirmBlockInventory(block.getBlockId()); + }); + Assert.assertTrue(confirmationStarted.await(5, TimeUnit.SECONDS)); + try { + confirmation.get(100, TimeUnit.MILLISECONDS); + Assert.fail("Confirmation must wait for the eligible inventory to be registered"); + } catch (TimeoutException expected) { + // The receipt observed the old head; it must be registered before confirmation finishes. + } + finishRegistration.countDown(); + registration.get(5, TimeUnit.SECONDS); + confirmation.get(5, TimeUnit.SECONDS); + + Assert.assertEquals(100, advertiser.getLastInteractiveTime()); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvInvReceive().getIfPresent(item)); + Assert.assertNull(advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } finally { + finishRegistration.countDown(); + executor.shutdownNow(); + Assert.assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private void record(long receivedAt) { + adv.recordInventory(advertiser, item, receivedAt); + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertEquals(Long.valueOf(receivedAt), + advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } + + private void acceptBlock() { + head.set(block.getBlockId()); + accepted.set(true); + } + + private void processBlock() throws Exception { + blocks.processMessage(provider, new BlockMessage(block)); + } + + private InventoryMessage inventory(Item value) { + return new InventoryMessage(Collections.singletonList(value.getHash()), value.getType()); + } + + private void assertInvalidBlockDoesNotConfirm(TypeEnum type) throws Exception { + record(100); + when(delegate.validBlock(block)).thenThrow(new P2pException(type, "invalid block")); + try { + processBlock(); + Assert.fail("Expected block validation to fail"); + } catch (P2pException e) { + Assert.assertEquals(type, e.getType()); + } + assertUnconfirmed(); + } + + private void assertUnconfirmed() { + Assert.assertEquals(1, advertiser.getLastInteractiveTime()); + Assert.assertEquals(0, advertiser.getBlockRcvTime()); + Assert.assertEquals(Long.valueOf(100), advertiser.getAdvBlockInvReceive().getIfPresent(item)); + } +} diff --git a/framework/src/test/java/org/tron/core/net/peer/PeerConnectionTest.java b/framework/src/test/java/org/tron/core/net/peer/PeerConnectionTest.java index cc30fb70b0b..81bd9fa70dc 100644 --- a/framework/src/test/java/org/tron/core/net/peer/PeerConnectionTest.java +++ b/framework/src/test/java/org/tron/core/net/peer/PeerConnectionTest.java @@ -70,6 +70,8 @@ public void testOnDisconnect() { Long time = System.currentTimeMillis(); BlockCapsule.BlockId blockId = new BlockCapsule.BlockId(); peerConnection.getAdvInvReceive().put(item, time); + peerConnection.getAdvBlockInvReceive().put( + new Item(blockId, Protocol.Inventory.InventoryType.BLOCK), time); peerConnection.getAdvInvSpread().put(item, time); peerConnection.getSyncBlockIdCache().put(item.getHash(), time); peerConnection.getSyncBlockToFetch().add(blockId); @@ -79,6 +81,7 @@ public void testOnDisconnect() { peerConnection.onDisconnect(); Assert.assertEquals(0, peerConnection.getAdvInvReceive().size()); + Assert.assertEquals(0, peerConnection.getAdvBlockInvReceive().size()); Assert.assertEquals(0, peerConnection.getAdvInvSpread().size()); Assert.assertEquals(0, peerConnection.getSyncBlockIdCache().size()); Assert.assertEquals(0, peerConnection.getSyncBlockToFetch().size()); @@ -86,6 +89,16 @@ public void testOnDisconnect() { Assert.assertEquals(0, peerConnection.getSyncBlockInProcess().size()); } + @Test + public void testLastInteractiveTimeOnlyAdvances() { + PeerConnection peer = new PeerConnection(); + peer.updateLastInteractiveTime(200); + peer.updateLastInteractiveTime(100); + Assert.assertEquals(200, peer.getLastInteractiveTime()); + peer.updateLastInteractiveTime(300); + Assert.assertEquals(300, peer.getLastInteractiveTime()); + } + @Test public void testIsIdle() { PeerConnection peerConnection = new PeerConnection(); diff --git a/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java index a07d849720c..27595669d5e 100644 --- a/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java +++ b/framework/src/test/java/org/tron/core/net/service/fetchblock/FetchBlockRetryTest.java @@ -200,6 +200,153 @@ public void testBackupFetchStillPrefersLowerLatency() throws Exception { Assert.assertSame(faster, ReflectUtils.getFieldObject(state(), "peer")); } + @Test + public void testBackupRequestsAreLimitedPerBlock() throws Exception { + for (int i = 0; i < 4; i++) { + peers.add(PeerBlockTestSupport.peer(18890 + i)); + } + peers.forEach(peer -> peer.getAdvInvReceive().put(item, System.currentTimeMillis())); + beginAgedFetch(2_000); + Long originalRequestTime = first.getAdvInvRequest().get(item); + + try (MockedStatic metrics = mockBlockLatencies(peers, + 2_000L, 2_000L, 2_000L, 2_000L, 2_000L, 2_000L)) { + tick(state()); + Assert.assertSame(second, ReflectUtils.getFieldObject(state(), "peer")); + ageState(2_000); + Object backupState = state(); + Long backupRequestTime = second.getAdvInvRequest().get(item); + for (int i = 0; i < 10; i++) { + // An eligible third provider and an expired short timeout must not bypass the cap. + Assert.assertTrue(fetch.canFetchBlock(peers.get(2), item, System.currentTimeMillis())); + tick(state()); + } + Assert.assertSame(backupState, state()); + Assert.assertEquals(backupRequestTime, second.getAdvInvRequest().get(item)); + } + + assertBlockRequests(second, item); + verify(first, never()).sendMessage(any()); + for (int i = 2; i < peers.size(); i++) { + verify(peers.get(i), never()).sendMessage(any()); + } + Assert.assertEquals(2L, peers.stream() + .filter(peer -> peer.getAdvInvRequest().containsKey(item)).count()); + Assert.assertEquals(originalRequestTime, first.getAdvInvRequest().get(item)); + } + + @Test + public void testSlowBackupWaitsForActualTimeout() throws Exception { + // Leave enough time for assertions without relying on sleeps or a sub-second test run. + ReflectUtils.setFieldValue(fetch, "fetchTimeOut", 10_000L); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + + try (MockedStatic metrics = mockBlockLatencies(first, 15_000L, second, 15_000L)) { + beginAgedFetch(100); + Object originalState = state(); + tick(state()); + + Assert.assertSame(originalState, state()); + Assert.assertFalse(second.getAdvInvRequest().containsKey(item)); + verify(second, never()).sendMessage(any()); + + ageState(11_000); + Long originalRequestTime = first.getAdvInvRequest().get(item); + tick(state()); + + assertBlockRequests(second, item); + Assert.assertSame(second, ReflectUtils.getFieldObject(state(), "peer")); + Assert.assertEquals(originalRequestTime, first.getAdvInvRequest().get(item)); + } + } + + @Test + public void testFasterBackupCanStartBeforeTimeout() throws Exception { + ReflectUtils.setFieldValue(fetch, "fetchTimeOut", 10_000L); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + + try (MockedStatic metrics = mockBlockLatencies(first, 8_000L, second, 100L)) { + beginAgedFetch(100); + Long originalRequestTime = first.getAdvInvRequest().get(item); + tick(state()); + + assertBlockRequests(second, item); + Assert.assertTrue(System.currentTimeMillis() - originalRequestTime < 10_000); + Assert.assertEquals(originalRequestTime, first.getAdvInvRequest().get(item)); + Assert.assertSame(second, ReflectUtils.getFieldObject(state(), "peer")); + } + } + + @Test + public void testOtherRequestsDoNotUseBlockBackupLimit() throws Exception { + PeerConnection otherProvider = PeerBlockTestSupport.peer(18890); + peers.add(otherProvider); + Item otherBlock = new Item(block.getParentBlockId(), InventoryType.BLOCK); + Item transaction = new Item(item.getHash(), InventoryType.TRX); + long now = System.currentTimeMillis(); + otherProvider.getAdvInvRequest().put(otherBlock, now); + second.getAdvInvRequest().put(transaction, now); + second.getAdvInvReceive().put(item, now); + beginAgedFetch(2_000); + + tick(state()); + + assertBlockRequests(second, item); + Assert.assertTrue(second.getAdvInvRequest().containsKey(item)); + Assert.assertEquals(Long.valueOf(now), second.getAdvInvRequest().get(transaction)); + Assert.assertEquals(Long.valueOf(now), otherProvider.getAdvInvRequest().get(otherBlock)); + verify(otherProvider, never()).sendMessage(any()); + } + + @Test + public void testOriginalDisconnectAllowsReplacementBackup() throws Exception { + PeerConnection replacement = PeerBlockTestSupport.peer(18890); + peers.add(replacement); + beginAgedFetch(2_000); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + ageState(2_000); + Long backupRequestTime = second.getAdvInvRequest().get(item); + replacement.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + verify(replacement, never()).sendMessage(any()); + + disconnect(first); + tick(state()); + + assertBlockRequests(second, item); + assertBlockRequests(replacement, item); + Assert.assertEquals(backupRequestTime, second.getAdvInvRequest().get(item)); + Assert.assertTrue(replacement.getAdvInvRequest().containsKey(item)); + Assert.assertSame(replacement, ReflectUtils.getFieldObject(state(), "peer")); + } + + @Test + public void testBackupDisconnectAllowsReplacementBackup() throws Exception { + PeerConnection replacement = PeerBlockTestSupport.peer(18890); + peers.add(replacement); + beginAgedFetch(2_000); + Long originalRequestTime = first.getAdvInvRequest().get(item); + second.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + ageState(2_000); + replacement.getAdvInvReceive().put(item, System.currentTimeMillis()); + tick(state()); + verify(replacement, never()).sendMessage(any()); + + disconnect(second); + Assert.assertSame(first, ReflectUtils.getFieldObject(state(), "peer")); + Assert.assertEquals(originalRequestTime, ReflectUtils.getFieldObject(state(), "time")); + tick(state()); + + assertBlockRequests(second, item); + assertBlockRequests(replacement, item); + verify(first, never()).sendMessage(any()); + Assert.assertEquals(originalRequestTime, first.getAdvInvRequest().get(item)); + Assert.assertTrue(replacement.getAdvInvRequest().containsKey(item)); + Assert.assertSame(replacement, ReflectUtils.getFieldObject(state(), "peer")); + } + @Test public void testQueuedBlocksAreFetchedInHeightOrder() throws Exception { List blocks = queueBlocksInReverseHeightOrder(); @@ -501,11 +648,15 @@ private void assertBlockRequests(PeerConnection peer, Item... blocks) { private MockedStatic mockBlockLatencies(PeerConnection left, Long leftLatency, PeerConnection right, Long rightLatency) { + return mockBlockLatencies(Arrays.asList(left, right), leftLatency, rightLatency); + } + + private MockedStatic mockBlockLatencies(List candidates, + Long... latencies) { + Assert.assertEquals(candidates.size(), latencies.length); MetricRegistry registry = new MetricRegistry(); - PeerConnection[] candidates = {left, right}; - Long[] latencies = {leftLatency, rightLatency}; - for (int i = 0; i < candidates.length; i++) { - PeerConnection peer = candidates[i]; + for (int i = 0; i < candidates.size(); i++) { + PeerConnection peer = candidates.get(i); InetSocketAddress address = new InetSocketAddress("127.0.0." + (i + 2), peer.getInetSocketAddress().getPort()); when(peer.getChannel().getInetAddress()).thenReturn(address.getAddress()); @@ -524,7 +675,8 @@ private MockedStatic mockBlockLatencies(PeerConnection left, Long l private void ageState(long age) { long time = System.currentTimeMillis() - age; ReflectUtils.setFieldValue(state(), "time", time); - first.getAdvInvRequest().put(item, time); + PeerConnection provider = (PeerConnection) ReflectUtils.getFieldObject(state(), "peer"); + provider.getAdvInvRequest().put(item, time); } private Object state() { diff --git a/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java index ec3c4b9cbaa..9216d93e7d7 100644 --- a/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java +++ b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java @@ -20,12 +20,16 @@ import org.tron.core.net.PeerBlockTestSupport; import org.tron.core.net.TronNetDelegate; import org.tron.core.net.messagehandler.PbftDataSyncHandler; +import org.tron.core.net.peer.Item; import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.service.adv.AdvService; +import org.tron.protos.Protocol.Inventory.InventoryType; import org.tron.protos.Protocol.ReasonCode; public class SyncBlockContributionTest { private SyncService sync; + private AdvService adv; private TronNetDelegate delegate; private PeerConnection provider; private PeerConnection other; @@ -34,6 +38,7 @@ public class SyncBlockContributionTest { @Before public void setUp() { sync = new SyncService(); + adv = new AdvService(); delegate = mock(TronNetDelegate.class); provider = PeerBlockTestSupport.peer(18888); other = PeerBlockTestSupport.peer(18889); @@ -43,12 +48,15 @@ public void setUp() { when(delegate.getActivePeer()).thenReturn(Arrays.asList(provider, other)); when(delegate.getHeadBlockId()).thenReturn(block.getParentBlockId()); ReflectUtils.setFieldValue(sync, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(adv, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(sync, "advService", adv); ReflectUtils.setFieldValue(sync, "pbftDataSyncHandler", mock(PbftDataSyncHandler.class)); } @After public void tearDown() { sync.close(); + adv.close(); } @Test @@ -62,6 +70,7 @@ public void testOnlyValidatedProviderGetsTimestamps() throws Exception { @Test public void testInvalidSignatureOnlyBlamesProvider() throws Exception { + adv.recordInventory(other, new Item(block.getBlockId(), InventoryType.BLOCK), 100); doThrow(new P2pException(TypeEnum.BLOCK_SIGN_INVALID, "bad signature")) .when(delegate).validSignature(block); process(); @@ -69,27 +78,45 @@ public void testInvalidSignatureOnlyBlamesProvider() throws Exception { verify(other, never()).disconnect(any()); Assert.assertEquals(1, provider.getLastInteractiveTime()); Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertEquals(1, other.getLastInteractiveTime()); } @Test - public void testStateFailureDoesNotBanPeersOrImproveTimestamps() throws Exception { + public void testStateFailureUsesBadBlockWithoutImprovingTimestamps() throws Exception { + adv.recordInventory(other, new Item(block.getBlockId(), InventoryType.BLOCK), 100); doThrow(new P2pException(TypeEnum.BAD_BLOCK, "state failure")) .when(delegate).processBlock(block, true); process(); - verify(provider).disconnect(ReasonCode.SYNC_FAIL); - verify(other).disconnect(ReasonCode.SYNC_FAIL); - verify(provider, never()).disconnect(ReasonCode.BAD_BLOCK); - verify(other, never()).disconnect(ReasonCode.BAD_BLOCK); + verify(provider).disconnect(ReasonCode.BAD_BLOCK); + verify(other).disconnect(ReasonCode.BAD_BLOCK); + verify(provider, never()).disconnect(ReasonCode.SYNC_FAIL); + verify(other, never()).disconnect(ReasonCode.SYNC_FAIL); Assert.assertEquals(1, provider.getLastInteractiveTime()); Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertEquals(1, other.getLastInteractiveTime()); } @Test public void testShutdownDoesNotImproveTimestamps() throws Exception { + adv.recordInventory(other, new Item(block.getBlockId(), InventoryType.BLOCK), 100); when(delegate.isHitDown()).thenReturn(true); process(); Assert.assertEquals(1, provider.getLastInteractiveTime()); Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertEquals(1, other.getLastInteractiveTime()); + } + + @Test + public void testSyncBlockConfirmsEligibleInventoryWithoutContribution() throws Exception { + Item item = new Item(block.getBlockId(), InventoryType.BLOCK); + adv.recordInventory(other, item, 100); + Assert.assertEquals(1, other.getLastInteractiveTime()); + + process(); + + Assert.assertEquals(100, other.getLastInteractiveTime()); + Assert.assertEquals(0, other.getBlockRcvTime()); + Assert.assertNull(other.getAdvBlockInvReceive().getIfPresent(item)); } @Test