Skip to content
Open
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
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 on lines 56 to 59

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] A single known block can still prematurely terminate synchronization despite being inconsistent with the peer’s advertised height.

If the peer advertises head=200 in HELLO and the local node sends summary [0,50,100], a response containing only [50] with remainNum=0 still passes check(). It then clears syncChainRequested, sets SYNC_COMPLETED, and sets needSyncFromPeer=false. When needSyncFromUs=false, isSyncFinish() also becomes true, so the connection immediately exits both the response-timeout check and the 30-second sync-no-progress check.

Please handle terminal responses that are inconsistent with the height previously advertised by the peer, while preserving appropriate no-progress handling and compatibility with normal responses from honestly lagging peers. Simply rejecting responses that end before the last block in the summary is insufficient: [100], remainNum=0 produces the same result.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Addressed in 4e46755.

check() now rejects single-block responses below the peer’s advertised HELLO head before clearing the pending request or changing synchronization state. Both [50], remainNum=0 and [100], remainNum=0 therefore trigger a SYNC_FAIL disconnect when HELLO advertised height 200.

Normal responses from lagging peers remain accepted—for example, HELLO=50 followed by [50], 0 when our local head has advanced to 100. A response observed during a temporary rollback can also trigger disconnection; using SYNC_FAIL allows recovery through a fresh handshake without the one-hour BAD_PROTOCOL ban.

Regression tests cover rejection, unchanged synchronization state, the actual disconnect reason, and compatibility with lagging peers. All 58 related tests and both Checkstyle checks passed locally.

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,20 @@ 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 (blockIds.size() == 1) {

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] Apply the HELLO head check to multi-block terminal responses as well.

The head check is currently limited to single-block responses, which still allows a peer to return remainNum=0 with 2–3 consecutive, locally known blocks whose final height is below the HELLO head. As long as the first block comes from the current summary, the response passes validation.

The known blocks are then popped and passed to setBlockBothHave(), refreshing the timestamp used by the no-progress timeout. Once the queue is empty, syncNext() sends another request and refreshes the inventory request time. Repeating this can keep the sync connection alive without any height progress while continuously triggering summary construction.

Keep the existing single-block rule, but for all responses with remainNum == 0, require the last block height to be at least the HELLO head height, and reject invalid responses before mutating peer state.

Add regression cases for two- and three-block responses, asserting SYNC_FAILED and that syncNext() is not called. Normal paginated responses with remainNum > 0 should still be allowed to end below the HELLO head.

long lastNum = blockIds.get(0).getNum();
long helloHeadNum = hello.getHeadBlockId().getNum();
if (lastNum < helloHeadNum) {
throw new P2pException(TypeEnum.SYNC_FAILED,
"Single-block response height " + lastNum + " 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