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/service/sync/SyncService.java b/framework/src/main/java/org/tron/core/net/service/sync/SyncService.java index bd656d9c41e..1200d2e8a8a 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: {}", 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/service/sync/SyncBlockContributionTest.java b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java new file mode 100644 index 00000000000..c9acda29d7d --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/service/sync/SyncBlockContributionTest.java @@ -0,0 +1,118 @@ +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()); + Assert.assertEquals(1, other.getLastInteractiveTime()); + } + + @Test + public void testStateFailureUsesBadBlockWithoutImprovingTimestamps() throws Exception { + 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); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertEquals(1, other.getLastInteractiveTime()); + } + + @Test + public void testShutdownDoesNotImproveTimestamps() throws Exception { + when(delegate.isHitDown()).thenReturn(true); + process(); + Assert.assertEquals(1, provider.getLastInteractiveTime()); + Assert.assertEquals(0, provider.getBlockRcvTime()); + Assert.assertEquals(1, other.getLastInteractiveTime()); + } + + @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); + } +}