Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 24 additions & 2 deletions framework/src/main/java/org/tron/core/net/peer/PeerConnection.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -158,7 +157,7 @@ public class PeerConnection {
private volatile Pair<Deque<BlockId>, Long> syncChainRequested = null;
@Setter
@Getter
private Set<BlockId> syncBlockInProcess = new HashSet<>();
private Set<BlockId> syncBlockInProcess = ConcurrentHashMap.newKeySet();
@Setter
@Getter
private volatile boolean needSyncFromPeer = true;
Expand Down Expand Up @@ -189,10 +188,33 @@ public void setBlockBothHave(BlockId blockId) {
this.blockBothHaveUpdateTime = System.currentTimeMillis();
}

/**
* Returns whether there are no outstanding inventory or sync requests.
*
* <p>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.
*
* <p>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.
*
* <p>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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<PeerConnection> peers = tronNetDelegate.getActivePeer().stream()
.filter(peer -> peer.isIdle())
.filter(peer -> !peer.isDisconnect())
.collect(Collectors.toList());
Collection<PeerConnection> 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<Item> 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();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<PeerConnection> 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()) {
Expand All @@ -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;
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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;
}
}
Original file line number Diff line number Diff line change
@@ -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());
}
}
Loading
Loading