Conversation
34517ce to
d1dcc82
Compare
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.
d1dcc82 to
112c83b
Compare
…olve-pr6965 # Conflicts: # framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java
waynercheung
left a comment
There was a problem hiding this comment.
[SHOULD] Please document in both framework/src/main/resources/config.conf and common/src/main/resources/reference.conf that positive sendqueuelength values now take effect, and that non-positive values use 1000. The latter preserves the previous effective socket behavior, but differs from ZeroMQ's native meaning of zero. Also clarify in the PR description that PUB delivery remains best-effort; making the configured HWM effective does not provide an end-to-end delivery guarantee.
[NIT] When preparing the final merge commit message, please remove the claim in 112c83b that this change introduces synchronous block/transaction trigger delivery. That behavior already landed in #6833; this PR strengthens its tests.
| @@ -56,13 +56,13 @@ public void close() { | |||
| } | |||
|
|
|||
| public void add(Event event) { | |||
There was a problem hiding this comment.
[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.
| ZContext context = contexts.constructed().get(0); | ||
|
|
||
| // ZContext applies its defaults when creating the socket, so ordering matters. | ||
| InOrder startup = inOrder(context, publisher); |
There was a problem hiding this comment.
[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 publisher.getSndHWM() == 2000 after startup with 2000 (e.g. via ReflectUtils.getFieldValue(queue, "publisher")), and verify the 1000 fallback for zero and negative values. This would cover the observable behavior requested in #6929.
| Mockito.when(plugin.isBusy()).thenReturn(pluginBusy); | ||
| Mockito.when(realtime.isBusy()).thenReturn(realtimeBusy); | ||
| blockEventLoad = Mockito.spy(blockEventLoad); | ||
| Mockito.doNothing().when(blockEventLoad).load(); |
There was a problem hiding this comment.
[SHOULD] Stubbing load() leaves the batch and recovery behavior untested. Please add a regression using the real loader and RealtimeEventService, mocking only block retrieval and external dependencies: start with 499 queued events, load a rollback + forward batch, and verify the complete event order and the cache head. Then run the same scheduled task while busy, and again after work() drains the queue, to verify that loading pauses and resumes correctly.
What does this PR do?
Realtime event loading dropped new events when downstream processing reached the soft queue limit, while the native publisher applied its configured send high-water mark after socket creation.
Pause block event loading at 500 queued events and retain pending events for later delivery. Normalize the native queue settings and apply the sendQueueLength before creating the publisher so explicit values take effect while non-positive values preserve the default behavior.
Fixes #6929
Why are these changes required?
This PR has been tested by:
Follow up
Extra details