From 112c83b1776bb829d1f972541dd5cfb6025dbf05 Mon Sep 17 00:00:00 2001 From: wb Date: Thu, 17 Sep 2026 17:59:57 +0800 Subject: [PATCH] fix(event): improve event delivery reliability Apply the publisher queue limit before socket creation and retain realtime events under backpressure. Deliver block and transaction triggers synchronously so rollback flags are captured before cached capsules are reused. --- .../nativequeue/NativeMessageQueue.java | 17 +-- .../core/services/event/BlockEventLoad.java | 2 +- .../services/event/RealtimeEventService.java | 10 +- .../logsfilter/NativeMessageQueueTest.java | 116 ++++++++------- .../tron/core/event/BlockEventLoadTest.java | 56 +++++++ .../core/event/RealtimeEventServiceTest.java | 139 ++++++++++++++---- 6 files changed, 244 insertions(+), 96 deletions(-) diff --git a/framework/src/main/java/org/tron/common/logsfilter/nativequeue/NativeMessageQueue.java b/framework/src/main/java/org/tron/common/logsfilter/nativequeue/NativeMessageQueue.java index 7d97e6f4ba9..7337f3df175 100644 --- a/framework/src/main/java/org/tron/common/logsfilter/nativequeue/NativeMessageQueue.java +++ b/framework/src/main/java/org/tron/common/logsfilter/nativequeue/NativeMessageQueue.java @@ -27,22 +27,21 @@ public static NativeMessageQueue getInstance() { } public boolean start(int bindPort, int sendQueueLength) { - context = new ZContext(); - publisher = context.createSocket(SocketType.PUB); - - if (Objects.isNull(publisher)) { - return false; - } - - if (bindPort == 0 || bindPort < 0) { + if (bindPort <= 0) { bindPort = DEFAULT_BIND_PORT; } - if (sendQueueLength < 0) { + if (sendQueueLength <= 0) { sendQueueLength = DEFAULT_QUEUE_LENGTH; } + context = new ZContext(); context.setSndHWM(sendQueueLength); + publisher = context.createSocket(SocketType.PUB); + + if (Objects.isNull(publisher)) { + return false; + } String bindAddress = String.format("tcp://*:%d", bindPort); return publisher.bind(bindAddress); diff --git a/framework/src/main/java/org/tron/core/services/event/BlockEventLoad.java b/framework/src/main/java/org/tron/core/services/event/BlockEventLoad.java index 7efc10a8ed5..65c2adf573d 100644 --- a/framework/src/main/java/org/tron/core/services/event/BlockEventLoad.java +++ b/framework/src/main/java/org/tron/core/services/event/BlockEventLoad.java @@ -37,7 +37,7 @@ public class BlockEventLoad { public void init() { executor.scheduleWithFixedDelay(() -> { try { - if (!instance.isBusy()) { + if (!instance.isBusy() && !realtimeEventService.isBusy()) { load(); } } catch (Exception e) { diff --git a/framework/src/main/java/org/tron/core/services/event/RealtimeEventService.java b/framework/src/main/java/org/tron/core/services/event/RealtimeEventService.java index cef16cd81c1..d2ec6a83c36 100644 --- a/framework/src/main/java/org/tron/core/services/event/RealtimeEventService.java +++ b/framework/src/main/java/org/tron/core/services/event/RealtimeEventService.java @@ -29,7 +29,7 @@ public class RealtimeEventService { private static BlockingQueue queue = new LinkedBlockingQueue<>(); - private int maxEventSize = 10000; + private static final int BUSY_EVENT_SIZE = 500; private final ScheduledExecutorService executor = ExecutorServiceManager .newSingleThreadScheduledExecutor("realtime-event"); @@ -56,13 +56,13 @@ public void close() { } public void add(Event event) { - if (queue.size() >= maxEventSize) { - logger.warn("Add event failed, blockId {}.", event.getBlockEvent().getBlockId().getString()); - return; - } queue.offer(event); } + public boolean isBusy() { + return queue.size() >= BUSY_EVENT_SIZE; + } + public synchronized void work() { while (queue.size() > 0) { Event event = queue.poll(); 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..1c7844a9bd8 100644 --- a/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java +++ b/framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java @@ -1,10 +1,17 @@ package org.tron.common.logsfilter; -import java.util.concurrent.ExecutorService; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockConstruction; +import static org.mockito.Mockito.when; + import org.junit.After; import org.junit.Assert; +import org.junit.Before; import org.junit.Test; -import org.tron.common.es.ExecutorServiceManager; +import org.mockito.InOrder; +import org.mockito.MockedConstruction; import org.tron.common.logsfilter.nativequeue.NativeMessageQueue; import org.zeromq.SocketType; import org.zeromq.ZContext; @@ -12,76 +19,75 @@ public class NativeMessageQueueTest { - public int bindPort = 5555; - public String dataToSend = "################"; - public String topic = "testTopic"; - - private ExecutorService subscriberExecutor; - private final String zmqSubscriber = "zmq-subscriber"; + private NativeMessageQueue queue; + private ZMQ.Socket publisher; + private MockedConstruction contexts; + + @Before + public void setUp() { + publisher = mock(ZMQ.Socket.class); + when(publisher.bind(anyString())).thenReturn(true); + contexts = mockConstruction(ZContext.class, (context, construction) -> + when(context.createSocket(SocketType.PUB)).thenReturn(publisher)); + queue = new NativeMessageQueue(); + } @After public void tearDown() { - ExecutorServiceManager.shutdownAndAwaitTermination(subscriberExecutor, zmqSubscriber); - subscriberExecutor = null; + try { + if (queue != null) { + queue.stop(); + } + } finally { + if (contexts != null) { + contexts.close(); + } + } } @Test - public void invalidBindPort() { - boolean bRet = NativeMessageQueue.getInstance().start(-1111, 0); - Assert.assertEquals(true, bRet); - NativeMessageQueue.getInstance().stop(); + public void configuredSendQueueLengthIsAppliedBeforeSocketCreation() { + assertStartup(6000, 2000, 6000, 2000); } @Test - public void invalidSendLength() { - boolean bRet = NativeMessageQueue.getInstance().start(0, -2222); - Assert.assertEquals(true, bRet); - NativeMessageQueue.getInstance().stop(); + public void invalidBindPortUsesDefaultPort() { + assertStartup(-1111, 1000, 5555, 1000); } @Test - public void publishTrigger() { - - int sendLength = 0; - boolean bRet = NativeMessageQueue.getInstance().start(bindPort, sendLength); - Assert.assertEquals(true, bRet); - - startSubscribeThread(); - - try { - Thread.sleep(1000); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - - NativeMessageQueue.getInstance().publishTrigger(dataToSend, topic); - - try { - Thread.sleep(1000); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } + public void negativeSendQueueLengthUsesDefaultSndHWM() { + assertStartup(6000, -1, 6000, 1000); + } - NativeMessageQueue.getInstance().stop(); + @Test + public void zeroSendQueueLengthUsesDefaultSndHWM() { + assertStartup(6000, 0, 6000, 1000); } - public void startSubscribeThread() { - subscriberExecutor = ExecutorServiceManager.newSingleThreadExecutor(zmqSubscriber); - subscriberExecutor.execute(() -> { - try (ZContext context = new ZContext()) { - ZMQ.Socket subscriber = context.createSocket(SocketType.SUB); + @Test + public void publishTriggerSendsTopicBeforePayload() { + Assert.assertTrue(queue.start(6000, 1000)); - Assert.assertTrue(subscriber.connect(String.format("tcp://localhost:%d", bindPort))); - Assert.assertTrue(subscriber.subscribe(topic)); + queue.publishTrigger("payload", "topic"); - while (!Thread.currentThread().isInterrupted()) { - byte[] message = subscriber.recv(); - String triggerMsg = new String(message); + InOrder delivery = inOrder(publisher); + delivery.verify(publisher).bind("tcp://*:6000"); + delivery.verify(publisher).sendMore("topic"); + delivery.verify(publisher).send("payload"); + delivery.verifyNoMoreInteractions(); + } - Assert.assertTrue(triggerMsg.contains(dataToSend) || triggerMsg.contains(topic)); - } - // ZMQ.Socket will be automatically closed when ZContext is closed - } - }); + private void assertStartup(int port, int queueLength, int expectedPort, int expectedQueueLength) { + Assert.assertTrue(queue.start(port, queueLength)); + Assert.assertEquals(1, contexts.constructed().size()); + ZContext context = contexts.constructed().get(0); + + // ZContext applies its defaults when creating the socket, so ordering matters. + InOrder startup = inOrder(context, publisher); + startup.verify(context).setSndHWM(expectedQueueLength); + startup.verify(context).createSocket(SocketType.PUB); + startup.verify(publisher).bind("tcp://*:" + expectedPort); + startup.verifyNoMoreInteractions(); } } diff --git a/framework/src/test/java/org/tron/core/event/BlockEventLoadTest.java b/framework/src/test/java/org/tron/core/event/BlockEventLoadTest.java index 991133fee78..61d8ae43937 100644 --- a/framework/src/test/java/org/tron/core/event/BlockEventLoadTest.java +++ b/framework/src/test/java/org/tron/core/event/BlockEventLoadTest.java @@ -5,9 +5,14 @@ import java.lang.reflect.Field; import java.lang.reflect.Method; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import org.junit.After; import org.junit.Assert; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.mockito.Mockito; +import org.tron.common.logsfilter.EventPluginLoader; import org.tron.common.utils.ReflectUtils; import org.tron.core.ChainBaseManager; import org.tron.core.capsule.BlockCapsule; @@ -23,6 +28,57 @@ public class BlockEventLoadTest { BlockEventLoad blockEventLoad = new BlockEventLoad(); + @After + public void tearDown() throws Exception { + getExecutor().shutdownNow(); + } + + @Test + public void shouldNotLoadWhenRealtimeEventServiceIsBusy() throws Exception { + verifyScheduledLoad(false, true, 0); + } + + @Test + public void shouldNotLoadWhenPluginIsBusy() throws Exception { + verifyScheduledLoad(true, false, 0); + } + + @Test + public void shouldLoadWhenBothConsumersAreReady() throws Exception { + verifyScheduledLoad(false, false, 1); + } + + private void verifyScheduledLoad(boolean pluginBusy, boolean realtimeBusy, int loadCalls) + throws Exception { + EventPluginLoader plugin = mock(EventPluginLoader.class); + RealtimeEventService realtime = mock(RealtimeEventService.class); + ScheduledExecutorService scheduler = mock(ScheduledExecutorService.class); + // Replace only scheduling and loading; execute the real init() task synchronously. + getExecutor().shutdownNow(); + ReflectUtils.setFieldValue(blockEventLoad, "executor", scheduler); + ReflectUtils.setFieldValue(blockEventLoad, "instance", plugin); + ReflectUtils.setFieldValue(blockEventLoad, "realtimeEventService", realtime); + Mockito.when(plugin.isBusy()).thenReturn(pluginBusy); + Mockito.when(realtime.isBusy()).thenReturn(realtimeBusy); + blockEventLoad = Mockito.spy(blockEventLoad); + Mockito.doNothing().when(blockEventLoad).load(); + + blockEventLoad.init(); + ArgumentCaptor task = ArgumentCaptor.forClass(Runnable.class); + Mockito.verify(scheduler).scheduleWithFixedDelay(task.capture(), Mockito.anyLong(), + Mockito.anyLong(), Mockito.eq(TimeUnit.MILLISECONDS)); + task.getValue().run(); + + Mockito.verify(blockEventLoad, Mockito.times(loadCalls)).load(); + Mockito.verify(blockEventLoad, Mockito.never()).close(); + } + + private ScheduledExecutorService getExecutor() throws ReflectiveOperationException { + Field field = BlockEventLoad.class.getDeclaredField("executor"); + field.setAccessible(true); + return (ScheduledExecutorService) field.get(blockEventLoad); + } + @Test public void test() throws Exception { Method method = blockEventLoad.getClass().getDeclaredMethod("load"); diff --git a/framework/src/test/java/org/tron/core/event/RealtimeEventServiceTest.java b/framework/src/test/java/org/tron/core/event/RealtimeEventServiceTest.java index f58f725195c..a95c5fc536c 100644 --- a/framework/src/test/java/org/tron/core/event/RealtimeEventServiceTest.java +++ b/framework/src/test/java/org/tron/core/event/RealtimeEventServiceTest.java @@ -3,10 +3,15 @@ import static org.mockito.Mockito.mock; import com.google.protobuf.ByteString; +import java.lang.reflect.Field; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import java.util.concurrent.BlockingQueue; import org.junit.Assert; import org.junit.Test; +import org.mockito.InOrder; +import org.mockito.MockedStatic; import org.mockito.Mockito; import org.tron.common.logsfilter.EventPluginLoader; import org.tron.common.logsfilter.capsule.BlockLogTriggerCapsule; @@ -26,6 +31,38 @@ public class RealtimeEventServiceTest { RealtimeEventService realtimeEventService = new RealtimeEventService(); + @Test + public void shouldBecomeBusyAt500EventsAndRetainLaterEvents() throws Exception { + Field queueField = RealtimeEventService.class.getDeclaredField("queue"); + queueField.setAccessible(true); + BlockingQueue queue = (BlockingQueue) queueField.get(null); + queue.clear(); + + try { + Event event = new Event(new BlockEvent(), false); + for (int i = 0; i < 499; i++) { + realtimeEventService.add(event); + } + + Assert.assertFalse(realtimeEventService.isBusy()); + + realtimeEventService.add(event); + Assert.assertTrue(realtimeEventService.isBusy()); + + Event laterEvent = new Event(new BlockEvent(), true); + realtimeEventService.add(laterEvent); + Assert.assertEquals(501, queue.size()); + for (int i = 0; i < 500; i++) { + Assert.assertSame(event, queue.remove()); + } + Assert.assertSame(laterEvent, queue.remove()); + } finally { + queue.clear(); + } + + Assert.assertFalse(realtimeEventService.isBusy()); + } + @Test public void test() throws Exception { BlockEvent be1 = new BlockEvent(); @@ -54,10 +91,13 @@ public void test() throws Exception { BlockCapsule blockCapsule = new BlockCapsule(0L, Sha256Hash.ZERO_HASH, 0L, ByteString.copyFrom(BlockEventCacheTest.getBlockId())); - // spy so processTrigger() is a no-op (does not reach the real EventPluginLoader), - // while setRemoved() still mutates the real trigger so the removed flag can be asserted. + // Capture the removed flag at delivery time, before the cached capsule is reused. BlockLogTriggerCapsule blockCap = Mockito.spy(new BlockLogTriggerCapsule(blockCapsule)); - Mockito.doNothing().when(blockCap).processTrigger(); + List deliveredRemovedFlags = new ArrayList<>(); + Mockito.doAnswer(invocation -> { + deliveredRemovedFlags.add(blockCap.getBlockLogTrigger().isRemoved()); + return null; + }).when(blockCap).processTrigger(); be2.setBlockLogTriggerCapsule(blockCap); Mockito.when(instance.isBlockLogTriggerEnable()).thenReturn(true); Mockito.when(instance.isBlockLogTriggerSolidified()).thenReturn(false); @@ -70,7 +110,7 @@ public void test() throws Exception { realtimeEventService.flush(be2, false); Assert.assertFalse(blockCap.getBlockLogTrigger().isRemoved()); // posted directly to the plugin both times, never via the async queue - Mockito.verify(blockCap, Mockito.times(2)).processTrigger(); + Assert.assertEquals(Arrays.asList(true, false), deliveredRemovedFlags); be2.setBlockLogTriggerCapsule(null); @@ -83,32 +123,79 @@ public void test() throws Exception { // rollback: tx trigger posted synchronously with removed=true realtimeEventService.flush(be2, true); - Mockito.verify(txCap).setRemoved(true); - Mockito.verify(txCap).processTrigger(); - - be2.setTransactionLogTriggerCapsules(null); + realtimeEventService.flush(be2, false); + InOrder delivery = Mockito.inOrder(txCap); + delivery.verify(txCap).setRemoved(true); + delivery.verify(txCap).processTrigger(); + delivery.verify(txCap).setRemoved(false); + delivery.verify(txCap).processTrigger(); + delivery.verifyNoMoreInteractions(); - SmartContractTrigger contractTrigger = new SmartContractTrigger(); - be2.setSmartContractTrigger(contractTrigger); + } - contractTrigger.getContractEventTriggers().add(mock(ContractEventTrigger.class)); - Mockito.when(instance.isContractLogTriggerEnable()).thenReturn(true); - try { - realtimeEventService.flush(be2, event.isRemove()); - } catch (Exception e) { - Assert.assertTrue(e instanceof NullPointerException); + @Test + public void shouldDeliverContractEventsOnlyWhenEnabledWithCurrentRemovedFlag() { + EventPluginLoader plugin = mock(EventPluginLoader.class); + ReflectUtils.setFieldValue(realtimeEventService, "instance", plugin); + BlockEvent block = new BlockEvent(); + block.setBlockId(new BlockCapsule.BlockId(BlockEventCacheTest.getBlockId(), 1)); + SmartContractTrigger triggers = new SmartContractTrigger(); + ContractEventTrigger event = new ContractEventTrigger(); + event.setTriggerName("staleName"); + triggers.getContractEventTriggers().add(event); + block.setSmartContractTrigger(triggers); + List deliveredFlags = new ArrayList<>(); + Mockito.doAnswer(invocation -> { + Assert.assertEquals("contractEventTrigger", event.getTriggerName()); + deliveredFlags.add(event.isRemoved()); + return null; + }).when(plugin).postContractEventTrigger(event); + + try (MockedStatic loader = Mockito.mockStatic(EventPluginLoader.class)) { + loader.when(EventPluginLoader::getInstance).thenReturn(plugin); + realtimeEventService.flush(block, true); + Mockito.verify(plugin, Mockito.never()).postContractEventTrigger(Mockito.any()); + + Mockito.when(plugin.isContractEventTriggerEnable()).thenReturn(true); + realtimeEventService.flush(block, true); + realtimeEventService.flush(block, false); + + Assert.assertEquals(Arrays.asList(true, false), deliveredFlags); + Mockito.verify(plugin, Mockito.times(2)).postContractEventTrigger(event); + Mockito.verify(plugin, Mockito.never()).postContractLogTrigger(Mockito.any()); } + } - contractTrigger.getContractEventTriggers().clear(); - - realtimeEventService.flush(be2, event.isRemove()); - - contractTrigger.getContractLogTriggers().add(mock(ContractLogTrigger.class)); - Mockito.when(instance.isContractEventTriggerEnable()).thenReturn(true); - try { - realtimeEventService.flush(be2, event.isRemove()); - } catch (Exception e) { - Assert.assertTrue(e instanceof NullPointerException); + @Test + public void shouldDeliverContractLogsOnlyWhenEnabledWithCurrentRemovedFlag() { + EventPluginLoader plugin = mock(EventPluginLoader.class); + ReflectUtils.setFieldValue(realtimeEventService, "instance", plugin); + BlockEvent block = new BlockEvent(); + block.setBlockId(new BlockCapsule.BlockId(BlockEventCacheTest.getBlockId(), 1)); + SmartContractTrigger triggers = new SmartContractTrigger(); + ContractLogTrigger log = new ContractLogTrigger(); + log.setTriggerName("staleName"); + triggers.getContractLogTriggers().add(log); + block.setSmartContractTrigger(triggers); + List deliveredFlags = new ArrayList<>(); + Mockito.doAnswer(invocation -> { + Assert.assertEquals("contractLogTrigger", log.getTriggerName()); + deliveredFlags.add(log.isRemoved()); + return null; + }).when(plugin).postContractLogTrigger(log); + + try (MockedStatic loader = Mockito.mockStatic(EventPluginLoader.class)) { + loader.when(EventPluginLoader::getInstance).thenReturn(plugin); + realtimeEventService.flush(block, true); + Mockito.verify(plugin, Mockito.never()).postContractLogTrigger(Mockito.any()); + + Mockito.when(plugin.isContractLogTriggerEnable()).thenReturn(true); + realtimeEventService.flush(block, true); + realtimeEventService.flush(block, false); + + Assert.assertEquals(Arrays.asList(true, false), deliveredFlags); + Mockito.verify(plugin, Mockito.times(2)).postContractLogTrigger(log); + Mockito.verify(plugin, Mockito.never()).postContractEventTrigger(Mockito.any()); } } }