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 b32f1c22d39..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,90 +1,93 @@ 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.tron.common.utils.PublicMethod; import org.zeromq.SocketType; import org.zeromq.ZContext; import org.zeromq.ZMQ; public class NativeMessageQueueTest { - // 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"; - - 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()); } } }