feat(state): DeadLetterStore + TaskStore Meta backends (issue #5301 Sub-PR C, fixes #5292 fully, unblocks #5302) (#5312)
* feat(state): DeadLetterStore + TaskStore Meta backends (issue #5301 Sub-PR C, fixes #5292 fully, unblocks #5302)
Two production implementations on top of the Sub-PR A interfaces:
- MetaBackedDeadLetterStore: idempotent ledger under /em/dlq/<deliveryId> via
MetaStore.putIfAbsent CAS; first recorder wins, subsequent writers are no-ops
that still return true so the dispatcher proceeds to retire.
- MetaBackedTaskStore: per-task record under /em/tasks/<taskId> with a self-describing
v1|<base64> wire envelope (all variable-width fields base64-encoded so opaque caller
input containing '|' or newlines cannot corrupt the format). Status updates use
MetaStore.tryAcquire as an epoch CAS, rejecting stale writers (issue #5291 style).
ReliableDispatcher 9-arg ctor: same as 8-arg + a DeadLetterStore that is invoked on
every confirmed DLQ transition. The 8-arg ctor (legacy Sub-PR A/B) is untouched and
the DLQ path skips the ledger call when no store is wired. Backward compatible.
Tests: MetaBackedDeadLetterStoreTest, MetaBackedTaskStoreTest, ReliableDispatcherDlqLedgerTest.
Interface contract tests DeadLetterStoreTest and TaskStoreTest (Sub-PR A) are unchanged
and continue to pass against the new production implementations.
Scope: fixes #5292 fully; unblocks #5302. A2A dispatch mode (#5302) and the
TaskExpirer reaper are follow-up scope in Sub-PR D (issue #5297).
* fix(state): allow deadLetterStore to be re-assigned in 9-arg ctor chain (issue #5301 Sub-PR C)
The 8-arg ctor assigns `deadLetterStore = null` and the 9-arg ctor chains to the 8-arg
ctor and then reassigns with the supplied ledger; this hit
`variable deadLetterStore might already have been assigned` because the field was
`final`. Drop `final` (with a Javadoc note that it is effectively final after
construction) so the chain compiles. No runtime semantics change.
* fix(state): fix imports + long line + test fully-qualified refs (Sub-PR C, issue #5301)
- ReliableDispatcher.java (CRLF): import order (DeadLetterStore before DeliveryStateStore),
drop the duplicate DeliveryStateStore import, split the 274-char log.warn into
three concatenated string fragments to satisfy LineLength <= 150.
- ReliableDispatcherDlqLedgerTest.java (LF): import AckCallback / DeadLetterSink / PushChannel
from `org.apache.eventmesh.runtime.delivery` (not `connector` / `push`); drop the
fully-qualified `org.apache.eventmesh.runtime.push.AckCallback` in the FakeChannel
override (now imports correctly).
All three compile errors and the three checkstyle violations from run 33057816390 are
addressed. No semantic change.
* fix(state): test checkstyle cleanups (Sub-PR C, issue #5301)
- MetaBackedDeadLetterStoreTest: drop unused MetaStore import (only InMemoryMetaStore is used).
- MetaBackedTaskStoreTest: rename local `bView` -> `bViewTwo` to satisfy LocalVariableName
(lowercase-first camelCase pattern).
- ReliableDispatcherDlqLedgerTest: drop three same-package `org.apache.eventmesh.runtime.delivery.*`
imports (the test lives in that package), and move `org.junit.jupiter.api.Test` to the top
import block (it was below `io.*` and broke ImportOrder).
All six checkstyle violations from run 33062276583 addressed; no semantic change.
* fix(state): second pass at test checkstyle (Sub-PR C, issue #5301)
- MetaBackedTaskStoreTest: rename local `bViewTwo` -> `taskB`. The checkstyle
LocalVariableName regex is `^[a-z]([a-z0-9][a-zA-Z0-9]*)?$` -- the second char
must be lowercase or digit, so `bViewTwo` (b + V) and `bView` (b + V) both fail.
`taskB` is the clean alternative.
- ReliableDispatcherDlqLedgerTest: import order was wrong. ASCII order requires
`org.apache.*` (a) before `org.junit.*` (j) before `java.*` (a again, but `j`
group sorts before `java`). The previous amendment placed `org.junit` ABOVE
`org.apache.eventmesh.*`, which violates ImportOrder.
* fix(state): reorder test imports for Apache checkstyle ImportOrder (#5301 Sub-PR C)
Apache checkstyle's ImportOrder rule uses a configured group order
(org.apache.eventmesh | org.apache | java | javax | org | io | net |
junit | com | lombok), not strict ASCII order. 'java.*' is group 4, so
it must appear AFTER 'org.apache.*' (group 3) and BEFORE 'org.junit.*'
(group 9). The previous amendment put 'org.junit.*' before 'java.*',
tripping checkstyleTest with 'Wrong order for java.net.URI import'.
Move all 'org.junit.*' imports to the LAST non-static group, after any
'java.*' or 'io.*' imports, in:
- ReliableDispatcherDlqLedgerTest
- MetaBackedTaskStoreTest
- MetaBackedDeadLetterStoreTest (no java imports; ordering unchanged)
* fix(state): reorder test imports for Apache checkstyle ImportOrder (#5301 Sub-PR C)
Apache checkstyle's ImportOrder rule uses a configured group order
(org.apache.eventmesh | org.apache | java | javax | org | io | net |
junit | com | lombok), not strict ASCII order. 'java.*' is group 4, so
it must appear AFTER 'org.apache.*' (group 3) and BEFORE 'org.junit.*'
(group 9). The previous amendment put 'org.junit.*' before 'java.*',
tripping checkstyleTest with 'Wrong order for java.net.URI import'.
Move all 'org.junit.*' imports to the LAST non-static group, after any
'java.*' or 'io.*' imports, in:
- ReliableDispatcherDlqLedgerTest
- MetaBackedTaskStoreTest
- MetaBackedDeadLetterStoreTest (no java imports; ordering unchanged)
* fix(state): put org.junit.* before io.* in test imports (Sub-PR C)
Apache checkstyle's ImportOrder rule uses a *first-match* prefix
algorithm against the `groups` config. `org.junit.*` matches the
`org` prefix (index 4) before the `junit` prefix (index 7), so
org.junit imports belong to the `org` group, NOT the `junit` group.
Required order: org(4) -> io(5). Previous fix had io(5) before
org.junit(4), tripping checkstyleTest with 'Wrong order for
org.junit.jupiter.api.Test' on line 42.
Reorder imports in ReliableDispatcherDlqLedgerTest.java to match
the pattern used by UniAdminServerTest.java (which passes):
static
org.apache.eventmesh.*
java.*
org.junit.* (org group, index 4)
io.* (io group, index 5)
* fix(state): correct test logic for DLQ exhaustion + stale expiry (Sub-PR C)
Two Sub-PR C tests had runtime logic bugs that compile/checkstyle did not
catch:
1. ReliableDispatcherDlqLedgerTest (both tests). The original test loop
used MAX_ATTEMPTS=2 iterations of (nack, clock+=1s, tick) which is one
iteration short: nack moves the next-attempt deadline into the future
(clock + backoff(1) = 1s) but the clock only advances 1s per tick,
so the second tick sees the delivery as not yet expired and skips it.
The test never reaches the DLQ branch (attempt >= maxAttempts).
Switch to timeout-driven exhaustion (matching
ReliableDispatcherTest.exhaustedRetriesGoToDLQ): bump MAX_ATTEMPTS to 3
and run 3 ticks each preceded by clock.addAndGet(ACK_TIMEOUT), so
attempt 1 -> 2 -> 3 -> DLQ. The first nack is no longer needed.
2. MetaBackedTaskStoreTest.expireStaleRemovesOldRecords. The test calls
expireStale(1L) immediately after updateStatus; both records are >1ms
old by then, so both end up in the expired list and the size==1
assertion fails.
Add a 3ms sleep between the last updateStatus and expireStale so the
just-updated 'new' record's updatedAtMs is within 1ms of 'now' and
survives, while 'old' (created several ms earlier) is removed.
These are test-only changes; production code (MetaBackedTaskStore,
ReliableDispatcher, MetaBackedDeadLetterStore) is unchanged.
* fix(state): declare InterruptedException in expireStale test (Sub-PR C)
Thread.sleep(3L) inside expireStaleRemovesOldRecords throws
InterruptedException; the test method must declare `throws Exception`
(or wrap the call in try/catch) for the test source to compile.
* fix(state): assert sink receives source topic, not source_DLQ (Sub-PR C)
In ReliableDispatcher.tick, the dlqSink is invoked with rec.topic
(the SOURCE topic, e.g. 'orders'). The '_DLQ' suffix is only
applied when recording on the durable ledger via
deadLetterStore.recordDeadLetter. The previous assertion
channel.dlqTopics.contains("orders_DLQ") is therefore wrong;
the sink sees the source topic unchanged.
Update the assertion to check for "orders" (the source topic),
with a comment explaining the source-vs-ledger distinction so the
next reader doesn't trip on the same misunderstanding.
* fix(state): bump expireStale test sleep from 3ms to 50ms (Sub-PR C)
The previous 3ms sleep worked locally but was too tight for CI
scheduler noise: on a busy host the 1ms window can collapse, so
both 'old' and 'new' records end up older than 1ms by the time
expireStale(1L) is called.
Bump the sleep to 50ms (still fast at < 100ms total per test)
so the just-updated 'new' record's updatedAtMs is well within 1ms
of 'now' (it survives), while 'old' (created ~50ms earlier) is
removed. The assertion logic is unchanged.
* fix(state): move updateStatus after sleep in expireStale test (Sub-PR C)
The previous test (and two follow-up attempts) updated 'new'
BEFORE the sleep, so by the time expireStale(1L) ran both
records were >1ms old and both got removed (expired.size() == 2).
Correct order:
1. createTask('old')
2. createTask('new')
3. Thread.sleep -- makes 'old' stale
4. updateStatus('new') -- refreshes 'new' to current time
5. expireStale(1L) -- removes only 'old' (now >1ms old)
Bump sleep to 30ms (was 50ms) -- still well under 100ms test budget,
and >1ms so 'old' is reliably stale even on a busy CI host.
đĻ Documentation | đ Examples | âī¸ Roadmap | đ įŽäŊ䏿
Apache EventMesh is a new generation serverless event middleware for building distributed event-driven applications.
EventMesh adopts a unified CloudEvents-over-MQ architecture. The message queue (MQ) acts as a pure write-ahead log (WAL) for durable storage only â there are no consumer groups, no tags, and no broker-side subscription semantics. Instead, the stateless EventMesh Runtime owns all delivery logic: its SubscriptionManager maintains the subscription registry and offset tracking, and dispatches events with load-balance, broadcast, and multicast semantics. Applications interact through a lightweight HTTP + CloudEvents 1.0 SDK (publish / subscribe / unsubscribe), while integration with external systems runs in a standalone Connector Runtime via the connector SPI.
Apache EventMesh is packed with features that help users build event-driven applications with ease. Here are the highlights that set EventMesh apart:
Core architecture
publish, subscribe, unsubscribe); no heavyweight client, no vendor lock-in.Extensibility & ecosystem
This section of the guide will show you the steps to deploy EventMesh from Local, Docker, K8s.
This section guides the launch of EventMesh according to the default configuration, if you need more detailed EventMesh deployment steps, please visit the EventMesh official document.
Use the following command line to download the latest version of EventMesh:
sudo docker pull apache/eventmesh:latest
Use the following command to start the EventMesh container:
sudo docker run -d --name eventmesh -p 8080:8080 -p 8081:8081 -p 8082:8082 -t apache/eventmesh:latest
Ports:
8080= traffic HTTP (/events/*CloudEvents API),8081= admin HTTP (/admin/*),8082= WebSocket push (opt-in). The connector runtime image exposes8083for its optional admin port.
Applications talk to EventMesh over the /events/* HTTP endpoints using standard CloudEvents 1.0. Publish an event with an HTTP POST (202 Accepted means the event was written to the WAL):
POST /events/publish HTTP/1.1 Host: localhost:8080 Content-Type: application/cloudevents+json { "specversion": "1.0", "id": "89010a5a-3c6f-4a1e-9b2d-0f7c1f2e3a4b", "source": "/example/producer", "type": "com.example.order.created", "subject": "orders", "datacontenttype": "application/json", "data": { "content": "Hello, EventMesh!" } }
Subscriptions are registered with EventMesh (there are no consumer groups or tags). Provide a clientId, the topic, and a distribution mode (LOAD_BALANCE, BROADCAST, MULTICAST, or LOAD_BALANCE_STICKY). The response returns a subscriptionId:
POST /events/subscribe HTTP/1.1 Host: localhost:8080 Content-Type: application/json { "clientId": "order-svc", "topic": "orders", "mode": "LOAD_BALANCE" }
Subscribers can receive dispatched events through three transports. Pick whichever fits your workload â all of them deliver the same CloudEvents and honor EventMesh's ACK / at-least-once semantics.
a) HTTP Long-Polling â pull events; the request blocks until events arrive or the timeout elapses:
GET /events/poll?clientId=order-svc&topics=orders&timeout=30000 HTTP/1.1 Host: localhost:8080
After processing the delivered events, acknowledge them so the offset advances (at-least-once delivery):
POST /events/ack HTTP/1.1 Host: localhost:8080 Content-Type: application/json { "subId": "sub-123", "clientId": "order-svc", "topic": "orders", "partition": 0, "offset": 42 }
b) Server-Sent Events (SSE) â server push over a long-lived HTTP connection; the client does not poll:
GET /events/stream?clientId=order-svc&topics=orders HTTP/1.1 Host: localhost:8080 Accept: text/event-stream
c) WebSocket â full-duplex server push over a dedicated WebSocket port:
GET /events/stream HTTP/1.1 Host: localhost:8082 Upgrade: websocket Connection: Upgrade
Long-polling, SSE, and WebSocket are interchangeable delivery transports â a subscriber chooses one. SSE and WebSocket are pushed by the server (no polling loop), while long-polling is client-driven. Request-reply (
POST /events/request+POST /events/reply) is also supported. TheCloudEventsClientJava SDK wraps all of these (subscribe/subscribeSse/subscribeWs); see the CloudEvents client guide.
When you no longer need to receive events for a topic, unsubscribe by clientId (optionally with the subscriptionId):
POST /events/unsubscribe HTTP/1.1 Host: localhost:8080 Content-Type: application/json { "clientId": "order-svc", "topic": "orders" }
Each contributor has played an important role in promoting the robust development of Apache EventMesh. We sincerely appreciate all contributors who have contributed code and documents.
Apache EventMesh enriches the CNCF Cloud Native Landscape.
Apache EventMesh is licensed under the Apache License, Version 2.0.
| WeChat Assistant | WeChat Public Account | Slack |
|---|---|---|
| Join Slack Chat(Please open an issue if this link is expired) |
Bi-weekly meeting : #Tencent meeting : 346-6926-0133
Bi-weekly meeting record : bilibili
| Name | Description | Subscribe | Unsubscribe | Archive |
|---|---|---|---|---|
| Users | User discussion | Subscribe | Unsubscribe | Mail Archives |
| Development | Development discussion (Design Documents, Issues, etc.) | Subscribe | Unsubscribe | Mail Archives |
| Commits | Commits to related repositories | Subscribe | Unsubscribe | Mail Archives |
| Issues | Issues or PRs comments and reviews | Subscribe | Unsubscribe | Mail Archives |