From fe8d2a8602cd85568f9bb0eed22b1922ed740dd5 Mon Sep 17 00:00:00 2001 From: 0xbigapple Date: Tue, 18 Aug 2026 14:31:27 +0800 Subject: [PATCH 1/3] refactor(framework): decouple Manager from TronJsonRpcImpl --- .../common/application/ApplicationImpl.java | 12 ++ .../logsfilter/queue/FilterCapsuleQueue.java | 38 ++++++ .../main/java/org/tron/core/db/Manager.java | 47 +------ .../services/jsonrpc/TronJsonRpcImpl.java | 51 +++++++- .../tron/common/runtime/vm/Create2Test.java | 3 +- .../java/org/tron/core/db/ManagerTest.java | 13 +- .../core/jsonrpc/ConcurrentHashMapTest.java | 2 +- .../jsonrpc/FilterPipelineDeliveryTest.java | 100 +++++++++++++++ .../jsonrpc/FilterPipelineShutdownTest.java | 120 ++++++++++++++++++ .../core/jsonrpc/HandleLogsFilterTest.java | 2 +- .../JsonRpcCallAndEstimateGasTest.java | 9 +- .../tron/core/jsonrpc/JsonrpcServiceTest.java | 3 +- .../tron/core/jsonrpc/WalletCursorTest.java | 15 +-- 13 files changed, 336 insertions(+), 79 deletions(-) create mode 100644 framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java create mode 100644 framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineDeliveryTest.java create mode 100644 framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineShutdownTest.java diff --git a/framework/src/main/java/org/tron/common/application/ApplicationImpl.java b/framework/src/main/java/org/tron/common/application/ApplicationImpl.java index bab95d299ab..f1bba10937a 100644 --- a/framework/src/main/java/org/tron/common/application/ApplicationImpl.java +++ b/framework/src/main/java/org/tron/common/application/ApplicationImpl.java @@ -1,5 +1,6 @@ package org.tron.common.application; +import java.io.IOException; import java.util.concurrent.CountDownLatch; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -10,6 +11,7 @@ import org.tron.core.db.Manager; import org.tron.core.net.TronNetService; import org.tron.core.services.event.EventService; +import org.tron.core.services.jsonrpc.TronJsonRpcImpl; import org.tron.program.SolidityNode; @Slf4j(topic = "app") @@ -37,6 +39,9 @@ public class ApplicationImpl implements Application { @Autowired(required = false) private SolidityNode solidityNode; + @Autowired + private TronJsonRpcImpl tronJsonRpc; + private final CountDownLatch shutdown = new CountDownLatch(1); /** @@ -62,6 +67,13 @@ public void shutdown() { if (solidityNode != null) { solidityNode.close(); } + // producers are stopped; stop the json-rpc filter consumer before the DB closes + // (idempotent — Spring bean destruction may call close() again) + try { + tronJsonRpc.close(); + } catch (IOException e) { + logger.warn("Closing TronJsonRpcImpl failed.", e); + } dbManager.close(); shutdown.countDown(); } diff --git a/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java b/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java new file mode 100644 index 00000000000..7b5ba0e550d --- /dev/null +++ b/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java @@ -0,0 +1,38 @@ +package org.tron.common.logsfilter.queue; + +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.stream.Stream; +import org.springframework.stereotype.Component; +import org.tron.common.logsfilter.capsule.FilterTriggerCapsule; + +/** + * Queue between the block-processing producer (Manager) and the json-rpc filter + * consumer (TronJsonRpcImpl), so that neither side references the other. + */ +@Component +public class FilterCapsuleQueue { + + private final BlockingQueue queue = new LinkedBlockingQueue<>(); + + public boolean offer(FilterTriggerCapsule capsule) { + return queue.offer(capsule); + } + + public FilterTriggerCapsule poll(long timeout, TimeUnit unit) throws InterruptedException { + return queue.poll(timeout, unit); + } + + public int size() { + return queue.size(); + } + + public void clear() { + queue.clear(); + } + + public Stream stream() { + return queue.stream(); + } +} diff --git a/framework/src/main/java/org/tron/core/db/Manager.java b/framework/src/main/java/org/tron/core/db/Manager.java index 9d7a7c979b9..162325573dd 100644 --- a/framework/src/main/java/org/tron/core/db/Manager.java +++ b/framework/src/main/java/org/tron/core/db/Manager.java @@ -48,7 +48,6 @@ import org.apache.commons.collections4.CollectionUtils; import org.bouncycastle.util.encoders.Hex; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; import org.tron.api.GrpcAPI; import org.tron.api.GrpcAPI.TransactionInfoList; @@ -62,11 +61,11 @@ import org.tron.common.logsfilter.capsule.BlockFilterCapsule; import org.tron.common.logsfilter.capsule.BlockLogTriggerCapsule; import org.tron.common.logsfilter.capsule.ContractTriggerCapsule; -import org.tron.common.logsfilter.capsule.FilterTriggerCapsule; import org.tron.common.logsfilter.capsule.LogsFilterCapsule; import org.tron.common.logsfilter.capsule.SolidityTriggerCapsule; import org.tron.common.logsfilter.capsule.TransactionLogTriggerCapsule; import org.tron.common.logsfilter.capsule.TriggerCapsule; +import org.tron.common.logsfilter.queue.FilterCapsuleQueue; import org.tron.common.logsfilter.trigger.ContractEventTrigger; import org.tron.common.logsfilter.trigger.ContractLogTrigger; import org.tron.common.logsfilter.trigger.ContractTrigger; @@ -143,7 +142,6 @@ import org.tron.core.service.MortgageService; import org.tron.core.service.RewardViCalService; import org.tron.core.services.event.exception.EventException; -import org.tron.core.services.jsonrpc.TronJsonRpcImpl; import org.tron.core.store.AccountAssetStore; import org.tron.core.store.AccountIdIndexStore; import org.tron.core.store.AccountIndexStore; @@ -253,8 +251,8 @@ public class Manager { @Getter private BlockingQueue triggerCapsuleQueue; // log filter - private boolean isRunFilterProcessThread = true; - private BlockingQueue filterCapsuleQueue; + @Autowired + private FilterCapsuleQueue filterCapsuleQueue; @Getter private volatile long latestSolidityNumShutDown; @@ -273,16 +271,10 @@ public class Manager { private static final String rePushEsName = "repush"; private ExecutorService triggerEs; private static final String triggerEsName = "event-trigger"; - private ExecutorService filterEs; - private static final String filterEsName = "filter"; @Autowired private RewardViCalService rewardViCalService; - @Lazy - @Autowired - private TronJsonRpcImpl tronJsonRpcImpl; - /** * Cycle thread to rePush Transactions */ @@ -334,26 +326,6 @@ public class Manager { } }; - private Runnable filterProcessLoop = - () -> { - while (isRunFilterProcessThread) { - try { - FilterTriggerCapsule filterCapsule = filterCapsuleQueue.poll(1, TimeUnit.SECONDS); - if (filterCapsule instanceof LogsFilterCapsule) { - tronJsonRpcImpl.handleLogsFilter((LogsFilterCapsule) filterCapsule); - } else if (filterCapsule instanceof BlockFilterCapsule) { - tronJsonRpcImpl.handleBLockFilter((BlockFilterCapsule) filterCapsule); - } - } catch (InterruptedException e) { - logger.error("FilterProcessLoop get InterruptedException, error is {}.", - e.getMessage()); - Thread.currentThread().interrupt(); - } catch (Throwable throwable) { - logger.error("Unknown throwable happened in filterProcessLoop. ", throwable); - } - } - }; - private Comparator downComparator = (Comparator) (o1, o2) -> Long .compare(o2.getOrder(), o1.getOrder()); @@ -476,11 +448,6 @@ public void stopRePushTriggerThread() { ExecutorServiceManager.shutdownAndAwaitTermination(triggerEs, triggerEsName); } - public void stopFilterProcessThread() { - isRunFilterProcessThread = false; - ExecutorServiceManager.shutdownAndAwaitTermination(filterEs, filterEsName); - } - public void stopValidateSignThread() { ExecutorServiceManager.shutdownAndAwaitTermination(validateSignService, "validate-sign"); } @@ -510,7 +477,6 @@ public void init() { this.rePushTransactions = new LinkedBlockingQueue<>(); } this.triggerCapsuleQueue = new LinkedBlockingQueue<>(); - this.filterCapsuleQueue = new LinkedBlockingQueue<>(); chainBaseManager.setMerkleContainer(getMerkleContainer()); chainBaseManager.setMortgageService(mortgageService); this.initGenesis(); @@ -584,12 +550,6 @@ public void init() { ExecutorServiceManager.submit(triggerEs, triggerCapsuleProcessLoop); } - // start json rpc filter process - if (CommonParameter.getInstance().isJsonRpcFilterEnabled()) { - filterEs = ExecutorServiceManager.newSingleThreadExecutor(filterEsName); - ExecutorServiceManager.submit(filterEs, filterProcessLoop); - } - //initStoreFactory prepareStoreFactory(); //initActuatorCreator @@ -2658,7 +2618,6 @@ public void close() { stopRePushThread(); stopRePushTriggerThread(); EventPluginLoader.getInstance().stopPlugin(); - stopFilterProcessThread(); stopValidateSignThread(); chainBaseManager.shutdown(); revokingStore.shutdown(); diff --git a/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java b/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java index 6be47886117..44deee06859 100644 --- a/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java +++ b/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java @@ -36,7 +36,9 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.ForkJoinPool; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.regex.Pattern; +import javax.annotation.PostConstruct; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; @@ -53,7 +55,9 @@ import org.tron.common.es.ExecutorServiceManager; import org.tron.common.logsfilter.ContractEventParser; import org.tron.common.logsfilter.capsule.BlockFilterCapsule; +import org.tron.common.logsfilter.capsule.FilterTriggerCapsule; import org.tron.common.logsfilter.capsule.LogsFilterCapsule; +import org.tron.common.logsfilter.queue.FilterCapsuleQueue; import org.tron.common.parameter.CommonParameter; import org.tron.common.runtime.vm.DataWord; import org.tron.common.utils.ByteArray; @@ -193,20 +197,49 @@ public enum RequestSource { private final ExecutorService sectionExecutor; private final NodeInfoService nodeInfoService; private final Wallet wallet; - @Autowired - private Manager manager; + private final Manager manager; private final String esName = "query-section"; @Autowired - public TronJsonRpcImpl(@Autowired NodeInfoService nodeInfoService, @Autowired Wallet wallet) { + private FilterCapsuleQueue filterCapsuleQueue; + private ExecutorService filterEs; + private static final String filterEsName = "filter"; + private final AtomicBoolean closed = new AtomicBoolean(false); + + @Autowired + public TronJsonRpcImpl(NodeInfoService nodeInfoService, Wallet wallet, Manager manager) { this.nodeInfoService = nodeInfoService; this.wallet = wallet; + this.manager = manager; this.sectionExecutor = ExecutorServiceManager.newFixedThreadPool(esName, 5); } - @VisibleForTesting - public void setManager(Manager manager) { - this.manager = manager; + @PostConstruct + private void start() { + if (CommonParameter.getInstance().isJsonRpcFilterEnabled()) { + filterEs = ExecutorServiceManager.newSingleThreadExecutor(filterEsName, true); + ExecutorServiceManager.submit(filterEs, this::filterProcessLoop); + } + } + + private void filterProcessLoop() { + while (!closed.get()) { + try { + FilterTriggerCapsule filterCapsule = filterCapsuleQueue.poll(1, TimeUnit.SECONDS); + if (filterCapsule instanceof LogsFilterCapsule) { + handleLogsFilter((LogsFilterCapsule) filterCapsule); + } else if (filterCapsule instanceof BlockFilterCapsule) { + handleBLockFilter((BlockFilterCapsule) filterCapsule); + } else if (filterCapsule != null) { + logger.warn("Unknown FilterTriggerCapsule: {}", filterCapsule.getClass().getName()); + } + } catch (InterruptedException e) { + logger.error("FilterProcessLoop get InterruptedException, error is {}.", e.getMessage()); + Thread.currentThread().interrupt(); + } catch (Throwable throwable) { + logger.error("Unknown throwable happened in filterProcessLoop. ", throwable); + } + } } @VisibleForTesting @@ -1614,6 +1647,12 @@ public Object[] getFilterResult(String filterId, Map queue = - ReflectUtils.getFieldValue(dbManager, "filterCapsuleQueue"); + FilterCapsuleQueue queue = context.getBean(FilterCapsuleQueue.class); queue.clear(); // old branch: A carries a transfer; applied via the normal extend path @@ -1868,7 +1867,7 @@ private BlockCapsule blockWithTransfer(long time, long number, ByteString parent return blockCapsule; } - private boolean hasLogsFilterCapsule(BlockingQueue queue, BlockCapsule b, + private boolean hasLogsFilterCapsule(FilterCapsuleQueue queue, BlockCapsule b, boolean removed) { String blockHash = b.getBlockId().toString(); return queue.stream() @@ -1878,7 +1877,7 @@ private boolean hasLogsFilterCapsule(BlockingQueue queue, && blockHash.equals(c.getBlockHash())); } - private boolean hasBlockFilterCapsule(BlockingQueue queue, + private boolean hasBlockFilterCapsule(FilterCapsuleQueue queue, BlockCapsule b) { String blockHash = b.getBlockId().toString(); return queue.stream() diff --git a/framework/src/test/java/org/tron/core/jsonrpc/ConcurrentHashMapTest.java b/framework/src/test/java/org/tron/core/jsonrpc/ConcurrentHashMapTest.java index 2fcb624002e..855aae3f22d 100644 --- a/framework/src/test/java/org/tron/core/jsonrpc/ConcurrentHashMapTest.java +++ b/framework/src/test/java/org/tron/core/jsonrpc/ConcurrentHashMapTest.java @@ -23,7 +23,7 @@ @Slf4j public class ConcurrentHashMapTest { private static final String EXECUTOR_NAME = "jsonrpc-concurrent-map-test"; - private final TronJsonRpcImpl jsonRpc = new TronJsonRpcImpl(null, null); + private final TronJsonRpcImpl jsonRpc = new TronJsonRpcImpl(null, null, null); private static int randomInt(int minInt, int maxInt) { return (int) round(random(true) * (maxInt - minInt) + minInt, true); diff --git a/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineDeliveryTest.java b/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineDeliveryTest.java new file mode 100644 index 00000000000..97bd4e03c21 --- /dev/null +++ b/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineDeliveryTest.java @@ -0,0 +1,100 @@ +package org.tron.core.jsonrpc; + +import java.util.Collections; +import java.util.List; +import java.util.function.BooleanSupplier; +import javax.annotation.Resource; +import org.junit.Assert; +import org.junit.Test; +import org.tron.common.BaseTest; +import org.tron.common.TestConstants; +import org.tron.common.logsfilter.capsule.BlockFilterCapsule; +import org.tron.common.logsfilter.capsule.LogsFilterCapsule; +import org.tron.common.logsfilter.queue.FilterCapsuleQueue; +import org.tron.common.parameter.CommonParameter; +import org.tron.common.runtime.vm.DataWord; +import org.tron.common.runtime.vm.LogInfo; +import org.tron.core.config.args.Args; +import org.tron.core.services.jsonrpc.TronJsonRpc.FilterRequest; +import org.tron.core.services.jsonrpc.TronJsonRpc.LogFilterElement; +import org.tron.core.services.jsonrpc.TronJsonRpcImpl; +import org.tron.core.services.jsonrpc.filters.BlockFilterAndResult; +import org.tron.core.services.jsonrpc.filters.LogFilterAndResult; +import org.tron.protos.Protocol.TransactionInfo; + +/** + * End-to-end coverage of the decoupled filter pipeline: capsules offered to the + * FilterCapsuleQueue bean are delivered to registered filters by the consumer + * thread that TronJsonRpcImpl starts at context startup. + */ +public class FilterPipelineDeliveryTest extends BaseTest { + + static { + Args.setParam(new String[] {"--output-directory", dbPath()}, TestConstants.TEST_CONF); + // isJsonRpcFilterEnabled() must hold at context startup so the consumer thread starts + CommonParameter.getInstance().setJsonRpcHttpFullNodeEnable(true); + } + + @Resource + private TronJsonRpcImpl tronJsonRpc; + @Resource + private FilterCapsuleQueue filterCapsuleQueue; + + private static TransactionInfo buildTxInfoWithLog() { + LogInfo logInfo = new LogInfo(new byte[20], + Collections.singletonList(new DataWord(new byte[32])), new byte[0]); + return TransactionInfo.newBuilder().addLog(LogInfo.buildLog(logInfo)).build(); + } + + private static void await(BooleanSupplier condition, String message) + throws InterruptedException { + long deadline = System.currentTimeMillis() + 10_000; + while (System.currentTimeMillis() < deadline) { + if (condition.getAsBoolean()) { + return; + } + Thread.sleep(50); + } + Assert.fail(message); + } + + @Test + public void consumerThreadStartedOnce() { + long count = Thread.getAllStackTraces().keySet().stream() + .filter(t -> "filter".equals(t.getName())).count(); + Assert.assertEquals(1, count); + } + + @Test + public void blockFilterDeliveredOnFullAndSolidityPaths() throws Exception { + BlockFilterAndResult full = new BlockFilterAndResult(); + BlockFilterAndResult solidity = new BlockFilterAndResult(); + tronJsonRpc.getBlockFilter2ResultFull().put("pipeline-block-full", full); + tronJsonRpc.getBlockFilter2ResultSolidity().put("pipeline-block-solidity", solidity); + try { + filterCapsuleQueue.offer(new BlockFilterCapsule("e2e-full-hash", false)); + filterCapsuleQueue.offer(new BlockFilterCapsule("e2e-solidity-hash", true)); + await(() -> full.getResult().size() == 1, "full-path block hash not delivered"); + await(() -> solidity.getResult().size() == 1, "solidity-path block hash not delivered"); + } finally { + tronJsonRpc.getBlockFilter2ResultFull().remove("pipeline-block-full"); + tronJsonRpc.getBlockFilter2ResultSolidity().remove("pipeline-block-solidity"); + } + } + + @Test + public void reorgLogsDeliveredWithRemovedFlag() throws Exception { + LogFilterAndResult filter = new LogFilterAndResult(new FilterRequest(), 100L, null); + tronJsonRpc.getEventFilter2ResultFull().put("pipeline-log-removed", filter); + try { + filterCapsuleQueue.offer(new LogsFilterCapsule(150L, "0xreorg", null, + Collections.singletonList(buildTxInfoWithLog()), false, true)); + await(() -> !filter.getResult().isEmpty(), "removed log not delivered"); + List elements = filter.popAll(); + Assert.assertEquals(1, elements.size()); + Assert.assertTrue(elements.get(0).isRemoved()); + } finally { + tronJsonRpc.getEventFilter2ResultFull().remove("pipeline-log-removed"); + } + } +} diff --git a/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineShutdownTest.java b/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineShutdownTest.java new file mode 100644 index 00000000000..b58e95a5784 --- /dev/null +++ b/framework/src/test/java/org/tron/core/jsonrpc/FilterPipelineShutdownTest.java @@ -0,0 +1,120 @@ +package org.tron.core.jsonrpc; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.util.Collections; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import javax.annotation.Resource; +import org.junit.Assert; +import org.junit.FixMethodOrder; +import org.junit.Test; +import org.junit.runners.MethodSorters; +import org.tron.common.BaseTest; +import org.tron.common.TestConstants; +import org.tron.common.logsfilter.capsule.LogsFilterCapsule; +import org.tron.common.logsfilter.queue.FilterCapsuleQueue; +import org.tron.common.parameter.CommonParameter; +import org.tron.common.runtime.vm.DataWord; +import org.tron.common.runtime.vm.LogInfo; +import org.tron.core.config.args.Args; +import org.tron.core.services.jsonrpc.TronJsonRpc.FilterRequest; +import org.tron.core.services.jsonrpc.TronJsonRpcImpl; +import org.tron.core.services.jsonrpc.filters.LogFilterAndResult; +import org.tron.protos.Protocol.TransactionInfo; + +/** + * Shutdown semantics of the decoupled filter pipeline. Method order matters: + * test1 closes the shared TronJsonRpcImpl bean, test2 asserts the second close + * is a guarded no-op, so they run in name order. + */ +@FixMethodOrder(MethodSorters.NAME_ASCENDING) +public class FilterPipelineShutdownTest extends BaseTest { + + static { + Args.setParam(new String[] {"--output-directory", dbPath()}, TestConstants.TEST_CONF); + // isJsonRpcFilterEnabled() must hold at context startup so the consumer thread starts + CommonParameter.getInstance().setJsonRpcHttpFullNodeEnable(true); + } + + @Resource + private TronJsonRpcImpl tronJsonRpc; + @Resource + private FilterCapsuleQueue filterCapsuleQueue; + + private static TransactionInfo buildTxInfoWithLog() { + LogInfo logInfo = new LogInfo(new byte[20], + Collections.singletonList(new DataWord(new byte[32])), new byte[0]); + return TransactionInfo.newBuilder().addLog(LogInfo.buildLog(logInfo)).build(); + } + + /** + * close() while the consumer is mid-capsule on the parallel (logsFilterPool) path: + * the in-flight capsule must complete instead of dying on a RejectedExecutionException, + * which is exactly the filterEs-before-logsFilterPool ordering constraint. + */ + @Test + public void test1CloseDuringParallelProcessingLosesNoEvent() throws Exception { + tronJsonRpc.setFilterParallelThreshold(0); + LogFilterAndResult filter = new LogFilterAndResult(new FilterRequest(), 100L, null); + tronJsonRpc.getEventFilter2ResultFull().put("shutdown-race", filter); + + CountDownLatch entered = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + LogsFilterCapsule capsule = new LogsFilterCapsule(150L, "0xrace", null, + Collections.singletonList(buildTxInfoWithLog()), false, false) { + @Override + public boolean isSolidified() { + // first call is the head of handleLogsFilter, on the consumer thread + entered.countDown(); + try { + release.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + return false; + } + }; + filterCapsuleQueue.offer(capsule); + Assert.assertTrue("consumer did not pick up the capsule", + entered.await(10, TimeUnit.SECONDS)); + + Thread closer = new Thread(() -> { + try { + tronJsonRpc.close(); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + }, "test-closer"); + closer.start(); + // close() must be parked awaiting the consumer, with logsFilterPool still open + Thread.sleep(500); + Assert.assertTrue("close() finished while a capsule was in flight", closer.isAlive()); + + release.countDown(); + closer.join(150_000); + Assert.assertFalse("close() did not finish", closer.isAlive()); + Assert.assertEquals("in-flight capsule was lost during close()", + 1, filter.getResult().size()); + } + + @Test + public void test2RepeatedCloseIsIdempotentAndThreadStopped() throws Exception { + long deadline = System.currentTimeMillis() + 10_000; + while (System.currentTimeMillis() < deadline && filterThreadAlive()) { + Thread.sleep(50); + } + Assert.assertFalse("consumer thread still alive after close()", filterThreadAlive()); + + long t0 = System.nanoTime(); + tronJsonRpc.close(); + long elapsedMs = (System.nanoTime() - t0) / 1_000_000; + Assert.assertTrue("repeated close should be a fast no-op, took " + elapsedMs + "ms", + elapsedMs < 1_000); + } + + private static boolean filterThreadAlive() { + return Thread.getAllStackTraces().keySet().stream() + .anyMatch(t -> "filter".equals(t.getName()) && t.isAlive()); + } +} diff --git a/framework/src/test/java/org/tron/core/jsonrpc/HandleLogsFilterTest.java b/framework/src/test/java/org/tron/core/jsonrpc/HandleLogsFilterTest.java index 33835c482fe..be3098095d5 100644 --- a/framework/src/test/java/org/tron/core/jsonrpc/HandleLogsFilterTest.java +++ b/framework/src/test/java/org/tron/core/jsonrpc/HandleLogsFilterTest.java @@ -27,7 +27,7 @@ public class HandleLogsFilterTest { @Before public void setUp() { - jsonRpc = new TronJsonRpcImpl(null, null); + jsonRpc = new TronJsonRpcImpl(null, null, null); } @After diff --git a/framework/src/test/java/org/tron/core/jsonrpc/JsonRpcCallAndEstimateGasTest.java b/framework/src/test/java/org/tron/core/jsonrpc/JsonRpcCallAndEstimateGasTest.java index 2ab455fa580..cc11b2e99f1 100644 --- a/framework/src/test/java/org/tron/core/jsonrpc/JsonRpcCallAndEstimateGasTest.java +++ b/framework/src/test/java/org/tron/core/jsonrpc/JsonRpcCallAndEstimateGasTest.java @@ -211,8 +211,7 @@ private static TronJsonRpcImpl newRpcWithMockedFailedCall(byte[] resData, Estima }); } - TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet); - rpc.setManager(mockManager); + TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet, mockManager); return rpc; } @@ -239,8 +238,7 @@ private static TronJsonRpcImpl newRpcWithMockedSuccessfulCall(byte[]... constant .build(); }); - TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet); - rpc.setManager(mockManager); + TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet, mockManager); return rpc; } @@ -276,8 +274,7 @@ private static TronJsonRpcImpl newRpcWithMockedEstimateGasSuccessfulCall(long en }); } - TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet); - rpc.setManager(mockManager); + TronJsonRpcImpl rpc = new TronJsonRpcImpl(mockNodeInfo, mockWallet, mockManager); return rpc; } } diff --git a/framework/src/test/java/org/tron/core/jsonrpc/JsonrpcServiceTest.java b/framework/src/test/java/org/tron/core/jsonrpc/JsonrpcServiceTest.java index e8d14ace060..a564396198d 100644 --- a/framework/src/test/java/org/tron/core/jsonrpc/JsonrpcServiceTest.java +++ b/framework/src/test/java/org/tron/core/jsonrpc/JsonrpcServiceTest.java @@ -216,8 +216,7 @@ public void init() { dbManager.getTransactionRetStore() .put(ByteArray.fromLong(blockCapsule2.getNum()), transactionRetCapsule2); - tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); } @Test diff --git a/framework/src/test/java/org/tron/core/jsonrpc/WalletCursorTest.java b/framework/src/test/java/org/tron/core/jsonrpc/WalletCursorTest.java index 24ca71a74bc..c26d81e1f99 100644 --- a/framework/src/test/java/org/tron/core/jsonrpc/WalletCursorTest.java +++ b/framework/src/test/java/org/tron/core/jsonrpc/WalletCursorTest.java @@ -60,8 +60,7 @@ public void init() { @Test public void testSource() { - TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); Assert.assertEquals(Cursor.HEAD, wallet.getCursor()); Assert.assertEquals(RequestSource.FULLNODE, tronJsonRpc.getSource()); @@ -92,8 +91,7 @@ public void testDisableInSolidity() { dbManager.setCursor(Cursor.SOLIDITY); - TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); try { tronJsonRpc.buildTransaction(buildArguments); tronJsonRpc.close(); @@ -115,8 +113,7 @@ public void testDisableInPBFT() { dbManager.setCursor(Cursor.PBFT); - TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); try { tronJsonRpc.buildTransaction(buildArguments); } catch (Exception e) { @@ -143,8 +140,7 @@ public void testEnableInFullNode() { buildArguments.setTo("0x548794500882809695a8a687866e76d4271a1abc"); buildArguments.setValue("0x1f4"); - TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); try { tronJsonRpc.buildTransaction(buildArguments); @@ -164,8 +160,7 @@ public void testNewFilter_exceedsCapThrowsException() throws Exception { int saved = Args.getInstance().getJsonRpcMaxLogFilterNum(); Args.getInstance().setJsonRpcMaxLogFilterNum(cap); FilterRequest fr = new FilterRequest(); - TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet); - tronJsonRpc.setManager(dbManager); + TronJsonRpcImpl tronJsonRpc = new TronJsonRpcImpl(nodeInfoService, wallet, dbManager); Map map = tronJsonRpc.getEventFilter2ResultFull(); List addedKeys = new ArrayList<>(); From 9f4895267af217560382a1c8e0b115e31f8ad5ec Mon Sep 17 00:00:00 2001 From: 0xbigapple Date: Wed, 19 Aug 2026 10:07:36 +0800 Subject: [PATCH 2/3] fix(framework): exit filter loop on interrupt, widen shutdown catch --- .../main/java/org/tron/common/application/ApplicationImpl.java | 3 +-- .../java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java | 1 + 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/framework/src/main/java/org/tron/common/application/ApplicationImpl.java b/framework/src/main/java/org/tron/common/application/ApplicationImpl.java index f1bba10937a..88c15ae26fc 100644 --- a/framework/src/main/java/org/tron/common/application/ApplicationImpl.java +++ b/framework/src/main/java/org/tron/common/application/ApplicationImpl.java @@ -1,6 +1,5 @@ package org.tron.common.application; -import java.io.IOException; import java.util.concurrent.CountDownLatch; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -71,7 +70,7 @@ public void shutdown() { // (idempotent — Spring bean destruction may call close() again) try { tronJsonRpc.close(); - } catch (IOException e) { + } catch (Exception e) { logger.warn("Closing TronJsonRpcImpl failed.", e); } dbManager.close(); diff --git a/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java b/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java index 44deee06859..890ffb9081a 100644 --- a/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java +++ b/framework/src/main/java/org/tron/core/services/jsonrpc/TronJsonRpcImpl.java @@ -236,6 +236,7 @@ private void filterProcessLoop() { } catch (InterruptedException e) { logger.error("FilterProcessLoop get InterruptedException, error is {}.", e.getMessage()); Thread.currentThread().interrupt(); + return; } catch (Throwable throwable) { logger.error("Unknown throwable happened in filterProcessLoop. ", throwable); } From d9db07054c4701a5c23754a377cd5c7f8967dc55 Mon Sep 17 00:00:00 2001 From: 0xbigapple Date: Wed, 19 Aug 2026 21:18:29 +0800 Subject: [PATCH 3/3] refactor(jsonrpc): narrow filter queue api Remove unused queue inspection and mutation methods, and mark the remaining stream accessor as test-only. Drop unreachable offer failure handling for the unbounded queue. --- .../common/logsfilter/queue/FilterCapsuleQueue.java | 10 ++-------- framework/src/main/java/org/tron/core/db/Manager.java | 8 ++------ .../src/test/java/org/tron/core/db/ManagerTest.java | 1 - 3 files changed, 4 insertions(+), 15 deletions(-) diff --git a/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java b/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java index 7b5ba0e550d..35f551966f1 100644 --- a/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java +++ b/framework/src/main/java/org/tron/common/logsfilter/queue/FilterCapsuleQueue.java @@ -1,5 +1,6 @@ package org.tron.common.logsfilter.queue; +import com.google.common.annotations.VisibleForTesting; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; @@ -24,14 +25,7 @@ public FilterTriggerCapsule poll(long timeout, TimeUnit unit) throws Interrupted return queue.poll(timeout, unit); } - public int size() { - return queue.size(); - } - - public void clear() { - queue.clear(); - } - + @VisibleForTesting public Stream stream() { return queue.stream(); } diff --git a/framework/src/main/java/org/tron/core/db/Manager.java b/framework/src/main/java/org/tron/core/db/Manager.java index 162325573dd..d4dbddf2768 100644 --- a/framework/src/main/java/org/tron/core/db/Manager.java +++ b/framework/src/main/java/org/tron/core/db/Manager.java @@ -2304,9 +2304,7 @@ private void reApplyBlockEvents(List newBranch) { private void postBlockFilter(final BlockCapsule blockCapsule, boolean solidified) { BlockFilterCapsule blockFilterCapsule = new BlockFilterCapsule(blockCapsule, solidified); - if (!filterCapsuleQueue.offer(blockFilterCapsule)) { - logger.info("Too many filters, block filter lost: {}.", blockCapsule.getBlockId()); - } + filterCapsuleQueue.offer(blockFilterCapsule); } private void postLogsFilter(final BlockCapsule blockCapsule, boolean solidified, @@ -2319,9 +2317,7 @@ private void postLogsFilter(final BlockCapsule blockCapsule, boolean solidified, blockCapsule.getBlockId().toString(), blockCapsule.getBloom(), transactionInfoList, solidified, removed); - if (!filterCapsuleQueue.offer(logsFilterCapsule)) { - logger.info("Too many filters, logs filter lost: {}.", blockNumber); - } + filterCapsuleQueue.offer(logsFilterCapsule); } } diff --git a/framework/src/test/java/org/tron/core/db/ManagerTest.java b/framework/src/test/java/org/tron/core/db/ManagerTest.java index 5a973889ec2..5a903925f19 100755 --- a/framework/src/test/java/org/tron/core/db/ManagerTest.java +++ b/framework/src/test/java/org/tron/core/db/ManagerTest.java @@ -1734,7 +1734,6 @@ public void switchForkShouldPostFullNodeFilterForNewBranch() throws Exception { long expiration = t + 1_000_000L; FilterCapsuleQueue queue = context.getBean(FilterCapsuleQueue.class); - queue.clear(); // old branch: A carries a transfer; applied via the normal extend path BlockCapsule a = blockWithTransfer(t + 6000, base + 2, p.getBlockId().getByteString(), keys,