Skip to content
Closed
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 Down Expand Up @@ -40,7 +41,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(requested, chainInventoryMessage);

peer.setFetchAble(false);

Expand All @@ -51,7 +53,11 @@ 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);
// This trusts the peer's claim that no blocks remain. A peer can echo an earlier
// known summary block while withholding newer blocks, so ending this download
// does not prove that we have caught up with the peer's actual chain. todo fix
peer.setNeedSyncFromPeer(false);
return;
}
Expand Down Expand Up @@ -98,8 +104,9 @@ 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(Pair<Deque<BlockId>, Long> requested, ChainInventoryMessage msg)
throws P2pException {
if (requested == null) {
throw new P2pException(TypeEnum.BAD_MESSAGE, "not send syncBlockChainMsg");
}

Expand All @@ -112,7 +119,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 +132,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,7 +145,7 @@ 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);
}
Expand Down
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
36 changes: 36 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,36 @@
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.peer.PeerConnection;
import org.tron.p2p.connection.Channel;

public final class PeerSyncTestSupport {

private PeerSyncTestSupport() {
}

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

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