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 @@ -262,7 +262,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 @@ -152,6 +152,9 @@ private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pExc

try {
tronNetDelegate.processBlock(block, false);
if (!tronNetDelegate.isHitDown()) {
advService.confirmBlockInventory(blockId);
}
peer.setBlockRcvTime(System.currentTimeMillis());
witnessProductBlockService.validWitnessProductTwoBlock(block);

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
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,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 @@ -175,7 +180,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,6 +194,12 @@ public void setBlockBothHave(BlockId blockId) {
this.blockBothHaveUpdateTime = System.currentTimeMillis();
}

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

public boolean isIdle() {
return advInvRequest.isEmpty() && isSyncIdle();
}
Expand Down Expand Up @@ -225,6 +236,7 @@ public void onDisconnect() {
syncService.onDisconnect(this);
advService.onDisconnect(this);
advInvReceive.invalidateAll();
advBlockInvReceive.invalidateAll();
advInvSpread.invalidateAll();
advInvRequest.clear();
syncBlockIdCache.invalidateAll();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,35 @@ public synchronized void addInvToCache(Item item) {
invToFetch.remove(item);
}

public void recordInventory(PeerConnection peer, Item item, long receivedAt) {
if (item.getType() != InventoryType.BLOCK) {
peer.getAdvInvReceive().put(item, receivedAt);
return;
}
synchronized (this) {
// Capture eligibility before fetching; validation may advance the head immediately.
if (peer.getAdvInvSpread().getIfPresent(item) == null
&& new BlockId(item.getHash()).getNum() > tronNetDelegate.getHeadBlockId().getNum()) {
peer.getAdvBlockInvReceive().put(item, receivedAt);
}
peer.getAdvInvReceive().put(item, receivedAt);
}
}

/**
* Confirms eligible announcements only after the block has been successfully processed.
*/
public synchronized void confirmBlockInventory(BlockId blockId) {
Item item = new Item(blockId, InventoryType.BLOCK);
// Share the registration lock so confirmation cannot miss an eligible announcement.
tronNetDelegate.getActivePeer().forEach(peer -> {
Long receivedAt = peer.getAdvBlockInvReceive().asMap().remove(item);
if (receivedAt != null) {
peer.updateLastInteractiveTime(receivedAt);
}
});
}

public boolean addInv(Item item) {
if (fastForward && item.getType().equals(InventoryType.TRX)) {
return false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.tron.core.net.messagehandler.PbftDataSyncHandler;
import org.tron.core.net.peer.PeerConnection;
import org.tron.core.net.peer.TronState;
import org.tron.core.net.service.adv.AdvService;
import org.tron.protos.Protocol.Inventory.InventoryType;
import org.tron.protos.Protocol.ReasonCode;

Expand All @@ -45,6 +46,9 @@ public class SyncService {
@Autowired
private PbftDataSyncHandler pbftDataSyncHandler;

@Autowired
private AdvService advService;

private Map<UnparsedBlock, PeerConnection> blockWaitToProcess = new ConcurrentHashMap<>();

private Map<UnparsedBlock, PeerConnection> blockJustReceived = new ConcurrentHashMap<>();
Expand Down Expand Up @@ -337,6 +341,9 @@ private void processSyncBlock(BlockCapsule block, PeerConnection peerConnection)
try {
tronNetDelegate.validSignature(block);
tronNetDelegate.processBlock(block, true);
if (!tronNetDelegate.isHitDown()) {
advService.confirmBlockInventory(blockId);
}
peerConnection.setBlockRcvTime(System.currentTimeMillis());
pbftDataSyncHandler.processPBFTCommitData(block);
} catch (P2pException p2pException) {
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