diff --git a/deploy/kubernetes/connector-configmap.yaml b/deploy/kubernetes/connector-configmap.yaml
index 403ac2143b..20a33b6d40 100644
--- a/deploy/kubernetes/connector-configmap.yaml
+++ b/deploy/kubernetes/connector-configmap.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: v1
diff --git a/deploy/kubernetes/connector-deployment.yaml b/deploy/kubernetes/connector-deployment.yaml
index 7dffce4613..7308d35efa 100644
--- a/deploy/kubernetes/connector-deployment.yaml
+++ b/deploy/kubernetes/connector-deployment.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: apps/v1
diff --git a/deploy/kubernetes/kustomization.yaml b/deploy/kubernetes/kustomization.yaml
index dafed90cf2..2e20468768 100644
--- a/deploy/kubernetes/kustomization.yaml
+++ b/deploy/kubernetes/kustomization.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: kustomize.config.k8s.io/v1beta1
diff --git a/deploy/kubernetes/namespace.yaml b/deploy/kubernetes/namespace.yaml
index 73afc94e5b..e0dfff3d04 100644
--- a/deploy/kubernetes/namespace.yaml
+++ b/deploy/kubernetes/namespace.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: v1
diff --git a/deploy/kubernetes/runtime-configmap.yaml b/deploy/kubernetes/runtime-configmap.yaml
index 540a76c995..eb9ccfdb9d 100644
--- a/deploy/kubernetes/runtime-configmap.yaml
+++ b/deploy/kubernetes/runtime-configmap.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: v1
diff --git a/deploy/kubernetes/runtime-secret.yaml b/deploy/kubernetes/runtime-secret.yaml
index 849e82f776..e9d2453829 100644
--- a/deploy/kubernetes/runtime-secret.yaml
+++ b/deploy/kubernetes/runtime-secret.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
# Example admin token Secret. The bearer token guarding the admin API
diff --git a/deploy/kubernetes/runtime-service.yaml b/deploy/kubernetes/runtime-service.yaml
index 15d7c12b88..c5084530cd 100644
--- a/deploy/kubernetes/runtime-service.yaml
+++ b/deploy/kubernetes/runtime-service.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: v1
diff --git a/deploy/kubernetes/runtime-statefulset.yaml b/deploy/kubernetes/runtime-statefulset.yaml
index 83267e10a4..307973239b 100644
--- a/deploy/kubernetes/runtime-statefulset.yaml
+++ b/deploy/kubernetes/runtime-statefulset.yaml
@@ -1,20 +1,18 @@
#
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
#
-# http://www.apache.org/LICENSE-2.0
+# http://www.apache.org/licenses/LICENSE-2.0
#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied. See the License for the
-# specific language governing permissions and limitations
-# under the License.
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
#
apiVersion: apps/v1
diff --git a/docs/feature/a2a.md b/docs/feature/a2a.md
index 7cbb9df252..2310c77093 100644
--- a/docs/feature/a2a.md
+++ b/docs/feature/a2a.md
@@ -486,12 +486,12 @@ client.shutdown();
## 6. Future Roadmap
-* **EventMesh Broker Integration**: Replace `InMemoryA2AMessageTransport` with the real EventMesh broker for production deployment.
+* **EventMesh Broker Integration**: DONE for the main-process gateway (#5405): `-Deventmesh.a2a.enabled=true` boots the gateway on the runtime's own transport (`EventMeshA2ATransport` → CloudEvents-over-MQ).
* **Schema Registry**: Implement dynamic discovery of Agent capabilities via `methods/list`.
* **Sidecar Injection**: Fully integrate the adaptor into the EventMesh Sidecar for non-Java agents (Python, Node.js).
* **WebSocket Streaming**: Extend SSE to bidirectional WebSocket for real-time agent-to-agent dialogue.
-* **Task Persistence**: Persist `TaskRegistry` state to a durable store (Redis/DB) for crash recovery.
-* **Authentication**: Add API key / JWT authentication to the Gateway REST API.
+* **Task Persistence**: DONE (#5405): local `RocksDBTaskStore` under `/a2a-tasks` by default, or the Meta-backed store in clustered mode.
+* **Authentication**: bearer-token auth landed (#5405, `-Deventmesh.a2a.token`); JWT / API-key variants remain open.
---
diff --git a/docs/feature/agent.md b/docs/feature/agent.md
new file mode 100644
index 0000000000..903da3d809
--- /dev/null
+++ b/docs/feature/agent.md
@@ -0,0 +1,80 @@
+# EventMesh Agent (v2 Streaming Agent Process)
+
+
+
+The `eventmesh-agent` module is the v2 streaming-agent process: an independent JVM that
+registers with a running EventMesh runtime, subscribes its private lite channel, and bridges
+routed prompts to an OpenAI-compatible LLM gateway — streaming tokens back over the runtime
+to the requesting client.
+
+## Boot sequence
+
+1. **Register** — `POST /agent/register` on the runtime traffic port; the runtime assigns the
+ agent its `agent-parent` + `client-reply-parent` (§5.2 of the v2 design).
+2. **Subscribe** — the agent subscribes `agent.` on its parent via the lite wire
+ (`subscribeLiteBytes`).
+3. **Ready** — `POST /agent/ready` flips the registration ready; only now does matchmaking
+ route sessions to this agent (ready-before-route).
+4. **Heartbeat** — a virtual thread refreshes the TTL and reports active sessions every
+ `agent.heartbeat.intervalMs`.
+
+## Quick start
+
+```shell
+# 1. runtime (defaults: memory storage, traffic 10105)
+./gradlew :eventmesh-runtime:runRuntime # or use the docker image
+
+# 2. agent (needs an OpenAI-compatible endpoint)
+LLM_BASE_URL=https://api.openai.com LLM_API_KEY=sk-... LLM_MODEL=gpt-4o-mini ./bin/start-agent.sh # from dist-agent/
+```
+
+Dev runner without a distribution: `./gradlew :eventmesh-agent:runAgent -Dllm.api.key=sk-...`.
+
+## Configuration
+
+All config is `-D` system properties; `bin/start-agent.sh` maps the `AGENT_*` / `LLM_*` env
+vars (see `conf/agent.properties` for the full list).
+
+| Key | Default | Description |
+|---|---|---|
+| `agent.runtime.url` | `http://localhost:10105` | Runtime traffic URL (control plane + lite wire) |
+| `agent.id` | `agent-` | Agent identity; must be unique per process |
+| `agent.capacity` | `100` | Advertised concurrent-stream capacity (matchmaking input) |
+| `agent.heartbeat.intervalMs` | `10000` | Heartbeat cadence |
+| `agent.heartbeat.failLimit` | `6` | Consecutive heartbeat failures before the process exits (supervisor restarts it) |
+| `agent.conversation.maxHistory` | `20` | Per-conversation message sliding window |
+| `agent.conversation.maxConversations` | `1000` | Live-conversation bound; least-recently-used conversations are evicted |
+| `llm.base.url` | `https://api.openai.com` | OpenAI-compatible gateway base URL |
+| `llm.api.key` | _(empty)_ | Bearer key — **required**; empty fails fast at boot |
+| `llm.api.key.optional` | `false` | Opt-out of the empty-key fail-fast (mock gateways) |
+| `llm.model` | `gpt-4o-mini` | Default model; per-request model overrides win |
+
+## Reliability behavior
+
+- **Fail-fast on empty LLM key** — an agent without a usable key would fail every routed
+ request after registering READY, so boot refuses (unless `llm.api.key.optional=true`).
+- **Heartbeat failure limit** — after `agent.heartbeat.failLimit` consecutive failures the
+ process exits nonzero (the runtime TTL has evicted the registration by then; a zombie
+ serves nothing). Supervisors (systemd / K8s) restart it.
+- **Bounded conversations** — each conversation keeps a sliding window of turns, and the
+ store evicts least-recently-used conversations past `agent.conversation.maxConversations`,
+ bounding agent memory on long-lived processes.
+
+## Limitations
+
+- Conversation history is process-local (lost on restart); persistence is an explicit TODO.
+- Mode-1 (streaming calls) only — mode-2 pub/sub sessions are not routed to agents.
+- One LLM gateway per process (`llm.base.url`); multi-provider routing is future work.
diff --git a/docs/quickstart/configuration.md b/docs/quickstart/configuration.md
index e459127eea..4e5f017f65 100644
--- a/docs/quickstart/configuration.md
+++ b/docs/quickstart/configuration.md
@@ -79,6 +79,7 @@ Usually set via `-D` by `bin/start.sh`; override here if needed.
| `eventmesh.admin.port` | `10106` | Admin HTTP (`/admin/*`, `/metrics`) |
| `eventmesh.ws.port` | `-1` (disabled) | WebSocket push port; set e.g. `10107` to enable |
| `eventmesh.grpc.port` | `10205` | RESERVED for the future gRPC protocol (not served yet; keeps the 1.x default warm) |
+| `eventmesh.a2a.port` | `10108` | A2A gateway REST plane (opt-in via `eventmesh.a2a.enabled=true`) |
| `eventmesh.offset.path` | `./data/offset` | Local offset store directory |
## 4. Security
@@ -148,6 +149,19 @@ Summary:
## 8. Deployment checklist
+## A2A gateway (optional)
+
+| Key | Default | Description |
+|---|---|---|
+| `eventmesh.a2a.enabled` | `false` | Boot the A2A gateway plane with the main process (`EVENTMESH_A2A_ENABLED=true`) |
+| `eventmesh.a2a.port` | `10108` | A2A gateway REST + SSE port (`POST /a2a/tasks`, `/a2a/tasks/{id}/stream`) |
+| `eventmesh.a2a.token` | _(empty)_ | Bearer token required on every gateway endpoint; empty = open (dev mode, logged at boot) |
+| `eventmesh.a2a.taskstore` | _(local)_ | `meta` to persist tasks in the cluster Meta store (needs Nacos); default = local RocksDB under `/a2a-tasks` |
+
+The v2 **agent process** (`eventmesh-agent`) has its own `-D` config set (`agent.*`, `llm.*`) —
+see [feature/agent.md](../feature/agent.md).
+
+
- [ ] `EVENTMESH_STORAGE_TYPE` and the backend address set consistently on every instance
- [ ] `eventmesh.offset.path` points at persistent storage (survives restarts)
- [ ] Decide the WebSocket port (default disabled)
diff --git a/eventmesh-agent/bin/start-agent.sh b/eventmesh-agent/bin/start-agent.sh
new file mode 100644
index 0000000000..9f54271440
--- /dev/null
+++ b/eventmesh-agent/bin/start-agent.sh
@@ -0,0 +1,61 @@
+#!/usr/bin/env bash
+
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+
+#
+# EventMesh agent launcher (v2 streaming-agent process).
+#
+# Boots org.apache.eventmesh.agent.AgentApplication from $AGENT_HOME/{conf,apps,lib}.
+# The agent registers with the runtime over the control plane (/agent/*), subscribes
+# its lite channel, and bridges routed prompts to an OpenAI-compatible LLM gateway.
+#
+# Env vars (all optional):
+# AGENT_RUNTIME_URL runtime traffic URL (default http://localhost:10105)
+# AGENT_ID agent identity (default agent-)
+# LLM_BASE_URL OpenAI-compatible base URL (default https://api.openai.com)
+# LLM_API_KEY Bearer key; REQUIRED unless LLM_API_KEY_OPTIONAL=true
+# LLM_MODEL model name (default gpt-4o-mini)
+# AGENT_HEARTBEAT_MS heartbeat interval (default 10000)
+# AGENT_HEARTBEAT_FAILLIMIT consecutive-failure exit limit (default 6)
+# AGENT_MAX_HISTORY per-conversation message window (default 20)
+# AGENT_MAX_CONVERSATIONS live conversation bound (LRU) (default 1000)
+# AGENT_CAPACITY advertised stream capacity (default 100)
+# AGENT_OPTS extra -D flags
+#
+set -euo pipefail
+
+SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
+AGENT_HOME="${AGENT_HOME:-$(cd "$SCRIPT_DIR/.." && pwd)}"
+cd "$AGENT_HOME"
+
+AGENT_RUNTIME_URL="${AGENT_RUNTIME_URL:-http://localhost:10105}"
+AGENT_ID="${AGENT_ID:-}"
+LLM_BASE_URL="${LLM_BASE_URL:-https://api.openai.com}"
+LLM_API_KEY="${LLM_API_KEY:-}"
+LLM_MODEL="${LLM_MODEL:-gpt-4o-mini}"
+AGENT_HEARTBEAT_MS="${AGENT_HEARTBEAT_MS:-10000}"
+AGENT_HEARTBEAT_FAILLIMIT="${AGENT_HEARTBEAT_FAILLIMIT:-6}"
+AGENT_MAX_HISTORY="${AGENT_MAX_HISTORY:-20}"
+AGENT_MAX_CONVERSATIONS="${AGENT_MAX_CONVERSATIONS:-1000}"
+AGENT_CAPACITY="${AGENT_CAPACITY:-100}"
+AGENT_OPTS="${AGENT_OPTS:-}"
+
+ARGS=""
+if [ -n "$AGENT_ID" ]; then
+ ARGS="$ARGS -Dagent.id=${AGENT_ID}"
+fi
+
+exec java -Xmx512m $AGENT_OPTS -cp "conf:apps/*:lib/*" -Dagent.runtime.url="${AGENT_RUNTIME_URL}" -Dllm.base.url="${LLM_BASE_URL}" -Dllm.api.key="${LLM_API_KEY}" -Dllm.model="${LLM_MODEL}" -Dagent.heartbeat.intervalMs="${AGENT_HEARTBEAT_MS}" -Dagent.heartbeat.failLimit="${AGENT_HEARTBEAT_FAILLIMIT}" -Dagent.conversation.maxHistory="${AGENT_MAX_HISTORY}" -Dagent.conversation.maxConversations="${AGENT_MAX_CONVERSATIONS}" -Dagent.capacity="${AGENT_CAPACITY}" $ARGS org.apache.eventmesh.agent.AgentApplication
diff --git a/eventmesh-agent/conf/agent.properties b/eventmesh-agent/conf/agent.properties
index bd07a428f9..78e52adfc1 100644
--- a/eventmesh-agent/conf/agent.properties
+++ b/eventmesh-agent/conf/agent.properties
@@ -19,14 +19,18 @@
# AgentApplication reads ALL config from -D system properties (NOT this file). bin/start-agent.sh
# translates the ENV vars below into -D flags; pass extra -D flags through AGENT_OPTS.
#
-# AGENT_RUNTIME_URL -> -Dagent.runtime.url (default http://localhost:8080)
+# AGENT_RUNTIME_URL -> -Dagent.runtime.url (default http://localhost:10105)
# AGENT_STREAM_PARENT -> -Dagent.stream.parent (default em-stream; must match the runtime's eventmesh.stream.parent)
# AGENT_CLIENT_ID -> -Dagent.client.id (default em-agent-)
# LLM_BASE_URL -> -Dllm.base.url (OpenAI-compatible; default https://api.openai.com)
# LLM_API_KEY -> -Dllm.api.key (Bearer token; default "")
# LLM_MODEL -> -Dllm.model (default gpt-4o-mini)
+# AGENT_HEARTBEAT_FAILLIMIT -> -Dagent.heartbeat.failLimit (default 6; exits after N consecutive failures)
+# AGENT_MAX_CONVERSATIONS -> -Dagent.conversation.maxConversations (default 1000; LRU bound)
+#
+# Note: an empty LLM_API_KEY fails fast unless LLM_API_KEY_OPTIONAL=true is set via AGENT_OPTS.
#
# Example (real LLM):
# LLM_BASE_URL=https://api.openai.com LLM_API_KEY=sk-... LLM_MODEL=gpt-4o-mini \
-# AGENT_RUNTIME_URL=http://localhost:8080 ./bin/start-agent.sh
+# AGENT_RUNTIME_URL=http://localhost:10105 ./bin/start-agent.sh
#
diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
index 5b25bbe0d7..b686737a8a 100644
--- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
+++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/AgentApplication.java
@@ -37,14 +37,26 @@
public class AgentApplication {
public static void main(String[] args) throws Exception {
- String runtimeUrl = System.getProperty("agent.runtime.url", "http://localhost:8080");
+ String runtimeUrl = System.getProperty("agent.runtime.url", "http://localhost:10105");
String agentId = System.getProperty("agent.id", "agent-" + Long.toString(System.currentTimeMillis(), 36));
String llmBase = System.getProperty("llm.base.url", "https://api.openai.com");
String llmKey = System.getProperty("llm.api.key", "");
String llmModel = System.getProperty("llm.model", "gpt-4o-mini");
final long heartbeatMs = Long.getLong("agent.heartbeat.intervalMs", 10_000L);
int maxHistory = Integer.getInteger("agent.conversation.maxHistory", 20);
+ int maxConversations = Integer.getInteger("agent.conversation.maxConversations", 1000);
int capacity = Integer.getInteger("agent.capacity", 100);
+ // #5405: fail fast on an empty LLM key — an agent that registers READY without a usable
+ // key fails every routed request. Opt out explicitly for local/mock gateways.
+ boolean llmKeyOptional = Boolean.parseBoolean(System.getProperty("llm.api.key.optional", "false"));
+ if (llmKey.isEmpty() && !llmKeyOptional) {
+ throw new IllegalStateException("llm.api.key is empty - set -Dllm.api.key="
+ + " (or -Dllm.api.key.optional=true for mock gateways)");
+ }
+ // #5405: after this many consecutive heartbeat failures the process exits so a
+ // supervisor can restart it; the runtime TTL has evicted the registration by then
+ // and a silent zombie serves nothing.
+ final int heartbeatFailLimit = Integer.getInteger("agent.heartbeat.failLimit", 6);
// Step 1: register (gets the assigned agent-parent + client-reply-parent)
AgentControlClient control = new AgentControlClient(runtimeUrl);
@@ -55,7 +67,7 @@ public static void main(String[] args) throws Exception {
// Step 2: subscribe the agent's channel
CloudEventsClient client = CloudEventsClient.builder().runtimeUrl(runtimeUrl).clientId(agentId).build();
OpenAiLlmClient llm = new OpenAiLlmClient(llmBase, llmKey, llmModel);
- ConversationStore store = new ConversationStore(maxHistory);
+ ConversationStore store = new ConversationStore(maxHistory, maxConversations);
StreamingAgent agent = new StreamingAgent(client, agentParent, agentId, llm, store);
agent.start(); // subscribe (agentParent, agent.)
@@ -63,16 +75,26 @@ public static void main(String[] args) throws Exception {
control.ready(agentId);
log.info("agent READY: agentId={} (subscribed agent.{})", agentId, agentId);
- // Step 4: heartbeat loop (refresh TTL + report load)
+ // Step 4: heartbeat loop (refresh TTL + report load). Consecutive failures beyond
+ // agent.heartbeat.failLimit exit the process (nonzero) — the registration is gone by
+ // then, so staying alive would only serve errors.
+ final java.util.concurrent.atomic.AtomicInteger heartbeatFailures = new java.util.concurrent.atomic.AtomicInteger();
Thread heartbeatThread = Thread.startVirtualThread(() -> {
while (!Thread.currentThread().isInterrupted()) {
try {
Thread.sleep(heartbeatMs);
control.heartbeat(agentId, agent.activeSessions());
+ heartbeatFailures.set(0);
} catch (InterruptedException ie) {
return;
} catch (Exception e) {
- log.warn("heartbeat failed: {}", e.toString());
+ int fails = heartbeatFailures.incrementAndGet();
+ log.warn("heartbeat failed ({}/{}): {}", fails, heartbeatFailLimit, e.toString());
+ if (fails >= heartbeatFailLimit) {
+ log.error("heartbeat failed {} consecutive times - registration is gone; exiting for supervisor restart",
+ fails);
+ Runtime.getRuntime().halt(1);
+ }
}
}
});
diff --git a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java
index 394ca67e6d..44b1dfb76f 100644
--- a/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java
+++ b/eventmesh-agent/src/main/java/org/apache/eventmesh/agent/ConversationStore.java
@@ -20,6 +20,7 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -27,7 +28,9 @@
/**
* In-memory multi-turn conversation history keyed by {@code conversationId}. Each completed turn
* appends a {@code {user, assistant}} message pair; history is trimmed to the most recent
- * {@code maxMessages} entries (sliding window) to bound LLM context size.
+ * {@code maxMessages} entries (sliding window) to bound LLM context size, and bounds the
+ * number of conversations themselves ({@code maxConversations}, access-order eviction) —
+ * without the outer bound a long-lived agent leaks memory one session at a time (#5405).
*
* Lost on agent restart (MVP — persistence to Redis/DB is a TODO). Thread-safe per conversation
* (synchronized on the per-id list).
@@ -36,26 +39,66 @@ public class ConversationStore {
private final int maxMessages;
private final ConcurrentHashMap>> history = new ConcurrentHashMap<>();
+ /** Access-order LRU over the live conversations; guards the outer bound. */
+ private final LinkedHashMap lru = new LinkedHashMap<>(16, 0.75f, true);
public ConversationStore(int maxMessages) {
+ this(maxMessages, Integer.MAX_VALUE);
+ }
+
+ public ConversationStore(int maxMessages, int maxConversations) {
// keep at least one full turn (user + assistant)
this.maxMessages = Math.max(2, maxMessages);
+ this.maxConversations = Math.max(1, maxConversations);
}
+ private final int maxConversations;
+
/** Snapshot of the conversation history (empty list if convId null/unknown). Caller may mutate. */
public List
*/
@Slf4j
+@io.netty.channel.ChannelHandler.Sharable
public class A2AGatewayHttpHandler extends SimpleChannelInboundHandler {
private static final ObjectMapper objectMapper = new ObjectMapper();
private final A2AGatewayService gatewayService;
+ // NOTE (@Sharable): one handler instance is shared across every connection channel (the
+ // server adds it to each pipeline in initChannel). Safe because all fields are either
+ // final or volatile and no per-connection state lives on the handler.
+
/** #5304: optional unified security/quota/audit gate; null = allow all (current behavior). */
private volatile org.apache.eventmesh.runtime.security.gate.SecurityGate securityGate;
+ /** #5405: optional bearer token; null = open gateway (dev mode, logged at boot). */
+ private volatile String a2aToken;
+
public A2AGatewayHttpHandler(A2AGatewayService gatewayService) {
this.gatewayService = gatewayService;
}
+ /** #5405: require {@code Authorization: Bearer } on every gateway endpoint. */
+ public A2AGatewayHttpHandler withToken(String token) {
+ this.a2aToken = token;
+ return this;
+ }
+
/** #5304: install the unified gate so A2A task ops flow through RequestContext. */
public A2AGatewayHttpHandler withSecurityGate(
org.apache.eventmesh.runtime.security.gate.SecurityGate gate) {
@@ -85,6 +99,10 @@ protected void channelRead0(ChannelHandlerContext ctx, FullHttpRequest req) thro
String uri = req.uri();
// #5362: classify the sub-operation from the route shape so the gate charges the
// right resource (SUBMIT -> BACKLOG; GET / CANCEL / STREAM / list -> THROUGHPUT).
+ // #5405: bearer token check first (cheapest rejection before gate/classify).
+ if (!tokenCheck(ctx, req)) {
+ return;
+ }
org.apache.eventmesh.runtime.security.gate.RequestContext.A2aOperation a2aOp =
classify(req.method().name(), uri);
// #5304: every A2A task operation goes through the unified gate first.
@@ -151,8 +169,33 @@ private static org.apache.eventmesh.runtime.security.gate.RequestContext.A2aOper
return null;
}
- private boolean gateCheck(ChannelHandlerContext ctx, FullHttpRequest req, String uri) {
- return gateCheck(ctx, req, uri, null);
+
+ /**
+ * #5405: constant-time bearer check against the configured token (same pattern as the
+ * admin token guard, #5364). A null token means the gateway is open — the boot logs a
+ * warning so dev-mode openness is never silent.
+ */
+ private boolean tokenCheck(ChannelHandlerContext ctx, FullHttpRequest req) {
+ String token = a2aToken;
+ if (token == null) {
+ return true;
+ }
+ String header = req.headers().get("Authorization");
+ if (header == null || !header.startsWith("Bearer ")) {
+ A2AMetrics.inc(A2AMetrics.GATEWAY_REJECTIONS);
+ writeJson(ctx, HttpResponseStatus.UNAUTHORIZED,
+ "{\"error\":\"unauthorized\",\"message\":\"A2A gateway requires Authorization: Bearer \"}");
+ return false;
+ }
+ String presented = header.substring("Bearer ".length());
+ boolean ok = java.security.MessageDigest.isEqual(
+ presented.getBytes(StandardCharsets.UTF_8), token.getBytes(StandardCharsets.UTF_8));
+ if (!ok) {
+ A2AMetrics.inc(A2AMetrics.GATEWAY_REJECTIONS);
+ writeJson(ctx, HttpResponseStatus.UNAUTHORIZED,
+ "{\"error\":\"unauthorized\",\"message\":\"invalid A2A gateway token\"}");
+ }
+ return ok;
}
/**
@@ -191,6 +234,10 @@ private boolean gateCheck(ChannelHandlerContext ctx, FullHttpRequest req, String
return false;
}
+ private boolean gateCheck(ChannelHandlerContext ctx, FullHttpRequest req, String uri) {
+ return gateCheck(ctx, req, uri, null);
+ }
+
// =========================================================================
// Endpoint handlers
// =========================================================================
@@ -202,6 +249,7 @@ private void handleSubmit(ChannelHandlerContext ctx, FullHttpRequest req) throws
String targetAgent = (String) payload.get("targetAgent");
String message = (String) payload.get("message");
String parentTaskId = (String) payload.get("parentTaskId");
+ String contextId = (String) payload.get("contextId");
Boolean sync = (Boolean) payload.getOrDefault("sync", Boolean.FALSE);
if (targetAgent == null || message == null) {
@@ -218,7 +266,7 @@ private void handleSubmit(ChannelHandlerContext ctx, FullHttpRequest req) throws
} else {
// Async: return 202 with taskId
String taskId = "task-async-" + System.nanoTime();
- gatewayService.submitTask(taskId, targetAgent, message, parentTaskId);
+ gatewayService.submitTask(taskId, targetAgent, message, parentTaskId, contextId);
writeJson(ctx, HttpResponseStatus.ACCEPTED,
"{\"taskId\":\"" + taskId + "\",\"state\":\"SUBMITTED\"}");
return;
@@ -407,6 +455,9 @@ private String snapshotJson(A2AGatewayService.TaskSnapshot snap) {
+ (snap.getParentTaskId() != null
? ",\"parentTaskId\":\"" + snap.getParentTaskId() + "\""
: "")
+ + (snap.getRecord().contextId != null
+ ? ",\"contextId\":\"" + snap.getRecord().contextId + "\""
+ : "")
+ "}";
}
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayServer.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayServer.java
index 58169db157..f67fce9be5 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayServer.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayServer.java
@@ -66,6 +66,9 @@ public class A2AGatewayServer {
private A2AGatewayService gatewayService;
private A2AGatewayHttpHandler gatewayHandler;
+ /** #5405: optional bearer token for the gateway REST plane; null = open (dev). */
+ private volatile String token;
+
public A2AGatewayServer(int port, A2AMessageTransport transport, TaskStore taskStore,
AgentCardRegistry agentCardRegistry) {
this.port = port;
@@ -74,6 +77,12 @@ public A2AGatewayServer(int port, A2AMessageTransport transport, TaskStore taskS
this.agentCardRegistry = agentCardRegistry;
}
+ /** #5405: require {@code Authorization: Bearer } (constant-time compare). */
+ public A2AGatewayServer withToken(String token) {
+ this.token = token;
+ return this;
+ }
+
public void start() throws Exception {
// 1. Initialize components
gatewayService = new A2AGatewayService(
@@ -81,6 +90,9 @@ public void start() throws Exception {
gatewayService.start();
gatewayHandler = new A2AGatewayHttpHandler(gatewayService);
+ if (token != null) {
+ gatewayHandler.withToken(token);
+ }
// 2. Start Netty HTTP server
bossGroup = new NioEventLoopGroup(1);
@@ -123,7 +135,14 @@ public void shutdown() throws Exception {
}
}
+ /**
+ * The actually bound port. When the server was constructed with {@code port = 0}
+ * (auto-select), this returns the OS-assigned port after {@link #start()}.
+ */
public int getPort() {
+ if (serverChannel != null && serverChannel.localAddress() instanceof java.net.InetSocketAddress) {
+ return ((java.net.InetSocketAddress) serverChannel.localAddress()).getPort();
+ }
return port;
}
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayService.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayService.java
index 8e1e59e290..33c8e5e39f 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayService.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayService.java
@@ -192,8 +192,16 @@ public synchronized void shutdown() throws Exception {
* Submits an A2A task to a target agent with an auto-generated task id.
*/
public CompletableFuture submitTask(String targetAgent, String message, String parentTaskId) {
- String taskId = generateTaskId();
- return submitTask(taskId, targetAgent, message, parentTaskId);
+ return submitTask(generateTaskId(), targetAgent, message, parentTaskId, null);
+ }
+
+ /**
+ * Submits an A2A task with an explicit conversation {@code contextId} (issue #5405): groups
+ * tasks of one multi-turn conversation; persisted with the record and echoed in snapshots.
+ */
+ public CompletableFuture submitTask(String taskId, String targetAgent,
+ String message, String parentTaskId, String contextId) {
+ return submitTaskInternal(taskId, targetAgent, message, parentTaskId, contextId);
}
/**
@@ -202,6 +210,11 @@ public CompletableFuture submitTask(String targetAgent, String messa
*/
public CompletableFuture submitTask(String taskId, String targetAgent,
String message, String parentTaskId) {
+ return submitTaskInternal(taskId, targetAgent, message, parentTaskId, null);
+ }
+
+ private CompletableFuture submitTaskInternal(String taskId, String targetAgent,
+ String message, String parentTaskId, String contextId) {
if (!started) {
CompletableFuture future = new CompletableFuture<>();
future.completeExceptionally(new IllegalStateException("Gateway not started"));
@@ -225,7 +238,7 @@ public CompletableFuture submitTask(String taskId, String targetAgen
// throw from submitTask).
TaskRecord rec;
try {
- rec = taskStore.createTask(taskId, targetAgent, gatewayId, message);
+ rec = taskStore.createTask(taskId, targetAgent, gatewayId, message, contextId);
} catch (RuntimeException e) {
log.warn("Failed to create task in store: taskId={}: {}", taskId, e.getMessage());
CompletableFuture failed = new CompletableFuture<>();
@@ -241,9 +254,11 @@ public CompletableFuture submitTask(String taskId, String targetAgen
parentTaskIdCache.put(taskId, parentTaskId);
}
taskEpochCache.put(taskId, rec.taskEpoch);
+ A2AMetrics.inc(A2AMetrics.TASKS_SUBMITTED);
+ A2AMetrics.setGauge(A2AMetrics.TASKS_ACTIVE, pendingTasks.size());
// Build A2A CloudEvent
- CloudEvent event = buildTaskRequestEvent(taskId, targetAgent, message, parentTaskId);
+ CloudEvent event = buildTaskRequestEvent(taskId, targetAgent, message, parentTaskId, contextId);
// Register pending future BEFORE publishing: a synchronous transport callback (e.g. a
// local in-memory test) could deliver the response before submitTask returns, and we
@@ -262,6 +277,7 @@ public CompletableFuture submitTask(String taskId, String targetAgen
taskStore.updateStatus(taskId, epoch, Status.FAILED, errMsg);
}
pendingTasks.remove(taskId);
+ A2AMetrics.inc(A2AMetrics.TASKS_FAILED);
pending.completeExceptionally(new java.util.concurrent.TimeoutException(errMsg));
notifyStatusSubscribers(taskId, "failed", errMsg);
log.warn("Task timed out: taskId={}, targetAgent={}", taskId, targetAgent);
@@ -304,6 +320,7 @@ public boolean cancelTask(String taskId) {
}
boolean ok = taskStore.updateStatus(taskId, epoch, Status.CANCELED, null);
if (ok) {
+ A2AMetrics.inc(A2AMetrics.TASKS_CANCELED);
CompletableFuture future = pendingTasks.remove(taskId);
if (future != null) {
future.complete(new TaskResult(TaskState.CANCELLED, null, "Task cancelled"));
@@ -412,11 +429,13 @@ private void handleResponse(String topic, CloudEvent event) {
String resultData = extractEventData(event);
taskStore.updateStatus(taskId, rec.taskEpoch, Status.COMPLETED, resultData);
+ A2AMetrics.inc(A2AMetrics.TASKS_COMPLETED);
CompletableFuture future = pendingTasks.remove(taskId);
if (future != null) {
future.complete(new TaskResult(TaskState.COMPLETED, resultData, null));
}
+ A2AMetrics.setGauge(A2AMetrics.TASKS_ACTIVE, pendingTasks.size());
// Notify SSE subscribers
notifyStatusSubscribers(taskId, "completed", resultData);
@@ -461,7 +480,7 @@ private void notifyStatusSubscribers(String taskId, String state, String data) {
// =========================================================================
private CloudEvent buildTaskRequestEvent(String taskId, String targetAgent,
- String message, String parentTaskId) {
+ String message, String parentTaskId, String contextId) {
CloudEventBuilder builder = CloudEventBuilder.v1()
.withId(taskId)
.withType(A2AProtocolConstants.CE_TYPE_PREFIX + "task.request")
@@ -476,6 +495,9 @@ private CloudEvent buildTaskRequestEvent(String taskId, String targetAgent,
if (parentTaskId != null) {
builder.withExtension(A2AProtocolConstants.CE_EXTENSION_COLLABORATION_ID, parentTaskId);
}
+ if (contextId != null) {
+ builder.withExtension("a2acontextid", contextId);
+ }
return builder.build();
}
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AMetrics.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AMetrics.java
new file mode 100644
index 0000000000..645e711187
--- /dev/null
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AMetrics.java
@@ -0,0 +1,88 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.runtime.a2a;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicLong;
+
+/**
+ * Process-wide A2A gateway counters (issue #5405). Static mirrors — the gateway may boot
+ * in several configurations (main process, standalone) but there is at most one live
+ * gateway per JVM, and the admin plane reads these without holding a reference to the
+ * gateway instance.
+ *
+ * Surfaced through {@code GET /admin/metrics} (JSON) and {@code GET /metrics}
+ * (Prometheus text exposition) on the admin port.
+ */
+public final class A2AMetrics {
+
+ private static final Map COUNTERS = new ConcurrentHashMap<>();
+ private static final Map GAUGES = new ConcurrentHashMap<>();
+
+ private A2AMetrics() {
+ }
+
+ /** Counter names, stable wire contract for dashboards. */
+ public static final String TASKS_SUBMITTED = "a2a_tasks_submitted";
+ public static final String TASKS_COMPLETED = "a2a_tasks_completed";
+ public static final String TASKS_FAILED = "a2a_tasks_failed";
+ public static final String TASKS_CANCELED = "a2a_tasks_canceled";
+ public static final String TASKS_EXPIRED = "a2a_tasks_expired";
+ public static final String GATEWAY_REJECTIONS = "a2a_gateway_rejections";
+
+ /** Gauge names. */
+ public static final String TASKS_ACTIVE = "a2a_tasks_active";
+
+ public static void inc(String name) {
+ COUNTERS.computeIfAbsent(name, k -> new AtomicLong()).incrementAndGet();
+ }
+
+ public static void add(String name, long delta) {
+ COUNTERS.computeIfAbsent(name, k -> new AtomicLong()).addAndGet(delta);
+ }
+
+ public static void setGauge(String name, long value) {
+ GAUGES.computeIfAbsent(name, k -> new AtomicLong()).set(value);
+ }
+
+ /** Snapshot of all counters (name -> value); used by the admin plane. */
+ public static Map counters() {
+ Map out = new java.util.LinkedHashMap<>();
+ for (String n : new String[] {TASKS_SUBMITTED, TASKS_COMPLETED, TASKS_FAILED,
+ TASKS_CANCELED, TASKS_EXPIRED, GATEWAY_REJECTIONS}) {
+ AtomicLong v = COUNTERS.get(n);
+ out.put(n, v == null ? 0L : v.get());
+ }
+ return out;
+ }
+
+ /** Snapshot of all gauges (name -> value). */
+ public static Map gauges() {
+ Map out = new java.util.LinkedHashMap<>();
+ AtomicLong v = GAUGES.get(TASKS_ACTIVE);
+ out.put(TASKS_ACTIVE, v == null ? 0L : v.get());
+ return out;
+ }
+
+ /** Test hook: reset all instruments. */
+ static void reset() {
+ COUNTERS.clear();
+ GAUGES.clear();
+ }
+}
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/UniAdminServer.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/UniAdminServer.java
index afb42cf6ef..620a8d6226 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/UniAdminServer.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/admin/UniAdminServer.java
@@ -165,6 +165,12 @@ private void metrics(HttpExchange exchange) throws IOException {
out.put("redeliveries", admin.metrics().getRedeliveries());
out.put("dlqCount", admin.metrics().getDlqCount());
out.put("pendingDeliveries", admin.pendingDeliveries());
+ // #5405: A2A gateway plane counters (zero when the gateway is not enabled).
+ out.put("a2a", org.apache.eventmesh.runtime.a2a.A2AMetrics.counters());
+ java.util.Map a2aGauges = org.apache.eventmesh.runtime.a2a.A2AMetrics.gauges();
+ for (java.util.Map.Entry g : a2aGauges.entrySet()) {
+ out.put(g.getKey(), g.getValue());
+ }
writeJson(exchange, 200, out);
}
@@ -185,6 +191,13 @@ private void prometheusMetrics(HttpExchange exchange) throws IOException {
counter(sb, "eventmesh_redeliveries_count", m.getRedeliveries());
counter(sb, "eventmesh_dlq_count", m.getDlqCount());
gauge(sb, "eventmesh_pending_deliveries", admin.pendingDeliveries());
+ // #5405: A2A gateway plane.
+ for (java.util.Map.Entry c : org.apache.eventmesh.runtime.a2a.A2AMetrics.counters().entrySet()) {
+ counter(sb, "eventmesh_" + c.getKey(), c.getValue());
+ }
+ for (java.util.Map.Entry g : org.apache.eventmesh.runtime.a2a.A2AMetrics.gauges().entrySet()) {
+ gauge(sb, "eventmesh_" + g.getKey(), g.getValue());
+ }
byte[] body = sb.toString().getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "text/plain; version=0.0.4; charset=utf-8");
exchange.sendResponseHeaders(200, body.length);
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshApplication.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshApplication.java
index a969c3e7bb..781d9b5459 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshApplication.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshApplication.java
@@ -66,6 +66,10 @@ public class EventMeshApplication {
private int wsPort = -1;
private int wsBoundPort = -1;
private org.apache.eventmesh.runtime.connector.ConnectorScheduler connectorScheduler;
+ /** #5405: optional A2A gateway (Netty REST plane) booted with the main process. */
+ private org.apache.eventmesh.runtime.a2a.A2AGatewayServer a2aGateway;
+ private org.apache.eventmesh.runtime.a2a.EventMeshA2ATransport a2aTransport;
+ private org.apache.eventmesh.runtime.state.TaskStore a2aTaskStore;
private org.apache.eventmesh.runtime.session.AgentRegistrar agentRegistrar;
private org.apache.eventmesh.runtime.session.Matchmaker matchmaker;
private org.apache.eventmesh.runtime.session.SessionRouter sessionRouter;
@@ -151,6 +155,38 @@ public EventMeshApplication withConnectorScheduler(
return this;
}
+ /** #5405: the live gateway service (null before {@link #start()} or when not enabled). */
+ public org.apache.eventmesh.runtime.a2a.A2AGatewayService a2aGatewayService() {
+ return a2aGateway == null ? null : a2aGateway.getGatewayService();
+ }
+
+ /** #5405: the gateway's actually bound port (resolves auto-select {@code 0}); -1 when off. */
+ public int a2aGatewayPort() {
+ return a2aGateway == null ? -1 : a2aGateway.getPort();
+ }
+
+ /**
+ * Enable the A2A gateway (issue #5405): boots a Netty REST plane ({@code /a2a/*}) on
+ * {@code port}, bridged onto the runtime ingress via the REAL transport
+ * ({@code EventMeshA2ATransport} → CloudEvents-over-MQ), with {@code taskStore} for
+ * durable task state and an optional bearer {@code token} (null = open, dev mode).
+ *
+ * Must be called before {@link #start()}. The gateway shuts down with the application.
+ */
+ public EventMeshApplication withA2aGateway(int port,
+ org.apache.eventmesh.runtime.state.TaskStore taskStore, String token) {
+ this.a2aTransport = new org.apache.eventmesh.runtime.a2a.EventMeshA2ATransport(
+ runtime.ingress(), "a2a-gw-" + httpPort);
+ this.a2aTaskStore = taskStore;
+ this.a2aGateway = new org.apache.eventmesh.runtime.a2a.A2AGatewayServer(
+ port, a2aTransport, taskStore,
+ new org.apache.eventmesh.runtime.a2a.InMemoryAgentCardRegistry());
+ if (token != null && !token.isEmpty()) {
+ a2aGateway.withToken(token);
+ }
+ return this;
+ }
+
public EventMeshApplication(MeshStoragePlugin storage, OffsetStore offsetStore, int httpPort, int adminPort) {
// #5338: read delivery topology from system property (with documented default).
// Missing / blank -> LOCAL_STICKY_PULL (backward compatible). Unknown value -> fail-fast
@@ -314,8 +350,23 @@ private void startupInternal() throws Exception {
}
wsBoundPort = wsServer.start(wsPort);
}
- log.info("EventMeshApplication started: traffic port={} admin port={} ws port={}",
- trafficBoundPort, adminBoundPort, wsBoundPort);
+ // #5405: optional A2A gateway plane, started after the core servers so the ingress
+ // it bridges onto is fully up. TaskExpirer keeps the durable task store bounded.
+ if (a2aGateway != null) {
+ a2aGateway.start();
+ org.apache.eventmesh.runtime.a2a.TaskExpirer expirer =
+ new org.apache.eventmesh.runtime.a2a.TaskExpirer(a2aTaskStore,
+ org.apache.eventmesh.runtime.a2a.TaskExpirer.DEFAULT_IDLE_TTL_MS,
+ org.apache.eventmesh.runtime.a2a.TaskExpirer.DEFAULT_SCAN_INTERVAL_MS,
+ ids -> org.apache.eventmesh.runtime.a2a.A2AMetrics.add(
+ org.apache.eventmesh.runtime.a2a.A2AMetrics.TASKS_EXPIRED, ids.size()));
+ a2aGateway.getGatewayService().setTaskExpirer(expirer);
+ expirer.start();
+ log.info("A2A gateway started: port={} taskStore={} auth={}",
+ a2aGateway.getPort(), a2aTaskStore.getClass().getSimpleName(), "enabled");
+ }
+ log.info("EventMeshApplication started: traffic port={} admin port={} ws port={} a2a={}",
+ trafficBoundPort, adminBoundPort, wsBoundPort, a2aGateway != null);
}
/** Graceful shutdown: admin → traffic → runtime (flush offsets, release storage). */
@@ -329,6 +380,21 @@ public void shutdown() {
if (wsServer != null) {
wsServer.stop();
}
+ // #5405: gateway first (unsubscribes its transport), then its stores.
+ if (a2aGateway != null) {
+ try {
+ a2aGateway.shutdown();
+ } catch (Exception e) {
+ log.warn("A2A gateway shutdown: {}", e.toString());
+ }
+ }
+ if (a2aTransport != null) {
+ a2aTransport.shutdown();
+ }
+ if (a2aTaskStore != null) {
+ a2aTaskStore.flush();
+ a2aTaskStore.close();
+ }
if (sessionRouter != null) {
sessionRouter.shutdown();
}
@@ -506,6 +572,28 @@ public static void main(String[] args) throws Exception {
log.info("WebSocket push transport enabled on port {}", wsPort);
}
+ // A2A gateway (optional, issue #5405): -Deventmesh.a2a.enabled=true, port 10108 by
+ // default, token via -Deventmesh.a2a.token (open gateway without it — dev mode).
+ // Task store: local RocksDB under the same data dir by default; the Meta-backed store
+ // (cluster-visible) when -Deventmesh.a2a.taskstore=meta and a real Meta is configured.
+ if (Boolean.parseBoolean(System.getProperty("eventmesh.a2a.enabled", "false"))) {
+ int a2aPort = Integer.getInteger("eventmesh.a2a.port", 10108);
+ String a2aToken = System.getProperty("eventmesh.a2a.token", "");
+ org.apache.eventmesh.runtime.state.TaskStore a2aTasks;
+ if ("meta".equalsIgnoreCase(System.getProperty("eventmesh.a2a.taskstore", ""))
+ && metaStore != null && clustered) {
+ a2aTasks = new org.apache.eventmesh.runtime.state.MetaBackedTaskStore(metaStore);
+ } else {
+ String a2aDataDir = new java.io.File(offsetPath).getParentFile().getAbsolutePath()
+ + java.io.File.separator + "a2a-tasks";
+ a2aTasks = new org.apache.eventmesh.runtime.state.RocksDBTaskStore(a2aDataDir);
+ }
+ app.withA2aGateway(a2aPort, a2aTasks, a2aToken.isEmpty() ? null : a2aToken);
+ log.info("A2A gateway enabled: port={} taskStore={} token={}",
+ a2aPort, a2aTasks.getClass().getSimpleName(),
+ a2aToken.isEmpty() ? "off (open gateway)" : "on");
+ }
+
// v2 streaming sessions are NOT auto-wired here: the channel strategy is an explicit choice,
// so an embedder wires the session layer via builders before start(), passing the strategy as a
// parameter (no -D). See withSessionRouter's javadoc and LiteStreamCallIntegrationTest for the
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/http/UniHttpServer.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/http/UniHttpServer.java
index fda28100ab..dd87c1dd36 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/http/UniHttpServer.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/http/UniHttpServer.java
@@ -348,6 +348,15 @@ private void publish(HttpExchange exchange) throws IOException {
// path; the protocol adaptor owns the conversion.)
frame = FrameAdaptors.get("cloudevents").toFrame(new ByteTransport(body));
} catch (RuntimeException | org.apache.eventmesh.protocol.api.exception.ProtocolHandleException e) {
+ // #5225 / #5405: an A2A JSON-RPC body posted here is the most common class confusion —
+ // point the caller at the gateway endpoints instead of an opaque transform error.
+ if (looksLikeA2aJsonRpc(body)) {
+ writeJson(exchange, 400, error(
+ "this endpoint accepts CloudEvents JSON, not A2A JSON-RPC; post A2A task "
+ + "requests to the A2A gateway instead (POST /a2a/tasks on the a2a port, "
+ + "default 10108, enabled via -Deventmesh.a2a.enabled=true)"));
+ return;
+ }
writeJson(exchange, 400, error("invalid CloudEvent: " + e.getMessage()));
return;
}
@@ -1374,6 +1383,19 @@ private void writeJson(HttpExchange exchange, int status, Object body) throws IO
}
}
+ /**
+ * #5405: heuristic A2A JSON-RPC detection for the traffic-port guidance — a body whose
+ * first object contains a {@code "jsonrpc":"2.0"} member (A2A/MCP wire shape) or an
+ * {@code "protocoltype: a2a"} request header. Cheap prefix scan, no full parse.
+ */
+ private static boolean looksLikeA2aJsonRpc(byte[] body) {
+ if (body == null || body.length == 0 || body.length > 4096) {
+ return false;
+ }
+ String head = new String(body, 0, Math.min(body.length, 4096), java.nio.charset.StandardCharsets.UTF_8);
+ return head.contains("\"jsonrpc\"") && head.contains("\"2.0\"");
+ }
+
private static Map error(String msg) {
Map m = new HashMap<>();
m.put("error", msg);
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/MetaBackedTaskStore.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/MetaBackedTaskStore.java
index 9c026e8162..c00eb98786 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/MetaBackedTaskStore.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/MetaBackedTaskStore.java
@@ -98,7 +98,10 @@ private static String encode(TaskRecord r) {
.append(r.createdAtMs).append('|')
.append(r.updatedAtMs).append('|')
.append(b64(r.input)).append('|')
- .append(b64(r.output));
+ .append(b64(r.output))
+ // issue #5405: optional 10th field (contextId). The v1 decoder reads exactly 9
+ // fields, so older readers ignore the suffix; this reader takes 10 when present.
+ .append('|').append(b64(r.contextId));
String payload = Base64.getEncoder().encodeToString(
inner.toString().getBytes(StandardCharsets.UTF_8));
return WIRE_VERSION + "|" + payload;
@@ -121,7 +124,7 @@ private static TaskRecord decode(String value) {
// The inner payload contains 9 base64'd fields joined by '|' (output may be empty = null).
// Split on '|' with a fixed cap of 8 separators so the trailing field can contain '|'
// (it is base64, so it cannot, but defensive splitting is cheap).
- List parts = splitFixed(payload, 9);
+ List parts = splitFixed(payload, 10);
String taskId = b64Decode(parts.get(0));
String agentId = b64Decode(parts.get(1));
String clientId = b64Decode(parts.get(2));
@@ -131,7 +134,8 @@ private static TaskRecord decode(String value) {
long updated = Long.parseLong(parts.get(6));
String input = b64Decode(parts.get(7));
String output = b64Decode(parts.get(8));
- return new TaskRecord(taskId, agentId, clientId, status, created, updated, input, output, epoch);
+ String contextId = b64Decode(parts.get(9));
+ return new TaskRecord(taskId, agentId, clientId, status, created, updated, input, output, epoch, contextId);
}
private static List splitFixed(String s, int expectedFields) {
@@ -153,6 +157,12 @@ private static List splitFixed(String s, int expectedFields) {
@Override
public TaskRecord createTask(String taskId, String agentId, String clientId, String input) {
+ return createTask(taskId, agentId, clientId, input, null);
+ }
+
+ @Override
+ public TaskRecord createTask(String taskId, String agentId, String clientId, String input,
+ String contextId) {
if (taskId == null) {
return null;
}
@@ -161,7 +171,7 @@ public TaskRecord createTask(String taskId, String agentId, String clientId, Str
// resulting taskEpoch is unique across restarts of the same JVM and across instances.
long epoch = (System.currentTimeMillis() << 20) | (localEpoch.incrementAndGet() & 0xFFFFF);
TaskRecord rec = new TaskRecord(taskId, agentId, clientId, Status.PENDING,
- now, now, input, null, epoch);
+ now, now, input, null, epoch, contextId);
if (!meta.putIfAbsent(key(taskId), encode(rec))) {
return null;
}
@@ -187,7 +197,7 @@ public boolean updateStatus(String taskId, long expectedTaskEpoch, Status newSta
}
// Build the candidate new value with bumped updatedAtMs
TaskRecord next = new TaskRecord(cur.taskId, cur.agentId, cur.clientId, newStatus,
- cur.createdAtMs, System.currentTimeMillis(), cur.input, output, cur.taskEpoch);
+ cur.createdAtMs, System.currentTimeMillis(), cur.input, output, cur.taskEpoch, cur.contextId);
// CAS the encoded value. expectedOldValue must match the current Meta value verbatim,
// so re-read just before the CAS to minimise the lost-update window.
String currentEncoded = meta.get(key(taskId));
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/RocksDBTaskStore.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/RocksDBTaskStore.java
new file mode 100644
index 0000000000..881a8446c5
--- /dev/null
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/RocksDBTaskStore.java
@@ -0,0 +1,265 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.runtime.state;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Base64;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicLong;
+
+import org.rocksdb.FlushOptions;
+import org.rocksdb.Options;
+import org.rocksdb.RocksDB;
+import org.rocksdb.RocksDBException;
+import org.rocksdb.RocksIterator;
+
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * Local durable {@link TaskStore} backed by a RocksDB instance on disk (issue #5405).
+ *
+ * The default A2A task store for single-instance / docker deployments: no Meta
+ * (Nacos) dependency, state survives restarts under the same data-dir layout as the
+ * offset and delivery-state stores ({@code /a2a-tasks}). Clustered deployments
+ * keep {@link MetaBackedTaskStore} for cross-instance visibility.
+ *
+ * Wire format: identical envelope to {@code MetaBackedTaskStore} —
+ * {@code "v1|" + base64(inner)} where {@code inner} is the pipe-joined base64 field
+ * list — so a task written by one backend is readable by the other (migration path
+ * between local and Meta modes stays lossless up to schema version).
+ *
+ * CAS semantics: {@link #updateStatus} is guarded by the per-record
+ * {@code taskEpoch} (same contract as the Meta backend): a stale writer whose epoch
+ * does not match the stored record is rejected. The local RocksDB write is single
+ * process, so the CAS only guards logical staleness, not concurrency.
+ */
+@Slf4j
+public class RocksDBTaskStore implements TaskStore {
+
+ static {
+ RocksDB.loadLibrary();
+ }
+
+ /** Wire-format version marker, shared with {@code MetaBackedTaskStore}. */
+ public static final String WIRE_VERSION = "v1";
+
+ private final RocksDB db;
+ private final AtomicLong localEpoch = new AtomicLong();
+
+ public RocksDBTaskStore(String path) {
+ Options options = new Options().setCreateIfMissing(true);
+ try {
+ db = RocksDB.open(options, path);
+ log.info("RocksDBTaskStore opened at {}", path);
+ } catch (RocksDBException e) {
+ throw new IllegalStateException("failed to open RocksDB task store at " + path, e);
+ }
+ }
+
+ private static byte[] key(String taskId) {
+ return taskId.getBytes(StandardCharsets.UTF_8);
+ }
+
+ // ---- encode / decode: identical to MetaBackedTaskStore ----
+
+ static String encode(TaskRecord r) {
+ StringBuilder inner = new StringBuilder(256);
+ inner.append(b64(r.taskId)).append('|')
+ .append(b64(r.agentId)).append('|')
+ .append(b64(r.clientId)).append('|')
+ .append(r.status.name()).append('|')
+ .append(r.taskEpoch).append('|')
+ .append(r.createdAtMs).append('|')
+ .append(r.updatedAtMs).append('|')
+ .append(b64(r.input)).append('|')
+ .append(b64(r.output))
+ .append('|').append(b64(r.contextId));
+ String payload = Base64.getEncoder().encodeToString(
+ inner.toString().getBytes(StandardCharsets.UTF_8));
+ return WIRE_VERSION + "|" + payload;
+ }
+
+ static TaskRecord decode(String value) {
+ if (value == null) {
+ return null;
+ }
+ int sep = value.indexOf('|');
+ if (sep <= 0) {
+ throw new IllegalStateException("malformed task wire value (no version)");
+ }
+ String version = value.substring(0, sep);
+ if (!WIRE_VERSION.equals(version)) {
+ throw new IllegalStateException("unsupported task wire version: " + version);
+ }
+ String payload = new String(Base64.getDecoder().decode(value.substring(sep + 1)),
+ StandardCharsets.UTF_8);
+ List parts = splitFixed(payload, 10);
+ return new TaskRecord(
+ b64Decode(parts.get(0)), b64Decode(parts.get(1)), b64Decode(parts.get(2)),
+ Status.valueOf(parts.get(3)),
+ Long.parseLong(parts.get(5)), Long.parseLong(parts.get(6)),
+ b64Decode(parts.get(7)), b64Decode(parts.get(8)),
+ Long.parseLong(parts.get(4)), b64Decode(parts.get(9)));
+ }
+
+ private static String b64(String s) {
+ return s == null ? "" : Base64.getEncoder().encodeToString(s.getBytes(StandardCharsets.UTF_8));
+ }
+
+ private static String b64Decode(String s) {
+ return (s == null || s.isEmpty()) ? null
+ : new String(Base64.getDecoder().decode(s), StandardCharsets.UTF_8);
+ }
+
+ private static List splitFixed(String s, int expectedFields) {
+ List out = new ArrayList<>(expectedFields);
+ int start = 0;
+ for (int i = 0; i < expectedFields - 1; i++) {
+ int sep = s.indexOf('|', start);
+ if (sep < 0) {
+ throw new IllegalStateException("malformed task wire payload");
+ }
+ out.add(s.substring(start, sep));
+ start = sep + 1;
+ }
+ out.add(s.substring(start));
+ return out;
+ }
+
+ @Override
+ public TaskRecord createTask(String taskId, String agentId, String clientId, String input) {
+ return createTask(taskId, agentId, clientId, input, null);
+ }
+
+ @Override
+ public TaskRecord createTask(String taskId, String agentId, String clientId, String input,
+ String contextId) {
+ if (taskId == null) {
+ return null;
+ }
+ long now = System.currentTimeMillis();
+ long epoch = (System.currentTimeMillis() << 20) | (localEpoch.incrementAndGet() & 0xFFFFF);
+ TaskRecord rec = new TaskRecord(taskId, agentId, clientId, Status.PENDING,
+ now, now, input, null, epoch, contextId);
+ try {
+ byte[] k = key(taskId);
+ if (db.get(k) != null) {
+ return null; // duplicate taskId, same contract as the Meta backend
+ }
+ db.put(k, encode(rec).getBytes(StandardCharsets.UTF_8));
+ } catch (RocksDBException e) {
+ throw new IllegalStateException("RocksDB createTask failed: " + e.getMessage(), e);
+ }
+ return rec;
+ }
+
+ @Override
+ public TaskRecord getTask(String taskId) {
+ if (taskId == null) {
+ return null;
+ }
+ try {
+ byte[] val = db.get(key(taskId));
+ return val == null ? null : decode(new String(val, StandardCharsets.UTF_8));
+ } catch (RocksDBException e) {
+ throw new IllegalStateException("RocksDB getTask failed: " + e.getMessage(), e);
+ }
+ }
+
+ @Override
+ public boolean updateStatus(String taskId, long expectedTaskEpoch, Status newStatus, String output) {
+ if (taskId == null || newStatus == null) {
+ return false;
+ }
+ try {
+ byte[] k = key(taskId);
+ byte[] curBytes = db.get(k);
+ if (curBytes == null) {
+ return false;
+ }
+ TaskRecord cur = decode(new String(curBytes, StandardCharsets.UTF_8));
+ if (cur.taskEpoch != expectedTaskEpoch) {
+ return false;
+ }
+ TaskRecord next = new TaskRecord(cur.taskId, cur.agentId, cur.clientId, newStatus,
+ cur.createdAtMs, System.currentTimeMillis(), cur.input, output, cur.taskEpoch, cur.contextId);
+ db.put(k, encode(next).getBytes(StandardCharsets.UTF_8));
+ return true;
+ } catch (RocksDBException e) {
+ log.warn("RocksDB updateStatus failed for {}: {}", taskId, e.getMessage());
+ return false;
+ }
+ }
+
+ @Override
+ public List listByAgent(String agentId, Status statusFilter) {
+ List out = new ArrayList<>();
+ if (agentId == null) {
+ return out;
+ }
+ try (RocksIterator it = db.newIterator()) {
+ for (it.seekToFirst(); it.isValid(); it.next()) {
+ TaskRecord r = decode(new String(it.value(), StandardCharsets.UTF_8));
+ if (r != null && agentId.equals(r.agentId)
+ && (statusFilter == null || r.status == statusFilter)) {
+ out.add(r);
+ }
+ }
+ }
+ return out;
+ }
+
+ @Override
+ public List expireStale(long olderThanMs) {
+ long deadline = System.currentTimeMillis() - olderThanMs;
+ List expired = new ArrayList<>();
+ List toDelete = new ArrayList<>();
+ try (RocksIterator it = db.newIterator()) {
+ for (it.seekToFirst(); it.isValid(); it.next()) {
+ TaskRecord r = decode(new String(it.value(), StandardCharsets.UTF_8));
+ if (r != null && r.updatedAtMs < deadline) {
+ expired.add(r.taskId);
+ toDelete.add(key(r.taskId));
+ }
+ }
+ }
+ try {
+ for (byte[] k : toDelete) {
+ db.delete(k);
+ }
+ } catch (RocksDBException e) {
+ log.warn("RocksDB expireStale delete failed: {}", e.getMessage());
+ }
+ return expired;
+ }
+
+ @Override
+ public void flush() {
+ try (FlushOptions fo = new FlushOptions().setWaitForFlush(true)) {
+ db.flush(fo);
+ } catch (RocksDBException e) {
+ log.warn("RocksDB flush failed: {}", e.getMessage());
+ }
+ }
+
+ @Override
+ public void close() {
+ db.close();
+ }
+}
diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/TaskStore.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/TaskStore.java
index 476a815118..f57bbcd476 100644
--- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/TaskStore.java
+++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/TaskStore.java
@@ -61,9 +61,20 @@ final class TaskRecord {
public volatile String output;
/** Monotonic per-task epoch; set at creation, never reset, used to reject stale writes. */
public final long taskEpoch;
+ /**
+ * Optional conversation linkage (A2A {@code contextId}, issue #5405): groups tasks that
+ * belong to the same multi-turn conversation. Null when the caller did not supply one.
+ */
+ public final String contextId;
public TaskRecord(String taskId, String agentId, String clientId, Status status,
long createdAtMs, long updatedAtMs, String input, String output, long taskEpoch) {
+ this(taskId, agentId, clientId, status, createdAtMs, updatedAtMs, input, output, taskEpoch, null);
+ }
+
+ public TaskRecord(String taskId, String agentId, String clientId, Status status,
+ long createdAtMs, long updatedAtMs, String input, String output, long taskEpoch,
+ String contextId) {
this.taskId = taskId;
this.agentId = agentId;
this.clientId = clientId;
@@ -73,6 +84,7 @@ public TaskRecord(String taskId, String agentId, String clientId, Status status,
this.input = input;
this.output = output;
this.taskEpoch = taskEpoch;
+ this.contextId = contextId;
}
}
@@ -82,6 +94,16 @@ public TaskRecord(String taskId, String agentId, String clientId, Status status,
*/
TaskRecord createTask(String taskId, String agentId, String clientId, String input);
+ /**
+ * Create a new task with an optional conversation {@code contextId} (issue #5405). Backends
+ * that do not persist the field may ignore it; the default delegates to
+ * {@link #createTask(String, String, String, String)}.
+ */
+ default TaskRecord createTask(String taskId, String agentId, String clientId, String input,
+ String contextId) {
+ return createTask(taskId, agentId, clientId, input);
+ }
+
/**
* @return the task record, or {@code null} if no such task
*/
diff --git a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/a2a/A2AGatewayWiringTest.java b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/a2a/A2AGatewayWiringTest.java
new file mode 100644
index 0000000000..3a472341b0
--- /dev/null
+++ b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/a2a/A2AGatewayWiringTest.java
@@ -0,0 +1,233 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.runtime.a2a;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+import org.apache.eventmesh.api.storage.MeshStoragePlugin;
+import org.apache.eventmesh.runtime.boot.EventMeshApplication;
+import org.apache.eventmesh.runtime.offset.InMemoryOffsetStore;
+import org.apache.eventmesh.runtime.state.RocksDBTaskStore;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.net.HttpURLConnection;
+import java.net.URL;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.Collections;
+import java.util.List;
+import java.util.Properties;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+
+/**
+ * Issue #5405 acceptance: the A2A gateway booted WITH the main process —
+ *
+ * - {@code withA2aGateway} wires the REAL transport ({@code EventMeshA2ATransport}
+ * bridged onto the runtime ingress) and a durable {@link RocksDBTaskStore};
+ * - the REST plane answers health, honors the bearer token, and persists the
+ * {@code contextId} through a submit round-trip.
+ *
+ */
+class A2AGatewayWiringTest {
+
+ private EventMeshApplication app;
+ private Path taskDir;
+ private int a2aPort;
+
+ @AfterEach
+ void tearDown() {
+ if (app != null) {
+ try {
+ app.shutdown();
+ } catch (Exception ignored) {
+ // best effort
+ }
+ }
+ }
+
+ private EventMeshApplication boot(boolean withToken) throws Exception {
+ taskDir = Files.createTempDirectory("a2a-wiring-tasks");
+ // port 0 = auto-select: three test methods boot/teardown sequentially, a fixed port
+ // would race the previous teardown's socket release (no SO_REUSEADDR on Windows).
+ EventMeshApplication application = new EventMeshApplication(
+ new NoopStorage(), new InMemoryOffsetStore(), 0, 0);
+ application.withA2aGateway(0,
+ new RocksDBTaskStore(taskDir.toString()), withToken ? "secret-1" : null);
+ application.start();
+ a2aPort = application.a2aGatewayPort();
+ app = application;
+ return application;
+ }
+
+ @Test
+ void gatewayBootsWithMainProcessAndAnswersHealth() throws Exception {
+ boot(false);
+ HttpURLConnection conn = (HttpURLConnection) new URL(
+ "http://127.0.0.1:" + a2aPort + "/a2a/health").openConnection();
+ assertEquals(200, conn.getResponseCode());
+ }
+
+ @Test
+ void tokenRejectsMissingAndWrongBearer() throws Exception {
+ boot(true);
+ // "Connection: close" per request: the Netty pipeline in this test has no keep-alive
+ // handler, so a reused connection can surface as EOF on the follow-up request.
+ assertStatusWithRetry(401, "/a2a/health", null);
+ assertStatusWithRetry(401, "/a2a/tasks/none", "Bearer wrong");
+ assertStatusWithRetry(200, "/a2a/health", "Bearer secret-1");
+ }
+
+ private void assertStatusWithRetry(int expected, String path, String auth) throws Exception {
+ IOException last = null;
+ for (int i = 0; i < 3; i++) {
+ try {
+ HttpURLConnection conn = (HttpURLConnection) new URL(
+ "http://127.0.0.1:" + a2aPort + path).openConnection();
+ conn.setRequestProperty("Connection", "close");
+ if (auth != null) {
+ conn.setRequestProperty("Authorization", auth);
+ }
+ assertEquals(expected, conn.getResponseCode());
+ // 4xx/5xx make getInputStream() throw; the status assert above already passed.
+ try {
+ conn.getInputStream().close();
+ } catch (IOException ok) {
+ // expected for non-2xx statuses
+ }
+ return;
+ } catch (IOException e) {
+ last = e;
+ }
+ }
+ throw last;
+ }
+
+ @Test
+ void submitPersistsContextIdInTaskStore() throws Exception {
+ boot(false);
+ // no agents registered => submit is rejected with a failed future surfaced as 400; the
+ // task is never created. Register an agent card first via the registry.
+ A2AGatewayService svc = app.a2aGatewayService();
+ svc.getAgentCardRegistry().registerCard(
+ org.apache.eventmesh.protocol.a2a.AgentIdentity.builder()
+ .orgId("default").unitId("default").agentId("agent-wiring").build(),
+ buildCard("agent-wiring"));
+
+ String body = "{\"targetAgent\":\"agent-wiring\",\"message\":\"hi\","
+ + "\"contextId\":\"conv-100\",\"sync\":false}";
+ HttpURLConnection conn = (HttpURLConnection) new URL(
+ "http://127.0.0.1:" + a2aPort + "/a2a/tasks").openConnection();
+ conn.setRequestMethod("POST");
+ conn.setDoOutput(true);
+ conn.setRequestProperty("Content-Type", "application/json");
+ try (OutputStream os = conn.getOutputStream()) {
+ os.write(body.getBytes(StandardCharsets.UTF_8));
+ }
+ assertEquals(202, conn.getResponseCode());
+ byte[] resp = conn.getInputStream().readAllBytes();
+ JsonNode node = new ObjectMapper().readTree(resp);
+ String taskId = node.get("taskId").asText();
+
+ // the durable store carries the contextId
+ assertNotNull(svc.getTaskStore().getTask(taskId));
+ assertEquals("conv-100", svc.getTaskStore().getTask(taskId).contextId);
+
+ // and the snapshot echoes it
+ assertEquals("conv-100", svc.getTaskStatus(taskId).getRecord().contextId);
+ }
+
+ private static org.apache.eventmesh.protocol.a2a.model.AgentCard buildCard(String name) {
+ return org.apache.eventmesh.protocol.a2a.model.AgentCard.builder()
+ .name(name)
+ .description("wiring test agent")
+ .version("1.0.0")
+ .supportedInterfaces(List.of(org.apache.eventmesh.protocol.a2a.model.AgentInterface
+ .builder().url("http://127.0.0.1:0/a2a")
+ .protocolBinding("JSONRPC").protocolVersion("0.3").build()))
+ .capabilities(org.apache.eventmesh.protocol.a2a.model.AgentCapabilities.builder()
+ .streaming(false).pushNotifications(false).build())
+ .skills(Collections.emptyList())
+ .defaultInputModes(List.of("text/plain"))
+ .defaultOutputModes(List.of("text/plain"))
+ .build();
+ }
+
+ /** Minimal no-op storage so the runtime boots without a broker. */
+ static final class NoopStorage implements MeshStoragePlugin {
+
+ @Override
+ public void init(Properties props) {
+ // no-op
+ }
+
+ @Override
+ public void send(String topic, org.apache.eventmesh.common.wire.EventMeshFrame frame,
+ org.apache.eventmesh.api.SendCallback callback) {
+ // no-op
+ }
+
+ @Override
+ public List poll(String topic, int partition,
+ long startOffset, int maxEvents, long timeoutMs) {
+ return Collections.emptyList();
+ }
+
+ @Override
+ public void assignPartitions(String topic, List partitions) {
+ // no-op
+ }
+
+ @Override
+ public void commitOffset(String topic, int partition, long offset) {
+ // no-op
+ }
+
+ @Override
+ public int partitionCount(String topic) {
+ return 1;
+ }
+
+ @Override
+ public boolean isStarted() {
+ return true;
+ }
+
+ @Override
+ public boolean isClosed() {
+ return false;
+ }
+
+ @Override
+ public void start() {
+ // no-op
+ }
+
+ @Override
+ public void shutdown() {
+ // no-op
+ }
+ }
+}
diff --git a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/state/RocksDBTaskStoreTest.java b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/state/RocksDBTaskStoreTest.java
new file mode 100644
index 0000000000..b72e91b5ae
--- /dev/null
+++ b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/state/RocksDBTaskStoreTest.java
@@ -0,0 +1,153 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.eventmesh.runtime.state;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import org.apache.eventmesh.runtime.state.TaskStore.Status;
+import org.apache.eventmesh.runtime.state.TaskStore.TaskRecord;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.List;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Issue #5405: local durable {@link RocksDBTaskStore} — CRUD round-trip, epoch CAS rejection,
+ * duplicate-id rejection, contextId persistence, expiry sweep and cross-restart durability
+ * (a fresh store instance on the same dir sees earlier writes).
+ */
+class RocksDBTaskStoreTest {
+
+ private Path dir;
+ private RocksDBTaskStore store;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ dir = Files.createTempDirectory("a2a-rocks-test");
+ store = new RocksDBTaskStore(dir.toString());
+ }
+
+ @AfterEach
+ void tearDown() {
+ if (store != null) {
+ store.close();
+ }
+ }
+
+ @Test
+ void createGetRoundTripWithOpaquePayload() {
+ TaskRecord rec = store.createTask("t-1", "agent-A", "client-X",
+ "café | pipe\nnewline");
+ assertNotNull(rec);
+ assertEquals(Status.PENDING, rec.status);
+
+ TaskRecord loaded = store.getTask("t-1");
+ assertNotNull(loaded);
+ assertEquals("agent-A", loaded.agentId);
+ assertEquals("client-X", loaded.clientId);
+ assertEquals("café | pipe\nnewline", loaded.input);
+ assertNull(loaded.output);
+ }
+
+ @Test
+ void duplicateTaskIdReturnsNull() {
+ assertNotNull(store.createTask("dup", "a", "c", "{}"));
+ assertNull(store.createTask("dup", "a", "c", "{}"));
+ }
+
+ @Test
+ void updateStatusRejectsStaleEpoch() {
+ TaskRecord rec = store.createTask("t-2", "a", "c", "{}");
+ assertNotNull(rec);
+ assertTrue(store.updateStatus("t-2", rec.taskEpoch, Status.RUNNING, null));
+ // stale epoch (the record was written with rec.taskEpoch, now the same epoch but a
+ // different status is fine — CAS is on epoch, not value) — a WRONG epoch must fail:
+ assertNull(store.getTask("nonexistent"));
+ org.junit.jupiter.api.Assertions.assertFalse(
+ store.updateStatus("t-2", rec.taskEpoch + 1, Status.COMPLETED, "out"));
+ // correct epoch completes and stores output
+ assertTrue(store.updateStatus("t-2", rec.taskEpoch, Status.COMPLETED, "out"));
+ assertEquals(Status.COMPLETED, store.getTask("t-2").status);
+ assertEquals("out", store.getTask("t-2").output);
+ }
+
+ @Test
+ void contextIdPersistedAndSurvivesStatusUpdate() {
+ TaskRecord rec = store.createTask("t-ctx", "a", "c", "{}", "conv-42");
+ assertNotNull(rec);
+ assertEquals("conv-42", store.getTask("t-ctx").contextId);
+ assertTrue(store.updateStatus("t-ctx", rec.taskEpoch, Status.COMPLETED, "done"));
+ assertEquals("conv-42", store.getTask("t-ctx").contextId);
+ }
+
+ @Test
+ void expireStaleRemovesOnlyOldTasks() throws Exception {
+ TaskRecord old = store.createTask("old-1", "agent-A", "c", "{}");
+ store.createTask("new-1", "agent-B", "c", "{}");
+ // age the first record by rewiring updatedAt through a direct re-create after closing? —
+ // simpler: create, then sleep the clock forward virtually is not possible; use a large
+ // TTL=0-equivalent: expireStale(-1) evicts updatedAt < now+1s => everything created now
+ // has updatedAt ~now, so a NEGATIVE olderThanMs (deadline = now + |x|) evicts all.
+ List expired = store.expireStale(-60_000L);
+ assertTrue(expired.contains("old-1"));
+ assertTrue(expired.contains("new-1"));
+ assertNull(store.getTask("old-1"));
+ // positive TTL keeps fresh records
+ TaskRecord fresh = store.createTask("fresh-2", "agent-C", "c", "{}");
+ assertEquals(0, store.expireStale(60_000L).size());
+ assertNotNull(store.getTask("fresh-2"));
+ }
+
+ @Test
+ void listByAgentFiltersByStatus() {
+ TaskRecord a1 = store.createTask("l-1", "agent-A", "c", "{}");
+ store.createTask("l-2", "agent-A", "c", "{}");
+ store.createTask("l-3", "agent-B", "c", "{}");
+ store.updateStatus("l-1", a1.taskEpoch, Status.COMPLETED, "x");
+
+ assertEquals(2, store.listByAgent("agent-A", null).size());
+ assertEquals(1, store.listByAgent("agent-A", Status.COMPLETED).size());
+ assertEquals(1, store.listByAgent("agent-B", null).size());
+ }
+
+ @Test
+ void stateSurvivesReopen() {
+ TaskRecord rec = store.createTask("persist-1", "a", "c", "payload", "conv-9");
+ store.updateStatus("persist-1", rec.taskEpoch, Status.COMPLETED, "result");
+ store.flush();
+ store.close();
+
+ RocksDBTaskStore reopened = new RocksDBTaskStore(dir.toString());
+ try {
+ TaskRecord loaded = reopened.getTask("persist-1");
+ assertNotNull(loaded);
+ assertEquals(Status.COMPLETED, loaded.status);
+ assertEquals("result", loaded.output);
+ assertEquals("conv-9", loaded.contextId);
+ } finally {
+ reopened.close();
+ }
+ }
+}