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.
6 files changed
tree: 6d832fcd7416a1332fe32f0aec9ff0bda9ae2433
  1. .github/
  2. docker/
  3. docs/
  4. eventmesh-agent/
  5. eventmesh-common/
  6. eventmesh-connector-plugin/
  7. eventmesh-connector-runtime/
  8. eventmesh-examples/
  9. eventmesh-protocol-plugin/
  10. eventmesh-runtime/
  11. eventmesh-sdks/
  12. eventmesh-spi/
  13. eventmesh-storage-plugin/
  14. gradle/
  15. resources/
  16. style/
  17. tools/
  18. .asf.yaml
  19. .dockerignore
  20. .gitattributes
  21. .gitignore
  22. .gitmodules
  23. .licenserc.yaml
  24. build.gradle
  25. gradle.properties
  26. gradlew
  27. gradlew.bat
  28. install.sh
  29. LICENSE
  30. maturity.md
  31. NOTICE
  32. README.md
  33. README.zh-CN.md
  34. settings.gradle
README.md




CI status CodeCov Code Scanning

License GitHub Release Slack Status

đŸ“Ļ Documentation | 📔 Examples | âš™ī¸ Roadmap | 🌐 įŽ€äŊ“中文

Apache EventMesh

Apache EventMesh is a new generation serverless event middleware for building distributed event-driven applications.

EventMesh Architecture

EventMesh Architecture

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.

Features

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

  • CloudEvents-native, end to end — built entirely around the CloudEvents 1.0 specification, so events stay vendor-neutral and portable.
  • Lightweight, language-agnostic SDK — just three operations over plain HTTP (publish, subscribe, unsubscribe); no heavyweight client, no vendor lock-in.
  • Runtime-owned subscription & dispatch — subscription state and delivery semantics (load-balance / broadcast / multicast) are managed by EventMesh itself, not the underlying MQ, giving you consistent behavior across any storage backend.
  • MQ as a pure write-ahead log (WAL) — append-only, no consumer groups, no tags; the broker is reduced to durable storage, dramatically simplifying operations.
  • Guaranteed at-least-once delivery — EventMesh owns reliability through self-managed offsets and explicit ACK.
  • Multiple delivery transports — subscribers choose HTTP long-polling, Server-Sent Events (SSE), or WebSocket push, with request-reply support.
  • Effortless horizontal scaling — stateless Runtime instances scale out seamlessly with no rebalancing cost.

Extensibility & ecosystem

  • Agent-to-Agent (A2A) collaboration — a built-in A2A protocol turns EventMesh into an agent collaboration bus, bridging synchronous MCP / JSON-RPC 2.0 tool calls and asynchronous event-driven pub/sub for LLM and multi-agent systems.
  • Pluggable storage layer — Apache RocketMQ, Apache Kafka, Apache Pulsar, RabbitMQ, Redis, and more.
  • Pluggable interconnector layer — connectors run as standalone processes acting as the source or sink of SaaS, CloudService, Database, etc.
  • Pluggable meta service — Consul, Nacos, ETCD, and Zookeeper.
  • Event schema management via catalog service.
  • Powerful event orchestration through the Serverless workflow engine.
  • Powerful event filtering and transformation.

Subprojects

Quick start

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.

1. Pull EventMesh Image

Use the following command line to download the latest version of EventMesh:

sudo docker pull apache/eventmesh:latest

2. Run EventMesh Runtime

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 exposes 8083 for its optional admin port.

3. Publishing a CloudEvent

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!"
  }
}

4. Subscribing to a Topic

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"
}

5. Receiving Events

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. The CloudEventsClient Java SDK wraps all of these (subscribe / subscribeSse / subscribeWs); see the CloudEvents client guide.

6. Unsubscribing

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"
}

Contributing

GitHub repo Good Issues for newbies GitHub Help Wanted issues GitHub Help Wanted PRs GitHub repo Issues

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.

CNCF Landscape

Apache EventMesh enriches the CNCF Cloud Native Landscape.

License

Apache EventMesh is licensed under the Apache License, Version 2.0.

Community

WeChat AssistantWeChat Public AccountSlack
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

Mailing List

NameDescriptionSubscribeUnsubscribeArchive
UsersUser discussionSubscribeUnsubscribeMail Archives
DevelopmentDevelopment discussion (Design Documents, Issues, etc.)SubscribeUnsubscribeMail Archives
CommitsCommits to related repositoriesSubscribeUnsubscribeMail Archives
IssuesIssues or PRs comments and reviewsSubscribeUnsubscribeMail Archives