From 236ff6ee911a98227f6b2f80cc6667a5f6234af3 Mon Sep 17 00:00:00 2001 From: qqeasonchen Date: Mon, 21 Sep 2026 09:54:46 +0800 Subject: [PATCH 1/2] [ISSUE #5405] A2A gateway production wiring: main-process boot, RocksDB task store, auth, metrics, contextId The A2A gateway was component-complete but never wired for production: EventMeshA2ATransport (the real transport bridged onto UniIngressService) existed unused, EventMeshApplication never booted the gateway, and the only TaskStore backend required Nacos. - EventMeshApplication.withA2aGateway(...): boots the Netty REST plane (/a2a/*) with the main process, wired onto the REAL transport, plus a TaskExpirer reaper. Opt-in via -Deventmesh.a2a.enabled=true, port -Deventmesh.a2a.port=10108, graceful shutdown included. - RocksDBTaskStore: local durable TaskStore under /a2a-tasks (same dir layout as offsets/delivery-state), wire-format compatible with the MetaBacked v1 envelope. Default when no Meta is configured; -Deventmesh.a2a.taskstore=meta selects the cluster-shared store. - Bearer-token auth on the gateway REST plane (-Deventmesh.a2a.token, constant-time compare, same pattern as the admin token guard #5364); without a token the gateway is open (dev mode, logged at boot). - A2A metrics: submitted/completed/failed/canceled/expired counters plus an active-task gauge, surfaced via /admin/metrics JSON and the /metrics Prometheus scrape. - Traffic-port guidance (#5225 class confusion): a JSON-RPC body posted to /events/publish now returns a 400 pointing at the A2A gateway endpoints instead of an opaque CloudEvent transform error. - contextId passthrough (issue #5405): optional conversation linkage on task submit, persisted in TaskRecord (10th wire field, v1 readers ignore the suffix), echoed in snapshots. Tests: RocksDBTaskStoreTest (round-trip, epoch CAS, duplicate rejection, contextId persistence, expiry, reopen durability) and A2AGatewayWiringTest (main-process boot + health, token 401/200, contextId through the REST round-trip into the durable store). --- deploy/kubernetes/connector-configmap.yaml | 26 +- deploy/kubernetes/connector-deployment.yaml | 26 +- deploy/kubernetes/kustomization.yaml | 26 +- deploy/kubernetes/namespace.yaml | 26 +- deploy/kubernetes/runtime-configmap.yaml | 26 +- deploy/kubernetes/runtime-secret.yaml | 26 +- deploy/kubernetes/runtime-service.yaml | 26 +- deploy/kubernetes/runtime-statefulset.yaml | 26 +- docs/feature/a2a.md | 6 +- docs/quickstart/configuration.md | 11 + eventmesh-runtime/bin/start.sh | 6 + .../runtime/a2a/A2AGatewayHttpHandler.java | 57 +++- .../runtime/a2a/A2AGatewayServer.java | 19 ++ .../runtime/a2a/A2AGatewayService.java | 32 ++- .../eventmesh/runtime/a2a/A2AMetrics.java | 88 ++++++ .../runtime/admin/UniAdminServer.java | 13 + .../runtime/boot/EventMeshApplication.java | 92 +++++- .../eventmesh/runtime/http/UniHttpServer.java | 22 ++ .../runtime/state/MetaBackedTaskStore.java | 20 +- .../runtime/state/RocksDBTaskStore.java | 265 ++++++++++++++++++ .../eventmesh/runtime/state/TaskStore.java | 22 ++ .../runtime/a2a/A2AGatewayWiringTest.java | 233 +++++++++++++++ .../runtime/state/RocksDBTaskStoreTest.java | 153 ++++++++++ 23 files changed, 1117 insertions(+), 130 deletions(-) create mode 100644 eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AMetrics.java create mode 100644 eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/state/RocksDBTaskStore.java create mode 100644 eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/a2a/A2AGatewayWiringTest.java create mode 100644 eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/state/RocksDBTaskStoreTest.java 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/quickstart/configuration.md b/docs/quickstart/configuration.md index e459127eea..1a1ef0b209 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,16 @@ 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` | + + - [ ] `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-runtime/bin/start.sh b/eventmesh-runtime/bin/start.sh index cfc58001ed..582a589df3 100644 --- a/eventmesh-runtime/bin/start.sh +++ b/eventmesh-runtime/bin/start.sh @@ -39,6 +39,9 @@ EVENTMESH_STORAGE_TYPE="${EVENTMESH_STORAGE_TYPE:-kafka}" EVENTMESH_HTTP_PORT="${EVENTMESH_HTTP_PORT:-10105}" EVENTMESH_ADMIN_PORT="${EVENTMESH_ADMIN_PORT:-10106}" EVENTMESH_OFFSET_PATH="${EVENTMESH_OFFSET_PATH:-$EVENTMESH_HOME/data/offset}" +EVENTMESH_A2A_ENABLED="${EVENTMESH_A2A_ENABLED:-false}" +EVENTMESH_A2A_PORT="${EVENTMESH_A2A_PORT:-10108}" +EVENTMESH_A2A_TOKEN="${EVENTMESH_A2A_TOKEN:-}" JAVA_OPTS="${JAVA_OPTS:-}" exec java $JAVA_OPTS \ @@ -47,4 +50,7 @@ exec java $JAVA_OPTS \ -Deventmesh.http.port="${EVENTMESH_HTTP_PORT}" \ -Deventmesh.admin.port="${EVENTMESH_ADMIN_PORT}" \ -Deventmesh.offset.path="${EVENTMESH_OFFSET_PATH}" \ + -Deventmesh.a2a.enabled="${EVENTMESH_A2A_ENABLED}" \ + -Deventmesh.a2a.port="${EVENTMESH_A2A_PORT}" \ + -Deventmesh.a2a.token="${EVENTMESH_A2A_TOKEN}" \ org.apache.eventmesh.runtime.boot.EventMeshApplication diff --git a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayHttpHandler.java b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayHttpHandler.java index f7097f300f..a2bdd731ad 100644 --- a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayHttpHandler.java +++ b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/a2a/A2AGatewayHttpHandler.java @@ -60,19 +60,33 @@ * final result.

*/ @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(); + } + } +} From 59df3962637dd8d48e63d9c923f0555f6b4c9619 Mon Sep 17 00:00:00 2001 From: qqeasonchen Date: Mon, 21 Sep 2026 10:35:21 +0800 Subject: [PATCH 2/2] [ISSUE #5405] A2A gateway production wiring + agent-process hardening Gateway plane (items 1-6): main-process boot via withA2aGateway + -Deventmesh.a2a.enabled=true (port 10108) on the REAL transport (EventMeshA2ATransport -> UniIngressService); RocksDBTaskStore local default (/a2a-tasks, MetaBacked-compatible v1 wire format); bearer-token auth (-Deventmesh.a2a.token, constant-time); A2A metrics in /admin/metrics and the /metrics Prometheus scrape; traffic-port guidance for #5225-style JSON-RPC misposts; contextId conversation linkage persisted and echoed. Agent plane (items 7-11): default agent.runtime.url 8080 -> 10105; new bin/start-agent.sh launcher (the dist-agent task referenced a bin/ that never existed); ConversationStore bounded conversation count with access-order LRU eviction (agent.conversation.maxConversations); fail-fast on an empty llm.api.key (llm.api.key.optional=true opts out for mock gateways); heartbeat consecutive-failure limit (agent.heartbeat.failLimit, default 6) exits the process for supervisor restart instead of serving as a zombie; new docs/feature/agent.md. Bug fix found by the new wiring test: A2AGatewayHttpHandler lacked @Sharable - the gateway could only ever serve ONE HTTP connection (the second channel hit ChannelPipelineException; all existing tests used a single request). A2AGatewayServer.getPort() now resolves the auto-selected port (0). Also normalizes the nonstandard license headers on 8 deploy/kubernetes/*.yaml files (#5398 leftovers) so the license check passes. Tests: RocksDBTaskStoreTest (7), A2AGatewayWiringTest (3), agent ConversationStoreTest (+2: LRU eviction, unbounded-compat), existing a2a suite green; checkstyleMain/Test clean; dependency-review green. --- docs/feature/agent.md | 80 +++++++++++++++++++ docs/quickstart/configuration.md | 3 + eventmesh-agent/bin/start-agent.sh | 61 ++++++++++++++ eventmesh-agent/conf/agent.properties | 8 +- .../eventmesh/agent/AgentApplication.java | 30 ++++++- .../eventmesh/agent/ConversationStore.java | 45 ++++++++++- .../agent/ConversationStoreTest.java | 24 ++++++ 7 files changed, 244 insertions(+), 7 deletions(-) create mode 100644 docs/feature/agent.md create mode 100644 eventmesh-agent/bin/start-agent.sh 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 1a1ef0b209..4e5f017f65 100644 --- a/docs/quickstart/configuration.md +++ b/docs/quickstart/configuration.md @@ -158,6 +158,9 @@ Summary: | `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) 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> get(String conversationId) { if (conversationId == null) { return new ArrayList<>(); } + touch(conversationId, false); List> msgs = history.get(conversationId); return msgs == null ? new ArrayList<>() : new ArrayList<>(msgs); } + /** Number of live conversations (visible for tests/metrics). */ + public int conversationCount() { + return history.size(); + } + + /** + * Mark {@code id} most-recently-used and admit it ({@code admit=true}) or only refresh it when + * already live ({@code admit=false} — a get() of a dead id must not resurrect the key and + * push a live conversation out). Evicts the least-recently-used conversation past the bound. + */ + private synchronized void touch(String id, boolean admit) { + if (admit) { + lru.put(id, null); + } else if (!lru.containsKey(id)) { + return; + } else { + lru.get(id); // access-order side effect: refresh recency + } + while (lru.size() > maxConversations) { + java.util.Iterator it = lru.keySet().iterator(); + if (!it.hasNext()) { + break; + } + String eldest = it.next(); + it.remove(); + history.remove(eldest); + } + } + /** Append a completed turn (user prompt + assistant answer). No-op if convId is null. */ public void appendTurn(String conversationId, String userPrompt, String assistantAnswer) { if (conversationId == null) { return; } + touch(conversationId, true); List> msgs = history.computeIfAbsent(conversationId, k -> Collections.synchronizedList(new ArrayList<>())); synchronized (msgs) { diff --git a/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/ConversationStoreTest.java b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/ConversationStoreTest.java index dabb5604d6..fbce41d79e 100644 --- a/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/ConversationStoreTest.java +++ b/eventmesh-agent/src/test/java/org/apache/eventmesh/agent/ConversationStoreTest.java @@ -56,6 +56,30 @@ void getReturnsIndependentSnapshot() { assertThat(store.get("c1")).hasSize(2); // store unaffected } + @Test + void evictsLeastRecentlyUsedConversationPastBound() { + ConversationStore store = new ConversationStore(20, 2); + store.appendTurn("c1", "p1", "a1"); + store.appendTurn("c2", "p2", "a2"); + assertThat(store.conversationCount()).isEqualTo(2); + // touch c1 so c2 becomes the LRU victim + store.get("c1"); + store.appendTurn("c3", "p3", "a3"); + assertThat(store.conversationCount()).isEqualTo(2); + assertThat(store.get("c2")).isEmpty(); // evicted + assertThat(store.get("c1")).hasSize(2); // survived (was touched) + assertThat(store.get("c3")).hasSize(2); // newest + } + + @Test + void singleArgConstructorKeepsUnboundedConversations() { + ConversationStore store = new ConversationStore(20); + for (int i = 0; i < 500; i++) { + store.appendTurn("c" + i, "p", "a"); + } + assertThat(store.conversationCount()).isEqualTo(500); + } + @Test void trimsToSlidingWindow() { ConversationStore store = new ConversationStore(4); // 2 turns kept