Skip to content

fix(event): improve event delivery reliability - #6965

Open
xxo1shine wants to merge 2 commits into
tronprotocol:release_v4.8.3from
xxo1shine:fix/event-delivery-backpressure
Open

xxo1shine wants to merge 2 commits into
tronprotocol:release_v4.8.3from
xxo1shine:fix/event-delivery-backpressure

Conversation

@xxo1shine

Copy link
Copy Markdown
Collaborator

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:

  • Unit Tests
  • Manual Testing

Follow up

Extra details

@github-actions
github-actions Bot requested a review from 0xbigapple September 11, 2026 07:20
@halibobo1205 halibobo1205 added this to the GreatVoyage-v4.8.3 milestone Sep 14, 2026
@halibobo1205 halibobo1205 added the topic:event subscribe transaction trigger, block trigger, contract event, contract log label Sep 14, 2026
@xxo1shine
xxo1shine changed the base branch from develop to release_v4.8.3 September 14, 2026 03:25
lxcmyf

This comment was marked as low quality.

@xxo1shine
xxo1shine force-pushed the fix/event-delivery-backpressure branch 2 times, most recently from 34517ce to d1dcc82 Compare September 17, 2026 10:05
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.
…olve-pr6965

# Conflicts:
#	framework/src/test/java/org/tron/common/logsfilter/NativeMessageQueueTest.java

@waynercheung waynercheung left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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) {

Copy link
Copy Markdown
Collaborator

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.

ZContext context = contexts.constructed().get(0);

// ZContext applies its defaults when creating the socket, so ordering matters.
InOrder startup = inOrder(context, publisher);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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 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();

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

topic:event subscribe transaction trigger, block trigger, contract event, contract log

Projects

Status: No status

Development

Successfully merging this pull request may close these issues.

Improve event delivery and ZeroMQ queue configuration

5 participants