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..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 @@ -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; @@ -189,10 +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(); } + /** + * 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/adv/AdvService.java b/framework/src/main/java/org/tron/core/net/service/adv/AdvService.java index 2241ce8d1ce..7bc7fcdc904 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 @@ -26,6 +26,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; @@ -268,36 +269,79 @@ public void onDisconnect(PeerConnection peer) { } private void consumerInvToFetch() { + // Snapshot connected peers; block fetching can still use peers with pending TRX requests. 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()); + // 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; } long now = System.currentTimeMillis(); + List blocks = new ArrayList<>(); invToFetch.forEach((item, time) -> { - if (time < now - TIMEOUT) { + 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; } - peers.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 -> { + if (item.getType() == InventoryType.BLOCK) { + // Defer block allocation until all queued blocks can be ordered by height. + blocks.add(item); + return; + } + // 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); } 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 -> { + // 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; + } + // 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)) { + invSender.add(item, peer); + invToFetch.remove(item); + } + }); + }); + // Items without an eligible provider remain queued for a later pass. } + // Send outside the scheduling lock. 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 bda2646abbc..6a55134385f 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 @@ -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; @@ -88,17 +89,23 @@ public void blockFetchSuccess(Sha256Hash sha256Hash) { this.fetchBlockInfo = null; } + 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 void fetchBlockProcess(FetchBlockInfo fetchBlock) { if (null == fetchBlock) { return; } + long now = System.currentTimeMillis(); 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) + .filter(peer -> canFetchBlock(peer, item, now)) .min(Comparator.comparingDouble(this::getPeerTop75)); if (optionalPeerConnection.isPresent()) { @@ -123,7 +130,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/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/peer/PeerBlockIdleTest.java b/framework/src/test/java/org/tron/core/net/peer/PeerBlockIdleTest.java new file mode 100644 index 00000000000..69d4aef4fe7 --- /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.isBlockFetchIdle()); + + Item block = new Item(new BlockId(), InventoryType.BLOCK); + peer.getAdvInvRequest().put(block, System.currentTimeMillis()); + Assert.assertFalse(peer.isBlockFetchIdle()); + peer.getAdvInvRequest().remove(block); + Assert.assertTrue(peer.isBlockFetchIdle()); + } + + @Test + public void testSyncRequestAndProcessingExcludeBlockProvider() { + PeerConnection peer = new PeerConnection(); + peer.getSyncBlockRequested().put(new BlockId(), System.currentTimeMillis()); + Assert.assertFalse(peer.isBlockFetchIdle()); + peer.getSyncBlockRequested().clear(); + peer.setSyncChainRequested(new Pair<>(new ArrayDeque<>(), System.currentTimeMillis())); + Assert.assertFalse(peer.isBlockFetchIdle()); + peer.setSyncChainRequested(null); + peer.getSyncBlockInProcess().add(new BlockId()); + Assert.assertFalse(peer.isBlockFetchIdle()); + peer.getSyncBlockInProcess().clear(); + Assert.assertTrue(peer.isBlockFetchIdle()); + } +} diff --git a/framework/src/test/java/org/tron/core/net/service/fetchblock/BlockFetchSchedulingTest.java b/framework/src/test/java/org/tron/core/net/service/fetchblock/BlockFetchSchedulingTest.java new file mode 100644 index 00000000000..7421ef7e902 --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/service/fetchblock/BlockFetchSchedulingTest.java @@ -0,0 +1,376 @@ +package org.tron.core.net.service.fetchblock; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +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; +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.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.FetchInvDataMessage; +import org.tron.core.net.message.adv.InventoryMessage; +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.service.adv.AdvService; +import org.tron.protos.Protocol.Inventory.InventoryType; + +public class BlockFetchSchedulingTest { + + 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 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 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)); + } + + @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.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)); + } + } + + @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.assertFalse(fetch.canFetchBlock(first, item, System.currentTimeMillis())); + Assert.assertTrue(fetch.canFetchBlock(second, item, System.currentTimeMillis())); + second.setNeedSyncFromPeer(true); + 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.canFetchBlock(second, item, System.currentTimeMillis())); + } + + @Test + 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.canFetchBlock(first, item, System.currentTimeMillis())); + } + + private InventoryMessage inventory() { + return new InventoryMessage(Collections.singletonList(item.getHash()), InventoryType.BLOCK); + } + + private void beginAgedFetch(long age) { + first.getAdvInvRequest().put(item, System.currentTimeMillis() - age); + fetch.fetchBlock(Collections.singletonList(item.getHash()), first); + Assert.assertNotNull(state()); + // The base registers tracking at send time; control elapsed time independently of retries. + ageState(age); + } + + 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) { + 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(); + 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()); + 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); + PeerConnection provider = (PeerConnection) ReflectUtils.getFieldObject(state(), "peer"); + provider.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); + } + +}