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
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,6 @@ private void updateLastInteractiveTime(PeerConnection peer, TronMessage msg) {
switch (type) {
case SYNC_BLOCK_CHAIN:
case BLOCK_CHAIN_INVENTORY:
case BLOCK:
flag = true;
break;
case FETCH_INV_DATA:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import org.springframework.stereotype.Component;
import org.tron.common.prometheus.MetricKeys;
import org.tron.common.prometheus.Metrics;
import org.tron.common.utils.Sha256Hash;
import org.tron.core.Constant;
import org.tron.core.capsule.BlockCapsule;
import org.tron.core.capsule.BlockCapsule.BlockId;
Expand Down Expand Up @@ -77,6 +78,12 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep
check(peer, blockMessage);
}

if (blockCapsule.getNum() <= 0
|| blockCapsule.getParentHashStr().size() != Sha256Hash.LENGTH
|| new BlockId(blockCapsule.getParentHash()).getNum() != blockCapsule.getNum() - 1) {
throw new P2pException(TypeEnum.BAD_BLOCK, "block number does not follow parent");
}

blockMessage.sanitize();

if (peer.getSyncBlockRequested().containsKey(blockId)) {
Expand All @@ -89,7 +96,12 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep
if (peer.isRelayPeer()) {
peer.getAdvInvSpread().put(item, now);
}
Long time = peer.getAdvInvRequest().remove(item);
Long time = peer.getAdvInvRequest().get(item);
long interval = blockId.getNum() - tronNetDelegate.getHeadBlockId().getNum();
BlockResult result = processBlock(peer, blockMessage.getBlockCapsule());
if (result == BlockResult.ACCEPTED || result == BlockResult.SYNC_REQUIRED) {
peer.setBlockRcvTime(System.currentTimeMillis());
}
if (null != time) {
MetricsUtil.histogramUpdateUnCheck(MetricsKey.NET_LATENCY_FETCH_BLOCK
+ peer.getInetAddress(), now - time);
Expand All @@ -98,9 +110,6 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep
}
Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_RECEIVE_DELAY,
(now - blockMessage.getBlockCapsule().getTimeStamp()) / Metrics.MILLISECONDS_PER_SECOND);
fetchBlockService.blockFetchSuccess(blockId);
long interval = blockId.getNum() - tronNetDelegate.getHeadBlockId().getNum();
processBlock(peer, blockMessage.getBlockCapsule());
logger.info(
"Receive block/interval {}/{} from {} fetch/delay {}/{}ms, "
+ "txs/process {}/{}ms, witness: {}",
Expand All @@ -125,46 +134,60 @@ private void check(PeerConnection peer, BlockMessage msg) throws P2pException {
}
}

private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pException {
private BlockResult processBlock(PeerConnection peer, BlockCapsule block) throws P2pException {
BlockId blockId = block.getBlockId();
boolean flag = tronNetDelegate.validBlock(block);
if (!flag) {
logger.warn("Receive a bad block from {}, {}, {}",
boolean activeWitness = tronNetDelegate.validBlock(block);
// Keep the request until validation succeeds so disconnect can retry invalid data.
fetchBlockService.blockFetchSuccess(blockId);
peer.getAdvInvRequest().remove(new Item(blockId, InventoryType.BLOCK));
if (!activeWitness) {
logger.warn("Receive a block from an inactive witness, peer {}, block {}, witness {}",
peer.getInetSocketAddress(), blockId.getString(),
Hex.toHexString(block.getWitnessAddress().toByteArray()));
return;
syncService.startSync(peer);
return BlockResult.STATE_FAILED;
}

peer.setLastInteractiveTime(System.currentTimeMillis());
long headNum = tronNetDelegate.getHeadBlockId().getNum();
if (block.getNum() < headNum || tronNetDelegate.containBlock(blockId)) {
logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum);
return BlockResult.IGNORED;
}

if (!tronNetDelegate.containBlock(block.getParentBlockId())) {
logger.warn("Get unlink block {} from {}, head is {}", blockId.getString(),
peer.getInetAddress(), tronNetDelegate.getHeadBlockId().getString());
syncService.startSync(peer);
return;
}

long headNum = tronNetDelegate.getHeadBlockId().getNum();
if (block.getNum() < headNum) {
logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum);
return;
return BlockResult.SYNC_REQUIRED;
}

broadcast(new BlockMessage(block));

try {
tronNetDelegate.processBlock(block, false);
peer.setBlockRcvTime(System.currentTimeMillis());
witnessProductBlockService.validWitnessProductTwoBlock(block);

Item item = new Item(blockId, InventoryType.BLOCK);
tronNetDelegate.getActivePeer().forEach(p -> {
if (p.getAdvInvReceive().getIfPresent(item) != null) {
p.setBlockBothHave(blockId);
}
});
} catch (Exception e) {
logger.warn("Process adv block {} from peer {} failed. reason: {}",
blockId, peer.getInetAddress(), e.getMessage());
syncService.startSync(peer);
return BlockResult.STATE_FAILED;
}
if (tronNetDelegate.isHitDown()) {
return BlockResult.IGNORED;
}
witnessProductBlockService.validWitnessProductTwoBlock(block);

Item item = new Item(blockId, InventoryType.BLOCK);
tronNetDelegate.getActivePeer().forEach(p -> {
if (p.getAdvInvReceive().getIfPresent(item) != null) {
p.setBlockBothHave(blockId);
}
});
return BlockResult.ACCEPTED;
}

private enum BlockResult {
ACCEPTED, SYNC_REQUIRED, IGNORED, STATE_FAILED
}

private void broadcast(BlockMessage blockMessage) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -336,8 +336,16 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection)
BlockId blockId = block.getBlockId();
try {
tronNetDelegate.validSignature(block);
boolean useful = block.getNum() >= tronNetDelegate.getHeadBlockId().getNum()
&& !tronNetDelegate.containBlock(blockId);
tronNetDelegate.processBlock(block, true);
peerConnection.setBlockRcvTime(System.currentTimeMillis());
if (tronNetDelegate.isHitDown()) {
return;
}
peerConnection.setLastInteractiveTime(System.currentTimeMillis());
if (useful) {
peerConnection.setBlockRcvTime(System.currentTimeMillis());
}
pbftDataSyncHandler.processPBFTCommitData(block);
} catch (P2pException p2pException) {
logger.error("Process sync block {} failed, type: {}",
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;
}
}
Loading
Loading