configMap = new HashMap<>();
+ configMap.put("node.dns.publish", true);
+ configMap.put("node.dns.dnsDomain", "nodes.example.org");
+ configMap.put("node.dns.dnsPrivate",
+ "1234567890123456789012345678901234567890123456789012345678901234");
+ configMap.put("node.dns.serverType", "aliyun");
+ configMap.put("node.dns.accessKeyId", "access-key-id");
+ configMap.put("node.dns.accessKeySecret", "access-key-secret");
+ configMap.put("node.dns.aliyunDnsEndpoint", "dns.aliyuncs.com");
+ configMap.put(key, value);
+ return ConfigFactory.parseMap(configMap)
+ .withFallback(ConfigFactory.defaultReference());
+ }
+
@Test
public void testRpcMaxMessageSizeExceedsIntMax() {
// HOCON's Config.getInt() throws when a numeric value exceeds int range.
From 4d6c24085ab151d8c7ef821ff0e7ad2b7f168733 Mon Sep 17 00:00:00 2001
From: Jeremy Zhang <50477615+warku123@users.noreply.github.com>
Date: Tue, 22 Sep 2026 17:38:13 +0800
Subject: [PATCH 06/25] fix(test): isolate cross-test state leaks and stabilize
flaky suites (#6974)
---
.../java/org/tron/common/BaseMethodTest.java | 5 +
.../test/java/org/tron/common/BaseTest.java | 11 +
.../java/org/tron/common/VMConfigRule.java | 52 ++++
.../tron/common/backup/BackupManagerTest.java | 98 +++---
.../tron/common/backup/BackupServerTest.java | 34 +-
.../tron/common/backup/BackupTestUtils.java | 72 +++++
.../common/logsfilter/EventLoaderTest.java | 11 +-
.../logsfilter/NativeMessageQueueTest.java | 5 +-
.../tron/common/prometheus/SRMetricsTest.java | 2 +
.../runtime/vm/AllowTvmCompatibleEvmTest.java | 1 +
.../common/runtime/vm/AllowTvmLondonTest.java | 4 +-
.../utils/PeerManagerStateResetter.java | 99 ++++++
.../utils/PeerManagerStateResetterTest.java | 79 +++++
.../tron/core/event/BlockEventGetTest.java | 6 +-
.../prometheus/PrometheusApiServiceTest.java | 1 +
.../messagehandler/MessageHandlerTest.java | 2 +
.../messagehandler/PbftMsgHandlerTest.java | 2 +
.../TransactionsMsgHandlerTest.java | 169 ++++++----
.../tron/core/net/peer/PeerManagerTest.java | 2 +
.../net/services/HandShakeServiceTest.java | 2 +
.../org/tron/core/services/WalletApiTest.java | 2 +
.../tron/core/zksnark/SendCoinShieldTest.java | 10 +
.../core/zksnark/ShieldedReceiveTest.java | 294 +++++++++---------
23 files changed, 717 insertions(+), 246 deletions(-)
create mode 100644 framework/src/test/java/org/tron/common/VMConfigRule.java
create mode 100644 framework/src/test/java/org/tron/common/backup/BackupTestUtils.java
create mode 100644 framework/src/test/java/org/tron/common/utils/PeerManagerStateResetter.java
create mode 100644 framework/src/test/java/org/tron/common/utils/PeerManagerStateResetterTest.java
diff --git a/framework/src/test/java/org/tron/common/BaseMethodTest.java b/framework/src/test/java/org/tron/common/BaseMethodTest.java
index 9ee1dfa3b36..c91310681a1 100644
--- a/framework/src/test/java/org/tron/common/BaseMethodTest.java
+++ b/framework/src/test/java/org/tron/common/BaseMethodTest.java
@@ -10,6 +10,7 @@
import org.tron.common.application.Application;
import org.tron.common.application.ApplicationFactory;
import org.tron.common.application.TronApplicationContext;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.core.ChainBaseManager;
import org.tron.core.config.DefaultConfig;
import org.tron.core.config.args.Args;
@@ -42,6 +43,9 @@ public abstract class BaseMethodTest {
@Rule
public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+ @Rule
+ public final VMConfigRule vmConfigRule = new VMConfigRule();
+
protected TronApplicationContext context;
protected Application appT;
protected Manager dbManager;
@@ -57,6 +61,7 @@ protected String configFile() {
@Before
public final void initContext() throws IOException {
+ PeerManagerStateResetter.reset();
String[] baseArgs = new String[]{
"--output-directory", temporaryFolder.newFolder().toString()};
String[] allArgs = mergeArgs(baseArgs, extraArgs());
diff --git a/framework/src/test/java/org/tron/common/BaseTest.java b/framework/src/test/java/org/tron/common/BaseTest.java
index 6d075a2d6aa..471aaa3d383 100644
--- a/framework/src/test/java/org/tron/common/BaseTest.java
+++ b/framework/src/test/java/org/tron/common/BaseTest.java
@@ -7,7 +7,9 @@
import lombok.extern.slf4j.Slf4j;
import org.junit.AfterClass;
import org.junit.Assert;
+import org.junit.Before;
import org.junit.ClassRule;
+import org.junit.Rule;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
import org.springframework.test.annotation.DirtiesContext;
@@ -17,6 +19,7 @@
import org.tron.common.crypto.ECKey;
import org.tron.common.parameter.CommonParameter;
import org.tron.common.utils.Commons;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.Sha256Hash;
import org.tron.consensus.base.Param;
import org.tron.core.ChainBaseManager;
@@ -64,6 +67,9 @@ public abstract class BaseTest {
@ClassRule
public static final TemporaryFolder temporaryFolder = new TemporaryFolder();
+ @Rule
+ public final VMConfigRule vmConfigRule = new VMConfigRule();
+
@Resource
protected Manager dbManager;
@Resource
@@ -75,6 +81,11 @@ public abstract class BaseTest {
private static Application appT1;
+ @Before
+ public void resetPeerManagerState() {
+ PeerManagerStateResetter.reset();
+ }
+
@PostConstruct
private void prepare() {
appT1 = appT;
diff --git a/framework/src/test/java/org/tron/common/VMConfigRule.java b/framework/src/test/java/org/tron/common/VMConfigRule.java
new file mode 100644
index 00000000000..2ee7374d816
--- /dev/null
+++ b/framework/src/test/java/org/tron/common/VMConfigRule.java
@@ -0,0 +1,52 @@
+package org.tron.common;
+
+import java.lang.reflect.Field;
+import org.junit.rules.ExternalResource;
+import org.tron.common.parameter.CommonParameter;
+import org.tron.core.vm.config.ConfigLoader;
+import org.tron.core.vm.config.VMConfig;
+
+/**
+ * Restores VM flags after each test, including failed setup and assertion paths.
+ *
+ * Snapshotting enumerates {@link VMConfig.Snapshot} fields reflectively, so any static flag
+ * not mirrored there is outside this rule's protection: when adding a static flag to
+ * {@link VMConfig} or {@link ConfigLoader}, it must also be mirrored into
+ * {@code VMConfig.Snapshot} (before/after save and restore) or it will leak across tests.
+ *
+ *
This is a method-level rule: the baseline is captured before every test method, so global
+ * flags written from class-level {@code @BeforeClass} code are not covered — such classes must
+ * add their own {@code @AfterClass} to reset them manually (a leaked London hard-fork flag from
+ * class-level setup is an instance of exactly this gap).
+ */
+public class VMConfigRule extends ExternalResource {
+
+ private VMConfig.Snapshot savedSnapshot;
+ private boolean savedLoaderDisabled;
+ private boolean savedHardFork;
+ private boolean savedTrace;
+
+ @Override
+ protected void before() throws Exception {
+ Field global = VMConfig.class.getDeclaredField("globalSnapshot");
+ global.setAccessible(true);
+ VMConfig.Snapshot current = (VMConfig.Snapshot) global.get(null);
+ savedSnapshot = new VMConfig.Snapshot();
+ // init* methods mutate the snapshot in place, so saving only its reference is insufficient.
+ for (Field flag : VMConfig.Snapshot.class.getFields()) {
+ flag.set(savedSnapshot, flag.get(current));
+ }
+ savedLoaderDisabled = ConfigLoader.disable;
+ savedHardFork = CommonParameter.ENERGY_LIMIT_HARD_FORK;
+ savedTrace = VMConfig.vmTrace();
+ VMConfig.clearLocalSnapshot();
+ }
+
+ @Override
+ protected void after() {
+ VMConfig.setGlobalSnapshot(savedSnapshot);
+ ConfigLoader.disable = savedLoaderDisabled;
+ VMConfig.initVmHardFork(savedHardFork);
+ VMConfig.setVmTrace(savedTrace);
+ }
+}
diff --git a/framework/src/test/java/org/tron/common/backup/BackupManagerTest.java b/framework/src/test/java/org/tron/common/backup/BackupManagerTest.java
index 5ff02fc8cb5..0efbb13a481 100644
--- a/framework/src/test/java/org/tron/common/backup/BackupManagerTest.java
+++ b/framework/src/test/java/org/tron/common/backup/BackupManagerTest.java
@@ -1,5 +1,6 @@
package org.tron.common.backup;
+import io.netty.channel.Channel;
import java.lang.reflect.Field;
import java.lang.reflect.Method;
import java.net.InetAddress;
@@ -9,8 +10,8 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
-import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import org.junit.After;
import org.junit.Assert;
@@ -25,6 +26,7 @@
import org.tron.common.backup.socket.UdpEvent;
import org.tron.common.parameter.CommonParameter;
import org.tron.common.utils.PublicMethod;
+import org.tron.common.utils.ReflectUtils;
import org.tron.core.config.args.Args;
import org.tron.core.config.args.InetUtil;
@@ -34,7 +36,8 @@ public class BackupManagerTest {
public TemporaryFolder temporaryFolder = new TemporaryFolder();
private BackupManager manager;
private BackupServer backupServer;
- private BiFunction savedLookup;
+ private BiFunction previousDnsLookup;
+ private boolean backupServerClosed;
@Before
public void setUp() throws Exception {
@@ -43,13 +46,36 @@ public void setUp() throws Exception {
CommonParameter.getInstance().setBackupPort(PublicMethod.chooseRandomPort());
manager = new BackupManager();
backupServer = new BackupServer(manager);
- savedLookup = InetUtil.dnsLookup;
+ previousDnsLookup = InetUtil.dnsLookup;
}
@After
- public void tearDown() {
- InetUtil.dnsLookup = savedLookup;
+ public void tearDown() throws Exception {
+ List errors = new ArrayList<>();
+ Channel channel = null;
+ if (backupServer != null) {
+ try {
+ channel = BackupTestUtils.getChannel(backupServer);
+ } catch (Throwable t) {
+ errors.add(t);
+ }
+ }
+ if (!backupServerClosed && backupServer != null) {
+ BackupTestUtils.runQuietly(errors, backupServer::close);
+ }
+ if (manager != null) {
+ BackupTestUtils.runQuietly(errors, manager::stop);
+ }
+ Channel captured = channel;
+ BackupTestUtils.runQuietly(errors, () -> {
+ if (captured != null) {
+ Assert.assertFalse("backup channel must close", captured.isOpen());
+ }
+ BackupTestUtils.assertExecutorsTerminated(manager, backupServer);
+ });
+ InetUtil.dnsLookup = previousDnsLookup;
Args.clearParam();
+ BackupTestUtils.throwIfAnyError(errors);
}
@Test
@@ -121,7 +147,7 @@ public void test() throws Exception {
}
@Test
- public void testSendKeepAliveMessage() throws Exception {
+ public void testBackupServerLifecycleDuringKeepAliveInterval() throws Exception {
CommonParameter parameter = CommonParameter.getInstance();
parameter.setBackupPriority(8);
List members = new ArrayList<>();
@@ -134,21 +160,19 @@ public void testSendKeepAliveMessage() throws Exception {
Assert.assertEquals(manager.getStatus(), BackupManager.BackupStatusEnum.MASTER);
backupServer.initServer();
+ awaitBackupServerReady();
manager.init();
-
- Thread.sleep(parameter.getKeepAliveInterval() + 1000);//test send KeepAliveMessage
-
- field = manager.getClass().getDeclaredField("executorService");
- field.setAccessible(true);
- ScheduledExecutorService executorService = (ScheduledExecutorService) field.get(manager);
- executorService.shutdown();
-
- Field field2 = backupServer.getClass().getDeclaredField("executor");
- field2.setAccessible(true);
- ExecutorService executorService2 = (ExecutorService) field2.get(backupServer);
- executorService2.shutdown();
+ long keepAliveDeadline = System.nanoTime()
+ + TimeUnit.MILLISECONDS.toNanos(parameter.getKeepAliveInterval() + 1000L);
+ BackupTestUtils.awaitCondition("keep-alive interval",
+ () -> System.nanoTime() >= keepAliveDeadline);
Assert.assertEquals(BackupManager.BackupStatusEnum.INIT, manager.getStatus());
+ Channel channel = BackupTestUtils.getChannel(backupServer);
+ backupServer.close();
+ backupServerClosed = true;
+ Assert.assertFalse("backup channel must close", channel.isOpen());
+ BackupTestUtils.assertExecutorsTerminated(manager, backupServer);
}
// ===== domain-handling tests for init() =====
@@ -161,8 +185,8 @@ public void testInitResolvesDomainsToMembers() throws Exception {
InetUtil.dnsLookup = (host, ipv4) ->
("node.example.com".equals(host) && ipv4) ? resolved : null;
manager.init();
- Set members = getField(manager, "members");
- Map cache = getField(manager, "domainIpCache");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
+ Map cache = ReflectUtils.getFieldValue(manager, "domainIpCache");
Assert.assertTrue(members.contains("1.2.3.4"));
Assert.assertEquals("1.2.3.4", cache.get("node.example.com"));
manager.stop();
@@ -174,8 +198,8 @@ public void testInitSkipsUnresolvableDomain() throws Exception {
Collections.singletonList("bad.invalid.domain"));
InetUtil.dnsLookup = (host, ipv4) -> null;
manager.init();
- Set members = getField(manager, "members");
- Map cache = getField(manager, "domainIpCache");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
+ Map cache = ReflectUtils.getFieldValue(manager, "domainIpCache");
Assert.assertTrue("unresolvable domain should be silently dropped", members.isEmpty());
Assert.assertTrue(cache.isEmpty());
manager.stop();
@@ -190,7 +214,7 @@ public void testInitSkipsDomainResolvingToLocalIp() throws Exception {
InetUtil.dnsLookup = (host, ipv4) ->
("self.local.host".equals(host) && ipv4) ? selfAddr : null;
manager.init();
- Set members = getField(manager, "members");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
Assert.assertFalse("domain resolving to local IP should not be in members",
members.contains(localIp));
manager.stop();
@@ -200,8 +224,8 @@ public void testInitSkipsDomainResolvingToLocalIp() throws Exception {
@Test(timeout = 5000)
public void testRefreshMemberIpsIpChanged() throws Exception {
- Set members = getField(manager, "members");
- Map cache = getField(manager, "domainIpCache");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
+ Map cache = ReflectUtils.getFieldValue(manager, "domainIpCache");
members.add("1.1.1.1");
cache.put("peer.tron.network", "1.1.1.1");
@@ -216,8 +240,8 @@ public void testRefreshMemberIpsIpChanged() throws Exception {
@Test(timeout = 5000)
public void testRefreshMemberIpsIpUnchanged() throws Exception {
- Set members = getField(manager, "members");
- Map cache = getField(manager, "domainIpCache");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
+ Map cache = ReflectUtils.getFieldValue(manager, "domainIpCache");
members.add("1.1.1.1");
cache.put("peer.tron.network", "1.1.1.1");
@@ -231,8 +255,8 @@ public void testRefreshMemberIpsIpUnchanged() throws Exception {
@Test(timeout = 5000)
public void testRefreshMemberIpsDnsFailure() throws Exception {
- Set members = getField(manager, "members");
- Map cache = getField(manager, "domainIpCache");
+ Set members = ReflectUtils.getFieldValue(manager, "members");
+ Map cache = ReflectUtils.getFieldValue(manager, "domainIpCache");
members.add("1.1.1.1");
cache.put("peer.tron.network", "1.1.1.1");
@@ -242,16 +266,18 @@ public void testRefreshMemberIpsDnsFailure() throws Exception {
Assert.assertEquals("1.1.1.1", cache.get("peer.tron.network"));
}
- @SuppressWarnings("unchecked")
- private T getField(Object obj, String name) throws Exception {
- Field f = obj.getClass().getDeclaredField(name);
- f.setAccessible(true);
- return (T) f.get(obj);
- }
-
private void invokeRefreshMemberIps(BackupManager mgr) throws Exception {
Method m = mgr.getClass().getDeclaredMethod("refreshMemberIps");
m.setAccessible(true);
m.invoke(mgr);
}
+
+ private void awaitBackupServerReady() throws Exception {
+ BackupTestUtils.awaitCondition("backup channel to become active",
+ () -> BackupTestUtils.getChannel(backupServer) != null
+ && BackupTestUtils.getChannel(backupServer).isActive());
+ BackupTestUtils.awaitCondition("backup message handler assignment",
+ () -> ReflectUtils.getFieldObject(manager, "messageHandler") != null);
+ }
+
}
diff --git a/framework/src/test/java/org/tron/common/backup/BackupServerTest.java b/framework/src/test/java/org/tron/common/backup/BackupServerTest.java
index 50778970d87..b1d60d5d38b 100644
--- a/framework/src/test/java/org/tron/common/backup/BackupServerTest.java
+++ b/framework/src/test/java/org/tron/common/backup/BackupServerTest.java
@@ -1,8 +1,10 @@
package org.tron.common.backup;
+import io.netty.channel.Channel;
import java.util.ArrayList;
import java.util.List;
import org.junit.After;
+import org.junit.Assert;
import org.junit.Before;
import org.junit.Rule;
import org.junit.Test;
@@ -23,6 +25,8 @@ public class BackupServerTest {
@Rule
public Timeout globalTimeout = Timeout.seconds(60);
private BackupServer backupServer;
+ private BackupManager backupManager;
+ private boolean backupServerClosed;
@Before
public void setUp() throws Exception {
@@ -32,21 +36,39 @@ public void setUp() throws Exception {
List members = new ArrayList<>();
members.add("127.0.0.2");
CommonParameter.getInstance().setBackupMembers(members);
- BackupManager backupManager = new BackupManager();
+ backupManager = new BackupManager();
backupManager.init();
backupServer = new BackupServer(backupManager);
}
@After
- public void tearDown() {
- backupServer.close();
+ public void tearDown() throws Exception {
+ List errors = new ArrayList<>();
+ if (!backupServerClosed && backupServer != null) {
+ BackupTestUtils.runQuietly(errors, backupServer::close);
+ }
+ if (backupManager != null) {
+ BackupTestUtils.runQuietly(errors, backupManager::stop);
+ }
+ BackupTestUtils.runQuietly(errors,
+ () -> BackupTestUtils.assertExecutorsTerminated(backupManager, backupServer));
Args.clearParam();
+ BackupTestUtils.throwIfAnyError(errors);
}
@Test(timeout = 60_000)
- public void test() throws InterruptedException {
+ public void test() throws Exception {
backupServer.initServer();
- // wait for the server to start so channel is assigned before close() is called
- Thread.sleep(1000);
+ BackupTestUtils.awaitCondition("backup channel to become active",
+ () -> BackupTestUtils.getChannel(backupServer) != null
+ && BackupTestUtils.getChannel(backupServer).isActive());
+ Channel channel = BackupTestUtils.getChannel(backupServer);
+ Assert.assertTrue("backup channel must be active after startup", channel.isActive());
+
+ backupServer.close();
+ backupServerClosed = true;
+
+ Assert.assertFalse("backup channel must close", channel.isOpen());
+ BackupTestUtils.assertExecutorsTerminated(backupManager, backupServer);
}
}
diff --git a/framework/src/test/java/org/tron/common/backup/BackupTestUtils.java b/framework/src/test/java/org/tron/common/backup/BackupTestUtils.java
new file mode 100644
index 00000000000..45f2cae596f
--- /dev/null
+++ b/framework/src/test/java/org/tron/common/backup/BackupTestUtils.java
@@ -0,0 +1,72 @@
+package org.tron.common.backup;
+
+import io.netty.channel.Channel;
+import java.util.List;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.function.BooleanSupplier;
+import org.junit.Assert;
+import org.junit.function.ThrowingRunnable;
+import org.tron.common.backup.socket.BackupServer;
+import org.tron.common.utils.ReflectUtils;
+
+/**
+ * Shared reflection/await/cleanup helpers for backup tests. Assertion messages and timeout
+ * parameters mirror the helpers they replace, so failure output is unchanged.
+ */
+public final class BackupTestUtils {
+
+ private BackupTestUtils() {
+ }
+
+ public static Channel getChannel(BackupServer server) {
+ try {
+ return (Channel) ReflectUtils.getFieldObject(server, "channel");
+ } catch (Exception e) {
+ throw new AssertionError("cannot inspect backup channel", e);
+ }
+ }
+
+ public static void awaitCondition(String description, BooleanSupplier condition)
+ throws Exception {
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
+ while (System.nanoTime() < deadline) {
+ if (condition.getAsBoolean()) {
+ return;
+ }
+ Thread.sleep(20);
+ }
+ Assert.fail("timed out waiting for " + description);
+ }
+
+ public static void assertExecutorsTerminated(BackupManager manager, BackupServer server)
+ throws Exception {
+ if (manager == null || server == null) {
+ return;
+ }
+ ExecutorService managerExecutor =
+ (ExecutorService) ReflectUtils.getFieldObject(manager, "executorService");
+ Assert.assertTrue("backup manager executor must terminate", managerExecutor.isTerminated());
+ ExecutorService serverExecutor =
+ (ExecutorService) ReflectUtils.getFieldObject(server, "executor");
+ if (serverExecutor != null) {
+ Assert.assertTrue("backup server executor must terminate", serverExecutor.isTerminated());
+ }
+ }
+
+ public static void runQuietly(List errors, ThrowingRunnable step) {
+ try {
+ step.run();
+ } catch (Throwable t) {
+ errors.add(t);
+ }
+ }
+
+ public static void throwIfAnyError(List errors) {
+ if (!errors.isEmpty()) {
+ AssertionError failure = new AssertionError("backup test cleanup failed");
+ errors.forEach(failure::addSuppressed);
+ throw failure;
+ }
+ }
+}
diff --git a/framework/src/test/java/org/tron/common/logsfilter/EventLoaderTest.java b/framework/src/test/java/org/tron/common/logsfilter/EventLoaderTest.java
index 958af4f7b7b..0857b1b9391 100644
--- a/framework/src/test/java/org/tron/common/logsfilter/EventLoaderTest.java
+++ b/framework/src/test/java/org/tron/common/logsfilter/EventLoaderTest.java
@@ -15,6 +15,7 @@
import org.pf4j.PluginWrapper;
import org.tron.common.logsfilter.trigger.BlockLogTrigger;
import org.tron.common.logsfilter.trigger.TransactionLogTrigger;
+import org.tron.common.utils.PublicMethod;
public class EventLoaderTest {
@@ -22,7 +23,7 @@ public class EventLoaderTest {
public void launchNativeQueue() {
EventPluginConfig config = new EventPluginConfig();
config.setSendQueueLength(1000);
- config.setBindPort(5555);
+ config.setBindPort(PublicMethod.chooseRandomPort());
config.setUseNativeQueue(true);
config.setPluginPath("pluginPath");
config.setServerAddress("serverAddress");
@@ -48,9 +49,11 @@ public void launchNativeQueue() {
config.setTriggerConfigList(triggerConfigList);
- assertTrue(EventPluginLoader.getInstance().start(config));
-
- EventPluginLoader.getInstance().stopPlugin();
+ try {
+ assertTrue(EventPluginLoader.getInstance().start(config));
+ } finally {
+ EventPluginLoader.getInstance().stopPlugin();
+ }
}
@Test
diff --git a/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java b/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java
index 5219654977b..b32f1c22d39 100644
--- a/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java
+++ b/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java
@@ -6,13 +6,16 @@
import org.junit.Test;
import org.tron.common.es.ExecutorServiceManager;
import org.tron.common.logsfilter.nativequeue.NativeMessageQueue;
+import org.tron.common.utils.PublicMethod;
import org.zeromq.SocketType;
import org.zeromq.ZContext;
import org.zeromq.ZMQ;
public class NativeMessageQueueTest {
- public int bindPort = 5555;
+ // Random port avoids fixed 5555 conflicts; note invalidBindPort/invalidSendLength still
+ // remap to DEFAULT_BIND_PORT (5555) in production start() — known low-risk residual.
+ public int bindPort = PublicMethod.chooseRandomPort();
public String dataToSend = "################";
public String topic = "testTopic";
diff --git a/framework/src/test/java/org/tron/common/prometheus/SRMetricsTest.java b/framework/src/test/java/org/tron/common/prometheus/SRMetricsTest.java
index 4c2e9292d29..4c1404bb232 100644
--- a/framework/src/test/java/org/tron/common/prometheus/SRMetricsTest.java
+++ b/framework/src/test/java/org/tron/common/prometheus/SRMetricsTest.java
@@ -13,6 +13,7 @@
import org.junit.Test;
import org.tron.common.BaseTest;
import org.tron.common.TestConstants;
+import org.tron.common.utils.PublicMethod;
import org.tron.common.utils.StringUtil;
import org.tron.consensus.dpos.MaintenanceManager;
import org.tron.core.capsule.AccountCapsule;
@@ -38,6 +39,7 @@ public class SRMetricsTest extends BaseTest {
Args.setParam(new String[]{"-d", dbPath()}, TestConstants.TEST_CONF);
Args.getInstance().setNodeListenPort(20000 + PORT.incrementAndGet());
Args.getInstance().setMetricsPrometheusEnable(true);
+ Args.getInstance().setMetricsPrometheusPort(PublicMethod.chooseRandomPort());
Metrics.init();
}
diff --git a/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmCompatibleEvmTest.java b/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmCompatibleEvmTest.java
index 74d44dfca7d..a3711ba8de7 100644
--- a/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmCompatibleEvmTest.java
+++ b/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmCompatibleEvmTest.java
@@ -306,6 +306,7 @@ public static void afterClass() {
VMConfig.initAllowTvmConstantinople(0);
VMConfig.initAllowTvmSolidity059(0);
VMConfig.initAllowTvmIstanbul(0);
+ VMConfig.initAllowTvmLondon(0);
VMConfig.initAllowTvmCompatibleEvm(0);
}
diff --git a/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmLondonTest.java b/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmLondonTest.java
index e93eca39092..11a02e615db 100644
--- a/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmLondonTest.java
+++ b/framework/src/test/java/org/tron/common/runtime/vm/AllowTvmLondonTest.java
@@ -74,8 +74,8 @@ public void testBaseFee() throws ContractExeException, ReceiptCheckErrException,
factoryAddress, Hex.decode(hexInput), 0, feeLimit, manager, null);
byte[] returnValue = result.getRuntime().getResult().getHReturn();
Assert.assertNull(result.getRuntime().getRuntimeError());
- Assert.assertArrayEquals(returnValue,
- longTo32Bytes(manager.getDynamicPropertiesStore().getEnergyFee()));
+ Assert.assertArrayEquals(longTo32Bytes(manager.getDynamicPropertiesStore().getEnergyFee()),
+ returnValue);
}
@Test
diff --git a/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetter.java b/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetter.java
new file mode 100644
index 00000000000..b4766115f2d
--- /dev/null
+++ b/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetter.java
@@ -0,0 +1,99 @@
+package org.tron.common.utils;
+
+import java.lang.reflect.Field;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.springframework.util.ReflectionUtils;
+import org.tron.common.es.ExecutorServiceManager;
+import org.tron.core.net.peer.PeerConnection;
+import org.tron.core.net.peer.PeerManager;
+import org.tron.protos.Protocol.ReasonCode;
+
+/**
+ * Test source-set utility: restores {@link PeerManager} to a cold-start state.
+ *
+ * {@link PeerManager#close()} neither clears the raw static peers/counters nor rebuilds the
+ * static executor, and tests share one JVM across Spring contexts.
+ *
+ *
The old executor is drained to termination before touching any state: {@code
+ * check()} is not synchronized (it snapshots peers, then removes entries and decrements
+ * counters), so a task left running would decrement the counters zeroed in step 4 and leave
+ * them negative; shutdown alone does not guarantee that, so reset fails if termination cannot
+ * be confirmed.
+ * Residual peers are disconnected before the raw list is cleared because {@code close()} may
+ * fail midway and leave live channels that a bare {@code clear()} would orphan; each peer is
+ * handled defensively (null channel tolerated, per-peer catch). A fresh executor is then
+ * installed (lazy thread, no tasks until the next {@code init()}).
+ *
+ *
Wired broadly from BaseTest/BaseMethodTest against unknown prior pollution; the reset is
+ * idempotent and cheap for tests that never use PeerManager. Remove this utility once
+ * production {@code close()}/{@code init()} is restart-safe.
+ */
+public final class PeerManagerStateResetter {
+
+ private static final String EXECUTOR_NAME = "peer-manager";
+
+ private PeerManagerStateResetter() {
+ }
+
+ public static synchronized void reset() {
+ // 1) Drain the old executor first: let running/queued check() tasks die out so they
+ // cannot interleave with the list/counter reset below. Gate on isTerminated(), not
+ // isShutdown(): shutdown() still lets a running check() finish asynchronously, and
+ // check() decrements the counters even when its peers.remove() is a no-op. If
+ // termination cannot be confirmed, fail the setup instead of resetting anyway.
+ ScheduledExecutorService executor = getFieldValue("executor");
+ if (executor != null && !executor.isTerminated()) {
+ ExecutorServiceManager.shutdownAndAwaitTermination(executor, EXECUTOR_NAME);
+ if (!executor.isTerminated()) {
+ throw new IllegalStateException(
+ "peer-manager executor did not terminate; refusing to reset shared state");
+ }
+ }
+ // 2) Unconditionally install a fresh executor (the old one may be shut down or null);
+ // its thread is created lazily.
+ setFieldValue("executor",
+ ExecutorServiceManager.newSingleThreadScheduledExecutor(EXECUTOR_NAME));
+
+ // 3) Release residual live connections before clearing the raw list.
+ List peers = getFieldValue("peers");
+ if (peers == null) {
+ setFieldValue("peers", Collections.synchronizedList(new ArrayList()));
+ } else {
+ for (PeerConnection peer : new ArrayList<>(peers)) {
+ try {
+ if (!peer.isDisconnect()) {
+ peer.disconnect(ReasonCode.PEER_QUITING);
+ if (peer.getChannel() != null) {
+ peer.getChannel().close();
+ }
+ }
+ } catch (Exception e) {
+ // best effort: a single corrupted leftover peer must not fail the reset
+ }
+ }
+ peers.clear();
+ }
+
+ // 4) Zero the counters; old tasks can no longer decrement them at this point.
+ AtomicInteger active = PeerManager.getActivePeersCount();
+ AtomicInteger passive = PeerManager.getPassivePeersCount();
+ active.set(0);
+ passive.set(0);
+ }
+
+ private static T getFieldValue(String fieldName) {
+ Field field = ReflectionUtils.findField(PeerManager.class, fieldName);
+ ReflectionUtils.makeAccessible(field);
+ return (T) ReflectionUtils.getField(field, null);
+ }
+
+ private static void setFieldValue(String fieldName, Object value) {
+ Field field = ReflectionUtils.findField(PeerManager.class, fieldName);
+ ReflectionUtils.makeAccessible(field);
+ ReflectionUtils.setField(field, null, value);
+ }
+}
diff --git a/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetterTest.java b/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetterTest.java
new file mode 100644
index 00000000000..d7316c06333
--- /dev/null
+++ b/framework/src/test/java/org/tron/common/utils/PeerManagerStateResetterTest.java
@@ -0,0 +1,79 @@
+package org.tron.common.utils;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.Modifier;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.junit.Test;
+import org.mockito.Mockito;
+import org.tron.core.net.peer.PeerConnection;
+import org.tron.core.net.peer.PeerManager;
+
+/**
+ * Pins the cold-start guarantees of {@link PeerManagerStateResetter#reset()}: after a reset the
+ * shared peers list is empty, both counters are zero, and the installed executor is fresh —
+ * including on the path where the old executor was already shut down but might still be
+ * draining a {@code check()} task.
+ */
+public class PeerManagerStateResetterTest {
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testResetRestoresColdStartState() throws Exception {
+ Field peersField = PeerManager.class.getDeclaredField("peers");
+ peersField.setAccessible(true);
+ List peers = (List) peersField.get(null);
+ peers.clear();
+ PeerConnection stalePeer = Mockito.mock(PeerConnection.class);
+ Mockito.when(stalePeer.isDisconnect()).thenReturn(true);
+ peers.add(stalePeer);
+ AtomicInteger active = PeerManager.getActivePeersCount();
+ AtomicInteger passive = PeerManager.getPassivePeersCount();
+ active.set(7);
+ passive.set(3);
+
+ PeerManagerStateResetter.reset();
+
+ assertEquals(0, peers.size());
+ assertEquals(0, active.get());
+ assertEquals(0, passive.get());
+ Field executorField = PeerManager.class.getDeclaredField("executor");
+ executorField.setAccessible(true);
+ ScheduledExecutorService executor = (ScheduledExecutorService) executorField.get(null);
+ assertFalse(executor.isShutdown());
+ }
+
+ // pin: reset() clears PeerManager statics by hardcoded field names — a new static field
+ // silently leaks across tests unless it gets resetter coverage or an allowlist entry here.
+ @Test
+ public void resetterCoversAllPeerManagerStaticFields() {
+ // fields reset() actually drains/rebuilds/clears/zeroes
+ Set handled = new HashSet<>(Arrays.asList(
+ "peers", "executor", "activePeersCount", "passivePeersCount"));
+ // fields intentionally untouched: constants / config that never mutates across tests
+ Set allowed = new HashSet<>(Arrays.asList(
+ "esName", "DISCONNECTION_TIME_OUT", "logger"));
+
+ List unclassified = new java.util.ArrayList<>();
+ for (Field field : PeerManager.class.getDeclaredFields()) {
+ // skip compiler/JaCoCo-generated synthetic fields (e.g. $jacocoData) — not business state
+ if (!Modifier.isStatic(field.getModifiers()) || field.isSynthetic()) {
+ continue;
+ }
+ String name = field.getName();
+ if (!handled.contains(name) && !allowed.contains(name)) {
+ unclassified.add(name + " (" + field.getType().getSimpleName() + ")");
+ }
+ }
+ assertEquals("new static field needs resetter coverage or explicit allowlist entry: "
+ + unclassified, Collections.emptyList(), unclassified);
+ }
+}
diff --git a/framework/src/test/java/org/tron/core/event/BlockEventGetTest.java b/framework/src/test/java/org/tron/core/event/BlockEventGetTest.java
index e2815e46063..6df63c6c04e 100644
--- a/framework/src/test/java/org/tron/core/event/BlockEventGetTest.java
+++ b/framework/src/test/java/org/tron/core/event/BlockEventGetTest.java
@@ -125,6 +125,10 @@ public void before() throws IOException {
@AfterClass
public static void after() throws IOException {
+ // stopPlugin() is safe when never started: it null-checks pluginManager, and
+ // NativeMessageQueue.stop() null-checks publisher/context. Ensures the native
+ // queue socket bound in test() is released even when assertions fail earlier.
+ EventPluginLoader.getInstance().stopPlugin();
context.destroy();
Args.clearParam();
}
@@ -174,7 +178,7 @@ public void test() throws Exception {
EventPluginConfig config = new EventPluginConfig();
config.setSendQueueLength(1000);
- config.setBindPort(5555);
+ config.setBindPort(PublicMethod.chooseRandomPort());
config.setUseNativeQueue(true);
config.setTriggerConfigList(new ArrayList<>());
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..4722caec08d 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
@@ -62,6 +62,7 @@ public class PrometheusApiServiceTest extends BaseTest {
Args.setParam(new String[] {"-d", dbPath()}, TestConstants.TEST_CONF);
Args.getInstance().setNodeListenPort(10000 + port.incrementAndGet());
initParameter(Args.getInstance());
+ Args.getInstance().setMetricsPrometheusPort(PublicMethod.chooseRandomPort());
Metrics.init();
}
diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/MessageHandlerTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/MessageHandlerTest.java
index be843674632..c8205b6b721 100644
--- a/framework/src/test/java/org/tron/core/net/messagehandler/MessageHandlerTest.java
+++ b/framework/src/test/java/org/tron/core/net/messagehandler/MessageHandlerTest.java
@@ -16,6 +16,7 @@
import org.tron.common.ClassLevelAppContextFixture;
import org.tron.common.TestConstants;
import org.tron.common.application.TronApplicationContext;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.ReflectUtils;
import org.tron.common.utils.Sha256Hash;
import org.tron.consensus.pbft.message.PbftMessage;
@@ -45,6 +46,7 @@ public class MessageHandlerTest {
@BeforeClass
public static void init() throws Exception {
+ PeerManagerStateResetter.reset();
Args.setParam(new String[] {"--output-directory",
temporaryFolder.newFolder().toString(), "--debug"}, TestConstants.TEST_CONF);
context = APP_FIXTURE.createContext();
diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/PbftMsgHandlerTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/PbftMsgHandlerTest.java
index 65a8f615bfe..15d7107b58f 100644
--- a/framework/src/test/java/org/tron/core/net/messagehandler/PbftMsgHandlerTest.java
+++ b/framework/src/test/java/org/tron/core/net/messagehandler/PbftMsgHandlerTest.java
@@ -17,6 +17,7 @@
import org.tron.common.crypto.SignInterface;
import org.tron.common.crypto.SignUtils;
import org.tron.common.utils.FileUtil;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.PublicMethod;
import org.tron.common.utils.ReflectUtils;
import org.tron.common.utils.Sha256Hash;
@@ -46,6 +47,7 @@ public class PbftMsgHandlerTest {
@BeforeClass
public static void init() {
+ PeerManagerStateResetter.reset();
Args.setParam(new String[] {"--output-directory", dbPath, "--debug"},
TestConstants.TEST_CONF);
context = new TronApplicationContext(DefaultConfig.class);
diff --git a/framework/src/test/java/org/tron/core/net/messagehandler/TransactionsMsgHandlerTest.java b/framework/src/test/java/org/tron/core/net/messagehandler/TransactionsMsgHandlerTest.java
index 78af06e64bc..282c80f9f6d 100644
--- a/framework/src/test/java/org/tron/core/net/messagehandler/TransactionsMsgHandlerTest.java
+++ b/framework/src/test/java/org/tron/core/net/messagehandler/TransactionsMsgHandlerTest.java
@@ -10,11 +10,13 @@
import java.util.Map;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.TimeUnit;
-import lombok.Getter;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
@@ -23,7 +25,6 @@
import org.tron.common.TestConstants;
import org.tron.common.runtime.TvmTestUtils;
import org.tron.common.utils.ByteArray;
-import org.tron.common.utils.ReflectUtils;
import org.tron.core.ChainBaseManager;
import org.tron.core.config.args.Args;
import org.tron.core.exception.P2pException;
@@ -48,6 +49,7 @@ public static void init() {
@Test
public void testProcessMessage() {
TransactionsMsgHandler transactionsMsgHandler = new TransactionsMsgHandler();
+ ExecutorService originalPool = null;
try {
transactionsMsgHandler.init();
@@ -80,17 +82,29 @@ public void testProcessMessage() {
Item item = new Item(new TransactionMessage(trx).getMessageId(),
Protocol.Inventory.InventoryType.TRX);
advInvRequest.put(item, 0L);
+ // The non-executing pool must be installed before the first submission so no
+ // real-pool worker can touch the peer mock while the test re-stubs it (Mockito
+ // stubbing is not thread-safe). The latch counts down only for off-thread callers,
+ // which after the replacement is exactly the smart-contract scheduler.
+ CountDownLatch smartContractSubmitted = new CountDownLatch(1);
+ Thread testThread = Thread.currentThread();
+ ExecutorService mockPool = Mockito.mock(ExecutorService.class);
+ Future> submittedTask = Mockito.mock(Future.class);
+ Mockito.when(mockPool.submit(Mockito.any(Runnable.class))).thenAnswer(invocation -> {
+ if (Thread.currentThread() != testThread) {
+ smartContractSubmitted.countDown();
+ }
+ return submittedTask;
+ });
+ originalPool = replaceTrxHandlePool(transactionsMsgHandler, mockPool);
+
Mockito.when(peer.getAdvInvRequest()).thenReturn(advInvRequest);
List transactionList = new ArrayList<>();
transactionList.add(trx);
transactionsMsgHandler.processMessage(peer, new TransactionsMessage(transactionList));
Assert.assertNull(advInvRequest.get(item));
- //Thread.sleep(10);
- BlockingQueue smartContractQueue =
- new LinkedBlockingQueue(2);
- smartContractQueue.offer(new TrxEvent(null, null));
- smartContractQueue.offer(new TrxEvent(null, null));
+ BlockingQueue> smartContractQueue = new LinkedBlockingQueue<>(1);
Field field1 = TransactionsMsgHandler.class.getDeclaredField("smartContractQueue");
field1.setAccessible(true);
field1.set(transactionsMsgHandler, smartContractQueue);
@@ -99,15 +113,27 @@ public void testProcessMessage() {
ByteArray.fromHexString("121212a9cf"),
ByteArray.fromHexString("123456"),
100, 100000000, 0, 0);
+ Protocol.Transaction trx3 = TvmTestUtils.generateTriggerSmartContractAndGetTransaction(
+ ByteArray.fromHexString("121212a9cf"),
+ ByteArray.fromHexString("121212a9cf"),
+ ByteArray.fromHexString("123457"),
+ 100, 100000000, 0, 0);
Map- advInvRequest1 = new ConcurrentHashMap<>();
Item item1 = new Item(new TransactionMessage(trx1).getMessageId(),
Protocol.Inventory.InventoryType.TRX);
advInvRequest1.put(item1, 0L);
+ Item item3 = new Item(new TransactionMessage(trx3).getMessageId(),
+ Protocol.Inventory.InventoryType.TRX);
+ advInvRequest1.put(item3, 0L);
Mockito.when(peer.getAdvInvRequest()).thenReturn(advInvRequest1);
List transactionList1 = new ArrayList<>();
transactionList1.add(trx1);
+ transactionList1.add(trx3);
transactionsMsgHandler.processMessage(peer, new TransactionsMessage(transactionList1));
- Assert.assertNull(advInvRequest.get(item1));
+ Assert.assertNull(advInvRequest1.get(item1));
+ Assert.assertNull(advInvRequest1.get(item3));
+ Assert.assertTrue("smart-contract scheduler did not submit work",
+ smartContractSubmitted.await(3, TimeUnit.SECONDS));
// test 0 contract
Protocol.Transaction trx2 = Protocol.Transaction.newBuilder().setRawData(
@@ -132,37 +158,40 @@ public void testProcessMessage() {
Assert.assertTrue(true);
}
} catch (Exception e) {
- Assert.fail();
+ Assert.fail(e.getMessage());
} finally {
- transactionsMsgHandler.close();
+ closeHandlerAndOriginalPool(transactionsMsgHandler, originalPool);
}
}
@Test
public void testProcessMessageAfterClose() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
- handler.init();
- handler.close();
+ try {
+ handler.init();
+ handler.close();
- PeerConnection peer = Mockito.mock(PeerConnection.class);
- TransactionsMessage msg = Mockito.mock(TransactionsMessage.class);
+ PeerConnection peer = Mockito.mock(PeerConnection.class);
+ TransactionsMessage msg = Mockito.mock(TransactionsMessage.class);
- handler.processMessage(peer, msg);
+ handler.processMessage(peer, msg);
- Mockito.verify(msg, Mockito.never()).getTransactions();
- Mockito.verifyNoInteractions(peer);
+ Mockito.verify(msg, Mockito.never()).getTransactions();
+ Mockito.verifyNoInteractions(peer);
+ } finally {
+ handler.close();
+ }
}
@Test
public void testRejectedExecution() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
+ ExecutorService originalPool = null;
try {
ExecutorService mockPool = Mockito.mock(ExecutorService.class);
Mockito.when(mockPool.submit(Mockito.any(Runnable.class)))
.thenThrow(new RejectedExecutionException("pool closed"));
- Field poolField = TransactionsMsgHandler.class.getDeclaredField("trxHandlePool");
- poolField.setAccessible(true);
- poolField.set(handler, mockPool);
+ originalPool = replaceTrxHandlePool(handler, mockPool);
PeerConnection peer = Mockito.mock(PeerConnection.class);
TransactionsMessage msg = buildTransferMessage(2);
@@ -172,26 +201,26 @@ public void testRejectedExecution() throws Exception {
Mockito.verify(mockPool, Mockito.times(1)).submit(Mockito.any(Runnable.class));
} finally {
- handler.close();
+ closeHandlerAndOriginalPool(handler, originalPool);
}
}
@Test
public void testCloseDuringProcessing() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
+ ExecutorService originalPool = null;
try {
Field closedField = TransactionsMsgHandler.class.getDeclaredField("isClosed");
closedField.setAccessible(true);
ExecutorService mockPool = Mockito.mock(ExecutorService.class);
+ Future> submittedTask = Mockito.mock(Future.class);
// on the first submit, flip isClosed to true so the second iteration breaks
Mockito.when(mockPool.submit(Mockito.any(Runnable.class))).thenAnswer(inv -> {
closedField.set(handler, true);
- return null;
+ return submittedTask;
});
- Field poolField = TransactionsMsgHandler.class.getDeclaredField("trxHandlePool");
- poolField.setAccessible(true);
- poolField.set(handler, mockPool);
+ originalPool = replaceTrxHandlePool(handler, mockPool);
PeerConnection peer = Mockito.mock(PeerConnection.class);
TransactionsMessage msg = buildTransferMessage(2);
@@ -200,7 +229,7 @@ public void testCloseDuringProcessing() throws Exception {
Mockito.verify(mockPool, Mockito.times(1)).submit(Mockito.any(Runnable.class));
} finally {
- handler.close();
+ closeHandlerAndOriginalPool(handler, originalPool);
}
}
@@ -234,6 +263,34 @@ private void stubAdvInvRequest(PeerConnection peer, TransactionsMessage msg) {
Mockito.when(peer.getAdvInvRequest()).thenReturn(advInvRequest);
}
+ private ExecutorService replaceTrxHandlePool(TransactionsMsgHandler handler, ExecutorService pool)
+ throws Exception {
+ Field poolField = TransactionsMsgHandler.class.getDeclaredField("trxHandlePool");
+ poolField.setAccessible(true);
+ ExecutorService originalPool = (ExecutorService) poolField.get(handler);
+ poolField.set(handler, pool);
+ return originalPool;
+ }
+
+ private void closeHandlerAndOriginalPool(TransactionsMsgHandler handler,
+ ExecutorService originalPool) {
+ try {
+ handler.close();
+ } finally {
+ if (originalPool != null) {
+ originalPool.shutdown();
+ try {
+ if (!originalPool.awaitTermination(5, TimeUnit.SECONDS)) {
+ originalPool.shutdownNow();
+ }
+ } catch (InterruptedException e) {
+ originalPool.shutdownNow();
+ Thread.currentThread().interrupt();
+ }
+ }
+ }
+ }
+
@Test
public void testHandleTransaction() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
@@ -341,7 +398,20 @@ public void testDuplicateTransactionRejected() throws Exception {
public void testInvalidSigLength() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
handler.init();
+ ExecutorService originalPool = null;
try {
+ // Mock pool never executes submitted tasks: the async worker would invoke isBadPeer()
+ // on the stubbed peer concurrently with main-thread re-stubbing of getAdvInvRequest(),
+ // and Mockito's per-mock invocationForStubbing state is not thread-safe
+ // (intermittent WrongTypeOfReturnValue: ConcurrentHashMap cannot be returned by
+ // isBadPeer()). This test only asserts the synchronous check() length validation,
+ // so not running the worker is intentional.
+ ExecutorService mockPool = Mockito.mock(ExecutorService.class);
+ Future> submittedTask = Mockito.mock(Future.class);
+ Mockito.when(mockPool.submit(Mockito.any(Runnable.class)))
+ .thenAnswer(invocation -> submittedTask);
+ originalPool = replaceTrxHandlePool(handler, mockPool);
+
PeerConnection peer = Mockito.mock(PeerConnection.class);
BalanceContract.TransferContract transferContract = BalanceContract.TransferContract
@@ -418,45 +488,32 @@ public void testInvalidSigLength() throws Exception {
stubAdvInvRequest(peer, new TransactionsMessage(paddedList));
handler.processMessage(peer, new TransactionsMessage(paddedList));
} finally {
- handler.close();
+ closeHandlerAndOriginalPool(handler, originalPool);
}
}
@Test
public void testIsBusyWithCachedTransactions() throws Exception {
TransactionsMsgHandler handler = new TransactionsMsgHandler();
+ try {
+ int threshold = Args.getInstance().getMaxTrxCacheSize();
+ TronNetDelegate tronNetDelegateMock = Mockito.mock(TronNetDelegate.class);
+ Field field = TransactionsMsgHandler.class.getDeclaredField("tronNetDelegate");
+ field.setAccessible(true);
+ field.set(handler, tronNetDelegateMock);
- int threshold = Args.getInstance().getMaxTrxCacheSize();
- TronNetDelegate tronNetDelegateMock = Mockito.mock(TronNetDelegate.class);
- Field field = TransactionsMsgHandler.class.getDeclaredField("tronNetDelegate");
- field.setAccessible(true);
- field.set(handler, tronNetDelegateMock);
-
- // queue and smartContractQueue are empty, but cached size > threshold
- Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(threshold + 1);
- Assert.assertTrue(handler.isBusy());
-
- // boundary: cached size == threshold, isBusy() uses strict >, so not busy
- Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(threshold);
- Assert.assertFalse(handler.isBusy());
-
- Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(0);
- Assert.assertFalse(handler.isBusy());
- }
-
- class TrxEvent {
+ // queue and smartContractQueue are empty, but cached size > threshold
+ Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(threshold + 1);
+ Assert.assertTrue(handler.isBusy());
- @Getter
- private PeerConnection peer;
- @Getter
- private TransactionMessage msg;
- @Getter
- private long time;
+ // boundary: cached size == threshold, isBusy() uses strict >, so not busy
+ Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(threshold);
+ Assert.assertFalse(handler.isBusy());
- public TrxEvent(PeerConnection peer, TransactionMessage msg) {
- this.peer = peer;
- this.msg = msg;
- this.time = System.currentTimeMillis();
+ Mockito.when(tronNetDelegateMock.getCachedTransactionSize()).thenReturn(0);
+ Assert.assertFalse(handler.isBusy());
+ } finally {
+ handler.close();
}
}
}
diff --git a/framework/src/test/java/org/tron/core/net/peer/PeerManagerTest.java b/framework/src/test/java/org/tron/core/net/peer/PeerManagerTest.java
index ffba127a6fd..16e88b38584 100644
--- a/framework/src/test/java/org/tron/core/net/peer/PeerManagerTest.java
+++ b/framework/src/test/java/org/tron/core/net/peer/PeerManagerTest.java
@@ -17,6 +17,7 @@
import org.springframework.context.ApplicationContext;
import org.tron.common.TestConstants;
import org.tron.common.parameter.CommonParameter;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.ReflectUtils;
import org.tron.core.config.args.Args;
import org.tron.p2p.connection.Channel;
@@ -25,6 +26,7 @@ public class PeerManagerTest {
@BeforeClass
public static void initArgs() {
+ PeerManagerStateResetter.reset();
Args.setParam(new String[]{}, TestConstants.TEST_CONF);
CommonParameter.getInstance().setRateLimiterSyncBlockChain(10);
CommonParameter.getInstance().setRateLimiterFetchInvData(10);
diff --git a/framework/src/test/java/org/tron/core/net/services/HandShakeServiceTest.java b/framework/src/test/java/org/tron/core/net/services/HandShakeServiceTest.java
index b8b0d5f6deb..dce5ccb851f 100644
--- a/framework/src/test/java/org/tron/core/net/services/HandShakeServiceTest.java
+++ b/framework/src/test/java/org/tron/core/net/services/HandShakeServiceTest.java
@@ -19,6 +19,7 @@
import org.springframework.context.ApplicationContext;
import org.tron.common.TestConstants;
import org.tron.common.application.TronApplicationContext;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.ReflectUtils;
import org.tron.common.utils.Sha256Hash;
import org.tron.core.ChainBaseManager;
@@ -52,6 +53,7 @@ public class HandShakeServiceTest {
@BeforeClass
public static void init() throws Exception {
+ PeerManagerStateResetter.reset();
Args.setParam(new String[] {"--output-directory",
temporaryFolder.newFolder().toString(), "--debug"}, TestConstants.TEST_CONF);
context = new TronApplicationContext(DefaultConfig.class);
diff --git a/framework/src/test/java/org/tron/core/services/WalletApiTest.java b/framework/src/test/java/org/tron/core/services/WalletApiTest.java
index 4a55556afb1..25b21f30872 100644
--- a/framework/src/test/java/org/tron/core/services/WalletApiTest.java
+++ b/framework/src/test/java/org/tron/core/services/WalletApiTest.java
@@ -17,6 +17,7 @@
import org.tron.common.ClassLevelAppContextFixture;
import org.tron.common.TestConstants;
import org.tron.common.application.TronApplicationContext;
+import org.tron.common.utils.PeerManagerStateResetter;
import org.tron.common.utils.PublicMethod;
import org.tron.common.utils.TimeoutInterceptor;
import org.tron.core.config.args.Args;
@@ -38,6 +39,7 @@ public class WalletApiTest {
@BeforeClass
public static void init() throws IOException {
+ PeerManagerStateResetter.reset();
Args.setParam(new String[] {"-d", temporaryFolder.newFolder().toString(),
"--p2p-disable", "true"}, TestConstants.TEST_CONF);
Args.getInstance().setRpcPort(PublicMethod.chooseRandomPort());
diff --git a/framework/src/test/java/org/tron/core/zksnark/SendCoinShieldTest.java b/framework/src/test/java/org/tron/core/zksnark/SendCoinShieldTest.java
index 08de83ca8bf..efa60139b12 100644
--- a/framework/src/test/java/org/tron/core/zksnark/SendCoinShieldTest.java
+++ b/framework/src/test/java/org/tron/core/zksnark/SendCoinShieldTest.java
@@ -14,6 +14,7 @@
import java.util.Optional;
import javax.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
+import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.BeforeClass;
@@ -106,6 +107,7 @@ public class SendCoinShieldTest extends BaseTest {
private static final int VOTE_SCORE = 2;
private static final String DESCRIPTION = "TRX";
private static final String URL = "https://tron.network";
+ private long previousAllowShieldedTransaction;
@Resource
private Wallet wallet;
@@ -130,6 +132,8 @@ public static void initZksnarkParams() {
*/
@Before
public void init() {
+ previousAllowShieldedTransaction = dbManager.getDynamicPropertiesStore()
+ .getAllowShieldedTransaction();
if (init) {
return;
}
@@ -155,6 +159,12 @@ public void init() {
init = true;
}
+ @After
+ public void restoreAllowShieldedTransaction() {
+ dbManager.getDynamicPropertiesStore()
+ .saveAllowShieldedTransaction(previousAllowShieldedTransaction);
+ }
+
private void addZeroValueOutputNote(ZenTransactionBuilder builder) throws ZksnarkException {
SpendingKey spendingKey = SpendingKey.random();
FullViewingKey fullViewingKey = spendingKey.fullViewingKey();
diff --git a/framework/src/test/java/org/tron/core/zksnark/ShieldedReceiveTest.java b/framework/src/test/java/org/tron/core/zksnark/ShieldedReceiveTest.java
index 5854b731e97..e62396bc046 100755
--- a/framework/src/test/java/org/tron/core/zksnark/ShieldedReceiveTest.java
+++ b/framework/src/test/java/org/tron/core/zksnark/ShieldedReceiveTest.java
@@ -8,7 +8,6 @@
import com.google.protobuf.Any;
import com.google.protobuf.ByteString;
import com.google.protobuf.InvalidProtocolBufferException;
-import java.lang.reflect.Field;
import java.security.SignatureException;
import java.util.Arrays;
import java.util.HashSet;
@@ -21,6 +20,7 @@
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
+import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.BeforeClass;
@@ -47,7 +47,6 @@
import org.tron.common.zksnark.LibrustzcashParam.OutputProofParams;
import org.tron.common.zksnark.LibrustzcashParam.SpendSigParams;
import org.tron.consensus.dpos.DposSlot;
-import org.tron.consensus.dpos.DposTask;
import org.tron.core.Wallet;
import org.tron.core.actuator.Actuator;
import org.tron.core.actuator.ActuatorCreator;
@@ -126,6 +125,7 @@ public class ShieldedReceiveTest extends BaseTest {
"librustzcashSaplingCheckSpend error",
"Rt is invalid."
));
+ private long previousAllowShieldedTransaction;
private static final String FROM_ADDRESS;
private static final String ADDRESS_ONE_PRIVATE_KEY;
@@ -143,13 +143,11 @@ public class ShieldedReceiveTest extends BaseTest {
@Resource
private ConsensusService consensusService;
@Resource
- private DposTask dposTask;
- @Resource
private Wallet wallet;
@Resource
private DposSlot dposSlot;
-
private static boolean init;
+ private static boolean consensusScheduleInitialized;
static {
Args.setParam(new String[] {"--output-directory", dbPath(), "-w"}, SHIELD_CONF);
@@ -167,14 +165,21 @@ public static void initZksnarkParams() {
*/
@Before
public void init() {
+ previousAllowShieldedTransaction = chainBaseManager.getDynamicPropertiesStore()
+ .getAllowShieldedTransaction();
if (init) {
return;
}
- consensusService.start();
chainBaseManager.getDynamicPropertiesStore().saveTotalShieldedPoolValue(10_000_000_000L);
init = true;
}
+ @After
+ public void restoreAllowShieldedTransaction() {
+ chainBaseManager.getDynamicPropertiesStore()
+ .saveAllowShieldedTransaction(previousAllowShieldedTransaction);
+ }
+
private static byte[] randomUint256() {
return org.tron.keystore.Wallet.generateRandomBytes(32);
}
@@ -254,9 +259,28 @@ private void updateTotalShieldedPoolValue(long valueBalance) {
@Test
public void testIsMining() {
+ initializeActiveWitnessSchedule();
Assert.assertTrue(wallet.isMining());
}
+ private void initializeActiveWitnessSchedule() {
+ synchronized (ShieldedReceiveTest.class) {
+ if (consensusScheduleInitialized) {
+ return;
+ }
+ boolean started = false;
+ try {
+ consensusService.start();
+ started = true;
+ } finally {
+ if (started) {
+ consensusService.stop();
+ }
+ }
+ consensusScheduleInitialized = true;
+ }
+ }
+
/*
* test of change ShieldedTransactionFee proposal
*/
@@ -2407,144 +2431,134 @@ public void pushSameSkAndScanAndSpend() throws Exception {
assert ecKey != null;
byte[] witnessAddress = ecKey.getAddress();
WitnessCapsule witnessCapsule = new WitnessCapsule(ByteString.copyFrom(witnessAddress));
- // Stop the consensus task before modifying the witness schedule: DposTask uses the same
- // localwitness key and would otherwise race to produce blocks at the same slot,
- // triggering fork resolution and making the test slow.
- consensusService.stop();
- try {
- chainBaseManager.addWitness(ByteString.copyFrom(witnessAddress));
-
- long time = nextScheduledTime(witnessCapsule.getAddress());
- Block block = getSignedBlock(witnessCapsule.getAddress(), time, privateKey);
- dbManager.pushBlock(new BlockCapsule(block));
-
- //create transactions
- chainBaseManager.getDynamicPropertiesStore().saveAllowShieldedTransaction(1);
- chainBaseManager.getDynamicPropertiesStore().saveTotalShieldedPoolValue(1000 * 1000000L);
- ZenTransactionBuilder builder = new ZenTransactionBuilder(wallet);
-
- // generate spend proof
- SpendingKey sk = SpendingKey
- .decode("ff2c06269315333a9207f817d2eca0ac555ca8f90196976324c7756504e7c9ee");
- ExpandedSpendingKey expsk = sk.expandedSpendingKey();
- byte[] senderOvk = expsk.getOvk();
- PaymentAddress address = sk.defaultAddress();
- Note note = new Note(address, 1000 * 1000000L);
- IncrementalMerkleVoucherContainer voucher = createSimpleMerkleVoucherContainer(note.cm());
- byte[] anchor = voucher.root().getContent().toByteArray();
- chainBaseManager.getMerkleContainer()
- .putMerkleTreeIntoStore(anchor, voucher.getVoucherCapsule().getTree());
- builder.addSpend(expsk, note, anchor, voucher);
-
- // generate output proof
- SpendingKey sk2 = SpendingKey.random();
- FullViewingKey fullViewingKey = sk2.fullViewingKey();
- IncomingViewingKey incomingViewingKey = fullViewingKey.inViewingKey();
-
- byte[] memo = org.tron.keystore.Wallet.generateRandomBytes(512);
-
- //send coin to 2 different address generated by same sk
- DiversifierT d1 = DiversifierT.random();
- PaymentAddress paymentAddress1 = incomingViewingKey.address(d1).get();
- builder.addOutput(senderOvk, paymentAddress1,
- (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2, memo);
-
- DiversifierT d2 = DiversifierT.random();
- PaymentAddress paymentAddress2 = incomingViewingKey.address(d2).get();
- builder.addOutput(senderOvk, paymentAddress2,
- (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2, memo);
+ // Initialize the same schedule as DPoS startup without starting its producer thread.
+ // Manual block production below therefore cannot race the background producer.
+ initializeActiveWitnessSchedule();
+ chainBaseManager.addWitness(ByteString.copyFrom(witnessAddress));
- TransactionCapsule transactionCap = builder.build();
+ long time = nextScheduledTime(witnessCapsule.getAddress());
+ Block block = getSignedBlock(witnessCapsule.getAddress(), time, privateKey);
+ dbManager.pushBlock(new BlockCapsule(block));
- byte[] trxId = transactionCap.getTransactionId().getBytes();
- boolean ok = dbManager.pushTransaction(transactionCap);
- Assert.assertTrue(ok);
-
- Thread.sleep(500);
- //package transaction to block
- long expectedBlockNum = chainBaseManager.getDynamicPropertiesStore()
- .getLatestBlockHeaderNumber() + 1;
- block = getSignedBlock(witnessCapsule.getAddress(),
- nextScheduledTime(witnessCapsule.getAddress()), privateKey);
- dbManager.pushBlock(new BlockCapsule(block));
-
- BlockCapsule blockCapsule3 = new BlockCapsule(wallet.getNowBlock());
- Assert.assertEquals("unexpected block number", expectedBlockNum, blockCapsule3.getNum());
-
- block = getSignedBlock(witnessCapsule.getAddress(),
- nextScheduledTime(witnessCapsule.getAddress()), privateKey);
- dbManager.pushBlock(new BlockCapsule(block));
-
- // scan note by ivk
- byte[] receiverIvk = incomingViewingKey.getValue();
- DecryptNotes notes1 = wallet.scanNoteByIvk(0, 100, receiverIvk);
- Assert.assertEquals(2, notes1.getNoteTxsCount());
-
- // scan note by ivk and mark
- DecryptNotesMarked notes3 = wallet.scanAndMarkNoteByIvk(0, 100, receiverIvk,
- fullViewingKey.getAk(), fullViewingKey.getNk());
- Assert.assertEquals(2, notes3.getNoteTxsCount());
-
- // scan note by ovk
- DecryptNotes notes2 = wallet.scanNoteByOvk(0, 100, senderOvk);
- Assert.assertEquals(2, notes2.getNoteTxsCount());
-
- // to spend received note above.
- ZenTransactionBuilder builder2 = new ZenTransactionBuilder(wallet);
-
- //query merkleinfo
- OutputPointInfo.Builder request = OutputPointInfo.newBuilder();
- for (int i = 0; i < notes1.getNoteTxsCount(); i++) {
- OutputPoint.Builder outPointBuild = OutputPoint.newBuilder();
- outPointBuild.setHash(ByteString.copyFrom(trxId));
- outPointBuild.setIndex(i);
- request.addOutPoints(outPointBuild.build());
- }
- request.setBlockNum(1);
- IncrementalMerkleVoucherInfo merkleVoucherInfo = wallet
- .getMerkleTreeVoucherInfo(request.build());
-
- //build spend proof. allow only one note in spend
- ExpandedSpendingKey expsk2 = sk2.expandedSpendingKey();
- for (int i = 0; i < 1; i++) {
- org.tron.api.GrpcAPI.Note grpcNote = notes1.getNoteTxs(i).getNote();
- PaymentAddress paymentAddress = KeyIo.decodePaymentAddress(grpcNote.getPaymentAddress());
- Note note2 = new Note(paymentAddress.getD(),
- paymentAddress.getPkD(),
- grpcNote.getValue(),
- grpcNote.getRcm().toByteArray()
- );
-
- IncrementalMerkleVoucherContainer voucher2 =
- new IncrementalMerkleVoucherContainer(
- new IncrementalMerkleVoucherCapsule(merkleVoucherInfo.getVouchers(i)));
- byte[] anchor2 = voucher2.root().getContent().toByteArray();
- builder2.addSpend(expsk2, note2, anchor2, voucher2);
- }
+ //create transactions
+ chainBaseManager.getDynamicPropertiesStore().saveAllowShieldedTransaction(1);
+ chainBaseManager.getDynamicPropertiesStore().saveTotalShieldedPoolValue(1000 * 1000000L);
+ ZenTransactionBuilder builder = new ZenTransactionBuilder(wallet);
+
+ // generate spend proof
+ SpendingKey sk = SpendingKey
+ .decode("ff2c06269315333a9207f817d2eca0ac555ca8f90196976324c7756504e7c9ee");
+ ExpandedSpendingKey expsk = sk.expandedSpendingKey();
+ byte[] senderOvk = expsk.getOvk();
+ PaymentAddress address = sk.defaultAddress();
+ Note note = new Note(address, 1000 * 1000000L);
+ IncrementalMerkleVoucherContainer voucher = createSimpleMerkleVoucherContainer(note.cm());
+ byte[] anchor = voucher.root().getContent().toByteArray();
+ chainBaseManager.getMerkleContainer()
+ .putMerkleTreeIntoStore(anchor, voucher.getVoucherCapsule().getTree());
+ builder.addSpend(expsk, note, anchor, voucher);
+
+ // generate output proof
+ SpendingKey sk2 = SpendingKey.random();
+ FullViewingKey fullViewingKey = sk2.fullViewingKey();
+ IncomingViewingKey incomingViewingKey = fullViewingKey.inViewingKey();
+
+ byte[] memo = org.tron.keystore.Wallet.generateRandomBytes(512);
+
+ //send coin to 2 different address generated by same sk
+ DiversifierT d1 = DiversifierT.random();
+ PaymentAddress paymentAddress1 = incomingViewingKey.address(d1).get();
+ builder.addOutput(senderOvk, paymentAddress1,
+ (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2, memo);
+
+ DiversifierT d2 = DiversifierT.random();
+ PaymentAddress paymentAddress2 = incomingViewingKey.address(d2).get();
+ builder.addOutput(senderOvk, paymentAddress2,
+ (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2, memo);
- //build output proof
- SpendingKey sk3 = SpendingKey.random();
- FullViewingKey fvk3 = sk3.fullViewingKey();
- IncomingViewingKey ivk3 = fvk3.inViewingKey();
-
- DiversifierT d3 = DiversifierT.random();
- PaymentAddress paymentAddress3 = incomingViewingKey.address(d3).get();
- byte[] memo3 = org.tron.keystore.Wallet.generateRandomBytes(512);
- builder2.addOutput(expsk2.getOvk(), paymentAddress3,
- (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2 - wallet
- .getShieldedTransactionFee(), memo3);
-
- TransactionCapsule transactionCap2 = builder2.build();
- boolean ok2 = dbManager.pushTransaction(transactionCap2);
- Assert.assertTrue(ok2);
- } finally {
- // DposTask.init() does not reset isRunning (it stays false after stop()), so force it back
- // to true via reflection before restarting.
- Field isRunning = DposTask.class.getDeclaredField("isRunning");
- isRunning.setAccessible(true);
- isRunning.set(dposTask, true);
- consensusService.start();
+ TransactionCapsule transactionCap = builder.build();
+
+ byte[] trxId = transactionCap.getTransactionId().getBytes();
+ boolean ok = dbManager.pushTransaction(transactionCap);
+ Assert.assertTrue(ok);
+
+ Thread.sleep(500);
+ //package transaction to block
+ long expectedBlockNum = chainBaseManager.getDynamicPropertiesStore()
+ .getLatestBlockHeaderNumber() + 1;
+ block = getSignedBlock(witnessCapsule.getAddress(),
+ nextScheduledTime(witnessCapsule.getAddress()), privateKey);
+ dbManager.pushBlock(new BlockCapsule(block));
+
+ BlockCapsule blockCapsule3 = new BlockCapsule(wallet.getNowBlock());
+ Assert.assertEquals("unexpected block number", expectedBlockNum, blockCapsule3.getNum());
+
+ block = getSignedBlock(witnessCapsule.getAddress(),
+ nextScheduledTime(witnessCapsule.getAddress()), privateKey);
+ dbManager.pushBlock(new BlockCapsule(block));
+
+ // scan note by ivk
+ byte[] receiverIvk = incomingViewingKey.getValue();
+ DecryptNotes notes1 = wallet.scanNoteByIvk(0, 100, receiverIvk);
+ Assert.assertEquals(2, notes1.getNoteTxsCount());
+
+ // scan note by ivk and mark
+ DecryptNotesMarked notes3 = wallet.scanAndMarkNoteByIvk(0, 100, receiverIvk,
+ fullViewingKey.getAk(), fullViewingKey.getNk());
+ Assert.assertEquals(2, notes3.getNoteTxsCount());
+
+ // scan note by ovk
+ DecryptNotes notes2 = wallet.scanNoteByOvk(0, 100, senderOvk);
+ Assert.assertEquals(2, notes2.getNoteTxsCount());
+
+ // to spend received note above.
+ ZenTransactionBuilder builder2 = new ZenTransactionBuilder(wallet);
+
+ //query merkleinfo
+ OutputPointInfo.Builder request = OutputPointInfo.newBuilder();
+ for (int i = 0; i < notes1.getNoteTxsCount(); i++) {
+ OutputPoint.Builder outPointBuild = OutputPoint.newBuilder();
+ outPointBuild.setHash(ByteString.copyFrom(trxId));
+ outPointBuild.setIndex(i);
+ request.addOutPoints(outPointBuild.build());
+ }
+ request.setBlockNum(1);
+ IncrementalMerkleVoucherInfo merkleVoucherInfo = wallet
+ .getMerkleTreeVoucherInfo(request.build());
+
+ //build spend proof. allow only one note in spend
+ ExpandedSpendingKey expsk2 = sk2.expandedSpendingKey();
+ for (int i = 0; i < 1; i++) {
+ org.tron.api.GrpcAPI.Note grpcNote = notes1.getNoteTxs(i).getNote();
+ PaymentAddress paymentAddress = KeyIo.decodePaymentAddress(grpcNote.getPaymentAddress());
+ Note note2 = new Note(paymentAddress.getD(),
+ paymentAddress.getPkD(),
+ grpcNote.getValue(),
+ grpcNote.getRcm().toByteArray()
+ );
+
+ IncrementalMerkleVoucherContainer voucher2 =
+ new IncrementalMerkleVoucherContainer(
+ new IncrementalMerkleVoucherCapsule(merkleVoucherInfo.getVouchers(i)));
+ byte[] anchor2 = voucher2.root().getContent().toByteArray();
+ builder2.addSpend(expsk2, note2, anchor2, voucher2);
}
+
+ //build output proof
+ SpendingKey sk3 = SpendingKey.random();
+ FullViewingKey fvk3 = sk3.fullViewingKey();
+ IncomingViewingKey ivk3 = fvk3.inViewingKey();
+
+ DiversifierT d3 = DiversifierT.random();
+ PaymentAddress paymentAddress3 = incomingViewingKey.address(d3).get();
+ byte[] memo3 = org.tron.keystore.Wallet.generateRandomBytes(512);
+ builder2.addOutput(expsk2.getOvk(), paymentAddress3,
+ (1000 * 1000000L - wallet.getShieldedTransactionFee()) / 2 - wallet
+ .getShieldedTransactionFee(), memo3);
+
+ TransactionCapsule transactionCap2 = builder2.build();
+ boolean ok2 = dbManager.pushTransaction(transactionCap2);
+ Assert.assertTrue(ok2);
}
// Returns the earliest timestamp at which witnessAddr is the DPoS-scheduled producer,
From 11d555d3d9cda3d01605ef1936b9db1553546f77 Mon Sep 17 00:00:00 2001
From: halibobo1205 <82020050+halibobo1205@users.noreply.github.com>
Date: Thu, 24 Sep 2026 16:48:25 +0800
Subject: [PATCH 07/25] refactor(vm): remove unreachable trace compression path
(#6997)
---
.../org/tron/core/actuator/VMActuator.java | 4 --
.../main/java/org/tron/core/vm/VMUtils.java | 45 -------------------
.../org/tron/core/vm/config/VMConfig.java | 6 ---
3 files changed, 55 deletions(-)
diff --git a/actuator/src/main/java/org/tron/core/actuator/VMActuator.java b/actuator/src/main/java/org/tron/core/actuator/VMActuator.java
index e0a721db28d..9a2cc8231da 100644
--- a/actuator/src/main/java/org/tron/core/actuator/VMActuator.java
+++ b/actuator/src/main/java/org/tron/core/actuator/VMActuator.java
@@ -309,10 +309,6 @@ public void execute(Object object) throws ContractExeException {
.error(result.getException())
.toString();
- if (VMConfig.vmTraceCompressed()) {
- traceContent = VMUtils.zipAndEncode(traceContent);
- }
-
String txHash = Hex.toHexString(rootInternalTx.getHash());
VMUtils.saveProgramTraceFile(txHash, traceContent);
}
diff --git a/actuator/src/main/java/org/tron/core/vm/VMUtils.java b/actuator/src/main/java/org/tron/core/vm/VMUtils.java
index 2f469e0579a..cbd6f61bf14 100644
--- a/actuator/src/main/java/org/tron/core/vm/VMUtils.java
+++ b/actuator/src/main/java/org/tron/core/vm/VMUtils.java
@@ -1,20 +1,14 @@
package org.tron.core.vm;
import static java.lang.String.format;
-import static org.apache.commons.codec.binary.Base64.encodeBase64String;
import static org.tron.common.math.Maths.addExact;
-import java.io.ByteArrayInputStream;
-import java.io.ByteArrayOutputStream;
import java.io.Closeable;
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
-import java.io.InputStream;
import java.io.OutputStream;
import java.util.Arrays;
-import java.util.zip.Deflater;
-import java.util.zip.DeflaterOutputStream;
import lombok.extern.slf4j.Slf4j;
import org.tron.common.utils.ByteArray;
import org.tron.common.utils.ByteUtil;
@@ -29,8 +23,6 @@
@Slf4j(topic = "VM")
public final class VMUtils {
- private static final int BUF_SIZE = 4096;
-
private VMUtils() {
}
@@ -97,43 +89,6 @@ public static void saveProgramTraceFile(String txHash, String content) {
}
}
- private static void write(InputStream in, OutputStream out, int bufSize) throws IOException {
- try {
- byte[] buf = new byte[bufSize];
- for (int count = in.read(buf); count != -1; count = in.read(buf)) {
- out.write(buf, 0, count);
- }
- } finally {
- closeQuietly(in);
- closeQuietly(out);
- }
- }
-
- public static byte[] compress(byte[] bytes) throws IOException {
- ByteArrayOutputStream baos = new ByteArrayOutputStream();
-
- ByteArrayInputStream in = new ByteArrayInputStream(bytes);
- DeflaterOutputStream out = new DeflaterOutputStream(baos, new Deflater(), BUF_SIZE);
-
- write(in, out, BUF_SIZE);
-
- return baos.toByteArray();
- }
-
- public static byte[] compress(String content) throws IOException {
- return compress(content.getBytes("UTF-8"));
- }
-
- public static String zipAndEncode(String content) {
- try {
- return encodeBase64String(compress(content));
- } catch (Exception e) {
- logger.error("Cannot zip or encode: ", e);
- return content;
- }
- }
-
-
public static boolean validateForSmartContract(Repository deposit, byte[] ownerAddress,
byte[] toAddress, long amount) throws ContractValidateException {
if (!DecodeUtil.addressValid(ownerAddress)) {
diff --git a/common/src/main/java/org/tron/core/vm/config/VMConfig.java b/common/src/main/java/org/tron/core/vm/config/VMConfig.java
index 304ced33698..9762ca31a6c 100644
--- a/common/src/main/java/org/tron/core/vm/config/VMConfig.java
+++ b/common/src/main/java/org/tron/core/vm/config/VMConfig.java
@@ -8,8 +8,6 @@
*/
public class VMConfig {
- private static boolean vmTraceCompressed = false;
-
@Setter
private static boolean vmTrace = false;
@@ -89,10 +87,6 @@ public static boolean vmTrace() {
return vmTrace;
}
- public static boolean vmTraceCompressed() {
- return vmTraceCompressed;
- }
-
public static void initVmHardFork(boolean pass) {
CommonParameter.ENERGY_LIMIT_HARD_FORK = pass;
}
From 9f3c9ae5b9731fd7a8d7a059fff56c7261c2dc39 Mon Sep 17 00:00:00 2001
From: barbatos2011 <162298485+barbatos2011@users.noreply.github.com>
Date: Tue, 29 Sep 2026 14:53:21 +0800
Subject: [PATCH 08/25] chore(p2p): internalize libp2p v2.2.9 as a local `p2p`
module (#6992)
---
build.gradle | 11 +
common/build.gradle | 18 +-
framework/build.gradle | 27 +-
gradle/verification-metadata.xml | 21 +
p2p/.gitignore | 2 +
p2p/README.md | 453 ++++++++++++++++
p2p/build.gradle | 196 +++++++
p2p/src/main/java/org/tron/p2p/P2pConfig.java | 37 ++
.../java/org/tron/p2p/P2pEventHandler.java | 20 +
.../main/java/org/tron/p2p/P2pService.java | 90 +++
.../main/java/org/tron/p2p/base/Constant.java | 16 +
.../java/org/tron/p2p/base/Parameter.java | 74 +++
.../java/org/tron/p2p/connection/Channel.java | 204 +++++++
.../tron/p2p/connection/ChannelManager.java | 307 +++++++++++
.../connection/business/MessageProcess.java | 8 +
.../business/detect/NodeDetectService.java | 229 ++++++++
.../connection/business/detect/NodeStat.java | 25 +
.../business/handshake/DisconnectCode.java | 30 +
.../business/handshake/HandshakeService.java | 90 +++
.../business/keepalive/KeepAliveService.java | 70 +++
.../business/pool/ConnPoolService.java | 346 ++++++++++++
.../business/upgrade/UpgradeController.java | 38 ++
.../tron/p2p/connection/message/Message.java | 78 +++
.../p2p/connection/message/MessageType.java | 42 ++
.../message/base/P2pDisconnectMessage.java | 39 ++
.../message/detect/StatusMessage.java | 61 +++
.../message/handshake/HelloMessage.java | 64 +++
.../message/keepalive/PingMessage.java | 33 ++
.../message/keepalive/PongMessage.java | 33 ++
.../p2p/connection/socket/MessageHandler.java | 90 +++
.../socket/P2pChannelInitializer.java | 62 +++
.../P2pProtobufVarint32FrameDecoder.java | 98 ++++
.../p2p/connection/socket/PeerClient.java | 103 ++++
.../p2p/connection/socket/PeerServer.java | 79 +++
.../tron/p2p/discover/DiscoverService.java | 25 +
.../main/java/org/tron/p2p/discover/Node.java | 197 +++++++
.../org/tron/p2p/discover/NodeManager.java | 47 ++
.../tron/p2p/discover/message/Message.java | 69 +++
.../p2p/discover/message/MessageType.java | 41 ++
.../discover/message/kad/FindNodeMessage.java | 55 ++
.../p2p/discover/message/kad/KadMessage.java | 35 ++
.../message/kad/NeighborsMessage.java | 80 +++
.../p2p/discover/message/kad/PingMessage.java | 59 ++
.../p2p/discover/message/kad/PongMessage.java | 53 ++
.../discover/protocol/kad/DiscoverTask.java | 87 +++
.../p2p/discover/protocol/kad/KadService.java | 219 ++++++++
.../discover/protocol/kad/NodeHandler.java | 250 +++++++++
.../kad/table/DistanceComparator.java | 26 +
.../protocol/kad/table/KademliaOptions.java | 12 +
.../protocol/kad/table/NodeBucket.java | 53 ++
.../protocol/kad/table/NodeEntry.java | 88 +++
.../protocol/kad/table/NodeTable.java | 129 +++++
.../protocol/kad/table/TimeComparator.java | 19 +
.../p2p/discover/socket/DiscoverServer.java | 93 ++++
.../p2p/discover/socket/EventHandler.java | 12 +
.../p2p/discover/socket/MessageHandler.java | 67 +++
.../p2p/discover/socket/P2pPacketDecoder.java | 51 ++
.../tron/p2p/discover/socket/UdpEvent.java | 32 ++
.../java/org/tron/p2p/dns/DnsManager.java | 79 +++
.../main/java/org/tron/p2p/dns/DnsNode.java | 95 ++++
.../org/tron/p2p/dns/lookup/LookUpTxt.java | 204 +++++++
.../java/org/tron/p2p/dns/sync/Client.java | 188 +++++++
.../org/tron/p2p/dns/sync/ClientTree.java | 195 +++++++
.../java/org/tron/p2p/dns/sync/LinkCache.java | 82 +++
.../org/tron/p2p/dns/sync/RandomIterator.java | 126 +++++
.../org/tron/p2p/dns/sync/SubtreeSync.java | 74 +++
.../java/org/tron/p2p/dns/tree/Algorithm.java | 150 +++++
.../org/tron/p2p/dns/tree/BranchEntry.java | 32 ++
.../java/org/tron/p2p/dns/tree/Entry.java | 10 +
.../java/org/tron/p2p/dns/tree/LinkEntry.java | 53 ++
.../org/tron/p2p/dns/tree/NodesEntry.java | 39 ++
.../java/org/tron/p2p/dns/tree/RootEntry.java | 113 ++++
.../main/java/org/tron/p2p/dns/tree/Tree.java | 247 +++++++++
.../org/tron/p2p/dns/update/AliClient.java | 341 ++++++++++++
.../org/tron/p2p/dns/update/AwsClient.java | 511 ++++++++++++++++++
.../java/org/tron/p2p/dns/update/DnsType.java | 23 +
.../java/org/tron/p2p/dns/update/Publish.java | 18 +
.../tron/p2p/dns/update/PublishConfig.java | 24 +
.../tron/p2p/dns/update/PublishService.java | 146 +++++
.../java/org/tron/p2p/example/StartApp.java | 406 ++++++++++++++
.../org/tron/p2p/exception/DnsException.java | 73 +++
.../org/tron/p2p/exception/P2pException.java | 59 ++
.../java/org/tron/p2p/stats/P2pStats.java | 15 +
.../java/org/tron/p2p/stats/StatsManager.java | 17 +
.../java/org/tron/p2p/stats/TrafficStats.java | 50 ++
.../java/org/tron/p2p/utils/ByteArray.java | 193 +++++++
.../org/tron/p2p/utils/CollectionUtils.java | 22 +
.../main/java/org/tron/p2p/utils/NetUtil.java | 294 ++++++++++
.../java/org/tron/p2p/utils/ProtoUtil.java | 48 ++
.../java/org/web3j/crypto/ECDSASignature.java | 61 +++
.../main/java/org/web3j/crypto/ECKeyPair.java | 112 ++++
p2p/src/main/java/org/web3j/crypto/Hash.java | 139 +++++
p2p/src/main/java/org/web3j/crypto/Sign.java | 356 ++++++++++++
.../exceptions/MessageDecodingException.java | 26 +
.../exceptions/MessageEncodingException.java | 26 +
.../main/java/org/web3j/utils/Assertions.java | 31 ++
.../main/java/org/web3j/utils/Numeric.java | 252 +++++++++
.../main/java/org/web3j/utils/Strings.java | 61 +++
p2p/src/main/proto/Connect.proto | 60 ++
p2p/src/main/proto/Discover.proto | 50 ++
.../java/org/tron/p2p/P2pServiceTest.java | 118 ++++
.../tron/p2p/connection/ChannelCoreTest.java | 173 ++++++
.../ChannelManagerAdmissionTest.java | 174 ++++++
.../p2p/connection/ChannelManagerTest.java | 185 +++++++
.../tron/p2p/connection/ChannelValueTest.java | 69 +++
.../p2p/connection/ConnPoolServiceTest.java | 186 +++++++
.../DisconnectReasonMappingTest.java | 51 ++
.../org/tron/p2p/connection/MessageTest.java | 87 +++
.../org/tron/p2p/connection/SocketTest.java | 82 +++
.../detect/NodeDetectServiceTest.java | 129 +++++
.../handshake/HandshakeServiceTest.java | 226 ++++++++
.../keepalive/KeepAliveServiceTest.java | 108 ++++
.../business/pool/ConnPoolLifecycleTest.java | 133 +++++
.../upgrade/UpgradeControllerTest.java | 77 +++
.../base/P2pDisconnectMessageTest.java | 31 ++
.../message/detect/StatusMessageTest.java | 70 +++
.../message/handshake/HelloMessageTest.java | 33 ++
.../connection/socket/MessageHandlerTest.java | 139 +++++
.../P2pProtobufVarint32FrameDecoderTest.java | 197 +++++++
.../tron/p2p/discover/NodeManagerTest.java | 25 +
.../java/org/tron/p2p/discover/NodeTest.java | 96 ++++
.../discover/message/DiscoverMessageTest.java | 128 +++++
.../discover/message/kad/KadMessagesTest.java | 158 ++++++
.../discover/protocol/kad/KadServiceTest.java | 53 ++
.../protocol/kad/NodeHandlerTest.java | 114 ++++
.../protocol/kad/table/NodeEntryTest.java | 69 +++
.../protocol/kad/table/NodeTableTest.java | 225 ++++++++
.../kad/table/TimeComparatorTest.java | 22 +
.../discover/socket/P2pPacketDecoderTest.java | 121 +++++
.../java/org/tron/p2p/dns/AlgorithmTest.java | 110 ++++
.../java/org/tron/p2p/dns/AwsRoute53Test.java | 169 ++++++
.../java/org/tron/p2p/dns/DnsManagerTest.java | 143 +++++
.../java/org/tron/p2p/dns/DnsNodeTest.java | 52 ++
.../java/org/tron/p2p/dns/LinkCacheTest.java | 35 ++
.../java/org/tron/p2p/dns/RandomTest.java | 36 ++
.../test/java/org/tron/p2p/dns/SyncTest.java | 33 ++
.../test/java/org/tron/p2p/dns/TreeTest.java | 248 +++++++++
.../tron/p2p/dns/lookup/LookUpTxtTest.java | 88 +++
.../tron/p2p/dns/tree/TreeSignAndTxtTest.java | 146 +++++
.../p2p/dns/update/AliClientDeployTest.java | 192 +++++++
.../tron/p2p/dns/update/AliClientTest.java | 344 ++++++++++++
.../p2p/dns/update/AwsClientBatchTest.java | 178 ++++++
.../p2p/dns/update/AwsClientChangeTest.java | 218 ++++++++
.../p2p/dns/update/AwsClientRecordsTest.java | 232 ++++++++
.../p2p/dns/update/PublishServiceTest.java | 208 +++++++
.../tron/p2p/example/ExampleUsageTest.java | 224 ++++++++
.../tron/p2p/example/StartAppArgsTest.java | 71 +++
.../tron/p2p/exception/DnsExceptionTest.java | 40 ++
.../tron/p2p/exception/P2pExceptionTest.java | 36 ++
.../org/tron/p2p/stats/StatsManagerTest.java | 32 ++
.../org/tron/p2p/stats/TrafficStatsTest.java | 75 +++
.../org/tron/p2p/utils/ByteArrayTest.java | 167 ++++++
.../tron/p2p/utils/NetUtilAddressTest.java | 103 ++++
.../java/org/tron/p2p/utils/NetUtilTest.java | 279 ++++++++++
.../org/tron/p2p/utils/ProtoUtilTest.java | 30 +
.../java/org/tron/p2p/utils/TestPort.java | 49 ++
.../org/web3j/crypto/ECDSASignatureTest.java | 38 ++
.../java/org/web3j/crypto/ECKeyPairTest.java | 53 ++
.../test/java/org/web3j/crypto/HashTest.java | 95 ++++
.../test/java/org/web3j/crypto/SignTest.java | 117 ++++
.../exceptions/MessageExceptionsTest.java | 29 +
.../java/org/web3j/utils/AssertionsTest.java | 22 +
.../java/org/web3j/utils/NumericTest.java | 163 ++++++
.../java/org/web3j/utils/StringsTest.java | 50 ++
plugins/build.gradle | 11 +-
protocol/build.gradle | 7 +-
settings.gradle | 1 +
167 files changed, 17365 insertions(+), 32 deletions(-)
create mode 100644 p2p/.gitignore
create mode 100644 p2p/README.md
create mode 100644 p2p/build.gradle
create mode 100644 p2p/src/main/java/org/tron/p2p/P2pConfig.java
create mode 100644 p2p/src/main/java/org/tron/p2p/P2pEventHandler.java
create mode 100644 p2p/src/main/java/org/tron/p2p/P2pService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/base/Constant.java
create mode 100644 p2p/src/main/java/org/tron/p2p/base/Parameter.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/Channel.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/ChannelManager.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/MessageProcess.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/detect/NodeDetectService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/detect/NodeStat.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/handshake/DisconnectCode.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/handshake/HandshakeService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/keepalive/KeepAliveService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/pool/ConnPoolService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/business/upgrade/UpgradeController.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/Message.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/MessageType.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/base/P2pDisconnectMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/detect/StatusMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/handshake/HelloMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/keepalive/PingMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/message/keepalive/PongMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/socket/MessageHandler.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/socket/P2pChannelInitializer.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/socket/P2pProtobufVarint32FrameDecoder.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/socket/PeerClient.java
create mode 100644 p2p/src/main/java/org/tron/p2p/connection/socket/PeerServer.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/DiscoverService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/Node.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/NodeManager.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/Message.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/MessageType.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/kad/FindNodeMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/kad/KadMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/kad/NeighborsMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/kad/PingMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/message/kad/PongMessage.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/DiscoverTask.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/KadService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/NodeHandler.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/DistanceComparator.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/KademliaOptions.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/NodeBucket.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/NodeEntry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/NodeTable.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/protocol/kad/table/TimeComparator.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/socket/DiscoverServer.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/socket/EventHandler.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/socket/MessageHandler.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/socket/P2pPacketDecoder.java
create mode 100644 p2p/src/main/java/org/tron/p2p/discover/socket/UdpEvent.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/DnsManager.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/DnsNode.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/lookup/LookUpTxt.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/sync/Client.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/sync/ClientTree.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/sync/LinkCache.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/sync/RandomIterator.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/sync/SubtreeSync.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/Algorithm.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/BranchEntry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/Entry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/LinkEntry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/NodesEntry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/RootEntry.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/tree/Tree.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/AliClient.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/AwsClient.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/DnsType.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/Publish.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/PublishConfig.java
create mode 100644 p2p/src/main/java/org/tron/p2p/dns/update/PublishService.java
create mode 100644 p2p/src/main/java/org/tron/p2p/example/StartApp.java
create mode 100644 p2p/src/main/java/org/tron/p2p/exception/DnsException.java
create mode 100644 p2p/src/main/java/org/tron/p2p/exception/P2pException.java
create mode 100644 p2p/src/main/java/org/tron/p2p/stats/P2pStats.java
create mode 100644 p2p/src/main/java/org/tron/p2p/stats/StatsManager.java
create mode 100644 p2p/src/main/java/org/tron/p2p/stats/TrafficStats.java
create mode 100644 p2p/src/main/java/org/tron/p2p/utils/ByteArray.java
create mode 100644 p2p/src/main/java/org/tron/p2p/utils/CollectionUtils.java
create mode 100644 p2p/src/main/java/org/tron/p2p/utils/NetUtil.java
create mode 100644 p2p/src/main/java/org/tron/p2p/utils/ProtoUtil.java
create mode 100644 p2p/src/main/java/org/web3j/crypto/ECDSASignature.java
create mode 100644 p2p/src/main/java/org/web3j/crypto/ECKeyPair.java
create mode 100644 p2p/src/main/java/org/web3j/crypto/Hash.java
create mode 100644 p2p/src/main/java/org/web3j/crypto/Sign.java
create mode 100644 p2p/src/main/java/org/web3j/exceptions/MessageDecodingException.java
create mode 100644 p2p/src/main/java/org/web3j/exceptions/MessageEncodingException.java
create mode 100644 p2p/src/main/java/org/web3j/utils/Assertions.java
create mode 100644 p2p/src/main/java/org/web3j/utils/Numeric.java
create mode 100644 p2p/src/main/java/org/web3j/utils/Strings.java
create mode 100644 p2p/src/main/proto/Connect.proto
create mode 100644 p2p/src/main/proto/Discover.proto
create mode 100644 p2p/src/test/java/org/tron/p2p/P2pServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/ChannelCoreTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/ChannelManagerAdmissionTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/ChannelManagerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/ChannelValueTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/ConnPoolServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/DisconnectReasonMappingTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/MessageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/SocketTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/business/detect/NodeDetectServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/business/handshake/HandshakeServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/business/keepalive/KeepAliveServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/business/pool/ConnPoolLifecycleTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/business/upgrade/UpgradeControllerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/message/base/P2pDisconnectMessageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/message/detect/StatusMessageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/message/handshake/HelloMessageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/socket/MessageHandlerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/connection/socket/P2pProtobufVarint32FrameDecoderTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/NodeManagerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/NodeTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/message/DiscoverMessageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/message/kad/KadMessagesTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/protocol/kad/KadServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/protocol/kad/NodeHandlerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/protocol/kad/table/NodeEntryTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/protocol/kad/table/NodeTableTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/protocol/kad/table/TimeComparatorTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/discover/socket/P2pPacketDecoderTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/AlgorithmTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/AwsRoute53Test.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/DnsManagerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/DnsNodeTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/LinkCacheTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/RandomTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/SyncTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/TreeTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/lookup/LookUpTxtTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/tree/TreeSignAndTxtTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/AliClientDeployTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/AliClientTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/AwsClientBatchTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/AwsClientChangeTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/AwsClientRecordsTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/dns/update/PublishServiceTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/example/ExampleUsageTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/example/StartAppArgsTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/exception/DnsExceptionTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/exception/P2pExceptionTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/stats/StatsManagerTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/stats/TrafficStatsTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/utils/ByteArrayTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/utils/NetUtilAddressTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/utils/NetUtilTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/utils/ProtoUtilTest.java
create mode 100644 p2p/src/test/java/org/tron/p2p/utils/TestPort.java
create mode 100644 p2p/src/test/java/org/web3j/crypto/ECDSASignatureTest.java
create mode 100644 p2p/src/test/java/org/web3j/crypto/ECKeyPairTest.java
create mode 100644 p2p/src/test/java/org/web3j/crypto/HashTest.java
create mode 100644 p2p/src/test/java/org/web3j/crypto/SignTest.java
create mode 100644 p2p/src/test/java/org/web3j/exceptions/MessageExceptionsTest.java
create mode 100644 p2p/src/test/java/org/web3j/utils/AssertionsTest.java
create mode 100644 p2p/src/test/java/org/web3j/utils/NumericTest.java
create mode 100644 p2p/src/test/java/org/web3j/utils/StringsTest.java
diff --git a/build.gradle b/build.gradle
index 04dee79fbae..7569c3bba32 100644
--- a/build.gradle
+++ b/build.gradle
@@ -7,6 +7,17 @@ plugins {
ext {
grpcVersion = "1.83.1"
+ // Netty 4.2 split io.netty.handler.codec.protobuf out of netty-codec into its
+ // own artifact. Both :framework and :p2p put the varint32 framing codecs on
+ // their channel pipelines and so must declare it explicitly. Netty itself
+ // arrives transitively through grpc-netty, so this version has to move with
+ // grpcVersion above — keeping it here makes that coupling visible instead of
+ // leaving two literals to drift apart.
+ nettyVersion = "4.2.15.Final"
+ // Shared by :protocol and :p2p, which both generate from .proto files.
+ protobufVersion = "3.25.8"
+ // Shared by :framework, :plugins and :p2p.
+ checkstyleVersion = "8.7"
}
allprojects {
diff --git a/common/build.gradle b/common/build.gradle
index 4b36d067b70..7d8922aab5c 100644
--- a/common/build.gradle
+++ b/common/build.gradle
@@ -23,23 +23,7 @@ dependencies {
api 'org.aspectj:aspectjrt:1.9.8'
api 'org.aspectj:aspectjweaver:1.9.8'
api 'org.aspectj:aspectjtools:1.9.8'
- api group: 'io.github.tronprotocol', name: 'libp2p', version: '2.2.9',{
- exclude group: 'io.grpc', module: 'grpc-context'
- exclude group: 'io.grpc', module: 'grpc-core'
- exclude group: 'io.grpc', module: 'grpc-netty'
- exclude group: 'com.google.protobuf', module: 'protobuf-java'
- exclude group: 'com.google.protobuf', module: 'protobuf-java-util'
- // https://github.com/dom4j/dom4j/pull/116
- // https://github.com/gradle/gradle/issues/13656
- // https://github.com/dom4j/dom4j/issues/99
- exclude group: 'jaxen', module: 'jaxen'
- exclude group: 'javax.xml.stream', module: 'stax-api'
- exclude group: 'net.java.dev.msv', module: 'xsdlib'
- exclude group: 'pull-parser', module: 'pull-parser'
- exclude group: 'xpp3', module: 'xpp3'
- exclude group: 'org.bouncycastle', module: 'bcprov-jdk18on'
- exclude group: 'org.bouncycastle', module: 'bcutil-jdk18on'
- }
+ api project(":p2p")
api project(":protocol")
api project(":platform")
}
diff --git a/framework/build.gradle b/framework/build.gradle
index 8255fc30d18..8edd6714f95 100644
--- a/framework/build.gradle
+++ b/framework/build.gradle
@@ -11,9 +11,7 @@ apply plugin: 'checkstyle'
mainClassName = 'org.tron.program.FullNode'
-def versions = [
- checkstyle: '8.7',
-]
+
@@ -40,7 +38,7 @@ dependencies {
// end local libraries
implementation group: 'com.beust', name: 'jcommander', version: '1.78'
implementation group: 'io.dropwizard.metrics', name: 'metrics-core', version: '3.1.2'
- implementation('io.netty:netty-codec-protobuf:4.2.15.Final') {
+ implementation("io.netty:netty-codec-protobuf:${rootProject.nettyVersion}") {
exclude group: 'com.google.protobuf'
exclude group: 'com.google.protobuf.nano'
}
@@ -61,17 +59,28 @@ dependencies {
testImplementation group: 'org.springframework', name: 'spring-test', version: "${springVersion}"
testImplementation group: 'javax.portlet', name: 'portlet-api', version: '3.0.1'
+
implementation group: 'org.zeromq', name: 'jeromq', version: '0.5.3'
api project(":chainbase")
api project(":protocol")
api project(":actuator")
api project(":consensus")
+ // org.tron.p2p is used directly in 17 files under src/main/java (org.tron.core.net
+ // and org.tron.core.config.args). It currently arrives only transitively, three
+ // hops away, because :common exposes it via `api project(":p2p")` -- and :common
+ // has to, since CommonParameter publishes P2pConfig/PublishConfig in its own API.
+ // Declare the direct use here as well, so framework keeps compiling if that
+ // transitive chain is ever narrowed. api, not implementation: framework does
+ // re-export p2p types -- P2pEventHandlerImpl extends org.tron.p2p.P2pEventHandler,
+ // HelloMessage.getFrom() returns org.tron.p2p.discover.Node, PeerManager takes
+ // org.tron.p2p.connection.Channel, Args.loadDnsPublishConfig returns PublishConfig.
+ api project(":p2p")
}
check.dependsOn 'lint'
checkstyle {
- toolVersion = "${versions.checkstyle}"
+ toolVersion = "${rootProject.checkstyleVersion}"
configFile = file("config/checkstyle/checkStyleAll.xml")
maxWarnings = 0
}
@@ -187,8 +196,14 @@ def binaryRelease(taskName, jarName, mainClass) {
}
// explicit_dependency
+ // :p2p is included because :common now exposes it via `api project(":p2p")`,
+ // so p2p-1.0.0.jar is on runtimeClasspath and gets zipped into the fat jar.
+ // Without it Gradle reports an implicit_dependency and disables execution
+ // optimizations, and a parallel build could assemble FullNode.jar before
+ // :p2p:jar has been written.
dependsOn (project(':actuator').jar, project(':consensus').jar, project(':chainbase').jar,
- project(':crypto').jar, project(':common').jar, project(':protocol').jar, project(':platform').jar)
+ project(':crypto').jar, project(':common').jar, project(':protocol').jar,
+ project(':platform').jar, project(':p2p').jar)
from {
configurations.runtimeClasspath.collect {
diff --git a/gradle/verification-metadata.xml b/gradle/verification-metadata.xml
index 2e30496116f..a73c4715fe1 100644
--- a/gradle/verification-metadata.xml
+++ b/gradle/verification-metadata.xml
@@ -448,6 +448,14 @@