Skip to content

[ISSUE #5411] Serve the legacy SDK gRPC protocol on the v2 runtime (port 10205 bridge) + connector plugin tests 23/23 - #5412

Merged
qqeasonchen merged 6 commits into
apache:developfrom
qqeasonchen:feat/grpc-legacy-bridge
Sep 23, 2026
Merged

qqeasonchen merged 6 commits into
apache:developfrom
qqeasonchen:feat/grpc-legacy-bridge

Conversation

@qqeasonchen

@qqeasonchen qqeasonchen commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

Fixes #5411.

The v2 runtime served HTTP + CloudEvents but kept eventmesh.grpc.port=10205 RESERVED — the legacy SDK gRPC surface (EventMeshGrpcProducer / EventMeshGrpcConsumer, protos in eventmesh-common) was fully preserved yet had nothing to connect to. This PR serves it as a compatibility bridge onto the v2 ingress path (no second messaging engine), and — per the request on the issue — folds in the connector-plugin test work that closes the "4 of 23 plugins carry the only unit tests" gap.

Changes

1. Legacy gRPC bridge (eventmesh-runtime/.../grpc/)

  • EventMeshGrpcServer — gRPC ServerBuilder on eventmesh.grpc.port (1.x default 10205, 0/unset = disabled, opt-in like the WS port), graceful shutdown with the app lifecycle.
    • publish / batchPublish / publishOneWay / batchPublishOneWay → proto CloudEvent mapped and persisted via UniIngressService.publish (topic = proto subject).
    • requestReply → UniIngressService.request, TTL attribute drives the timeout; replies map back to proto.
    • webhook subscribe → WebHookChannel push target + per-topic v2 subscriptions.
    • subscribeStream (bidi) → GrpcStreamChannel push target pumping the v2 dispatcher into the stream; ACKs ride back as stream replies (emdeliveryid echo).
    • heartbeat → GrpcClientRegistry TTL refresh; a reaper evicts (and unsubscribes) stale clients.
  • GrpcCloudEventMapper — pure proto ↔ v2 io.cloudevents.CloudEvent mapping + the legacy Response envelope (statuscode/responsemessage/time, the exact 3 keys the SDK parses).
  • Group semantics: CLUSTERING → LOAD_BALANCE, BROADCASTING → BROADCAST; clientId derived from consumerGroup+env+idc.
  • Wiring: EventMeshApplication.withGrpcBridge(port) + -Deventmesh.grpc.port in main(); runtime gains io.grpc:grpc-netty-shaded/protobuf/stub 1.68.0 (aligned with eventmesh-common).
  • Architecture guard: ruleGrpcProtocolHidden now exempts the sanctioned org.apache.eventmesh.runtime.grpc.. adapter (same precedent the TCP rule grants runtime..), package-info doc updated.

2. Tests

  • GrpcCloudEventMapperTest — hermetic: every proto attribute flavor round-trips; response envelope contract.
  • GrpcLegacyBridgeIntegrationTest — boots the bridge over in-memory storage and drives it with the real legacy SDK (EventMeshGrpcProducer publish / batch publish, EventMeshGrpcConsumer stream subscription round-trip). The SDK is its own conformance suite, per the issue.
  • 19 connector plugins → 38 test classes (23/23 plugins now have tests):
    • webhook sinks (dingtalk/http/knative/lark/slack/wechat/wecom/chatgpt): local capturing HTTP server, body/empty-batch/null-data contracts;
    • mcp (JSON-RPC 2.0 envelope + non-2xx throw), openfunction (Ce-Id/Type/Source headers), prometheus (single merged push), spring (EventForwarder fail-fast/in-order/error-surfacing);
    • webhook sources ×10 (dingtalk/lark/slack/wechat/wecom/http/knative/openfunction/chatgpt/mcp): lazy-bound ephemeral port, native callback payloads → poll() CloudEvents;
    • prometheus source (scrape + unreachable-URL degradation), spring source (ordered buffer drain);
    • external-client plugins (canal/jdbc/mongodb/rabbitmq/redis/s3/pravega): contract-level tests (init default-config parse where lazy, commit no-op) — their backends can't boot in CI.

3. Docs

Verification

  • gradlew test checkstyleMain checkstyleTest — runtime 340/340 (incl. 8 new gRPC tests), architecture-guard green, checkstyle 0 violations across all modules; 18/19 connector plugin modules green (81 plugin tests total).
  • eventmesh-connector-file module tests fail only on the dev machine with Failed to delete temp directory — a local DLP/ACL artifact (verified with a standalone Java probe: even never-opened temp files are delete-denied on this host). These tests are untouched by this PR and pass on CI/Linux.

Follow-ups (not in this PR)

  • Go/Rust SDK verification against the bridge (protos are shared; tracked separately per the issue's non-goals).
  • Stream-ACK correlation for request/reply routed over subscribeStream replies.

…ime (port 10205 bridge)

EventMeshGrpcServer binds PublisherService/ConsumerService/HeartbeatService on
eventmesh.grpc.port (opt-in, 0/-1 = off) and maps every call onto the v2
UniIngressService pipeline (WAL at-least-once, shared retry/DLQ):
- publish/batchPublish/publishOneWay/batchPublishOneWay -> ingress.publish (topic = subject)
- requestReply -> ingress.request with the legacy TTL attribute
- webhook subscribe -> WebHookChannel target; CLUSTERING->LOAD_BALANCE, BROADCASTING->BROADCAST
- subscribeStream (bidi) -> GrpcStreamChannel push target, ACKs ride back on the stream
- heartbeat -> GrpcClientRegistry TTL refresh + reaper unsubscribes stale clients
GrpcCloudEventMapper holds the proto<->v2 CloudEvent mapping + legacy Response envelope.
ruleGrpcProtocolHidden grants the sanctioned runtime.grpc.. adapter (TCP-rule precedent).
…real-SDK integration

GrpcCloudEventMapperTest covers every proto attribute flavor + the 3-key response
envelope. GrpcLegacyBridgeIntegrationTest boots the bridge over the in-memory
storage and drives it with the real EventMeshGrpcProducer/Consumer (the SDK as its
own conformance suite): publish, batch publish, stream subscription round-trip.
Hermetic where possible: webhook sinks (dingtalk/http/knative/lark/slack/wechat/
wecom/chatgpt) verified against a local capturing HTTP server; mcp (JSON-RPC
envelope), openfunction (Ce-* context headers), prometheus (single merged push),
spring (EventForwarder contract) exercised with their exact wire contracts;
webhook sources (10 plugins) POST their native callback payloads through the
lazy-bound hook port; prometheus/spring sources cover scrape + buffer drains;
external-client plugins (canal/jdbc/mongodb/rabbitmq/redis/s3/pravega) get
contract-level tests (init config-parse where lazy, commit no-op) since their
backends cannot boot in CI. Closes the "4 of 23 plugins carry the only unit
tests" gap tracked under the apache#5296 review.
…or Runtime test status

configuration.md flips eventmesh.grpc.port from RESERVED to served (opt-in, 1.x
default 10205); protocols.md gains the legacy-gRPC mapping table section; both
READMEs note the served bridge and update the Connector Runtime row to 23/23
plugin test coverage (drops the stale "4 of 23" wording).
Replace the vague "GA criteria tracked under the apache#5296 architecture review"
wording with the actual remaining GA blockers: per-plugin real-backend
integration tests and the connector-runtime HA story. Also reference apache#5412
for the 23/23 unit-test coverage.
…nd 23/23 tests

- deployment.md: connector-runtime status drops the stale "only 4 of 23
  plugins carry unit tests" (now 23/23 since apache#5412; real-backend ITs stay
  the GA gate)
- architecture/overview.md: protocol-status bullet + legacy-SDK paragraph
  mention the served opt-in gRPC bridge; Key-classes table gains the
  EventMeshGrpcServer entry
- quickstart/getting-started.md: port list mentions the opt-in 10205
  legacy gRPC bridge
@qqeasonchen
qqeasonchen merged commit c10fd2a into apache:develop Sep 23, 2026
9 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Serve the legacy SDK gRPC protocol on the v2 runtime (port 10205 compatibility bridge)

1 participant