Repository navigation
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.
There was a problem hiding this comment.
The responsibilities are already fairly clear here: isBusy() controls whether the next batch starts loading, while add() 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.
| 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.
There was a problem hiding this comment.
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 InOrder assertions would catch a regression that moves the configuration back after socket creation.
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.
| 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.
There was a problem hiding this comment.
Stubbing load() here is intended to isolate the scheduled task's backpressure checks. Batch processing is covered separately: test() in the same test class uses the real BlockEventLoad and RealtimeEventService to verify that rollback events are enqueued before forward events, as well as the final cache head.
RealtimeEventServiceTest separately verifies the busy-state boundary at 499/500 events, retention and ordering of events beyond the threshold, and the return to a non-busy state after the queue is emptied.
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 load(). I would prefer to retain the current test structure without adding this combined scenario.
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