Repository navigation
fix(event): improve event delivery reliability #6965
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: release_v4.8.3
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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<ZContext> 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); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [SHOULD] Both the context and the publisher are mocked, so this verifies initialization order without checking the actual socket HWM. Please add a real JeroMQ test using a dynamically allocated port: assert
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The current tests target the key aspects of this fix: setting the send HWM before creating the publisher, verifying the configured value of 2000, and verifying the fallback to 1000 for zero and negative values. The Reading the socket HWM from a real JeroMQ instance would provide additional integration coverage, but the existing tests already cover the parameter handling and call ordering addressed by this change. I would prefer to retain the current deterministic unit tests without adding a test that requires port binding in this PR. |
||
| startup.verify(context).setSndHWM(expectedQueueLength); | ||
| startup.verify(context).createSocket(SocketType.PUB); | ||
| startup.verify(publisher).bind("tcp://*:" + expectedPort); | ||
| startup.verifyNoMoreInteractions(); | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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(); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [SHOULD] Stubbing
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Stubbing
The suggested scenario would provide additional combined coverage, but these responsibilities already have separate tests, and this PR does not change the batch-processing logic inside |
||
|
|
||
| blockEventLoad.init(); | ||
| ArgumentCaptor<Runnable> 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"); | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
[SHOULD] Please document the soft-threshold contract here: normal loading checks
isBusy()before starting a batch, then enqueues the complete rollback/forward batch even if the queue crosses 500.add()does not enforce a capacity limit. This explains why future producers must participate in backpressure, and why rejecting events mid-batch would be incorrect.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The responsibilities are already fairly clear here:
isBusy()controls whether the next batch starts loading, whileadd()retains events from the current batch. The test also explicitly covers becoming busy at 500 events while continuing to accept subsequent events. I think the existing code and test sufficiently express this contract, so I would prefer not to add a comment here.