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 @@ -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 All @@ -262,7 +261,7 @@ private void updateLastInteractiveTime(PeerConnection peer, TronMessage msg) {
break;
}
if (flag) {
peer.setLastInteractiveTime(System.currentTimeMillis());
peer.updateLastInteractiveTime(System.currentTimeMillis());
}
}

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,61 @@ 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.updateLastInteractiveTime(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;
}
advService.confirmBlockInventory(blockId);
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 @@ -6,7 +6,6 @@
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.tron.common.utils.Sha256Hash;
import org.tron.core.capsule.BlockCapsule.BlockId;
import org.tron.core.config.args.Args;
import org.tron.core.exception.P2pException;
import org.tron.core.exception.P2pException.TypeEnum;
Expand Down Expand Up @@ -42,14 +41,8 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep

for (Sha256Hash id : inventoryMessage.getHashList()) {
Item item = new Item(id, type);
peer.getAdvInvReceive().put(item, System.currentTimeMillis());
advService.recordInventory(peer, item, System.currentTimeMillis());
advService.addInv(item);
if (type.equals(InventoryType.BLOCK) && peer.getAdvInvSpread().getIfPresent(item) == null) {
long headNum = tronNetDelegate.getHeadBlockId().getNum();
if (new BlockId(id).getNum() > headNum) {
peer.setLastInteractiveTime(System.currentTimeMillis());
}
}
}
}

Expand Down
40 changes: 37 additions & 3 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 @@ -122,6 +121,11 @@ public class PeerConnection {
private Cache<Item, Long> advInvReceive = CacheBuilder.newBuilder().maximumSize(invCacheSize)
.expireAfterWrite(1, TimeUnit.HOURS).recordStats().build();

// Eligible block INV receipt times, pending successful block processing before activity updates.
@Getter
private Cache<Item, Long> advBlockInvReceive = CacheBuilder.newBuilder().maximumSize(100)
.expireAfterWrite(1, TimeUnit.MINUTES).build();

@Setter
@Getter
private Cache<Item, Long> advInvSpread = CacheBuilder.newBuilder().maximumSize(invCacheSize)
Expand Down Expand Up @@ -158,7 +162,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 All @@ -175,7 +179,7 @@ public void setChannel(Channel channel) {
this.isRelayPeer = true;
}
this.nodeStatistics = TronStatsManager.getNodeStatistics(channel.getInetAddress());
lastInteractiveTime = System.currentTimeMillis();
updateLastInteractiveTime(System.currentTimeMillis());
p2pRateLimiter.register(SYNC_BLOCK_CHAIN.asByte(),
Args.getInstance().getRateLimiterSyncBlockChain());
p2pRateLimiter.register(FETCH_INV_DATA.asByte(),
Expand All @@ -189,10 +193,39 @@ public void setBlockBothHave(BlockId blockId) {
this.blockBothHaveUpdateTime = System.currentTimeMillis();
}

public synchronized void updateLastInteractiveTime(long time) {
if (time > lastInteractiveTime) {
lastInteractiveTime = time;
}
}

/**
* 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 Expand Up @@ -225,6 +258,7 @@ public void onDisconnect() {
syncService.onDisconnect(this);
advService.onDisconnect(this);
advInvReceive.invalidateAll();
advBlockInvReceive.invalidateAll();
advInvSpread.invalidateAll();
advInvRequest.clear();
syncBlockIdCache.invalidateAll();
Expand Down
Loading
Loading