Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import org.apache.commons.collections4.CollectionUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.tron.common.utils.Pair;
import org.tron.core.capsule.BlockCapsule.BlockId;
import org.tron.core.config.Parameter.ChainConstant;
import org.tron.core.config.Parameter.NetConstants;
Expand All @@ -18,6 +19,7 @@
import org.tron.core.exception.P2pException.TypeEnum;
import org.tron.core.net.TronNetDelegate;
import org.tron.core.net.message.TronMessage;
import org.tron.core.net.message.handshake.HelloMessage;
import org.tron.core.net.message.sync.ChainInventoryMessage;
import org.tron.core.net.peer.PeerConnection;
import org.tron.core.net.peer.TronState;
Expand All @@ -40,7 +42,8 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep

ChainInventoryMessage chainInventoryMessage = (ChainInventoryMessage) msg;

check(peer, chainInventoryMessage);
Pair<Deque<BlockId>, Long> requested = peer.getSyncChainRequested();
check(peer, requested, chainInventoryMessage);

peer.setFetchAble(false);

Expand All @@ -51,6 +54,7 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep
Deque<BlockId> blockIdWeGet = new LinkedList<>(chainInventoryMessage.getBlockIds());

if (blockIdWeGet.size() == 1 && tronNetDelegate.containBlock(blockIdWeGet.peek())) {
peer.setRemainNum(0);
peer.setTronState(TronState.SYNC_COMPLETED);
peer.setNeedSyncFromPeer(false);
Comment thread
3for marked this conversation as resolved.
return;
Expand Down Expand Up @@ -98,11 +102,17 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep
}
}

private void check(PeerConnection peer, ChainInventoryMessage msg) throws P2pException {
if (peer.getSyncChainRequested() == null) {
private void check(PeerConnection peer, Pair<Deque<BlockId>, Long> requested,
ChainInventoryMessage msg) throws P2pException {
if (requested == null) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "not send syncBlockChainMsg");
}

HelloMessage hello = peer.getHelloMessageReceive();
if (hello == null) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "hello message not received");
}

List<BlockId> blockIds = msg.getBlockIds();
if (CollectionUtils.isEmpty(blockIds)) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "blockIds is empty");
Expand All @@ -112,7 +122,8 @@ private void check(PeerConnection peer, ChainInventoryMessage msg) throws P2pExc
throw new P2pException(TypeEnum.BAD_MESSAGE, "big blockIds size: " + blockIds.size());
}

if (msg.getRemainNum() != 0 && blockIds.size() < NetConstants.SYNC_FETCH_BATCH_NUM) {
if (msg.getRemainNum() < 0
|| (msg.getRemainNum() != 0 && blockIds.size() < NetConstants.SYNC_FETCH_BATCH_NUM)) {
throw new P2pException(TypeEnum.BAD_MESSAGE,
"remain: " + msg.getRemainNum() + ", blockIds size: " + blockIds.size());
}
Expand All @@ -124,9 +135,9 @@ private void check(PeerConnection peer, ChainInventoryMessage msg) throws P2pExc
}
}

if (!peer.getSyncChainRequested().getKey().contains(blockIds.get(0))) {
if (!requested.getKey().contains(blockIds.get(0))) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "unlinked block, my head: "
+ peer.getSyncChainRequested().getKey().getLast().getString()
+ requested.getKey().getLast().getString()
+ ", peer: " + blockIds.get(0).getString());
}

Expand All @@ -137,11 +148,23 @@ private void check(PeerConnection peer, ChainInventoryMessage msg) throws P2pExc
long maxFutureNum =
maxRemainTime / BLOCK_PRODUCED_INTERVAL + tronNetDelegate.getSolidBlockId().getNum();
long lastNum = blockIds.get(blockIds.size() - 1).getNum();
if (lastNum + msg.getRemainNum() > maxFutureNum) {
if (lastNum > maxFutureNum || msg.getRemainNum() > maxFutureNum - lastNum) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "lastNum: " + lastNum + " + remainNum: "
+ msg.getRemainNum() + " > futureMaxNum: " + maxFutureNum);
}
}

if (msg.getRemainNum() == 0) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[SHOULD] Complete synchronization for terminal inventories containing only known blocks.

This guard still accepts a multi-block terminal response once its last block reaches the HELLO height. If every returned block is already known, processMessage() drains syncBlockToFetch but falls through to syncNext() at lines 95–101, leaving needSyncFromPeer=true.

A peer can repeatedly return a summary-linked suffix of known blocks with remainNum=0. Each round refreshes both the shared-block timestamp and the inventory-request timestamp, and rebuilds the summary under forkLock. For prolonged operation, the peer can periodically advance the suffix to blocks already obtained from honest peers before the summary’s lower bound overtakes it.

Please complete synchronization when remainNum == 0 && syncBlockToFetch.isEmpty() after removing known blocks, preserving needSyncFromUs. The existing testKnownMultiBlockResponseRequestsNextSummary currently asserts the redundant continuation and should instead assert completion and no further syncNext() call.

long lastNum = blockIds.get(blockIds.size() - 1).getNum();
long helloHeadNum = hello.getHeadBlockId().getNum();
if (lastNum < helloHeadNum) {
// In rare cases, a fork rollback can put an honest peer below its HELLO head.
// Accept the transient SYNC_FAIL disconnect; the peer can reconnect with a fresh
// HELLO after the default one-minute cooldown.
throw new P2pException(TypeEnum.SYNC_FAILED,
"lastNum " + lastNum + " (remainNum=0) is below hello head " + helloHeadNum);
}
}
}

}
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
package org.tron.core.net.peer;

import java.util.Deque;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.tron.common.es.ExecutorServiceManager;
import org.tron.common.utils.Pair;
import org.tron.core.capsule.BlockCapsule.BlockId;
import org.tron.core.config.Parameter.NetConstants;
import org.tron.core.net.TronNetDelegate;
import org.tron.protos.Protocol.ReasonCode;
Expand Down Expand Up @@ -64,6 +67,15 @@ public void statusCheck() {
}
}

if (!isDisconnected) {
Pair<Deque<BlockId>, Long> requested = peer.getSyncChainRequested();
isDisconnected = requested != null
&& requested.getValue() < now - NetConstants.SYNC_TIME_OUT;
if (isDisconnected) {
logger.warn("Peer {} get chain inventory timeout", peer.getInetAddress());
}
}

if (!isDisconnected) {
isDisconnected = peer.getSyncBlockRequested().values().stream()
.anyMatch(time -> time < now - NetConstants.SYNC_TIME_OUT);
Expand Down
45 changes: 45 additions & 0 deletions framework/src/test/java/org/tron/core/net/PeerSyncTestSupport.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
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 java.net.InetSocketAddress;
import org.tron.common.utils.ReflectUtils;
import org.tron.common.utils.Sha256Hash;
import org.tron.core.capsule.BlockCapsule.BlockId;
import org.tron.core.net.message.handshake.HelloMessage;
import org.tron.core.net.peer.PeerConnection;
import org.tron.p2p.connection.Channel;
import org.tron.protos.Protocol;

public final class PeerSyncTestSupport {

private PeerSyncTestSupport() {
}

public static BlockId blockId(long number) {
return new BlockId(Sha256Hash.ZERO_HASH, number);
}

public static HelloMessage helloMessage(long headNum) throws Exception {
return new HelloMessage(Protocol.HelloMessage.newBuilder()
.setHeadBlockId(Protocol.HelloMessage.BlockId.newBuilder()
.setHash(blockId(headNum).getByteString()).setNumber(headNum))
.build().toByteArray());
}

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);
doNothing().when(peer).sendMessage(any());
doNothing().when(peer).disconnect(any());
return peer;
}
}
Loading
Loading