From d90a81da419ad2af31aff84dbe803fe04273d306 Mon Sep 17 00:00:00 2001 From: 317787106 <317787106@qq.com> Date: Tue, 29 Sep 2026 13:30:59 +0800 Subject: [PATCH] fix(net): confirm block inventory activity after processing --- .../tron/core/net/P2pEventHandlerImpl.java | 2 +- .../net/messagehandler/BlockMsgHandler.java | 3 + .../messagehandler/InventoryMsgHandler.java | 9 +- .../tron/core/net/peer/PeerConnection.java | 14 +- .../tron/core/net/service/adv/AdvService.java | 29 ++ .../core/net/service/sync/SyncService.java | 7 + .../tron/core/net/PeerBlockTestSupport.java | 44 ++ .../BlockInventoryActivityTest.java | 408 ++++++++++++++++++ .../core/net/peer/PeerConnectionTest.java | 13 + .../sync/SyncBlockInventoryActivityTest.java | 120 ++++++ 10 files changed, 639 insertions(+), 10 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/BlockInventoryActivityTest.java create mode 100644 framework/src/test/java/org/tron/core/net/service/sync/SyncBlockInventoryActivityTest.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..ffa609acabe 100644 --- a/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java +++ b/framework/src/main/java/org/tron/core/net/P2pEventHandlerImpl.java @@ -262,7 +262,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 452209d575f..8717f669d21 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 @@ -152,6 +152,9 @@ private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pExc try { tronNetDelegate.processBlock(block, false); + if (!tronNetDelegate.isHitDown()) { + advService.confirmBlockInventory(blockId); + } peer.setBlockRcvTime(System.currentTimeMillis()); witnessProductBlockService.validWitnessProductTwoBlock(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 f96f7f0b0ff..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 @@ -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; @@ -42,14 +41,8 @@ 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); - 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..1c5adfaf626 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 @@ -122,6 +122,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) @@ -175,7 +180,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(), @@ -189,6 +194,12 @@ public void setBlockBothHave(BlockId blockId) { this.blockBothHaveUpdateTime = System.currentTimeMillis(); } + public synchronized void updateLastInteractiveTime(long time) { + if (time > lastInteractiveTime) { + lastInteractiveTime = time; + } + } + public boolean isIdle() { return advInvRequest.isEmpty() && isSyncIdle(); } @@ -225,6 +236,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 2241ce8d1ce..73c6d68b888 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 @@ -113,6 +113,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/sync/SyncService.java b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java index bd656d9c41e..4c495911153 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<>(); @@ -337,6 +341,9 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection) try { tronNetDelegate.validSignature(block); tronNetDelegate.processBlock(block, true); + if (!tronNetDelegate.isHitDown()) { + advService.confirmBlockInventory(blockId); + } peerConnection.setBlockRcvTime(System.currentTimeMillis()); pbftDataSyncHandler.processPBFTCommitData(block); } catch (P2pException p2pException) { 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/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/sync/SyncBlockInventoryActivityTest.java b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockInventoryActivityTest.java new file mode 100644 index 00000000000..5076ca4536e --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockInventoryActivityTest.java @@ -0,0 +1,120 @@ +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.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 SyncBlockInventoryActivityTest { + + private SyncService sync; + private AdvService adv; + private TronNetDelegate delegate; + private PeerConnection provider; + private PeerConnection other; + private BlockCapsule block; + + @Before + public void setUp() { + sync = new SyncService(); + adv = new AdvService(); + 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(adv, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(sync, "advService", adv); + ReflectUtils.setFieldValue(sync, "pbftDataSyncHandler", mock(PbftDataSyncHandler.class)); + } + + @After + public void tearDown() { + sync.close(); + adv.close(); + } + + @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(); + verify(provider).disconnect(ReasonCode.BAD_BLOCK); + verify(other, never()).disconnect(any()); + assertUnconfirmed(); + } + + @Test + public void testStateFailureDoesNotConfirmInventory() 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.BAD_BLOCK); + verify(other).disconnect(ReasonCode.BAD_BLOCK); + verify(provider, never()).disconnect(ReasonCode.SYNC_FAIL); + verify(other, never()).disconnect(ReasonCode.SYNC_FAIL); + assertUnconfirmed(); + } + + @Test + public void testShutdownDoesNotConfirmInventory() throws Exception { + adv.recordInventory(other, new Item(block.getBlockId(), InventoryType.BLOCK), 100); + when(delegate.isHitDown()).thenReturn(true); + process(); + assertUnconfirmed(); + } + + @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)); + } + + private void assertUnconfirmed() { + Assert.assertEquals(1, other.getLastInteractiveTime()); + Assert.assertEquals(0, other.getBlockRcvTime()); + Assert.assertEquals(Long.valueOf(100), other.getAdvBlockInvReceive() + .getIfPresent(new Item(block.getBlockId(), InventoryType.BLOCK))); + } + + private void process() throws Exception { + Method method = SyncService.class.getDeclaredMethod("processSyncBlock", + BlockCapsule.class, PeerConnection.class); + method.setAccessible(true); + method.invoke(sync, block, provider); + } +}