From 4e749e8bf6b4285fdaaae55b2d0b0c5c1d7d46a7 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 4 Aug 2026 14:43:13 +0800 Subject: [PATCH 01/20] refactor(metrics): use libp2p avg latency for fetch-block peer selection Replace the legacy Dropwizard per-peer histogram P75 with Channel.getAvgLatency() for fetch-block peer selection, and remove the per-peer histogram write in BlockMsgHandler. --- .../net/messagehandler/BlockMsgHandler.java | 4 ---- .../service/fetchblock/FetchBlockService.java | 23 ++++++++----------- 2 files changed, 10 insertions(+), 17 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index 452209d575f..6b6c9e293c8 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -15,8 +15,6 @@ import org.tron.core.config.args.Args; import org.tron.core.exception.P2pException; import org.tron.core.exception.P2pException.TypeEnum; -import org.tron.core.metrics.MetricsKey; -import org.tron.core.metrics.MetricsUtil; import org.tron.core.net.TronNetDelegate; import org.tron.core.net.message.TronMessage; import org.tron.core.net.message.adv.BlockMessage; @@ -91,8 +89,6 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep } Long time = peer.getAdvInvRequest().remove(item); if (null != time) { - MetricsUtil.histogramUpdateUnCheck(MetricsKey.NET_LATENCY_FETCH_BLOCK - + peer.getInetAddress(), now - time); Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_FETCH_LATENCY, (now - time) / Metrics.MILLISECONDS_PER_SECOND); } diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index bda2646abbc..dac03943256 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -16,8 +16,6 @@ import org.tron.common.utils.Sha256Hash; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.BlockCapsule; -import org.tron.core.metrics.MetricsKey; -import org.tron.core.metrics.MetricsUtil; import org.tron.core.net.TronNetDelegate; import org.tron.core.net.message.adv.FetchInvDataMessage; import org.tron.core.net.peer.Item; @@ -97,9 +95,9 @@ private void fetchBlockProcess(FetchBlockInfo fetchBlock) { .filter(PeerConnection::isIdle) .filter(filterPeer -> !filterPeer.equals(fetchBlock.getPeer())) .filter(filterPeer -> filterPeer.getAdvInvReceive().getIfPresent(item) != null) - .filter(filterPeer -> getPeerTop75(filterPeer) + .filter(filterPeer -> getPeerLatency(filterPeer) <= CommonParameter.getInstance().fetchBlockTimeout) - .min(Comparator.comparingDouble(this::getPeerTop75)); + .min(Comparator.comparingDouble(this::getPeerLatency)); if (optionalPeerConnection.isPresent()) { optionalPeerConnection.ifPresent(firstPeer -> { @@ -120,21 +118,20 @@ private void fetchBlockProcess(FetchBlockInfo fetchBlock) { } private boolean shouldFetchBlock(PeerConnection newPeer, FetchBlockInfo fetchBlock) { - double newPeerTop75 = getPeerTop75(newPeer); - double oldPeerTop75 = getPeerTop75(fetchBlock.getPeer()); + double newPeerLatency = getPeerLatency(newPeer); + double oldPeerLatency = getPeerLatency(fetchBlock.getPeer()); long oldPeerSpendTime = System.currentTimeMillis() - fetchBlock.getTime(); - if (oldPeerTop75 > fetchTimeOut || oldPeerSpendTime >= fetchTimeOut) { + if (oldPeerLatency > fetchTimeOut || oldPeerSpendTime >= fetchTimeOut) { return true; } - double oldPeerLeftTime = oldPeerTop75 - oldPeerSpendTime; - return newPeerTop75 < oldPeerLeftTime * BLOCK_FETCH_LEFT_TIME_PERCENT - && oldPeerSpendTime + newPeerTop75 < fetchTimeOut; + double oldPeerLeftTime = oldPeerLatency - oldPeerSpendTime; + return newPeerLatency < oldPeerLeftTime * BLOCK_FETCH_LEFT_TIME_PERCENT + && oldPeerSpendTime + newPeerLatency < fetchTimeOut; } - private double getPeerTop75(PeerConnection peerConnection) { - return MetricsUtil.getHistogram(MetricsKey.NET_LATENCY_FETCH_BLOCK - + peerConnection.getInetAddress()).getSnapshot().get75thPercentile(); + private double getPeerLatency(PeerConnection peerConnection) { + return peerConnection.getChannel().getAvgLatency(); } private static class FetchBlockInfo { From a1908f5481d7f24ec02556c6db6f6be43cf1025e Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 4 Aug 2026 16:56:46 +0800 Subject: [PATCH 02/20] test(metrics): add coverage for prometheus node_info, MetricsService and fetch-block peer selection Add tests to satisfy the changed-line coverage gate (>60%) that failed in the fork validation CI: - PrometheusApiServiceTest: testNodeInfoMetric verifies the tron:node_info Info collector is registered and exposes the node version; testNodeInfoUnknownKey exercises the null-guard branch in MetricsInfo.set; testApplyBlockDupWitness and testApplyBlockWithTxs cover the migrated MetricsService.applyBlock Prometheus-only path (dup-witness MINER counter and TXS counter). - FetchBlockServiceTest: testSelectLowestLatencyPeer verifies that fetchBlockProcess selects the idle peer with the lowest channel avg latency (the migrated replacement for the legacy per-IP histogram P75) and dispatches a FetchInvDataMessage; testSwitchOnOldPeerTimeout covers the fast-switch branch when the old peer exceeds fetchBlockTimeout. --- .../prometheus/PrometheusApiServiceTest.java | 29 +++- .../net/services/FetchBlockServiceTest.java | 161 ++++++++++++++++++ 2 files changed, 188 insertions(+), 2 deletions(-) create mode 100644 framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java diff --git a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java index dd260a1b869..800c3e76388 100644 --- a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java +++ b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java @@ -21,13 +21,13 @@ import org.tron.common.TestConstants; import org.tron.common.crypto.ECKey; import org.tron.common.parameter.CommonParameter; +import org.tron.common.prometheus.MetricKeys; import org.tron.common.prometheus.MetricLabels; import org.tron.common.prometheus.Metrics; import org.tron.common.utils.ByteArray; import org.tron.common.utils.PublicMethod; import org.tron.common.utils.Sha256Hash; -import org.tron.common.utils.StringUtil; -import org.tron.common.utils.Utils; +import org.tron.common.utils.StringUtil;import org.tron.common.utils.Utils; import org.tron.consensus.dpos.DposSlot; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.AccountCapsule; @@ -36,6 +36,7 @@ import org.tron.core.config.args.Args; import org.tron.core.consensus.ConsensusService; import org.tron.core.net.TronNetDelegate; +import org.tron.program.Version; import org.tron.protos.Protocol; @Slf4j(topic = "metric") @@ -206,4 +207,28 @@ private BlockCapsule createTestBlockCapsule(long time, return blockCapsule; } + @Test + public void testNodeInfoMetric() { + String version = Version.getVersion(); + Metrics.info(MetricKeys.Info.NODE_INFO, version); + // Prometheus Info collector appends "_info" to the sample name + Double value = CollectorRegistry.defaultRegistry.getSampleValue( + "tron:node_info_info", + new String[] {MetricLabels.Info.VERSION}, + new String[] {version}); + Assert.assertNotNull("tron:node_info_info sample should exist", value); + Assert.assertEquals(1.0, value, 0.0); + } + + @Test + public void testNodeInfoUnknownKey() { + // unknown key exercises the null-guard branch in MetricsInfo.set + Metrics.info("tron:unknown_info", "x"); + Double value = CollectorRegistry.defaultRegistry.getSampleValue( + "tron:unknown_info_info", + new String[] {"version"}, + new String[] {"x"}); + Assert.assertNull(value); + } + } \ No newline at end of file diff --git a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java new file mode 100644 index 00000000000..4bde4dc92b8 --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java @@ -0,0 +1,161 @@ +package org.tron.core.net.services; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.google.protobuf.ByteString; +import java.lang.reflect.Constructor; +import java.lang.reflect.Method; +import java.net.InetSocketAddress; +import java.util.ArrayList; +import java.util.List; +import org.junit.Assert; +import org.junit.Test; +import org.tron.common.BaseMethodTest; +import org.tron.common.utils.ReflectUtils; +import org.tron.common.utils.Sha256Hash; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.peer.Item; +import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.service.fetchblock.FetchBlockService; +import org.tron.p2p.connection.Channel; +import org.tron.protos.Protocol.Inventory.InventoryType; + +public class FetchBlockServiceTest extends BaseMethodTest { + + private FetchBlockService service; + private TronNetDelegate tronNetDelegate; + + @Override + protected void afterInit() { + service = context.getBean(FetchBlockService.class); + tronNetDelegate = mock(TronNetDelegate.class); + ReflectUtils.setFieldValue(service, "tronNetDelegate", tronNetDelegate); + } + + /** + * Verify that fetchBlockProcess selects the idle peer with the lowest avg latency + * (excluding the peer we are already fetching from) and sends it a FetchInvDataMessage. + * + *

Covers the migrated getPeerLatency / shouldFetchBlock path that replaced the + * legacy Dropwizard per-IP histogram P75 selection. + */ + @Test + public void testSelectLowestLatencyPeer() throws Exception { + InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.1", 10001); + InetSocketAddress newAddr = new InetSocketAddress("127.0.0.2", 10001); + + Channel oldChannel = mock(Channel.class); + when(oldChannel.getInetSocketAddress()).thenReturn(oldAddr); + when(oldChannel.getInetAddress()).thenReturn(oldAddr.getAddress()); + when(oldChannel.getAvgLatency()).thenReturn(200L); + doNothing().when(oldChannel).send(any(byte[].class)); + + Channel newChannel = mock(Channel.class); + when(newChannel.getInetSocketAddress()).thenReturn(newAddr); + when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); + when(newChannel.getAvgLatency()).thenReturn(50L); + doNothing().when(newChannel).send(any(byte[].class)); + + PeerConnection oldPeer = context.getBean(PeerConnection.class); + oldPeer.setChannel(oldChannel); + + PeerConnection newPeer = context.getBean(PeerConnection.class); + newPeer.setChannel(newChannel); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + + // both peers advertise having the block + oldPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + newPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + + List activePeers = new ArrayList<>(); + activePeers.add(oldPeer); + activePeers.add(newPeer); + when(tronNetDelegate.getActivePeer()).thenReturn(activePeers); + + // seed fetchBlockInfo via reflection (private static nested class) + Class fetchBlockInfoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = fetchBlockInfoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object fetchBlockInfo = constructor.newInstance( + hash, oldPeer, System.currentTimeMillis()); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", fetchBlockInfo); + + // invoke private fetchBlockProcess + Method method = FetchBlockService.class.getDeclaredMethod( + "fetchBlockProcess", fetchBlockInfoClass); + method.setAccessible(true); + method.invoke(service, fetchBlockInfo); + + // new peer (lowest latency) should receive the fetch request + verify(newChannel).send(any(byte[].class)); + // old peer should not be re-requested + verify(oldChannel, never()).send(any(byte[].class)); + // fetchBlockInfo should be cleared after successful dispatch + Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } + + /** + * When the old peer's latency exceeds fetchBlockTimeout, shouldFetchBlock returns + * true immediately (fast-switch path), covering the timeout branch. + */ + @Test + public void testSwitchOnOldPeerTimeout() throws Exception { + InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.3", 10001); + InetSocketAddress newAddr = new InetSocketAddress("127.0.0.4", 10001); + + Channel oldChannel = mock(Channel.class); + when(oldChannel.getInetSocketAddress()).thenReturn(oldAddr); + when(oldChannel.getInetAddress()).thenReturn(oldAddr.getAddress()); + // old peer latency exceeds default fetchBlockTimeout (500) + when(oldChannel.getAvgLatency()).thenReturn(600L); + doNothing().when(oldChannel).send(any(byte[].class)); + + Channel newChannel = mock(Channel.class); + when(newChannel.getInetSocketAddress()).thenReturn(newAddr); + when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); + when(newChannel.getAvgLatency()).thenReturn(50L); + doNothing().when(newChannel).send(any(byte[].class)); + + PeerConnection oldPeer = context.getBean(PeerConnection.class); + oldPeer.setChannel(oldChannel); + + PeerConnection newPeer = context.getBean(PeerConnection.class); + newPeer.setChannel(newChannel); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + + newPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + + List activePeers = new ArrayList<>(); + activePeers.add(oldPeer); + activePeers.add(newPeer); + when(tronNetDelegate.getActivePeer()).thenReturn(activePeers); + + Class fetchBlockInfoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = fetchBlockInfoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object fetchBlockInfo = constructor.newInstance( + hash, oldPeer, System.currentTimeMillis()); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", fetchBlockInfo); + + Method method = FetchBlockService.class.getDeclaredMethod( + "fetchBlockProcess", fetchBlockInfoClass); + method.setAccessible(true); + method.invoke(service, fetchBlockInfo); + + verify(newChannel).send(any(byte[].class)); + Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } +} From 20ee8a97f753a5a3bad594b6464a5098ef1bf0d7 Mon Sep 17 00:00:00 2001 From: warku123 Date: Fri, 14 Aug 2026 11:00:51 +0800 Subject: [PATCH 03/20] test(net): pin fetch-block failover semantics for unknown peer latency --- .../net/services/FetchBlockServiceTest.java | 121 ++++++++++++++++++ 1 file changed, 121 insertions(+) diff --git a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java index 4bde4dc92b8..af599f66638 100644 --- a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java +++ b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java @@ -158,4 +158,125 @@ public void testSwitchOnOldPeerTimeout() throws Exception { verify(newChannel).send(any(byte[].class)); Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); } + + /** + * When the old peer's avgLatency is 0 (unknown latency), failover must not happen: + * oldPeerLeftTime = 0 - 0 = 0, so 50 < 0 * 0.5 is false and shouldFetchBlock returns + * false. This pins the semantics that an unknown latency on the current peer + * suppresses switching to a candidate peer. + */ + @Test + public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { + InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.5", 10001); + InetSocketAddress newAddr = new InetSocketAddress("127.0.0.6", 10001); + + Channel oldChannel = mock(Channel.class); + when(oldChannel.getInetSocketAddress()).thenReturn(oldAddr); + when(oldChannel.getInetAddress()).thenReturn(oldAddr.getAddress()); + // old peer latency unknown + when(oldChannel.getAvgLatency()).thenReturn(0L); + doNothing().when(oldChannel).send(any(byte[].class)); + + Channel newChannel = mock(Channel.class); + when(newChannel.getInetSocketAddress()).thenReturn(newAddr); + when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); + when(newChannel.getAvgLatency()).thenReturn(50L); + doNothing().when(newChannel).send(any(byte[].class)); + + PeerConnection oldPeer = context.getBean(PeerConnection.class); + oldPeer.setChannel(oldChannel); + + PeerConnection newPeer = context.getBean(PeerConnection.class); + newPeer.setChannel(newChannel); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + + oldPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + newPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + + List activePeers = new ArrayList<>(); + activePeers.add(oldPeer); + activePeers.add(newPeer); + when(tronNetDelegate.getActivePeer()).thenReturn(activePeers); + + Class fetchBlockInfoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = fetchBlockInfoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object fetchBlockInfo = constructor.newInstance( + hash, oldPeer, System.currentTimeMillis()); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", fetchBlockInfo); + + Method method = FetchBlockService.class.getDeclaredMethod( + "fetchBlockProcess", fetchBlockInfoClass); + method.setAccessible(true); + method.invoke(service, fetchBlockInfo); + + // no failover: candidate peer must not receive a fetch request + verify(newChannel, never()).send(any(byte[].class)); + // in-flight fetchBlockInfo stays pending + Assert.assertNotNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } + + /** + * When the candidate peer's avgLatency is 0 (unknown latency), it is treated as the + * fastest peer: 0 < (200 - 0) * 0.5 holds, so failover to the candidate happens. + * This pins the semantics that an unknown latency on the candidate peer wins the + * latency comparison. + */ + @Test + public void testSwitchWhenCandidateLatencyUnknown() throws Exception { + InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.7", 10001); + InetSocketAddress newAddr = new InetSocketAddress("127.0.0.8", 10001); + + Channel oldChannel = mock(Channel.class); + when(oldChannel.getInetSocketAddress()).thenReturn(oldAddr); + when(oldChannel.getInetAddress()).thenReturn(oldAddr.getAddress()); + when(oldChannel.getAvgLatency()).thenReturn(200L); + doNothing().when(oldChannel).send(any(byte[].class)); + + Channel newChannel = mock(Channel.class); + when(newChannel.getInetSocketAddress()).thenReturn(newAddr); + when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); + // candidate peer latency unknown + when(newChannel.getAvgLatency()).thenReturn(0L); + doNothing().when(newChannel).send(any(byte[].class)); + + PeerConnection oldPeer = context.getBean(PeerConnection.class); + oldPeer.setChannel(oldChannel); + + PeerConnection newPeer = context.getBean(PeerConnection.class); + newPeer.setChannel(newChannel); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + + oldPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + newPeer.getAdvInvReceive().put(item, System.currentTimeMillis()); + + List activePeers = new ArrayList<>(); + activePeers.add(oldPeer); + activePeers.add(newPeer); + when(tronNetDelegate.getActivePeer()).thenReturn(activePeers); + + Class fetchBlockInfoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = fetchBlockInfoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object fetchBlockInfo = constructor.newInstance( + hash, oldPeer, System.currentTimeMillis()); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", fetchBlockInfo); + + Method method = FetchBlockService.class.getDeclaredMethod( + "fetchBlockProcess", fetchBlockInfoClass); + method.setAccessible(true); + method.invoke(service, fetchBlockInfo); + + // failover: candidate peer (unknown latency treated as fastest) receives the request + verify(newChannel).send(any(byte[].class)); + Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } } From e95c50c0f5b19298b3c5ace1fad1182095f0d2be Mon Sep 17 00:00:00 2001 From: warku123 Date: Fri, 4 Sep 2026 17:50:24 +0800 Subject: [PATCH 04/20] feat(net): bound fetch-block peer latency with per-connection EWMA estimator Replace the raw channel average latency used by fetch-block peer selection with a bounded per-connection estimator: - PeerConnection tracks a volatile fetchLatency EWMA (alpha = 0.1), seeded from the channel average latency on the first sample and clamped to [0, fetchBlockTimeout] to resist outliers - BlockMsgHandler feeds measured fetch durations into the estimator - FetchBlockService reads the estimator; the wall-clock hard timeout switches peers unconditionally while the latency-saturation gate requires a strictly better candidate to avoid 500v500 flapping Fetch latency stays observable via the unlabeled Prometheus histogram. --- .../net/messagehandler/BlockMsgHandler.java | 1 + .../tron/core/net/peer/PeerConnection.java | 27 +++ .../service/fetchblock/FetchBlockService.java | 16 +- .../net/services/FetchBlockServiceTest.java | 157 ++++++++++++++++-- 4 files changed, 178 insertions(+), 23 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index 6b6c9e293c8..096580410c9 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -89,6 +89,7 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep } Long time = peer.getAdvInvRequest().remove(item); if (null != time) { + peer.updateFetchLatency(now - time); Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_FETCH_LATENCY, (now - time) / Metrics.MILLISECONDS_PER_SECOND); } diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index 7d7457cf2fc..dfdd4bce39e 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -24,7 +24,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Scope; import org.springframework.stereotype.Component; +import org.tron.common.math.StrictMathWrapper; import org.tron.common.overlay.message.Message; +import org.tron.common.parameter.CommonParameter; import org.tron.common.prometheus.MetricKeys; import org.tron.common.prometheus.Metrics; import org.tron.common.utils.Pair; @@ -92,6 +94,11 @@ public class PeerConnection { @Getter private volatile long blockRcvTime; + @Getter + private volatile long fetchLatency; + + private volatile boolean fetchLatencySeeded; + @Getter @Setter private volatile TronState tronState = TronState.INIT; @@ -184,6 +191,26 @@ public void setChannel(Channel channel) { Args.getInstance().getRateLimiterDisconnect()); } + /** + * Updates the bounded fetch latency estimator with channel-latency initialization/fallback. + * The first measured fetch latency is intentionally blended with the channel estimate because + * the first fetch is often the slowest and the RTT prior is more robust against outliers. + * A single fetch worker reads this value while the channel event loop writes it; volatile is + * sufficient for this benign race and no lock should be added. + * + * @param latencyMillis measured fetch latency in milliseconds + */ + public void updateFetchLatency(long latencyMillis) { + if (!fetchLatencySeeded) { + fetchLatency = channel.getAvgLatency(); + fetchLatencySeeded = true; + } + fetchLatency = (fetchLatency * 9 + latencyMillis) / 10; + // Saturation intentionally makes the >= timeout gate in FetchBlockService trigger. + fetchLatency = StrictMathWrapper.max(0, + StrictMathWrapper.min(CommonParameter.getInstance().fetchBlockTimeout, fetchLatency)); + } + public void setBlockBothHave(BlockId blockId) { this.blockBothHave = blockId; this.blockBothHaveUpdateTime = System.currentTimeMillis(); diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index dac03943256..b0c525c73fe 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -95,8 +95,7 @@ private void fetchBlockProcess(FetchBlockInfo fetchBlock) { .filter(PeerConnection::isIdle) .filter(filterPeer -> !filterPeer.equals(fetchBlock.getPeer())) .filter(filterPeer -> filterPeer.getAdvInvReceive().getIfPresent(item) != null) - .filter(filterPeer -> getPeerLatency(filterPeer) - <= CommonParameter.getInstance().fetchBlockTimeout) + // Clamping bounds latency by timeout; min() ordering handles candidate selection. .min(Comparator.comparingDouble(this::getPeerLatency)); if (optionalPeerConnection.isPresent()) { @@ -121,7 +120,14 @@ private boolean shouldFetchBlock(PeerConnection newPeer, FetchBlockInfo fetchBlo double newPeerLatency = getPeerLatency(newPeer); double oldPeerLatency = getPeerLatency(fetchBlock.getPeer()); long oldPeerSpendTime = System.currentTimeMillis() - fetchBlock.getTime(); - if (oldPeerLatency > fetchTimeOut || oldPeerSpendTime >= fetchTimeOut) { + // Switch unconditionally on a hard timeout: an unseeded or saturated old peer must not + // permanently wedge fetchBlockInfo. + if (oldPeerSpendTime >= fetchTimeOut) { + return true; + } + + // Require a strictly better peer for the latency saturation gate to prevent 500v500 flapping. + if (oldPeerLatency >= fetchTimeOut && newPeerLatency < oldPeerLatency) { return true; } @@ -131,7 +137,7 @@ private boolean shouldFetchBlock(PeerConnection newPeer, FetchBlockInfo fetchBlo } private double getPeerLatency(PeerConnection peerConnection) { - return peerConnection.getChannel().getAvgLatency(); + return peerConnection.getFetchLatency(); } private static class FetchBlockInfo { @@ -156,4 +162,4 @@ public FetchBlockInfo(Sha256Hash hash, PeerConnection peer, long time) { } -} \ No newline at end of file +} diff --git a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java index af599f66638..39e75e8b26e 100644 --- a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java +++ b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java @@ -12,6 +12,7 @@ import java.lang.reflect.Method; import java.net.InetSocketAddress; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; import org.junit.Assert; import org.junit.Test; @@ -63,9 +64,10 @@ public void testSelectLowestLatencyPeer() throws Exception { PeerConnection oldPeer = context.getBean(PeerConnection.class); oldPeer.setChannel(oldChannel); - + oldPeer.updateFetchLatency(200L); PeerConnection newPeer = context.getBean(PeerConnection.class); newPeer.setChannel(newChannel); + newPeer.updateFetchLatency(50L); Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); Item item = new Item(hash, InventoryType.BLOCK); @@ -104,8 +106,7 @@ public void testSelectLowestLatencyPeer() throws Exception { } /** - * When the old peer's latency exceeds fetchBlockTimeout, shouldFetchBlock returns - * true immediately (fast-switch path), covering the timeout branch. + * A 600ms seed is clamped to 500ms, and >= timeout triggers fast-switch. */ @Test public void testSwitchOnOldPeerTimeout() throws Exception { @@ -115,7 +116,7 @@ public void testSwitchOnOldPeerTimeout() throws Exception { Channel oldChannel = mock(Channel.class); when(oldChannel.getInetSocketAddress()).thenReturn(oldAddr); when(oldChannel.getInetAddress()).thenReturn(oldAddr.getAddress()); - // old peer latency exceeds default fetchBlockTimeout (500) + // seed above timeout; estimator clamps this to default fetchBlockTimeout (500) when(oldChannel.getAvgLatency()).thenReturn(600L); doNothing().when(oldChannel).send(any(byte[].class)); @@ -126,10 +127,12 @@ public void testSwitchOnOldPeerTimeout() throws Exception { doNothing().when(newChannel).send(any(byte[].class)); PeerConnection oldPeer = context.getBean(PeerConnection.class); - oldPeer.setChannel(oldChannel); + ReflectUtils.setFieldValue(oldPeer, "channel", oldChannel); + oldPeer.updateFetchLatency(600L); PeerConnection newPeer = context.getBean(PeerConnection.class); newPeer.setChannel(newChannel); + newPeer.updateFetchLatency(50L); Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); Item item = new Item(hash, InventoryType.BLOCK); @@ -159,11 +162,44 @@ public void testSwitchOnOldPeerTimeout() throws Exception { Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); } + @Test + public void testSwitchOnHardTimeoutWhenOldPeerUnseeded() throws Exception { + PeerConnection oldPeer = context.getBean(PeerConnection.class); + PeerConnection candidate = context.getBean(PeerConnection.class); + Channel oldChannel = mock(Channel.class); + Channel candidateChannel = mock(Channel.class); + when(oldChannel.getAvgLatency()).thenReturn(0L); + when(candidateChannel.getAvgLatency()).thenReturn(50L); + ReflectUtils.setFieldValue(oldPeer, "channel", oldChannel); + ReflectUtils.setFieldValue(candidate, "channel", candidateChannel); + candidate.updateFetchLatency(50L); + doNothing().when(candidateChannel).send(any(byte[].class)); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + candidate.getAdvInvReceive().put(item, System.currentTimeMillis()); + when(tronNetDelegate.getActivePeer()).thenReturn(Arrays.asList(oldPeer, candidate)); + + Class fetchBlockInfoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = fetchBlockInfoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object fetchBlockInfo = constructor.newInstance( + hash, oldPeer, System.currentTimeMillis() - 600); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", fetchBlockInfo); + Method method = FetchBlockService.class.getDeclaredMethod( + "fetchBlockProcess", fetchBlockInfoClass); + method.setAccessible(true); + method.invoke(service, fetchBlockInfo); + + verify(candidateChannel).send(any(byte[].class)); + Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } + /** - * When the old peer's avgLatency is 0 (unknown latency), failover must not happen: - * oldPeerLeftTime = 0 - 0 = 0, so 50 < 0 * 0.5 is false and shouldFetchBlock returns - * false. This pins the semantics that an unknown latency on the current peer - * suppresses switching to a candidate peer. + * When the old peer is unseeded, its fetch latency is 0 and the candidate latency is 144 + * ((50 * 9 + 999) / 10), so failover must not happen while the fetch is still within timeout. */ @Test public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { @@ -188,6 +224,7 @@ public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { PeerConnection newPeer = context.getBean(PeerConnection.class); newPeer.setChannel(newChannel); + newPeer.updateFetchLatency(999L); Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); Item item = new Item(hash, InventoryType.BLOCK); @@ -221,13 +258,10 @@ public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { } /** - * When the candidate peer's avgLatency is 0 (unknown latency), it is treated as the - * fastest peer: 0 < (200 - 0) * 0.5 holds, so failover to the candidate happens. - * This pins the semantics that an unknown latency on the candidate peer wins the - * latency comparison. + * Both peers are unseeded, so zero latency suppresses failover. */ @Test - public void testSwitchWhenCandidateLatencyUnknown() throws Exception { + public void testNoSwitchWhenBothPeersUnseeded() throws Exception { InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.7", 10001); InetSocketAddress newAddr = new InetSocketAddress("127.0.0.8", 10001); @@ -240,13 +274,12 @@ public void testSwitchWhenCandidateLatencyUnknown() throws Exception { Channel newChannel = mock(Channel.class); when(newChannel.getInetSocketAddress()).thenReturn(newAddr); when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); - // candidate peer latency unknown + // candidate remains unseeded, so its fetch latency is zero when(newChannel.getAvgLatency()).thenReturn(0L); doNothing().when(newChannel).send(any(byte[].class)); PeerConnection oldPeer = context.getBean(PeerConnection.class); oldPeer.setChannel(oldChannel); - PeerConnection newPeer = context.getBean(PeerConnection.class); newPeer.setChannel(newChannel); @@ -275,8 +308,96 @@ public void testSwitchWhenCandidateLatencyUnknown() throws Exception { method.setAccessible(true); method.invoke(service, fetchBlockInfo); - // failover: candidate peer (unknown latency treated as fastest) receives the request - verify(newChannel).send(any(byte[].class)); + // no failover: both unseeded latencies are zero + verify(newChannel, never()).send(any(byte[].class)); + Assert.assertNotNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + } + + @Test + public void testSwitchToUnseededCandidateWhenOldPeerSeeded() throws Exception { + PeerConnection oldPeer = context.getBean(PeerConnection.class); + PeerConnection candidate = context.getBean(PeerConnection.class); + Channel oldChannel = mock(Channel.class); + Channel candidateChannel = mock(Channel.class); + when(oldChannel.getAvgLatency()).thenReturn(200L); + when(candidateChannel.getAvgLatency()).thenReturn(0L); + ReflectUtils.setFieldValue(oldPeer, "channel", oldChannel); + oldPeer.updateFetchLatency(200L); + ReflectUtils.setFieldValue(candidate, "channel", candidateChannel); + doNothing().when(candidateChannel).send(any(byte[].class)); + + Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); + Item item = new Item(hash, InventoryType.BLOCK); + candidate.getAdvInvReceive().put(item, System.currentTimeMillis()); + when(tronNetDelegate.getActivePeer()).thenReturn(Arrays.asList(oldPeer, candidate)); + + Class infoClass = Class.forName( + "org.tron.core.net.service.fetchblock.FetchBlockService$FetchBlockInfo"); + Constructor constructor = infoClass.getDeclaredConstructor( + Sha256Hash.class, PeerConnection.class, long.class); + constructor.setAccessible(true); + Object info = constructor.newInstance(hash, oldPeer, System.currentTimeMillis()); + ReflectUtils.setFieldValue(service, "fetchBlockInfo", info); + Method method = FetchBlockService.class.getDeclaredMethod("fetchBlockProcess", infoClass); + method.setAccessible(true); + method.invoke(service, info); + + verify(candidateChannel).send(any(byte[].class)); Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); } + + @Test + public void testFirstFetchLatencySampleSeedsFromChannel() { + Channel channel = mock(Channel.class); + when(channel.getAvgLatency()).thenReturn(123L); + PeerConnection peer = new PeerConnection(); + ReflectUtils.setFieldValue(peer, "channel", channel); + + // First call blends channel prior with measured fetch sample. + peer.updateFetchLatency(999L); + + Assert.assertEquals((123L * 9 + 999L) / 10, peer.getFetchLatency()); + } + + @Test + public void testFetchLatencyUsesEwma() { + Channel channel = mock(Channel.class); + when(channel.getAvgLatency()).thenReturn(100L); + PeerConnection peer = new PeerConnection(); + ReflectUtils.setFieldValue(peer, "channel", channel); + + peer.updateFetchLatency(100L); + peer.updateFetchLatency(200L); + + Assert.assertEquals(110L, peer.getFetchLatency()); + } + + @Test + public void testFetchLatencyIsClamped() { + Channel channel = mock(Channel.class); + when(channel.getAvgLatency()).thenReturn(100L); + PeerConnection peer = new PeerConnection(); + ReflectUtils.setFieldValue(peer, "channel", channel); + + peer.updateFetchLatency(100L); + peer.updateFetchLatency(9999L); + + Assert.assertEquals(500L, peer.getFetchLatency()); + } + + @Test + public void testFetchLatencyIsIsolatedAcrossConnections() { + Channel firstChannel = mock(Channel.class); + when(firstChannel.getAvgLatency()).thenReturn(100L); + PeerConnection first = new PeerConnection(); + ReflectUtils.setFieldValue(first, "channel", firstChannel); + first.updateFetchLatency(9999L); + + Channel secondChannel = mock(Channel.class); + when(secondChannel.getAvgLatency()).thenReturn(50L); + PeerConnection second = new PeerConnection(); + ReflectUtils.setFieldValue(second, "channel", secondChannel); + + Assert.assertEquals(0L, second.getFetchLatency()); + } } From 6f65b8b0a9dd905e1025acee00cf01cd2b90de0f Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 4 Aug 2026 15:09:12 +0800 Subject: [PATCH 05/20] feat(metrics): expose node version as a prometheus info metric Add tron:node_info{version="..."} so the node version that the legacy Monitor API used to report is still observable through prometheus. Node IP is intentionally not added; the prometheus instance label already identifies the source node. --- .../tron/common/prometheus/MetricKeys.java | 10 +++++ .../tron/common/prometheus/MetricLabels.java | 10 +++++ .../org/tron/common/prometheus/Metrics.java | 4 ++ .../tron/common/prometheus/MetricsInfo.java | 39 +++++++++++++++++++ .../main/java/org/tron/program/FullNode.java | 2 + 5 files changed, 65 insertions(+) create mode 100644 common/src/main/java/org/tron/common/prometheus/MetricsInfo.java diff --git a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java index 95a38c4b479..8ac35042b16 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java @@ -44,6 +44,16 @@ private Gauge() { } + // Info + public static class Info { + public static final String NODE_INFO = "tron:node_info"; + + private Info() { + throw new IllegalStateException("Info"); + } + + } + // Histogram public static class Histogram { public static final String HTTP_SERVICE_LATENCY = "tron:http_service_latency_seconds"; diff --git a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java index 1f0da214085..b27ea34bdd8 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java @@ -78,4 +78,14 @@ private Histogram() { } + // Info + public static class Info { + public static final String VERSION = "version"; + + private Info() { + throw new IllegalStateException("Info"); + } + + } + } diff --git a/common/src/main/java/org/tron/common/prometheus/Metrics.java b/common/src/main/java/org/tron/common/prometheus/Metrics.java index 6774dd7c315..d3506684231 100644 --- a/common/src/main/java/org/tron/common/prometheus/Metrics.java +++ b/common/src/main/java/org/tron/common/prometheus/Metrics.java @@ -65,4 +65,8 @@ public static void histogramObserve(Histogram.Timer startTimer) { public static void histogramObserve(String key, double amt, String... labels) { MetricsHistogram.observe(key, amt, labels); } + + public static void info(String key, String... labels) { + MetricsInfo.set(key, labels); + } } diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java new file mode 100644 index 00000000000..94d9ce15342 --- /dev/null +++ b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java @@ -0,0 +1,39 @@ +package org.tron.common.prometheus; + +import io.prometheus.client.Info; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import lombok.extern.slf4j.Slf4j; + +@Slf4j(topic = "metrics") +class MetricsInfo { + + private static final Map container = new ConcurrentHashMap<>(); + + static { + init(MetricKeys.Info.NODE_INFO, "tron node info .", MetricLabels.Info.VERSION); + } + + private MetricsInfo() { + throw new IllegalStateException("MetricsInfo"); + } + + private static void init(String name, String help, String... labels) { + container.put(name, Info.build() + .name(name) + .help(help) + .labelNames(labels) + .register()); + } + + static void set(String key, String... labels) { + if (Metrics.enabled()) { + Info info = container.get(key); + if (info == null) { + logger.info("{} not exist", key); + return; + } + info.labels(labels); + } + } +} \ No newline at end of file diff --git a/framework/src/main/java/org/tron/program/FullNode.java b/framework/src/main/java/org/tron/program/FullNode.java index 96b9f73d577..cf32262e1fa 100644 --- a/framework/src/main/java/org/tron/program/FullNode.java +++ b/framework/src/main/java/org/tron/program/FullNode.java @@ -10,6 +10,7 @@ import org.tron.common.exit.ExitManager; import org.tron.common.log.LogService; import org.tron.common.parameter.CommonParameter; +import org.tron.common.prometheus.MetricKeys; import org.tron.common.prometheus.Metrics; import org.tron.core.config.DefaultConfig; import org.tron.core.config.args.Args; @@ -50,6 +51,7 @@ public static void main(String[] args) { // init metrics first Metrics.init(); + Metrics.info(MetricKeys.Info.NODE_INFO, Version.getVersion()); DefaultListableBeanFactory beanFactory = new DefaultListableBeanFactory(); beanFactory.setAllowCircularReferences(false); From aeee8c7da91258f2b6ab366de4ca1f10a4540100 Mon Sep 17 00:00:00 2001 From: warku123 Date: Fri, 14 Aug 2026 11:00:50 +0800 Subject: [PATCH 06/20] feat(metrics): add chain id to node info metric --- .../java/org/tron/common/prometheus/MetricKeys.java | 2 +- .../java/org/tron/common/prometheus/MetricLabels.java | 1 + .../java/org/tron/common/prometheus/MetricsInfo.java | 3 ++- .../src/main/java/org/tron/program/FullNode.java | 4 +++- .../metrics/prometheus/PrometheusApiServiceTest.java | 11 ++++++----- 5 files changed, 13 insertions(+), 8 deletions(-) diff --git a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java index 8ac35042b16..cfe93818291 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java @@ -46,7 +46,7 @@ private Gauge() { // Info public static class Info { - public static final String NODE_INFO = "tron:node_info"; + public static final String NODE_INFO = "tron:node"; private Info() { throw new IllegalStateException("Info"); diff --git a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java index b27ea34bdd8..33b62358a42 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java @@ -81,6 +81,7 @@ private Histogram() { // Info public static class Info { public static final String VERSION = "version"; + public static final String CHAIN_ID = "chain_id"; private Info() { throw new IllegalStateException("Info"); diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java index 94d9ce15342..0c2c5e5697b 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java @@ -11,7 +11,8 @@ class MetricsInfo { private static final Map container = new ConcurrentHashMap<>(); static { - init(MetricKeys.Info.NODE_INFO, "tron node info .", MetricLabels.Info.VERSION); + init(MetricKeys.Info.NODE_INFO, "tron node info.", + MetricLabels.Info.VERSION, MetricLabels.Info.CHAIN_ID); } private MetricsInfo() { diff --git a/framework/src/main/java/org/tron/program/FullNode.java b/framework/src/main/java/org/tron/program/FullNode.java index cf32262e1fa..9385ba9a7df 100644 --- a/framework/src/main/java/org/tron/program/FullNode.java +++ b/framework/src/main/java/org/tron/program/FullNode.java @@ -51,7 +51,6 @@ public static void main(String[] args) { // init metrics first Metrics.init(); - Metrics.info(MetricKeys.Info.NODE_INFO, Version.getVersion()); DefaultListableBeanFactory beanFactory = new DefaultListableBeanFactory(); beanFactory.setAllowCircularReferences(false); @@ -62,6 +61,9 @@ public static void main(String[] args) { Application appT = ApplicationFactory.create(context); context.registerShutdownHook(); appT.startup(); + // chainId is only available after the context refresh (Manager.initGenesis) + Metrics.info(MetricKeys.Info.NODE_INFO, Version.getVersion(), + Args.getInstance().getChainId()); if (parameter.isSolidityNode()) { SolidityNode node = context.getBean(SolidityNode.class); node.run(); diff --git a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java index 800c3e76388..85ed38869ba 100644 --- a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java +++ b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java @@ -210,13 +210,14 @@ private BlockCapsule createTestBlockCapsule(long time, @Test public void testNodeInfoMetric() { String version = Version.getVersion(); - Metrics.info(MetricKeys.Info.NODE_INFO, version); + String testChainId = "00000000000000001ebf88508a03865c71d452e25f4d51194196a1d22b6653dc"; + Metrics.info(MetricKeys.Info.NODE_INFO, version, testChainId); // Prometheus Info collector appends "_info" to the sample name Double value = CollectorRegistry.defaultRegistry.getSampleValue( - "tron:node_info_info", - new String[] {MetricLabels.Info.VERSION}, - new String[] {version}); - Assert.assertNotNull("tron:node_info_info sample should exist", value); + "tron:node_info", + new String[] {MetricLabels.Info.VERSION, MetricLabels.Info.CHAIN_ID}, + new String[] {version, testChainId}); + Assert.assertNotNull("tron:node_info sample should exist", value); Assert.assertEquals(1.0, value, 0.0); } From bd4008db9315640fcb398db73dcce572b939a0a7 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:06:20 +0800 Subject: [PATCH 07/20] fix(net): replace estimator seeding with direct first-sample initialization The first real fetch latency now directly initializes the estimator (isomorphic to RFC 6298 SRTT initialization) instead of being blended with the channel average latency. The channel latency is demoted to a read-only fallback for the unsampled state via getFetchLatency(), and never enters the sample sequence. EWMA alpha=0.1 applies from the second sample onward; clamp keeps math-check compliance. --- .../tron/core/net/peer/PeerConnection.java | 41 ++++++++++++++----- 1 file changed, 31 insertions(+), 10 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index dfdd4bce39e..c738df204e1 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -94,7 +94,6 @@ public class PeerConnection { @Getter private volatile long blockRcvTime; - @Getter private volatile long fetchLatency; private volatile boolean fetchLatencySeeded; @@ -192,23 +191,45 @@ public void setChannel(Channel channel) { } /** - * Updates the bounded fetch latency estimator with channel-latency initialization/fallback. - * The first measured fetch latency is intentionally blended with the channel estimate because - * the first fetch is often the slowest and the RTT prior is more robust against outliers. - * A single fetch worker reads this value while the channel event loop writes it; volatile is - * sufficient for this benign race and no lock should be added. + * Bounded fetch latency estimator with an explicit unsampled state. + * + *

The channel's average latency is never part of the sample sequence; it is only a + * read fallback while the estimator is unsampled (see {@link #getFetchLatency()}). The + * first measured fetch latency directly replaces the unsampled state (isomorphic to + * RFC 6298 SRTT initialization), and subsequent samples are blended with an EWMA of + * alpha = 0.1. With integer division the EWMA has a fixed point, e.g. + * (499 * 9 + 500) / 10 = 499, which damps jitter around the saturation bound. + * + *

A single fetch worker reads this value while the channel event loop writes it; + * volatile is sufficient for this benign race and no lock should be added. * * @param latencyMillis measured fetch latency in milliseconds */ public void updateFetchLatency(long latencyMillis) { if (!fetchLatencySeeded) { - fetchLatency = channel.getAvgLatency(); + fetchLatency = clampFetchLatency(latencyMillis); fetchLatencySeeded = true; + } else { + fetchLatency = clampFetchLatency((fetchLatency * 9 + latencyMillis) / 10); + } + } + + /** + * Returns the bounded fetch latency estimate. While the estimator has not observed a + * real fetch sample yet, the channel's average latency is returned as a read fallback + * (an unknown peer is treated via its transport-level estimate instead of 0). + */ + public long getFetchLatency() { + if (!fetchLatencySeeded) { + return channel.getAvgLatency(); } - fetchLatency = (fetchLatency * 9 + latencyMillis) / 10; + return fetchLatency; + } + + private long clampFetchLatency(long latency) { // Saturation intentionally makes the >= timeout gate in FetchBlockService trigger. - fetchLatency = StrictMathWrapper.max(0, - StrictMathWrapper.min(CommonParameter.getInstance().fetchBlockTimeout, fetchLatency)); + return StrictMathWrapper.max(0, + StrictMathWrapper.min(CommonParameter.getInstance().fetchBlockTimeout, latency)); } public void setBlockBothHave(BlockId blockId) { From 833a7afccea1bcdb194a7958e0c57682118eaef3 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:10:05 +0800 Subject: [PATCH 08/20] test(net): recalculate fetch-block failover tests for unsampled read fallback Unsampled peers now read their channel avgLatency as a fallback instead of 0, so the both-unsampled quadrant flips from suppressing failover to allowing it (candidate 0 < (200 - 0) * 0.5). The first-sample test now asserts direct replacement with clamp instead of channel blending, and the isolation test asserts the fresh connection's channel fallback. --- .../net/services/FetchBlockServiceTest.java | 43 +++++++++++++------ 1 file changed, 31 insertions(+), 12 deletions(-) diff --git a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java index 39e75e8b26e..fb882827860 100644 --- a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java +++ b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java @@ -198,8 +198,9 @@ public void testSwitchOnHardTimeoutWhenOldPeerUnseeded() throws Exception { } /** - * When the old peer is unseeded, its fetch latency is 0 and the candidate latency is 144 - * ((50 * 9 + 999) / 10), so failover must not happen while the fetch is still within timeout. + * Old peer unsampled: getFetchLatency() falls back to its channel avgLatency (0), while + * the candidate's first real sample 999 clamps to 500. shouldFetchBlock: left time + * 0 - 0 = 0, so 500 < 0 * 0.5 = 0 is false — no failover while the fetch is within timeout. */ @Test public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { @@ -224,6 +225,7 @@ public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { PeerConnection newPeer = context.getBean(PeerConnection.class); newPeer.setChannel(newChannel); + // first real sample replaces the unsampled state, then clamps 999 to 500 newPeer.updateFetchLatency(999L); Sha256Hash hash = Sha256Hash.wrap(ByteString.copyFrom(new byte[32])); @@ -258,10 +260,12 @@ public void testNoSwitchWhenOldPeerLatencyUnknown() throws Exception { } /** - * Both peers are unseeded, so zero latency suppresses failover. + * Both peers are unsampled, so both reads fall back to their channel avgLatency + * (old = 200, candidate = 0). The candidate wins min() and shouldFetchBlock: + * left time 200 - 0 = 200, so 0 < 200 * 0.5 = 100 holds — failover happens. */ @Test - public void testNoSwitchWhenBothPeersUnseeded() throws Exception { + public void testSwitchWhenBothPeersUnseededReadsChannelFallback() throws Exception { InetSocketAddress oldAddr = new InetSocketAddress("127.0.0.7", 10001); InetSocketAddress newAddr = new InetSocketAddress("127.0.0.8", 10001); @@ -274,7 +278,7 @@ public void testNoSwitchWhenBothPeersUnseeded() throws Exception { Channel newChannel = mock(Channel.class); when(newChannel.getInetSocketAddress()).thenReturn(newAddr); when(newChannel.getInetAddress()).thenReturn(newAddr.getAddress()); - // candidate remains unseeded, so its fetch latency is zero + // candidate unsampled: read falls back to its channel avgLatency (0) when(newChannel.getAvgLatency()).thenReturn(0L); doNothing().when(newChannel).send(any(byte[].class)); @@ -308,11 +312,15 @@ public void testNoSwitchWhenBothPeersUnseeded() throws Exception { method.setAccessible(true); method.invoke(service, fetchBlockInfo); - // no failover: both unseeded latencies are zero - verify(newChannel, never()).send(any(byte[].class)); - Assert.assertNotNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); + // failover: unsampled candidate reads channel fallback 0 and wins the comparison + verify(newChannel).send(any(byte[].class)); + Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); } + /** + * Old peer seeded at 200; the unsampled candidate reads its channel fallback (0), wins + * min() and the left-time comparison (0 < 200 * 0.5) — failover to the candidate. + */ @Test public void testSwitchToUnseededCandidateWhenOldPeerSeeded() throws Exception { PeerConnection oldPeer = context.getBean(PeerConnection.class); @@ -346,17 +354,22 @@ public void testSwitchToUnseededCandidateWhenOldPeerSeeded() throws Exception { Assert.assertNull(ReflectUtils.getFieldObject(service, "fetchBlockInfo")); } + /** + * The first real fetch sample directly replaces the unsampled state (no blending with + * the channel prior, isomorphic to RFC 6298 SRTT initialization); the clamp still + * applies, so 999 saturates to the 500ms bound instead of the old blended (123*9+999)/10. + */ @Test - public void testFirstFetchLatencySampleSeedsFromChannel() { + public void testFirstFetchLatencySampleReplacesSeed() { Channel channel = mock(Channel.class); when(channel.getAvgLatency()).thenReturn(123L); PeerConnection peer = new PeerConnection(); ReflectUtils.setFieldValue(peer, "channel", channel); - // First call blends channel prior with measured fetch sample. + // First real sample replaces the unsampled state; channel prior (123) is not blended. peer.updateFetchLatency(999L); - Assert.assertEquals((123L * 9 + 999L) / 10, peer.getFetchLatency()); + Assert.assertEquals(500L, peer.getFetchLatency()); } @Test @@ -366,7 +379,9 @@ public void testFetchLatencyUsesEwma() { PeerConnection peer = new PeerConnection(); ReflectUtils.setFieldValue(peer, "channel", channel); + // first real sample initializes directly: clamp(100) = 100 peer.updateFetchLatency(100L); + // second sample onwards: EWMA alpha = 0.1 → (100 * 9 + 200) / 10 = 110 peer.updateFetchLatency(200L); Assert.assertEquals(110L, peer.getFetchLatency()); @@ -379,7 +394,9 @@ public void testFetchLatencyIsClamped() { PeerConnection peer = new PeerConnection(); ReflectUtils.setFieldValue(peer, "channel", channel); + // first real sample initializes directly: clamp(100) = 100 peer.updateFetchLatency(100L); + // EWMA: (100 * 9 + 9999) / 10 = 1089, then clamped to the 500ms bound peer.updateFetchLatency(9999L); Assert.assertEquals(500L, peer.getFetchLatency()); @@ -398,6 +415,8 @@ public void testFetchLatencyIsIsolatedAcrossConnections() { PeerConnection second = new PeerConnection(); ReflectUtils.setFieldValue(second, "channel", secondChannel); - Assert.assertEquals(0L, second.getFetchLatency()); + // the fresh connection is unsampled: it reads its own channel fallback (50), + // not the first connection's estimate (500) and not a cross-connection value + Assert.assertEquals(50L, second.getFetchLatency()); } } From 74d952e633fcfe5c88619adc8a30a3e2eeb05c93 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:11:18 +0800 Subject: [PATCH 09/20] feat(metrics): rename node_info chain_id label to genesis_block_id The label value is the chain id derived from the genesis block hash, so genesis_block_id describes what it identifies more accurately. --- .../main/java/org/tron/common/prometheus/MetricLabels.java | 4 +++- framework/src/main/java/org/tron/program/FullNode.java | 3 ++- .../core/metrics/prometheus/PrometheusApiServiceTest.java | 7 ++++--- 3 files changed, 9 insertions(+), 5 deletions(-) diff --git a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java index 33b62358a42..7f3d5fa6076 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricLabels.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricLabels.java @@ -81,7 +81,9 @@ private Histogram() { // Info public static class Info { public static final String VERSION = "version"; - public static final String CHAIN_ID = "chain_id"; + // identifies the genesis block: the label value is the chain id derived from the + // genesis block hash, so the label is named genesis_block_id + public static final String CHAIN_ID = "genesis_block_id"; private Info() { throw new IllegalStateException("Info"); diff --git a/framework/src/main/java/org/tron/program/FullNode.java b/framework/src/main/java/org/tron/program/FullNode.java index 9385ba9a7df..f7092a05a2e 100644 --- a/framework/src/main/java/org/tron/program/FullNode.java +++ b/framework/src/main/java/org/tron/program/FullNode.java @@ -61,7 +61,8 @@ public static void main(String[] args) { Application appT = ApplicationFactory.create(context); context.registerShutdownHook(); appT.startup(); - // chainId is only available after the context refresh (Manager.initGenesis) + // the genesis block id (chainId) is only available after the context refresh + // (Manager.initGenesis) Metrics.info(MetricKeys.Info.NODE_INFO, Version.getVersion(), Args.getInstance().getChainId()); if (parameter.isSolidityNode()) { diff --git a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java index 85ed38869ba..4ee5a16657c 100644 --- a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java +++ b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java @@ -210,13 +210,14 @@ private BlockCapsule createTestBlockCapsule(long time, @Test public void testNodeInfoMetric() { String version = Version.getVersion(); - String testChainId = "00000000000000001ebf88508a03865c71d452e25f4d51194196a1d22b6653dc"; - Metrics.info(MetricKeys.Info.NODE_INFO, version, testChainId); + String testGenesisBlockId = + "00000000000000001ebf88508a03865c71d452e25f4d51194196a1d22b6653dc"; + Metrics.info(MetricKeys.Info.NODE_INFO, version, testGenesisBlockId); // Prometheus Info collector appends "_info" to the sample name Double value = CollectorRegistry.defaultRegistry.getSampleValue( "tron:node_info", new String[] {MetricLabels.Info.VERSION, MetricLabels.Info.CHAIN_ID}, - new String[] {version, testChainId}); + new String[] {version, testGenesisBlockId}); Assert.assertNotNull("tron:node_info sample should exist", value); Assert.assertEquals(1.0, value, 0.0); } From 56ef29124049060ecb928e1a2b839339866d27e2 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:21:47 +0800 Subject: [PATCH 10/20] chore(protocol): mark legacy Monitor service and MetricsInfo as deprecated Add 'option deprecated = true;' to the Monitor gRPC service and the MetricsInfo message so generated classes carry @Deprecated. --- protocol/src/main/protos/api/api.proto | 2 ++ protocol/src/main/protos/core/Tron.proto | 2 ++ 2 files changed, 4 insertions(+) diff --git a/protocol/src/main/protos/api/api.proto b/protocol/src/main/protos/api/api.proto index f8d13a6bbd3..f2f51392b65 100644 --- a/protocol/src/main/protos/api/api.proto +++ b/protocol/src/main/protos/api/api.proto @@ -646,6 +646,8 @@ service Database { }; service Monitor { + option deprecated = true; + rpc GetStatsInfo (EmptyMessage) returns (MetricsInfo) { } } diff --git a/protocol/src/main/protos/core/Tron.proto b/protocol/src/main/protos/core/Tron.proto index a68e841bb60..50046903a5b 100644 --- a/protocol/src/main/protos/core/Tron.proto +++ b/protocol/src/main/protos/core/Tron.proto @@ -756,6 +756,8 @@ message NodeInfo { } message MetricsInfo { + option deprecated = true; + int64 interval = 1; NodeInfo node = 2; BlockChainInfo blockchain = 3; From 80740dd0a79b1665b2b4928400c1439f268b790b Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:27:15 +0800 Subject: [PATCH 11/20] chore(metrics): warn on legacy metrics stack usage at entry points Log a process-level warning once when the deprecated legacy metrics stack is used: node startup with node.metricsEnable, HTTP /monitor/getstatsinfo, and rpc Monitor.GetStatsInfo. The servlet and rpc warnings use independent once-flags. --- .../main/java/org/tron/core/services/RpcApiService.java | 7 +++++++ .../java/org/tron/core/services/http/MetricsServlet.java | 7 +++++++ framework/src/main/java/org/tron/program/FullNode.java | 5 +++++ 3 files changed, 19 insertions(+) diff --git a/framework/src/main/java/org/tron/core/services/RpcApiService.java b/framework/src/main/java/org/tron/core/services/RpcApiService.java index b9cb05a3b14..4949f62a7ab 100755 --- a/framework/src/main/java/org/tron/core/services/RpcApiService.java +++ b/framework/src/main/java/org/tron/core/services/RpcApiService.java @@ -10,6 +10,7 @@ import io.grpc.netty.NettyServerBuilder; import io.grpc.stub.StreamObserver; import java.util.Objects; +import java.util.concurrent.atomic.AtomicBoolean; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -2655,9 +2656,15 @@ public void getBlock(GrpcAPI.BlockReq request, public class MonitorApi extends MonitorGrpc.MonitorImplBase { + private final AtomicBoolean deprecatedWarned = new AtomicBoolean(false); + @Override public void getStatsInfo(EmptyMessage request, StreamObserver responseObserver) { + if (deprecatedWarned.compareAndSet(false, true)) { + logger.warn("rpc Monitor.GetStatsInfo is deprecated and will be removed in a " + + "future major release; migrate to the prometheus metrics endpoint"); + } responseObserver.onNext(metricsApiService.getMetricProtoInfo()); responseObserver.onCompleted(); } diff --git a/framework/src/main/java/org/tron/core/services/http/MetricsServlet.java b/framework/src/main/java/org/tron/core/services/http/MetricsServlet.java index aaaebb22146..0d14b3abc49 100644 --- a/framework/src/main/java/org/tron/core/services/http/MetricsServlet.java +++ b/framework/src/main/java/org/tron/core/services/http/MetricsServlet.java @@ -1,5 +1,6 @@ package org.tron.core.services.http; +import java.util.concurrent.atomic.AtomicBoolean; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; import lombok.extern.slf4j.Slf4j; @@ -13,10 +14,16 @@ @Slf4j(topic = "API") public class MetricsServlet extends RateLimiterServlet { + private static final AtomicBoolean deprecatedWarned = new AtomicBoolean(false); + @Autowired private MetricsApiService metricsApiService; protected void doGet(HttpServletRequest request, HttpServletResponse response) { + if (deprecatedWarned.compareAndSet(false, true)) { + logger.warn("HTTP /monitor/getstatsinfo is deprecated and will be removed in a " + + "future major release; migrate to the prometheus metrics endpoint"); + } try { MetricsInfo metricsInfo = metricsApiService.getMetricsInfo(); diff --git a/framework/src/main/java/org/tron/program/FullNode.java b/framework/src/main/java/org/tron/program/FullNode.java index f7092a05a2e..b7ff9f063d9 100644 --- a/framework/src/main/java/org/tron/program/FullNode.java +++ b/framework/src/main/java/org/tron/program/FullNode.java @@ -52,6 +52,11 @@ public static void main(String[] args) { // init metrics first Metrics.init(); + if (parameter.isNodeMetricsEnable()) { + logger.warn("legacy metrics stack (node.metricsEnable) is deprecated and will be " + + "removed in a future major release; migrate to node.metrics.prometheus.enable"); + } + DefaultListableBeanFactory beanFactory = new DefaultListableBeanFactory(); beanFactory.setAllowCircularReferences(false); TronApplicationContext context = From c8f6bf60db73b82c69a7ad5389cceca7ef7dd142 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:29:37 +0800 Subject: [PATCH 12/20] feat(metrics): add verification counters for fetch failover and duplicates Register two label-free prometheus counters: tron:block_fetch_secondary increments when the estimator-driven failover issues a secondary fetch; tron:block_duplicate increments when an adv block below head (already processed) is received. --- .../src/main/java/org/tron/common/prometheus/MetricKeys.java | 3 +++ .../main/java/org/tron/common/prometheus/MetricsCounter.java | 3 +++ .../java/org/tron/core/net/messagehandler/BlockMsgHandler.java | 1 + .../tron/core/net/service/fetchblock/FetchBlockService.java | 3 +++ 4 files changed, 10 insertions(+) diff --git a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java index cfe93818291..21a5c6fc262 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java @@ -21,6 +21,9 @@ public static class Counter { public static final String P2P_ERROR = "tron:p2p_error"; public static final String P2P_DISCONNECT = "tron:p2p_disconnect"; public static final String INTERNAL_SERVICE_FAIL = "tron:internal_service_fail"; + // verification counters for the bounded fetch latency estimator rollout + public static final String BLOCK_FETCH_SECONDARY = "tron:block_fetch_secondary"; + public static final String BLOCK_DUPLICATE = "tron:block_duplicate"; private Counter() { throw new IllegalStateException("Counter"); diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java index 7231baaba8f..593be2c0f5e 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java @@ -19,6 +19,9 @@ class MetricsCounter { init(MetricKeys.Counter.P2P_DISCONNECT, "tron p2p disconnect .", "type"); init(MetricKeys.Counter.INTERNAL_SERVICE_FAIL, "internal Service fail.", "class", "method"); + init(MetricKeys.Counter.BLOCK_FETCH_SECONDARY, + "secondary fetch requests issued by the fetch-block failover estimator."); + init(MetricKeys.Counter.BLOCK_DUPLICATE, "duplicate blocks received from peers."); } private MetricsCounter() { diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index 096580410c9..4f8ef9cd69c 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -141,6 +141,7 @@ private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pExc long headNum = tronNetDelegate.getHeadBlockId().getNum(); if (block.getNum() < headNum) { + Metrics.counterInc(MetricKeys.Counter.BLOCK_DUPLICATE, 1); logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); return; } diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index b0c525c73fe..e62b8faa8c5 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -13,6 +13,8 @@ import org.springframework.stereotype.Component; import org.tron.common.es.ExecutorServiceManager; import org.tron.common.parameter.CommonParameter; +import org.tron.common.prometheus.MetricKeys; +import org.tron.common.prometheus.Metrics; import org.tron.common.utils.Sha256Hash; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.BlockCapsule; @@ -104,6 +106,7 @@ private void fetchBlockProcess(FetchBlockInfo fetchBlock) { && firstPeer.checkAndPutAdvInvRequest(item, System.currentTimeMillis())) { firstPeer.sendMessage(new FetchInvDataMessage(Collections.singletonList(item.getHash()), item.getType())); + Metrics.counterInc(MetricKeys.Counter.BLOCK_FETCH_SECONDARY, 1); this.fetchBlockInfo = null; } }); From 10a8f2ad711d35c498d2cce8d358ed7e0f441c51 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 17:30:12 +0800 Subject: [PATCH 13/20] chore(metrics): drop unused NET_LATENCY_FETCH_BLOCK key Its read and write points were removed with the legacy fetch-block histogram. --- framework/src/main/java/org/tron/core/metrics/MetricsKey.java | 1 - .../tron/core/metrics/prometheus/PrometheusApiServiceTest.java | 3 ++- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/framework/src/main/java/org/tron/core/metrics/MetricsKey.java b/framework/src/main/java/org/tron/core/metrics/MetricsKey.java index 3ac7b5840d8..c4630994b3f 100644 --- a/framework/src/main/java/org/tron/core/metrics/MetricsKey.java +++ b/framework/src/main/java/org/tron/core/metrics/MetricsKey.java @@ -23,6 +23,5 @@ public class MetricsKey { public static final String NET_API_DETAIL_QPS = "net.api.detail.qps."; public static final String NET_API_DETAIL_FAIL_QPS = "net.api.detail.failQps."; public static final String NET_API_DETAIL_OUT_TRAFFIC = "net.api.detail.outTraffic."; - public static final String NET_LATENCY_FETCH_BLOCK = "net.latency.fetch.block."; } diff --git a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java index 4ee5a16657c..fba687fcdcc 100644 --- a/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java +++ b/framework/src/test/java/org/tron/core/metrics/prometheus/PrometheusApiServiceTest.java @@ -27,7 +27,8 @@ import org.tron.common.utils.ByteArray; import org.tron.common.utils.PublicMethod; import org.tron.common.utils.Sha256Hash; -import org.tron.common.utils.StringUtil;import org.tron.common.utils.Utils; +import org.tron.common.utils.StringUtil; +import org.tron.common.utils.Utils; import org.tron.consensus.dpos.DposSlot; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.AccountCapsule; From 1f1a54966184fcbba53106ee74baa57e210e243f Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 8 Sep 2026 23:06:31 +0800 Subject: [PATCH 14/20] chore(metrics): clarify clamp comment and add missing eof newline --- .../src/main/java/org/tron/common/prometheus/MetricsInfo.java | 2 +- .../tron/core/net/service/fetchblock/FetchBlockService.java | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java index 0c2c5e5697b..e3cd8b7ca65 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricsInfo.java @@ -37,4 +37,4 @@ static void set(String key, String... labels) { info.labels(labels); } } -} \ No newline at end of file +} diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index e62b8faa8c5..611fa543835 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -97,7 +97,8 @@ private void fetchBlockProcess(FetchBlockInfo fetchBlock) { .filter(PeerConnection::isIdle) .filter(filterPeer -> !filterPeer.equals(fetchBlock.getPeer())) .filter(filterPeer -> filterPeer.getAdvInvReceive().getIfPresent(item) != null) - // Clamping bounds latency by timeout; min() ordering handles candidate selection. + // Seeded estimates are clamped to the fetch timeout; the unseeded channel-latency + // fallback is not, but min() ordering and the saturation gate keep it safe. .min(Comparator.comparingDouble(this::getPeerLatency)); if (optionalPeerConnection.isPresent()) { From 2fe5584276cf2c727157ead861c0b3d894c6a348 Mon Sep 17 00:00:00 2001 From: warku123 Date: Wed, 9 Sep 2026 11:20:39 +0800 Subject: [PATCH 15/20] feat(metrics): add armed fetch counter and rename duplicate to block_already_known --- .../src/main/java/org/tron/common/prometheus/MetricKeys.java | 3 ++- .../main/java/org/tron/common/prometheus/MetricsCounter.java | 5 ++++- .../org/tron/core/net/messagehandler/BlockMsgHandler.java | 5 ++++- .../tron/core/net/service/fetchblock/FetchBlockService.java | 1 + 4 files changed, 11 insertions(+), 3 deletions(-) diff --git a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java index 21a5c6fc262..d392a47d200 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricKeys.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricKeys.java @@ -22,8 +22,9 @@ public static class Counter { public static final String P2P_DISCONNECT = "tron:p2p_disconnect"; public static final String INTERNAL_SERVICE_FAIL = "tron:internal_service_fail"; // verification counters for the bounded fetch latency estimator rollout + public static final String BLOCK_FETCH_ARMED = "tron:block_fetch_armed"; public static final String BLOCK_FETCH_SECONDARY = "tron:block_fetch_secondary"; - public static final String BLOCK_DUPLICATE = "tron:block_duplicate"; + public static final String BLOCK_ALREADY_KNOWN = "tron:block_already_known"; private Counter() { throw new IllegalStateException("Counter"); diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java index 593be2c0f5e..3a92358817a 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java @@ -19,9 +19,12 @@ class MetricsCounter { init(MetricKeys.Counter.P2P_DISCONNECT, "tron p2p disconnect .", "type"); init(MetricKeys.Counter.INTERNAL_SERVICE_FAIL, "internal Service fail.", "class", "method"); + init(MetricKeys.Counter.BLOCK_FETCH_ARMED, + "in-flight fetch requests armed by the fetch-block service."); init(MetricKeys.Counter.BLOCK_FETCH_SECONDARY, "secondary fetch requests issued by the fetch-block failover estimator."); - init(MetricKeys.Counter.BLOCK_DUPLICATE, "duplicate blocks received from peers."); + init(MetricKeys.Counter.BLOCK_ALREADY_KNOWN, + "received blocks already known (block num below head); best-effort signal."); } private MetricsCounter() { diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index 4f8ef9cd69c..d4470ca16f1 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -141,7 +141,10 @@ private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pExc long headNum = tronNetDelegate.getHeadBlockId().getNum(); if (block.getNum() < headNum) { - Metrics.counterInc(MetricKeys.Counter.BLOCK_DUPLICATE, 1); + // Best-effort signal: counts received blocks already known (block num < head); + // includes responses to secondary fetches and concurrent/redundant arrivals; + // cannot attribute specifically to a secondary fetch. + Metrics.counterInc(MetricKeys.Counter.BLOCK_ALREADY_KNOWN, 1); logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); return; } diff --git a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java index 611fa543835..4642cd0bdf4 100644 --- a/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java +++ b/framework/src/main/java/org/tron/core/net/service/fetchblock/FetchBlockService.java @@ -73,6 +73,7 @@ public void fetchBlock(List sha256HashList, PeerConnection peer) { .findFirst().ifPresent(sha256Hash -> { long now = System.currentTimeMillis(); fetchBlockInfo = new FetchBlockInfo(sha256Hash, peer, now); + Metrics.counterInc(MetricKeys.Counter.BLOCK_FETCH_ARMED, 1); logger.info("Set fetchBlockInfo, block: {}, peer: {}, time: {}", sha256Hash, peer.getInetAddress(), now); }); From f7c9adf871d3d29b5d431f89361d667d7425cf9e Mon Sep 17 00:00:00 2001 From: warku123 Date: Wed, 9 Sep 2026 11:20:39 +0800 Subject: [PATCH 16/20] refactor(metrics): extract EWMA divisor constant and add convergence tests --- .../tron/core/net/peer/PeerConnection.java | 16 +++++- .../net/services/FetchBlockServiceTest.java | 57 +++++++++++++++++++ 2 files changed, 71 insertions(+), 2 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index c738df204e1..dedf1c23f8e 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -94,6 +94,17 @@ public class PeerConnection { @Getter private volatile long blockRcvTime; + /** + * EWMA smoothing divisor for the fetch latency estimator: the previous estimate is + * weighted (EWMA_DIVISOR - 1) / EWMA_DIVISOR and the new sample 1 / EWMA_DIVISOR, + * i.e. alpha = 0.1. This trades off smoothing against responsiveness and sits in the + * same order of magnitude as TCP's SRTT gain (1/8, RFC 6298). Under a large + * degradation the relative ordering of two peers can flip within 1-2 samples, while + * the absolute value converges smoothly (e.g. seeded at 100, ten 500ms samples walk + * 140, 176, 208, 237, 263, 286, 307, 326, 343, 358 without ever hitting the clamp). + */ + private static final int EWMA_DIVISOR = 10; + private volatile long fetchLatency; private volatile boolean fetchLatencySeeded; @@ -197,7 +208,7 @@ public void setChannel(Channel channel) { * read fallback while the estimator is unsampled (see {@link #getFetchLatency()}). The * first measured fetch latency directly replaces the unsampled state (isomorphic to * RFC 6298 SRTT initialization), and subsequent samples are blended with an EWMA of - * alpha = 0.1. With integer division the EWMA has a fixed point, e.g. + * alpha = 1 / EWMA_DIVISOR = 0.1. With integer division the EWMA has a fixed point, e.g. * (499 * 9 + 500) / 10 = 499, which damps jitter around the saturation bound. * *

A single fetch worker reads this value while the channel event loop writes it; @@ -210,7 +221,8 @@ public void updateFetchLatency(long latencyMillis) { fetchLatency = clampFetchLatency(latencyMillis); fetchLatencySeeded = true; } else { - fetchLatency = clampFetchLatency((fetchLatency * 9 + latencyMillis) / 10); + fetchLatency = clampFetchLatency( + (fetchLatency * (EWMA_DIVISOR - 1) + latencyMillis) / EWMA_DIVISOR); } } diff --git a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java index fb882827860..b7a6ecdfcdc 100644 --- a/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java +++ b/framework/src/test/java/org/tron/core/net/services/FetchBlockServiceTest.java @@ -387,6 +387,63 @@ public void testFetchLatencyUsesEwma() { Assert.assertEquals(110L, peer.getFetchLatency()); } + /** + * Degradation: seeded at 100 (direct replacement), then ten 500ms samples. With + * alpha = 0.1 and integer division the estimate rises monotonically while converging + * smoothly and never exceeds the 500ms clamp bound. Hand-computed sequence: + * 140, 176, 208, 237, 263, 286, 307, 326, 343, 358. + */ + @Test + public void testFetchLatencyEwmaDegradationConverges() { + Channel channel = mock(Channel.class); + when(channel.getAvgLatency()).thenReturn(0L); + PeerConnection peer = new PeerConnection(); + ReflectUtils.setFieldValue(peer, "channel", channel); + + // first real sample initializes directly: clamp(100) = 100 + peer.updateFetchLatency(100L); + long previous = peer.getFetchLatency(); + long[] expected = {140, 176, 208, 237, 263, 286, 307, 326, 343, 358}; + for (long expectedValue : expected) { + peer.updateFetchLatency(500L); + long current = peer.getFetchLatency(); + Assert.assertEquals(expectedValue, current); + Assert.assertTrue(current > previous); + Assert.assertTrue(current <= 500L); + previous = current; + } + } + + /** + * Recovery: continuing from the degradation endpoint (358), ten 100ms samples pull the + * estimate back down monotonically and smoothly. Hand-computed sequence: + * 332, 308, 287, 268, 251, 235, 221, 208, 197, 187. + */ + @Test + public void testFetchLatencyEwmaRecoveryConverges() { + Channel channel = mock(Channel.class); + when(channel.getAvgLatency()).thenReturn(0L); + PeerConnection peer = new PeerConnection(); + ReflectUtils.setFieldValue(peer, "channel", channel); + + // replay the degradation phase to reach its endpoint + peer.updateFetchLatency(100L); + for (int i = 0; i < 10; i++) { + peer.updateFetchLatency(500L); + } + Assert.assertEquals(358L, peer.getFetchLatency()); + + long previous = peer.getFetchLatency(); + long[] expected = {332, 308, 287, 268, 251, 235, 221, 208, 197, 187}; + for (long expectedValue : expected) { + peer.updateFetchLatency(100L); + long current = peer.getFetchLatency(); + Assert.assertEquals(expectedValue, current); + Assert.assertTrue(current < previous); + previous = current; + } + } + @Test public void testFetchLatencyIsClamped() { Channel channel = mock(Channel.class); From b3e8dfd3f41f9f7ca9c33323ab4b3df7dad7178f Mon Sep 17 00:00:00 2001 From: warku123 Date: Wed, 9 Sep 2026 18:30:03 +0800 Subject: [PATCH 17/20] fix(metrics): count requested already-known blocks by exact id --- .../common/prometheus/MetricsCounter.java | 5 +- .../net/messagehandler/BlockMsgHandler.java | 14 +- .../BlockAlreadyKnownCounterTest.java | 179 ++++++++++++++++++ 3 files changed, 193 insertions(+), 5 deletions(-) create mode 100644 framework/src/test/java/org/tron/core/net/messagehandler/BlockAlreadyKnownCounterTest.java diff --git a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java index 3a92358817a..2e066064ed8 100644 --- a/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java +++ b/common/src/main/java/org/tron/common/prometheus/MetricsCounter.java @@ -24,7 +24,10 @@ class MetricsCounter { init(MetricKeys.Counter.BLOCK_FETCH_SECONDARY, "secondary fetch requests issued by the fetch-block failover estimator."); init(MetricKeys.Counter.BLOCK_ALREADY_KNOWN, - "received blocks already known (block num below head); best-effort signal."); + "adv block responses matched to an outstanding request whose exact block id " + + "was already known when the response was handled (best-effort: concurrent " + + "arrivals may be missed; a duplicate is not attributed to secondary " + + "fetches)."); } private MetricsCounter() { diff --git a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java index d4470ca16f1..b1631979913 100644 --- a/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java +++ b/framework/src/main/java/org/tron/core/net/messagehandler/BlockMsgHandler.java @@ -92,6 +92,16 @@ public void processMessage(PeerConnection peer, TronMessage msg) throws P2pExcep peer.updateFetchLatency(now - time); Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_FETCH_LATENCY, (now - time) / Metrics.MILLISECONDS_PER_SECOND); + // Best-effort duplicate signal: only responses matched to an outstanding adv + // request whose exact block id was already known before this response is + // processed (a concurrent or redundant arrival) are counted. The lookup is + // exact-id (block store + khaos), not a height comparison: an unknown fork + // block below head must not count. Concurrency can still let a simultaneous + // arrival slip through, and a duplicate is not attributed to a secondary + // fetch, so this is a lower-bound indicator rather than an exact count. + if (tronNetDelegate.containBlock(blockId)) { + Metrics.counterInc(MetricKeys.Counter.BLOCK_ALREADY_KNOWN, 1); + } } Metrics.histogramObserve(MetricKeys.Histogram.BLOCK_RECEIVE_DELAY, (now - blockMessage.getBlockCapsule().getTimeStamp()) / Metrics.MILLISECONDS_PER_SECOND); @@ -141,10 +151,6 @@ private void processBlock(PeerConnection peer, BlockCapsule block) throws P2pExc long headNum = tronNetDelegate.getHeadBlockId().getNum(); if (block.getNum() < headNum) { - // Best-effort signal: counts received blocks already known (block num < head); - // includes responses to secondary fetches and concurrent/redundant arrivals; - // cannot attribute specifically to a secondary fetch. - Metrics.counterInc(MetricKeys.Counter.BLOCK_ALREADY_KNOWN, 1); logger.warn("Receive a low block {}, head {}", blockId.getString(), headNum); return; } diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/BlockAlreadyKnownCounterTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/BlockAlreadyKnownCounterTest.java new file mode 100644 index 00000000000..64e77d77bba --- /dev/null +++ b/framework/src/test/java/org/tron/core/net/messagehandler/BlockAlreadyKnownCounterTest.java @@ -0,0 +1,179 @@ +package org.tron.core.net.messagehandler; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.fail; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import io.prometheus.client.CollectorRegistry; +import java.net.InetSocketAddress; +import java.util.Collections; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.tron.common.parameter.CommonParameter; +import org.tron.common.prometheus.MetricKeys; +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.exception.P2pException; +import org.tron.core.net.TronNetDelegate; +import org.tron.core.net.message.adv.BlockMessage; +import org.tron.core.net.peer.Item; +import org.tron.core.net.peer.PeerConnection; +import org.tron.core.net.service.adv.AdvService; +import org.tron.core.net.service.fetchblock.FetchBlockService; +import org.tron.core.net.service.relay.RelayService; +import org.tron.core.net.service.sync.SyncService; +import org.tron.core.services.WitnessProductBlockService; +import org.tron.p2p.connection.Channel; +import org.tron.protos.Protocol.Inventory.InventoryType; + +/** + * Focused unit tests for the tron:block_already_known counter semantics. + * + *

The counter only counts adv block responses that were matched to an outstanding + * adv request (the request entry is consumed exactly once) AND whose exact block id + * was already known before the response was processed. A height comparison is never + * used, because an unknown fork block can sit below head and must not count. + * + *

Pure unit tests (no Spring context): every BlockMsgHandler dependency is mocked + * and the counter value is read back from the default Prometheus registry as a delta, + * so leftover increments from earlier tests in the same JVM cannot interfere. + */ +public class BlockAlreadyKnownCounterTest { + + private static final long HEAD_NUM = 100L; + private static final String SAMPLE_NAME = MetricKeys.Counter.BLOCK_ALREADY_KNOWN + "_total"; + + private BlockMsgHandler handler; + private TronNetDelegate delegate; + private PeerConnection peer; + private boolean metricsEnabledBefore; + + @Before + public void setUp() { + metricsEnabledBefore = CommonParameter.getInstance().isMetricsPrometheusEnable(); + CommonParameter.getInstance().setMetricsPrometheusEnable(true); + delegate = mock(TronNetDelegate.class); + handler = new BlockMsgHandler(); + ReflectUtils.setFieldValue(handler, "tronNetDelegate", delegate); + ReflectUtils.setFieldValue(handler, "advService", mock(AdvService.class)); + ReflectUtils.setFieldValue(handler, "relayService", mock(RelayService.class)); + ReflectUtils.setFieldValue(handler, "syncService", mock(SyncService.class)); + ReflectUtils.setFieldValue(handler, "fetchBlockService", mock(FetchBlockService.class)); + ReflectUtils.setFieldValue(handler, "witnessProductBlockService", + mock(WitnessProductBlockService.class)); + // production default; pinned so a leaked fast-forward flag from another test + // class in the same JVM cannot skip the no-request validation below + ReflectUtils.setFieldValue(handler, "fastForward", false); + + peer = new PeerConnection(); + Channel channel = mock(Channel.class); + InetSocketAddress address = new InetSocketAddress("127.0.0.1", 18888); + when(channel.getInetSocketAddress()).thenReturn(address); + when(channel.getInetAddress()).thenReturn(address.getAddress()); + ReflectUtils.setFieldValue(peer, "channel", channel); + } + + @After + public void tearDown() { + CommonParameter.getInstance().setMetricsPrometheusEnable(metricsEnabledBefore); + } + + @Test + public void testMatchedRequestAlreadyKnownBlockIncrements() throws P2pException { + when(delegate.containBlock(any(BlockId.class))).thenReturn(true); + when(delegate.validBlock(any(BlockCapsule.class))).thenReturn(true); + when(delegate.getHeadBlockId()).thenReturn(new BlockId(Sha256Hash.ZERO_HASH, HEAD_NUM)); + + BlockMessage msg = newBlockMessage(1); + double before = sample(); + request(msg); + handler.processMessage(peer, msg); + assertEquals(before + 1, sample(), 0.0); + + // each matched delivery of an already-known id counts exactly once + double beforeSecond = sample(); + request(msg); + handler.processMessage(peer, msg); + assertEquals(beforeSecond + 1, sample(), 0.0); + } + + @Test + public void testMatchedRequestUnknownForkBelowHeadNotIncremented() throws P2pException { + // An unknown fork below head must not count. Only the block's own id is stubbed + // "unknown" while its parent is "known", so processing passes the unlink guard and + // reaches the low-height branch (num < head) without ever delegating: under the old + // height-based implementation that branch counted, so this test fails on it. + BlockMessage msg = newBlockMessage(1); + assertNotEquals(msg.getBlockId(), msg.getBlockCapsule().getParentBlockId()); + when(delegate.containBlock(msg.getBlockId())).thenReturn(false); + when(delegate.containBlock(msg.getBlockCapsule().getParentBlockId())).thenReturn(true); + when(delegate.validBlock(any(BlockCapsule.class))).thenReturn(true); + when(delegate.getHeadBlockId()).thenReturn(new BlockId(Sha256Hash.ZERO_HASH, HEAD_NUM)); + + double before = sample(); + request(msg); + handler.processMessage(peer, msg); + assertEquals(before, sample(), 0.0); + // the response-time exact-id sample ran and found the block unknown + verify(delegate).containBlock(msg.getBlockId()); + // the unlink guard really did check the parent before the low-height branch ... + verify(delegate).containBlock(msg.getBlockCapsule().getParentBlockId()); + // ... which returned without delegating to block processing + verify(delegate, never()).processBlock(any(BlockCapsule.class), anyBoolean()); + } + + @Test + public void testNoRequestKnownBlockNotIncremented() throws Exception { + when(delegate.containBlock(any(BlockId.class))).thenReturn(true); + + double before = sample(); + try { + handler.processMessage(peer, newBlockMessage(1)); + fail("expected P2pException for a block with no matching request"); + } catch (P2pException e) { + assertEquals("no request", e.getMessage()); + } + assertEquals(before, sample(), 0.0); + } + + @Test + public void testMatchedRequestHeadDuplicateIncrements() throws P2pException { + // re-delivery of the current head block id counts: the id is exactly known even + // though it is not below head. + when(delegate.containBlock(any(BlockId.class))).thenReturn(true); + when(delegate.validBlock(any(BlockCapsule.class))).thenReturn(true); + when(delegate.getHeadBlockId()).thenReturn(new BlockId(Sha256Hash.ZERO_HASH, HEAD_NUM)); + when(delegate.getActivePeer()).thenReturn(Collections.emptyList()); + + BlockMessage msg = newBlockMessage(HEAD_NUM); + double before = sample(); + request(msg); + handler.processMessage(peer, msg); + assertEquals(before + 1, sample(), 0.0); + } + + private BlockMessage newBlockMessage(long number) { + BlockCapsule capsule = new BlockCapsule(number, Sha256Hash.ZERO_HASH, + System.currentTimeMillis() - 60_000L, Sha256Hash.ZERO_HASH.getByteString()); + return new BlockMessage(capsule); + } + + private void request(BlockMessage msg) { + peer.getAdvInvRequest() + .put(new Item(msg.getBlockId(), InventoryType.BLOCK), System.currentTimeMillis()); + } + + private double sample() { + Double value = CollectorRegistry.defaultRegistry.getSampleValue(SAMPLE_NAME); + return value == null ? 0 : value; + } +} From 4443724c4f4a35103525c00fb7a37ae6ad632307 Mon Sep 17 00:00:00 2001 From: warku123 Date: Tue, 22 Sep 2026 11:00:28 +0800 Subject: [PATCH 18/20] docs(metrics): register phase 1 metrics in changelog --- METRICS_CHANGELOG.md | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/METRICS_CHANGELOG.md b/METRICS_CHANGELOG.md index 3c599796d7a..b897f6c5d26 100644 --- a/METRICS_CHANGELOG.md +++ b/METRICS_CHANGELOG.md @@ -3,6 +3,17 @@ Metrics Changelog This file tracks Prometheus metric additions, changes, and removals in java-tron. For the full set of metrics emitted today, see the references at the bottom. +**4.8.3** + +### New Metrics + +#### Core + +- `tron:node_info` (Info, labels `version`, `genesis_block_id`) — static node identity: the running node version string plus the full genesis block hash as the canonical chain identifier, for fleet dashboards and version/chain alert rules. ([#6923](https://github.com/tronprotocol/java-tron/issues/6923)) +- `tron:block_fetch_armed` (Counter) — incremented when `FetchBlockService` arms a fetch tracking for a block; denominator for secondary-fetch rates. ([#6923](https://github.com/tronprotocol/java-tron/issues/6923)) +- `tron:block_fetch_secondary` (Counter) — incremented when a secondary fetch request is sent to an alternate peer. ([#6923](https://github.com/tronprotocol/java-tron/issues/6923)) +- `tron:block_already_known` (Counter) — incremented when a block response matches an outstanding adv request whose exact block ID is already known before that response is processed (best-effort signal; concurrent arrivals may be missed). ([#6923](https://github.com/tronprotocol/java-tron/issues/6923)) + **4.8.2** ### New Metrics From feb287cf8762d5d59d14d48a3b6d3603c77b0e57 Mon Sep 17 00:00:00 2001 From: warku123 Date: Mon, 5 Oct 2026 21:03:28 +0800 Subject: [PATCH 19/20] refactor(metrics): rename EWMA divisor to sliding factor and trim comments --- .../tron/core/net/peer/PeerConnection.java | 31 +++++++------------ 1 file changed, 11 insertions(+), 20 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index dedf1c23f8e..41e3755ce25 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -95,15 +95,11 @@ public class PeerConnection { private volatile long blockRcvTime; /** - * EWMA smoothing divisor for the fetch latency estimator: the previous estimate is - * weighted (EWMA_DIVISOR - 1) / EWMA_DIVISOR and the new sample 1 / EWMA_DIVISOR, - * i.e. alpha = 0.1. This trades off smoothing against responsiveness and sits in the - * same order of magnitude as TCP's SRTT gain (1/8, RFC 6298). Under a large - * degradation the relative ordering of two peers can flip within 1-2 samples, while - * the absolute value converges smoothly (e.g. seeded at 100, ten 500ms samples walk - * 140, 176, 208, 237, 263, 286, 307, 326, 343, 358 without ever hitting the clamp). + * Sliding factor of the fetch latency EWMA: each new sample contributes + * 1 / FETCH_LATENCY_EWMA_FACTOR (alpha = 0.1), balancing smoothing against + * responsiveness, in the same order as TCP's SRTT gain (1/8, RFC 6298). */ - private static final int EWMA_DIVISOR = 10; + private static final int FETCH_LATENCY_EWMA_FACTOR = 10; private volatile long fetchLatency; @@ -202,17 +198,11 @@ public void setChannel(Channel channel) { } /** - * Bounded fetch latency estimator with an explicit unsampled state. - * - *

The channel's average latency is never part of the sample sequence; it is only a - * read fallback while the estimator is unsampled (see {@link #getFetchLatency()}). The - * first measured fetch latency directly replaces the unsampled state (isomorphic to - * RFC 6298 SRTT initialization), and subsequent samples are blended with an EWMA of - * alpha = 1 / EWMA_DIVISOR = 0.1. With integer division the EWMA has a fixed point, e.g. - * (499 * 9 + 500) / 10 = 499, which damps jitter around the saturation bound. - * - *

A single fetch worker reads this value while the channel event loop writes it; - * volatile is sufficient for this benign race and no lock should be added. + * Updates the fetch latency estimator, which has an explicit unsampled state: the + * channel's average latency is only a read fallback and never enters the sample sequence. + * The first measured sample directly replaces the placeholder; subsequent samples are + * blended with alpha = 1 / FETCH_LATENCY_EWMA_FACTOR. A single fetch worker reads while + * the channel event loop writes, so volatile suffices and no lock should be added. * * @param latencyMillis measured fetch latency in milliseconds */ @@ -222,7 +212,8 @@ public void updateFetchLatency(long latencyMillis) { fetchLatencySeeded = true; } else { fetchLatency = clampFetchLatency( - (fetchLatency * (EWMA_DIVISOR - 1) + latencyMillis) / EWMA_DIVISOR); + (fetchLatency * (FETCH_LATENCY_EWMA_FACTOR - 1) + latencyMillis) + / FETCH_LATENCY_EWMA_FACTOR); } } From e3a8b1275a8f7dce43196eb407777901607bcf90 Mon Sep 17 00:00:00 2001 From: warku123 Date: Mon, 5 Oct 2026 22:48:57 +0800 Subject: [PATCH 20/20] refactor(net): clarify fetch latency estimator naming Address PR #6988 review comments: - rename boolean field fetchLatencySeeded -> hasFetchLatencySampled (declaration and all reads/writes in updateFetchLatency/getFetchLatency) - spell out "Exponentially Weighted Moving Average (EWMA)" at its first occurrence in PeerConnection.java; later occurrences stay as EWMA --- .../java/org/tron/core/net/peer/PeerConnection.java | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java index 41e3755ce25..6ddc64ae76f 100644 --- a/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java +++ b/framework/src/main/java/org/tron/core/net/peer/PeerConnection.java @@ -95,7 +95,8 @@ public class PeerConnection { private volatile long blockRcvTime; /** - * Sliding factor of the fetch latency EWMA: each new sample contributes + * Sliding factor of the fetch latency Exponentially Weighted Moving Average (EWMA): + * each new sample contributes * 1 / FETCH_LATENCY_EWMA_FACTOR (alpha = 0.1), balancing smoothing against * responsiveness, in the same order as TCP's SRTT gain (1/8, RFC 6298). */ @@ -103,7 +104,7 @@ public class PeerConnection { private volatile long fetchLatency; - private volatile boolean fetchLatencySeeded; + private volatile boolean hasFetchLatencySampled; @Getter @Setter @@ -207,9 +208,9 @@ public void setChannel(Channel channel) { * @param latencyMillis measured fetch latency in milliseconds */ public void updateFetchLatency(long latencyMillis) { - if (!fetchLatencySeeded) { + if (!hasFetchLatencySampled) { fetchLatency = clampFetchLatency(latencyMillis); - fetchLatencySeeded = true; + hasFetchLatencySampled = true; } else { fetchLatency = clampFetchLatency( (fetchLatency * (FETCH_LATENCY_EWMA_FACTOR - 1) + latencyMillis) @@ -223,7 +224,7 @@ public void updateFetchLatency(long latencyMillis) { * (an unknown peer is treated via its transport-level estimate instead of 0). */ public long getFetchLatency() { - if (!fetchLatencySeeded) { + if (!hasFetchLatencySampled) { return channel.getAvgLatency(); } return fetchLatency;