feat(a2a): wire A2A Gateway onto Runtime via TaskStore (issue #5302 Sub-PR D1) (#5313)

* feat(a2a): wire A2A Gateway onto Runtime via TaskStore (issue #5302 Sub-PR D1)

Brings the A2A Gateway back into the Runtime, this time with the unified
control-plane stores from #5301 Sub-PR A/C as the durable backend instead
of the in-memory TaskRegistry that PR #5260 introduced.

What lands
----------
* A2AGatewayService — task lifecycle controller. Replaces the in-memory
  TaskRegistry with TaskStore (Sub-PR A/C). A small runtime cache holds
  the two fields TaskStore does not model: parentTaskId (used to render
  child-task queries) and taskEpoch (the per-task counter the store
  requires for stale-write rejection on updateStatus). Both caches are
  rebuildable from the store on a fresh JVM; the cache keys are
  taskIds, so the cache size tracks active tasks only.
* A2AGatewayHttpHandler — Netty HTTP handler for the REST + SSE API
  (/a2a/tasks, /a2a/tasks/{id}, /a2a/tasks/{id}/stream, /a2a/health).
  SSE registers a StatusSubscriber so intermediate PENDING -> RUNNING ->
  COMPLETED transitions stream to the client.
* A2AGatewayServer — Netty HTTP bootstrap. Production wires
  EventMeshA2ATransport (Runtime-bridged) for the A2AMessageTransport
  dependency; tests pass an in-process transport. The weather-agent
  demo from PR #5260 is removed — the demo client (A2AGatewayDemo) and
  the Meta-ized AgentCard registry land in Sub-PR D2.
* AgentCardRegistry + InMemoryAgentCardRegistry — minimal discovery
  surface (isAgentRegistered / registerCard / getCard). A Meta-backed
  implementation backed by SessionStore (Sub-PR A) lands in D2.

Status mapping (PR #5260 -> Sub-PR A/C)
---------------------------------------
  SUBMITTED -> PENDING, WORKING -> RUNNING, CANCELLED -> CANCELED (one L)
The legacy wire vocabulary stays available via A2AGatewayService.TaskState
and the toLegacyState() helper. JSON response payloads still emit
SUBMITTED/WORKING/CANCELLED so existing A2A clients do not break.

Documentation
-------------
* docs/a2a-protocol/README.md carries an EXPERIMENTAL banner that lists
  the three pieces still missing for production (TaskExpirer, Meta-ized
  AgentCard, Testcontainers E2E) — all scheduled for Sub-PR D2.

Acceptance against issue #5302
------------------------------
* "A2A transport over the Runtime delivery path" — A2AGatewayService
  takes an A2AMessageTransport; production wires EventMeshA2ATransport
  which already bridges A2A onto UniIngressService (issue #5301 Sub-PR
  A's runtime ingress). The parallel InMemoryA2AMessageTransport is no
  longer required by the gateway path.
* "A2A tasks persist through TaskStore and recover across restarts" —
  createTask / updateStatus / getTask all go through the persistent
  store. On a fresh JVM, a gateway that had pending tasks sees them
  with their last persisted status; in-flight work that was awaiting
  a response will be picked up on the next status callback.

Tests (7, all passing)
----------------------
* A2AGatewayServiceTest — 6 tests: create+complete, cancel, cancel-on-
  unknown, submit-to-unregistered, parent-child index, task-not-found.
  Uses an in-process TaskStore (mirrors Sub-PR A's test stub) and an
  in-process A2AMessageTransport that handles A2A's + single-segment
  wildcards.
* A2AGatewaySmokeTest — Netty HTTP /a2a/health loopback check.

Not in D1 (Sub-PR D2)
---------------------
* TaskExpirer reaper (periodic TaskStore.expireStale sweep).
* Meta-backed AgentCardRegistry (SessionStore).
* Testcontainers fault-injection E2E.

* fix(a2a): address checkstyle violations in A2A Gateway D1 (PR #5313)

A2AGatewayHttpHandler.java:
- Remove unused imports (A2AProtocolConstants, A2ATopicFactory, TaskState)
- Expand try-catch blocks for limit/offset parsing (EmptyCatchBlock, NeedBraces)
- Expand handleList loop control flow (NeedBraces)
- Collapse duplicate blank line between package and imports (EmptyLineSeparator)

A2AGatewayServiceTest.java:
- Expand single-line if-blocks in InProcessTaskStore.listByAgent
- Expand if-block in InProcessTransport.matches
- Expand if-blocks in tearDown
- Move taskId declaration closer to its use (VariableDeclarationUsageDistance)

A2AGatewaySmokeTest.java:
- Rewrite StubTaskStore / NoopTransport with multi-line method bodies
  (LeftCurly, RightCurlyAlone, NeedBraces, EmptyLineSeparator)
- Add missing static imports (assertEquals, assertNotNull)

Build PR #5313 fix.

* fix(a2a): test imports + variable usage distance (PR #5313 followup)

A2AGatewayServiceTest.java:
- Fix ImportOrder: AgentCapabilities before AgentCard (alphabetical)
- Fix VariableDeclarationUsageDistance for f: inline taskId into
  submitTask call so the declaration is closer to first use

A2AGatewaySmokeTest.java:
- Fix ImportOrder: static Assertions imports before non-static
  TaskStore import

Build PR #5313 followup.

* fix(a2a): correct cancelMarksTaskCanceled test body (PR #5313 followup)

The previous fixup broke the test logic by calling TaskResult.getTaskId()
which does not exist. Restore the original taskId declaration pattern
while keeping the f declaration within 3 lines of its first use:

- Declare taskId first
- Declare f immediately after
- Inline cancelTask result into the assertTrue (no intermediate boolean)
- Move f.get() call right after the cancel assertion so f's first use
  is within 3 lines of its declaration

Build PR #5313 followup.
8 files changed
tree: 55230dccd1d1a752beef26dd7b6fa7c869504302
  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