Skip to content
Open
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());
}
Comment on lines 253 to 265

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] Do not unconditionally update the activity timestamp for sync requests that make no progress.

A fully synchronized peer with remainNum=0 can repeatedly send a summary containing only the local node’s current head. The response contains the same single block with remainNum still at 0, so no synchronization progress is made, yet lastInteractiveTime is still updated here after processing. At the same time, check() bypasses the message rate limiter when remainNum=0.

As a result, even after tightening the accounting for INV/BLOCK, this path can still be used to avoid applicable inactivity filtering and reduce the peer’s likelihood of random eviction.

Please move the activity update to a point where response progress can be determined, avoid crediting single-block confirmation responses as useful activity, and apply rate limiting to requests with remainNum=0 as well.

@317787106 317787106 Oct 2, 2026 •

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.

Thanks for your review. Repeated single-BlockId SYNC_BLOCK_CHAIN requests can refresh peer activity without actual synchronization progress. We’ll address this known issue in a separate PR, while preserving legitimate startup/restart and sync-completion behavior. This PR will retain the existing SYNC_BLOCK_CHAIN handling.

}

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,62 @@ 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);
// Retain pending fetch/request state if validation throws so disconnect can retry.
// Otherwise, complete the fetch; inactive witnesses trigger sync recovery below.
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;
Comment on lines +144 to +149

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] Limit repeated no-progress synchronization triggered by inactive signers.

A false result from validBlock() only means that the signer is not in the local active-witness set; it does not prove that the local node needs to synchronize. After receiving a request via INV, a normal peer can provide a block from an inactive account with a valid Merkle root, signature, and parent-height structure, triggering startSync() here. It can then respond with a single locally known block in the summary and remainNum=0, causing the current ChainInventoryMsgHandler to complete synchronization. Repeating this with a new hash can trigger the same flow again.

Please consider adding a retry limit or cooldown for repeated recovery attempts that do not result in verified chain progress, while preserving the legitimate recovery path when the local node is actually behind or its witness set is stale.

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.

We prefer to keep startSync() here because an inactive signer may indicate that the local witness set is stale, and synchronization provides a recovery opportunity. An occasional attempt has limited cost and does not itself accept the block.

The concern about repeated recovery attempts without progress is valid. We’ll address retry limits and activity accounting in a separate PR while preserving this recovery path.

}

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