Skip to content

fix(router): Event Relays interventions (EventBus, Hermes, relay manager) - #1604

Merged
kkopanidis merged 9 commits into
feat/router-event-relaysfrom
cursor/event-relays-interventions-77ca
Sep 12, 2026
Merged

kkopanidis merged 9 commits into
feat/router-event-relaysfrom
cursor/event-relays-interventions-77ca

Conversation

@kkopanidis

Copy link
Copy Markdown
Contributor

Implements the Conduit backend slice of the Event Relays interventions plan on top of #1600.

Summary

  • EventBus: single Redis message dispatcher, callbacks keyed by subscriberId, SIGTERM/SIGINT quit, subscribe only after Redis ACK; unit tests for churn.
  • Hermes: engine global middleware installed once with namespace derived from handshake URL; no namespace broadcast when rooms and receivers are empty; hot-path logs at debug; recovered /events/ sockets re-run relay authorization; isSocketHandshake matches EIO+transport only (for admin ticket hardening in feat(database): stream live document updates over admin sockets #1601).
  • EventRelayManager: relays indexed by id; coalesced + 30s periodic reconcile; templates compiled at reconcile; 256KiB inbound cap; TTL cached emit-time ReBAC with fail-closed deny → leave-room; room eviction on deactivate/delete/permission or resourceType change; bounded local emit backpressure; metrics gauges/counters.
  • Admin: POST /router/event-relays/preview uses the runtime template renderer.
  • Docs: subscribe-only, dual-publish warning, polling stickiness, EventBus self-publish, room naming (er:<relayId>:<sha256(resourceId)>), no tenant/team in rooms.

Tests added/updated

  • EventBus.test.ts, isSocketHandshake.test.ts, router rebacCache.test.ts, interventions.test.ts (inbound cap + preview parity), existing event-relay unit suite updated for push flags.

Out of scope (follow-ups)

  • Conduit-UI (feat(admin): admin users expansion #327)
  • Integration tests: dual-router HA emit, live /events/ middleware E2E, deactivate relay E2E
  • Wiring isSocketHandshake into admin auth middleware (lands with database admin realtime PR)
Open in Web Open in Cursor 

EventBus uses a single Redis message dispatcher with subscriberId maps,
SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine
middleware once, skips empty-room namespace broadcasts, demotes hot-path
logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager
gains coalesced periodic reconcile, compiled templates, inbound caps,
TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and
metrics/docs aligned with the interventions plan.
Recovery re-auth keeps per-user subscription state across disconnect;
emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process
signal handlers and adds subscribeAck; manager subscribes only after Redis
ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local
sockets without emit acks; cache is bounded; preview caps sample JSON;
handshake matcher rejects ticket/sid/non-namespace paths; socket middleware
rebind avoids duplicate registration on first sockets enable.
Type the Socket.IO engine middleware runner as Express NextFunction so
recursive callbacks match registerGlobalMiddleware. Use Logger.info for
socket trace lines (IConduitLogger has no debug).
_rebindSocketGlobalMiddlewares passes handlers typed with fetch Response
because Response was not imported from express, failing tsc against Hermes
registerSocketGlobalMiddleware.
Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(),
registered on SIGTERM/SIGINT with process exit so Redis teardown is not
left to SDK signal handlers. Tighten registerGlobalMiddleware typing and
trim handshake helper comment (deslop).
Recovery re-auth keeps membership on Authorization UNAVAILABLE and only
leaves on deny; re-check subscriptions for restored er: rooms via room map.
Prune user subscriptions on disconnect when no other socket holds them.
Scope relay emits to receivers in the target room; use Engine.IO
writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs.
Store relay subs on socket.data; resolve recovered rooms from context and
TTL room map without last-writer userId. Prune room map when unused on
disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router
tests for room-scoped emit and queue metric.
@kkopanidis
kkopanidis marked this pull request as ready for review September 12, 2026 23:35
@kkopanidis
kkopanidis merged commit 9c7af8c into feat/router-event-relays Sep 12, 2026
11 of 12 checks passed
@kkopanidis
kkopanidis deleted the cursor/event-relays-interventions-77ca branch September 12, 2026 23:36
kkopanidis added a commit that referenced this pull request Sep 13, 2026
…, recovery) (#1605)

* fix(router): Event Relays interventions (EventBus, Hermes, relay manager) (#1604)

* fix(router): harden event relays for HA and auth lifetime

EventBus uses a single Redis message dispatcher with subscriberId maps,
SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine
middleware once, skips empty-room namespace broadcasts, demotes hot-path
logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager
gains coalesced periodic reconcile, compiled templates, inbound caps,
TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and
metrics/docs aligned with the interventions plan.

* fix(router): address #1604 re-review P1 and P2 follow-ups

Recovery re-auth keeps per-user subscription state across disconnect;
emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process
signal handlers and adds subscribeAck; manager subscribes only after Redis
ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local
sockets without emit acks; cache is bounded; preview caps sample JSON;
handshake matcher rejects ticket/sid/non-namespace paths; socket middleware
rebind avoids duplicate registration on first sockets enable.

* fix(hermes): restore production tsc for engine middleware chain

Type the Socket.IO engine middleware runner as Express NextFunction so
recursive callbacks match registerGlobalMiddleware. Use Logger.info for
socket trace lines (IConduitLogger has no debug).

* fix(router): import Express Response for socket global middleware

_rebindSocketGlobalMiddlewares passes handlers typed with fetch Response
because Response was not imported from express, failing tsc against Hermes
registerSocketGlobalMiddleware.

* fix(router): quit EventBus on module shutdown signals

Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(),
registered on SIGTERM/SIGINT with process exit so Redis teardown is not
left to SDK signal handlers. Tighten registerGlobalMiddleware typing and
trim handshake helper comment (deslop).

* fix(router): pass-3 recovery, scoped emit, and subscription hygiene

Recovery re-auth keeps membership on Authorization UNAVAILABLE and only
leaves on deny; re-check subscriptions for restored er: rooms via room map.
Prune user subscriptions on disconnect when no other socket holds them.
Scope relay emits to receivers in the target room; use Engine.IO
writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs.

* fix(router): pass-4 recovery map, backpressure, emit tests

Store relay subs on socket.data; resolve recovered rooms from context and
TTL room map without last-writer userId. Prune room map when unused on
disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router
tests for room-scoped emit and queue metric.

* fix(router): clear eventRelaySubs from socket.data on recovery deny

* fix(router): clear socket.data subs on recovery fail-closed leave

* fix(router): address CodeFactor findings on event relay validation

Split validateEventRelayInput into per-field parsers to reduce complexity.
Silence unused-parameter lint in EventBus test FakeRedis stub.

* fix(database): Mongo live-update interventions (watch, token, tickets, recovery)

Rebase #1601 onto current Event Relays Hermes and close the stack-review P1s:
watch $match/$project, serialized handleChange, persist resume token after
successful emit, keep token on CursorKilled 237, hermes isSocketHandshake,
restore /database/ Redis on recover, and shutdown the watch plus leader lock.

Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com>

* style: prettier on Mongo intervention tests

* fix(database): stop watch on emit failure; keep ReBAC membership

Failed socketPush/bus no longer advances the resume token. ReBAC
unavailable no longer removeUser. Auth middleware 401s a realtime
ticket on POST /realtime/ticket. Drop/rename/invalidate reopen the watch.

---------
kkopanidis added a commit that referenced this pull request Sep 13, 2026
…1606)

* fix(router): Event Relays interventions (EventBus, Hermes, relay manager) (#1604)

* fix(router): harden event relays for HA and auth lifetime

EventBus uses a single Redis message dispatcher with subscriberId maps,
SIGTERM/SIGINT shutdown, and subscribe-after-ACK. Hermes installs engine
middleware once, skips empty-room namespace broadcasts, demotes hot-path
logs, and re-authenticates recovered /events/ subscriptions. EventRelayManager
gains coalesced periodic reconcile, compiled templates, inbound caps,
TTL-cached emit-time ReBAC, room eviction, backpressure, preview API, and
metrics/docs aligned with the interventions plan.

* fix(router): address #1604 re-review P1 and P2 follow-ups

Recovery re-auth keeps per-user subscription state across disconnect;
emit-time ReBAC distinguishes unavailable vs deny; EventBus drops process
signal handlers and adds subscribeAck; manager subscribes only after Redis
ACK and broadcasts evictRelayIds on refresh; Hermes backpressure uses local
sockets without emit acks; cache is bounded; preview caps sample JSON;
handshake matcher rejects ticket/sid/non-namespace paths; socket middleware
rebind avoids duplicate registration on first sockets enable.

* fix(hermes): restore production tsc for engine middleware chain

Type the Socket.IO engine middleware runner as Express NextFunction so
recursive callbacks match registerGlobalMiddleware. Use Logger.info for
socket trace lines (IConduitLogger has no debug).

* fix(router): import Express Response for socket global middleware

_rebindSocketGlobalMiddlewares passes handlers typed with fetch Response
because Response was not imported from express, failing tsc against Hermes
registerSocketGlobalMiddleware.

* fix(router): quit EventBus on module shutdown signals

Stop event relays and call grpcSdk.bus.quit() from Router.shutdown(),
registered on SIGTERM/SIGINT with process exit so Redis teardown is not
left to SDK signal handlers. Tighten registerGlobalMiddleware typing and
trim handshake helper comment (deslop).

* fix(router): pass-3 recovery, scoped emit, and subscription hygiene

Recovery re-auth keeps membership on Authorization UNAVAILABLE and only
leaves on deny; re-check subscriptions for restored er: rooms via room map.
Prune user subscriptions on disconnect when no other socket holds them.
Scope relay emits to receivers in the target room; use Engine.IO
writeBuffer/writable for backpressure. Drop hot-path Hermes socket info logs.

* fix(router): pass-4 recovery map, backpressure, emit tests

Store relay subs on socket.data; resolve recovered rooms from context and
TTL room map without last-writer userId. Prune room map when unused on
disconnect. Backpressure uses writeBuffer depth only. Add Hermes/router
tests for room-scoped emit and queue metric.

* fix(router): clear eventRelaySubs from socket.data on recovery deny

* fix(router): clear socket.data subs on recovery fail-closed leave

* fix(router): address CodeFactor findings on event relay validation

Split validateEventRelayInput into per-field parsers to reduce complexity.
Silence unused-parameter lint in EventBus test FakeRedis stub.

* feat(database): restack Mongo live updates onto current Event Relays

Replay the Mongo delta onto feat/router-event-relays (9e6356e). Keep
#1604 Hermes (engine.use, er: recovery, writeBuffer) and graft only the
port-keyed adapter plus register-before-initSockets. Do not take old
#1601 Socket.ts wholesale.

Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com>

* fix(database): fence change-stream leadership and expire recovery Redis keys

Lock renew failure bumps a generation so a draining watch cannot emit.
Do not open a watch unless the leader lock extends after acquire.
Arm TTL on realtime socket/doc keys when a recoverable disconnect
never recovers; persist on restore. CMS-read deny fails closed at
emit without extra ReBAC round-trips and keeps membership; UNAVAILABLE
keeps membership. Cover dropDatabase reopen.

Co-authored-by: Konstantinos Kopanidis <kkopanidis@users.noreply.github.com>

* fix(database): serialize leader acquire so concurrent reconcile still watches

Two overlapping ensureLeader calls both bumped the lock generation while
the first openStream was still reading the resume token, so the watch
never attached. Hold an acquiring flag so only one opener runs.

---------
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.

2 participants