Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ status. See [docs](docs/) for the per-capability guides.
| [Kafka / RocketMQ storage](eventmesh-storage-plugin/) (4.x, 5.x) | **GA target** | Recommended — pluggable WAL backends, TCK-covered (`MeshStoragePluginTCK`) | Primary path |
| [Memory storage](eventmesh-storage-plugin/eventmesh-storage-memory/) (default) | **Beta** | Zero-dependency dev/CI/quick-start backend (`docker run apache/eventmesh` with no broker); state is process-local — not for production | Switch `EVENTMESH_STORAGE_TYPE` to kafka / rocketmq / rocketmq5 |
| SSE / WebSocket push | **Beta** | Usable — integration-tested; ACK-tracked redelivery + DLQ are shared with long-polling (same `ReliableDispatcher`), e2e-real-broker suite in [#5389](https://github.com/apache/eventmesh/pull/5389) | Unified push transports |
| Connector Runtime | **Beta** | Usable — working end-to-end; all 23 plugins fully implemented since [#5394](https://github.com/apache/eventmesh/pull/5394) (no template stubs left), data-loss hardening in [#5328](https://github.com/apache/eventmesh/pull/5328); 4 of 23 (file/kafka/pulsar/rocketmq) still carry the only unit tests | Remaining plugin tests + GA criteria tracked under the #5296 architecture review |
| Connector Runtime | **Beta** | Usable — working end-to-end; all 23 plugins fully implemented since [#5394](https://github.com/apache/eventmesh/pull/5394) (no template stubs left), data-loss hardening in [#5328](https://github.com/apache/eventmesh/pull/5328); unit tests cover all 23 plugins since [#5412](https://github.com/apache/eventmesh/pull/5412) | Promote to GA once the remaining #5296 review items land (real-backend integration tests per plugin, connector-runtime HA story) |
| [A2A / Agent Gateway](docs/feature/a2a.md) | **Beta** | Usable — task store + runtime bridge (#5302/#5304), task reaper + Meta-backed agent cards in [#5346](https://github.com/apache/eventmesh/pull/5346), quota classification in [#5373](https://github.com/apache/eventmesh/pull/5373); the Testcontainers E2E gate (#5340) closed in September 2026 | Unified Runtime A2A |
| [Agent tools & event triggers](docs/feature/agent-tools.md) | **Experimental** | Evaluate — `LlmClient`/`ConversationMemory`/`AgentTool` extension points, connector-backed tools and SPI plugin deployment landed in [#5408](https://github.com/apache/eventmesh/pull/5408)/[#5409](https://github.com/apache/eventmesh/pull/5409); runtime in `eventmesh-agent-runtime`, plugins in `eventmesh-agent-plugin/` | Agent tool ecosystem |
| TCP / gRPC / OpenMessaging SDKs | **Legacy-compatible** | Existing users only — kept so old clients run unmodified; not extended | [HTTP + CloudEvents](docs/feature/client-java.md) |
Expand All @@ -81,7 +81,10 @@ Status meanings:

> Migrating off TCP / gRPC SDKs? The legacy clients keep working against the current
> runtime; see the [client guide](docs/feature/client-java.md) for the
> HTTP + CloudEvents replacement (`CloudEventsClient`).
> HTTP + CloudEvents replacement (`CloudEventsClient`). The legacy gRPC SDK
> surface is served again since the #5411 bridge (opt-in via
> `eventmesh.grpc.port=10205`) — old `EventMeshGrpcProducer` /
> `EventMeshGrpcConsumer` clients run unmodified on the v2 runtime.

### Documentation

Expand Down
8 changes: 6 additions & 2 deletions README.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ Apache EventMesh 提供了丰富的能力,帮助用户轻松构建事件驱动
| [Kafka / RocketMQ 存储](eventmesh-storage-plugin/)(4.x、5.x) | **GA 目标** | 推荐——可插拔 WAL 后端,TCK 覆盖(`MeshStoragePluginTCK`) | 主路径 |
| [Memory 内存存储](eventmesh-storage-plugin/eventmesh-storage-memory/)(默认) | **Beta** | 零依赖的开发/CI/快速上手后端(`docker run apache/eventmesh` 无需 broker);状态仅存于进程内——不可用于生产 | 切换 `EVENTMESH_STORAGE_TYPE` 为 kafka / rocketmq / rocketmq5 |
| SSE / WebSocket 推送 | **Beta** | 可用——已有集成测试;ACK 追踪的重投递与 DLQ 已与长轮询共享同一 `ReliableDispatcher`,真实 broker 的 e2e 套件见 [#5389](https://github.com/apache/eventmesh/pull/5389) | 统一推送传输 |
| Connector Runtime | **Beta** | 可用——端到端工作;自 [#5394](https://github.com/apache/eventmesh/pull/5394) 起 23 个插件全部完整实现(不再有模板桩),数据丢失加固见 [#5328](https://github.com/apache/eventmesh/pull/5328);单元测试仍仅覆盖其中 4 个(file/kafka/pulsar/rocketmq) | 剩余插件测试与 GA 标准由 #5296 架构 review 跟踪 |
| Connector Runtime | **Beta** | 可用——端到端工作;自 [#5394](https://github.com/apache/eventmesh/pull/5394) 起 23 个插件全部完整实现(不再有模板桩),数据丢失加固见 [#5328](https://github.com/apache/eventmesh/pull/5328);自 [#5412](https://github.com/apache/eventmesh/pull/5412) 起单元测试覆盖全部 23 个插件 | 待 #5296 review 剩余项落地后升级 GA(各插件真实后端集成测试、connector-runtime HA 方案) |
| [A2A / Agent 网关](docs/feature/a2a.md) | **Beta** | 可用——TaskStore + Runtime 桥(#5302/#5304)、任务 reaper 与 Meta 化 AgentCard([#5346](https://github.com/apache/eventmesh/pull/5346))、配额分类([#5373](https://github.com/apache/eventmesh/pull/5373))均已落地;Testcontainers E2E 门槛(#5340)已于 2026 年 9 月关闭 | 统一 Runtime A2A |
| [智能体工具与事件触发](docs/feature/agent-tools.md) | **实验性** | 评估——`LlmClient`/`ConversationMemory`/`AgentTool` 扩展点、connector 工具化与 SPI 插件化部署已落地([#5408](https://github.com/apache/eventmesh/pull/5408)/[#5409](https://github.com/apache/eventmesh/pull/5409));运行时在 `eventmesh-agent-runtime`,插件在 `eventmesh-agent-plugin/` | 智能体工具生态 |
| TCP / gRPC / OpenMessaging SDK | **Legacy 兼容** | 仅存量用户——保持老客户端零改动运行;不再扩展 | [HTTP + CloudEvents](docs/feature/client-java.md) |
Expand All @@ -80,7 +80,11 @@ Apache EventMesh 提供了丰富的能力,帮助用户轻松构建事件驱动
- **Legacy 兼容** —— 仅为存量客户端零改动兼容而维护;只修缺陷、不加功能。新接入不要选这里。

> 正在从 TCP / gRPC SDK 迁移?Legacy 客户端在当前 Runtime 上继续可用;替代方案(HTTP +
> CloudEvents 的 `CloudEventsClient`)见[客户端指引](docs/feature/client-java.md)。
> CloudEvents 的 `CloudEventsClient`)

> #5411 之后 legacy gRPC SDK 接入面重新可用(通过 `eventmesh.grpc.port=10205`
> 显式开启)——旧的 `EventMeshGrpcProducer` / `EventMeshGrpcConsumer` 客户端
> 可以零改动运行在 v2 Runtime 上。见[客户端指引](docs/feature/client-java.md)。

### 文档导航

Expand Down
8 changes: 6 additions & 2 deletions docs/architecture/overview.md
Original file line number Diff line number Diff line change
Expand Up @@ -114,14 +114,16 @@ request-reply).
The legacy TCP / gRPC / OpenMessaging SDKs still work because the new
`EventMeshFrame` adaptor (`eventmesh-runtime/.../protocol/meshmessage`) and the
`UniTcpServer` (public/internal split per `#5297`) preserve the wire format and
semantics of the old `MeshMessage` / `OpenMessage` clients.
semantics of the old `MeshMessage` / `OpenMessage` clients; the legacy gRPC
SDK family is served by the opt-in `EventMeshGrpcServer` bridge (#5411).

### Key classes

| Component | File | Role |
| --- | --- | --- |
| HTTP entry | `eventmesh-runtime/.../http/UniHttpServer.java` | Netty HTTP/S; entry of all `/events/*` + `/a2a/*` traffic; `withSecurityGate(...)` wiring point |
| WebSocket entry | `eventmesh-runtime/.../http/UniWsServer.java` | WebSocket transport for subscribers |
| Legacy gRPC bridge | `eventmesh-runtime/.../grpc/EventMeshGrpcServer.java` | Opt-in port-10205 bridge serving the 1.x `PublisherService`/`ConsumerService`/`HeartbeatService` (#5411) |
| Ingress orchestrator | `eventmesh-runtime/.../ingress/UniIngressService.java` | Frame-typed facade; single protocol path (`EventMeshFrame`) used by both HTTP and the legacy TCP adaptor |
| Security gate | `eventmesh-runtime/.../security/gate/SecurityGate.java` | Opt-in unified gate (see §4) |
| Filter chain | `eventmesh-runtime/.../security/FilterChain.java` | Auth + ACL filters executed before the gate |
Expand Down Expand Up @@ -418,7 +420,9 @@ summarized there:
HTTP codec are still supported for backward compatibility but
are explicitly marked Legacy and receive only critical bug
fixes. See `docs/feature/protocols.md` for the migration path.
* gRPC framing is **Beta**: stable, but the API surface may shift.
* gRPC framing is **Beta**: stable, but the API surface may shift. The
legacy SDK gRPC surface is served again as an opt-in compatibility bridge
on `eventmesh.grpc.port` since #5411 (see `docs/feature/protocols.md` §1.1).
* A2A is **Experimental** until the readiness checks in #5340
pass.

Expand Down
6 changes: 4 additions & 2 deletions docs/feature/deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,8 +120,10 @@ Connector definitions are managed through `/admin/connectors` (CRUD) and
scheduled by the runtime's `ConnectorScheduler` with **generation
fencing** (#5382): every (re)assignment bumps a per-connector generation,
and a stale delayed start can never take over a connector running a newer
generation. Status: **Experimental** — only 4 of 23 plugins carry unit
tests today.
generation. Status: **Experimental** — all 23 plugins carry unit tests
since [#5412](https://github.com/apache/eventmesh/pull/5412) (hermetic
webhook tests + contract tests for external-client plugins); real-backend
integration tests remain the GA gate.

## Production checklist

Expand Down
29 changes: 28 additions & 1 deletion docs/feature/protocols.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,14 +20,41 @@ public surface.
| EventMeshFrame (internal) | **GA** | `org.apache.eventmesh.common.wire.EventMeshFrame` (Frame architecture) | runtime internal; producer / storage / push path | - |
| A2A (Agent-to-Agent) | **Experimental** | A2A JSON-RPC + SSE | `eventmesh-protocol-plugin/eventmesh-protocol-a2a`, A2A gateway on Runtime | - |
| MeshMessage TCP | **Legacy** | length-prefixed `MeshMessage` bytes | runtime `tcp/` subpackage, `eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/resolver/tcp` | CloudEvents HTTP, or A2A (for agent workloads) |
| gRPC (CloudEvents + EventMeshMessage) | **Beta** | protobuf over HTTP/2 | runtime `transport/grpc`, `eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/resolver/grpc` | CloudEvents HTTP for new clients |
| gRPC (CloudEvents + EventMeshMessage) | **Beta** | protobuf over HTTP/2 | runtime `grpc/` bridge (issue #5411), `eventmesh-protocol-plugin/eventmesh-protocol-meshmessage/resolver/grpc` | CloudEvents HTTP for new clients |
| OpenMessaging API (TCP) | **Legacy** | OMA spec, used by the legacy TCP client only | `eventmesh-sdks/eventmesh-sdk-java/client/tcp/impl/openmessage` | CloudEvents HTTP client |

> **GA** = production-ready and the recommended path. **Beta** = stable but
> the API surface may still shift. **Experimental** = subject to breaking
> change without notice. **Legacy** = still works, but is no longer the
> recommended choice and is being phased out.

### 1.1 Legacy gRPC bridge (served, opt-in — issue #5411)

The v2 runtime serves the legacy SDK gRPC protocol as a **compatibility
bridge** on `eventmesh.grpc.port` (1.x default `10205`; unset / `-1` =
disabled, opt-in like the WS port). There is no second messaging engine:
every call maps onto the same v2 ingress/delivery pipeline the HTTP plane
uses (WAL durability, at-least-once, shared retry/DLQ).

| Legacy call | v2 mapping |
| --- | --- |
| `publish` / `batchPublish` / `publishOneWay` / `batchPublishOneWay` | proto `CloudEvent` → v2 CloudEvent, persisted via `UniIngressService.publish` (topic = proto `subject` attribute) |
| `requestReply` | v2 request/reply correlation (`UniIngressService.request`), TTL attribute drives the timeout |
| `subscribe` (webhook `url`) | `WebHookChannel` push target + v2 subscription per topic |
| `subscribeStream` (bidi) | `GrpcStreamChannel` push target pumping the v2 dispatcher into the stream; ACKs ride back as stream replies |
| `unsubscribe` | v2 unsubscribe per topic + client deregistration |
| `heartbeat` | TTL refresh in the `GrpcClientRegistry`; a reaper evicts stale clients (unsubscribes them) |

Group semantics: the legacy `consumerGroup` maps onto a v2 subscription
group; the SDK's default `CLUSTERING` mode maps to `LOAD_BALANCE`
distribution, `BROADCASTING` maps to `BROADCAST`. The clientId is derived
from `consumerGroup` + `env` + `idc` (the 1.x triple).

Non-goals (follow-ups): a new gRPC-native v2 API (HTTP + CloudEvents stays
the primary path), the gRPC admin surface (stays unimplemented as on 1.x),
and Go/Rust SDK verification (the protos are shared; Java SDK is the
conformance suite — see `GrpcLegacyBridgeIntegrationTest`).

## 2. Server-side protocol plugins

The runtime discovers protocol adaptors via the
Expand Down
2 changes: 1 addition & 1 deletion docs/quickstart/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ Usually set via `-D` by `bin/start.sh`; override here if needed.
| `eventmesh.http.port` | `10105` | Traffic HTTP (`/events/*`, `/agent/*`, `/session/*`) |
| `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.grpc.port` | `-1` (disabled) | Legacy SDK gRPC bridge (Publisher/Consumer/Heartbeat services, issue #5411); set `10205` (the 1.x default) to serve old `EventMeshGrpcProducer`/`EventMeshGrpcConsumer` clients on the v2 runtime |
| `eventmesh.a2a.port` | `10108` | A2A gateway REST plane (opt-in via `eventmesh.a2a.enabled=true`) |
| `eventmesh.offset.path` | `./data/offset` | Local offset store directory |

Expand Down
3 changes: 2 additions & 1 deletion docs/quickstart/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,8 @@ dependency, ready for a smoke test as-is. For a real broker add
backend address keys shown below.

Ports: `10105` = traffic HTTP (`/events/*`), `10106` = admin HTTP (`/admin/*`). The WebSocket
push port (`10107`) is opt-in.
push port (`10107`) is opt-in, and the legacy SDK gRPC bridge (`10205`, for old
`EventMeshGrpcProducer`/`EventMeshGrpcConsumer` clients) is opt-in since #5411.

### Option B — From source

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,9 @@ public static JavaClasses loadProductionClasses() {
.and().resideOutsideOfPackage("org.apache.eventmesh.common..")
.and().resideOutsideOfPackage("org.apache.eventmesh.protocol.meshmessage..")
.and().resideOutsideOfPackage("org.apache.eventmesh.client..")
// #5411: the legacy gRPC bridge (runtime.grpc..) is the sanctioned server-side
// adapter for these types — same carve-out the TCP rule grants runtime..
.and().resideOutsideOfPackage("org.apache.eventmesh.runtime.grpc..")
.should().dependOnClassesThat().resideInAPackage("org.apache.eventmesh.common.protocol.grpc..");

public static ArchRule ruleTcpProtocolHidden = noClasses()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,9 @@

* <p><b>Forbidden (this package must NOT be imported by):</b>
* <ul>
* <li>{@code ANY module other than eventmesh-sdks:eventmesh-sdk-java and}
* {@code eventmesh-protocol-plugin:eventmesh-protocol-meshmessage (ruleGrpcProtocolHidden)}</li>
* <li>{@code ANY module other than eventmesh-sdks:eventmesh-sdk-java,}
* {@code eventmesh-protocol-plugin:eventmesh-protocol-meshmessage and the runtime gRPC bridge}
* {@code (org.apache.eventmesh.runtime.grpc.., issue #5411) (ruleGrpcProtocolHidden)}</li>
* </ul>
*/
package org.apache.eventmesh.common.protocol.grpc;
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
/*
* 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.connector.canal.sink;

import java.util.Collections;

import org.junit.jupiter.api.Test;

import io.cloudevents.CloudEvent;
import io.cloudevents.core.builder.CloudEventBuilder;

/**
* Contract-level unit test for CanalSinkConnector: the external system is NOT available in CI, so this covers
* the parts that must hold regardless — instantiation, the documented default-config surface of
* init() (where init only parses config), and the commit() no-op contract.
*/
class CanalSinkConnectorTest {

@Test
void commitIsNoOp() {
CanalSinkConnector connector = new CanalSinkConnector();
CloudEvent last = CloudEventBuilder.v1().withId("e1")
.withSource(java.net.URI.create("/test")).withType("test.event").build();
connector.commit(Collections.singletonList(last)); // must not throw
}

}
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
/*
* 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.connector.canal.source;

import java.util.Properties;

import org.junit.jupiter.api.Test;

import io.cloudevents.CloudEvent;
import io.cloudevents.core.builder.CloudEventBuilder;

/**
* Contract-level unit test for CanalSourceConnector: the external system is NOT available in CI, so this covers
* the parts that must hold regardless — instantiation, the documented default-config surface of
* init() (where init only parses config), and the commit() no-op contract.
*/
class CanalSourceConnectorTest {

@Test
void commitIsNoOp() {
CanalSourceConnector connector = new CanalSourceConnector();
CloudEvent last = CloudEventBuilder.v1().withId("e1")
.withSource(java.net.URI.create("/test")).withType("test.event").build();
connector.commit(last); // must not throw
}

@Test
void initParsesDocumentedDefaults() {
CanalSourceConnector connector = new CanalSourceConnector();
connector.init(new Properties()); // must not connect anything
}

}
Loading
Loading