Skip to content

Commit 239e98d

Browse files
committed
Document WebSocket relay HA, NATS vs AMQP backend selection, and operator troubleshooting.
Update architecture, ws-sessions guide, swagger WS/NATS creds paths, and changelog for relay production and NATS platform account naming.
1 parent d6d3dab commit 239e98d

3 files changed

Lines changed: 171 additions & 31 deletions

File tree

docs/architecture.md

Lines changed: 44 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -206,15 +206,45 @@ Full spec: [`.cursor/controllerv3.8/docs/15-fog-platform-reconcile.md`](../.curs
206206

207207
## WebSocket exec & log sessions
208208

209-
Interactive **exec** and **log streaming** use paired WebSocket sessions between operators (Bearer JWT), Controller, and Edgelet agents (fog token). Plan 16 hardens log sessions and shared WS infra (HA, drain, OTEL). **Plan 17** redesigns **microservice exec** to log-style multi-session flow (3 concurrent per MS, agent poll + session-scoped WS) — **Edgelet agent wire change required** for exec (see [edgelet-invariants.md §10.1](../.cursor/controllerv3.8/docs/edgelet-invariants.md)).
209+
Interactive **exec** and **log streaming** use paired WebSocket sessions between operators (Bearer JWT), Controller, and Edgelet agents (fog token). Plan 16 hardens log sessions and shared WS infra (HA, drain, OTEL). **Plan 17** redesigns **microservice exec** to log-style multi-session flow (3 concurrent per MS, agent poll + session-scoped WS). **Plan 18** production-hardens cross-replica relay via **`WsRelayTransport`** — AMQP pool + recovery when `nats.enabled=false`, NATS Core when `nats.enabled=true` (R102–R113). **Edgelet agent wire change required** for exec only (see [edgelet-invariants.md §10.1](../.cursor/controllerv3.8/docs/edgelet-invariants.md)).
210+
211+
```mermaid
212+
flowchart TB
213+
subgraph ws [WebSocketServer]
214+
U[User WS]
215+
A[Agent WS]
216+
end
217+
218+
subgraph factory [WsRelayTransportFactory]
219+
SEL{nats.enabled?}
220+
end
221+
222+
subgraph amqp [AmqpRelayTransport]
223+
POOL[RouterConnectionManager pool x8]
224+
Q[agent/user queues per sessionId]
225+
end
226+
227+
subgraph nats [NatsRelayTransport]
228+
NC[NatsRelayConnectionManager]
229+
SUB[Core pub/sub per sessionId]
230+
end
231+
232+
U --> ws
233+
A --> ws
234+
ws --> factory
235+
SEL -->|false| amqp
236+
SEL -->|true| nats
237+
POOL --> Q
238+
NC --> SUB
239+
```
210240

211241
```mermaid
212242
sequenceDiagram
213243
participant U as User WS
214244
participant C as Controller
215245
participant DB as MicroserviceExecSessions
216246
participant CT as Change tracking
217-
participant Q as AMQP
247+
participant R as WsRelayTransport
218248
participant A as Edgelet agent
219249
220250
Note over U,A: MS exec (R92–R98) — log-style
@@ -228,11 +258,11 @@ sequenceDiagram
228258
A->>C: GET /agent/exec/sessions
229259
A->>C: WS /agent/exec/microservice/:uuid/:sessionId
230260
C->>DB: ACTIVE agentConnected
231-
C->>Q: agent-{sessionId} user-{sessionId}
261+
C->>R: enable bridge if cross-replica
232262
U->>C: STDIN
233-
C->>Q->>A: relay
263+
C->>R->>A: relay
234264
A->>C: STDOUT
235-
C->>Q->>U: relay
265+
C->>R->>U: relay
236266
U-->>C: close
237267
C->>DB: DELETE row
238268
C->>CT: execSessions if needed
@@ -242,12 +272,12 @@ sequenceDiagram
242272
C->>C: provision debug system MS
243273
U->>C: WS /microservices/system/exec/:debugMsUuid
244274
245-
Note over U,A: Logs (R82–R83, R84) — unchanged Plan 16
275+
Note over U,A: Logs (R82–R83, R84/R112)
246276
U->>C: WS logs + tail params
247277
C->>DB: PENDING sessionId
248278
A->>C: WS agent/logs/:sessionId
249279
A->>C: LOG_LINE
250-
C->>Q->>U: logs-user-{sessionId}
280+
C->>R->>U: relay
251281
```
252282

253283
| Topic | Normative value |
@@ -263,21 +293,23 @@ sequenceDiagram
263293
| Log concurrency | **3** user log WS per microservice (or per fog for node logs) |
264294
| Log limits | Tail max **5,000** lines; **120s** pending; **2h** idle |
265295
| Log content | Live relay only — no log line persistence; audit connect/disconnect |
266-
| HA relay | Cross-replica sessions **require** AMQP (`WebSocketQueueService`); same-replica may use direct WS; **fail fast** when router down |
296+
| HA relay | Cross-replica sessions **require** a **relay backend** (R112): **AMQP** router queues when `nats.enabled=false`; **NATS Core** subjects on hub when `nats.enabled=true`. Same-replica may use direct WS; **fail fast** close **1013** when active backend unavailable |
267297
| Graceful drain | **30s** on SIGTERM / k8s `preStop` — CLOSE frames, queue cleanup, session row delete |
268298
| Security | Agent handlers validate fog token **before** message processing; **50** upgrades/min/IP; **100** active WS/IP; JWT in `?token=` (ingress log redaction required) |
269299
| Scale SLO | **500** concurrent WS per replica; **p99 pairing < 5s** |
270300
| Observability | OpenTelemetry: active/pending sessions, pairing latency, AMQP failures, router connectivity |
271301

272302
**OTEL metric names (R87):** `ws_exec_sessions_active`, `ws_log_sessions_active`, `ws_pending_pairings`, `ws_pairing_duration_ms` (histogram), `ws_amqp_publish_errors`, `ws_router_connected` (gauge). Emitted when `ENABLE_TELEMETRY=true`; see `src/websocket/ws-metrics.js`.
273303

274-
**HA config (`server.webSocket.ha`):** `crossReplicaRequiresAmqp` (default `true`), `failFastOnRouterUnavailable` (default `true`). Env: `WS_HA_CROSS_REPLICA_REQUIRES_AMQP`, `WS_HA_FAIL_FAST_ON_ROUTER_UNAVAILABLE`. Graceful drain timeout: `server.webSocket.session.drainTimeoutMs` (default **30s**, env `WS_DRAIN_TIMEOUT_MS`).
304+
**Relay transport (Plan 18, R102):** Selected once at startup from existing `nats.enabled` / `NATS_ENABLED`**no new relay env var**. `false` → AMQP pool (8 connections, sticky by `sessionId`); `true` → NATS Core on hub with **`controller`** NATS account. Relay connect is **lazy** — does not block Controller startup.
305+
306+
**HA config (`server.webSocket.ha`):** `crossReplicaRequiresAmqp` (default `true`; semantics: cross-replica requires **active relay backend** per R112), `failFastOnRouterUnavailable` (default `true`; applies to selected backend). Env: `WS_HA_CROSS_REPLICA_REQUIRES_AMQP`, `WS_HA_FAIL_FAST_ON_ROUTER_UNAVAILABLE`. Graceful drain timeout: `server.webSocket.session.drainTimeoutMs` (default **30s**, env `WS_DRAIN_TIMEOUT_MS`).
275307

276-
**Core modules:** `src/websocket/server.js`, `exec-session-manager.js` (Plan 17), `log-session-manager.js`, `src/services/websocket-queue-service.js`, `src/services/router-connection-service.js`.
308+
**Core modules:** `src/websocket/server.js`, `exec-session-manager.js` (Plan 17), `log-session-manager.js`, `ws-relay-transport-factory.js` / `amqp-relay-transport.js` / `nats-relay-transport.js` (Plan 18), `src/services/websocket-queue-service.js`, `src/services/router-connection-manager.js` (Plan 18), `src/services/nats-relay-connection-manager.js` (Plan 18).
277309

278-
**Operator guide:** [operations/ws-sessions.md](operations/ws-sessions.md) — ingress `?token=` log redaction, HTTPS/WSS, multi-replica AMQP requirement, k8s preStop drain, load SLO probe.
310+
**Operator guide:** [operations/ws-sessions.md](operations/ws-sessions.md) — ingress `?token=` log redaction, HTTPS/WSS, multi-replica relay backend (`nats.enabled`), k8s preStop drain, load SLO probe.
279311

280-
Full spec: Plan 16 [logs + shared infra](../.cursor/controllerv3.8/docs/16-ws-exec-log-hardening.md) · Plan 17 [MS exec](../.cursor/controllerv3.8/docs/17-multi-exec-sessions.md) · RFC R80–R91, R92–R101 · Edgelet contract: [edgelet-invariants.md §10–§10.1](../.cursor/controllerv3.8/docs/edgelet-invariants.md).
312+
Full spec: Plan 16 [logs + shared infra](../.cursor/controllerv3.8/docs/16-ws-exec-log-hardening.md) · Plan 17 [MS exec](../.cursor/controllerv3.8/docs/17-multi-exec-sessions.md) · Plan 18 [WS relay production](../.cursor/controllerv3.8/docs/18-ws-relay-production.md) · RFC R80–R91, R92–R101, R102–R113 · Edgelet contract: [edgelet-invariants.md §10–§10.1](../.cursor/controllerv3.8/docs/edgelet-invariants.md).
281313

282314
---
283315

docs/operations/ws-sessions.md

Lines changed: 43 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66

77
## Overview
88

9-
Controller exposes **interactive exec** and **log streaming** over WebSocket on the API port (default **51121**). Sessions pair an operator browser/CLI client (Bearer JWT) with an Edgelet agent (fog token). In multi-replica deployments, cross-replica relay requires the **Skupper-style AMQP router** microservice.
9+
Controller exposes **interactive exec** and **log streaming** over WebSocket on the API port (default **51121**). Sessions pair an operator browser/CLI client (Bearer JWT) with an Edgelet agent (fog token). In multi-replica deployments, cross-replica relay uses a **relay backend** selected at startup by **`nats.enabled`** (Plan 18, R102): **AMQP** router queues when `false`, **NATS Core** on the platform hub when `true`.
1010

1111
---
1212

@@ -38,18 +38,41 @@ Without redaction, long-lived bearer tokens may appear in load balancer logs.
3838

3939
## Multi-replica HA
4040

41+
Relay transport is selected **once at startup** from existing platform config — **no separate relay env var** (R102):
42+
43+
| `nats.enabled` | Cross-replica relay backend |
44+
|----------------|----------------------------|
45+
| `false` (default) | **AMQP** — Skupper-style router queues via `WebSocketQueueService` |
46+
| `true` | **NATS Core** — hub pub/sub subjects `controller.relay.v1.*` via `NatsRelayTransport` |
47+
48+
Set `NATS_ENABLED=true` only when the platform NATS hub is deployed and all Controller replicas share the same value.
49+
4150
| Setting | Default | Env |
4251
|---------|---------|-----|
43-
| Cross-replica requires AMQP | `true` | `WS_HA_CROSS_REPLICA_REQUIRES_AMQP` |
44-
| Fail fast when router down | `true` | `WS_HA_FAIL_FAST_ON_ROUTER_UNAVAILABLE` |
52+
| Cross-replica requires relay backend | `true` | `WS_HA_CROSS_REPLICA_REQUIRES_AMQP` |
53+
| Fail fast when relay backend down | `true` | `WS_HA_FAIL_FAST_ON_ROUTER_UNAVAILABLE |
4554

46-
**Requirements:**
55+
> Env names retain `AMQP`/`ROUTER` for backward compatibility; semantics apply to the **active relay backend** (AMQP or NATS) per R112.
4756
48-
1. Deploy the **router** system microservice and ensure Controller can reach AMQP (`RouterConnectionService`).
57+
### AMQP relay (`nats.enabled=false`)
58+
59+
1. Deploy the **router** system microservice and ensure Controller can reach AMQP (`RouterConnectionManager` pool).
4960
2. Run **2+ Controller replicas** behind a load balancer with **sticky sessions optional** — cross-replica exec/log uses AMQP queues (`agent-{sessionId}`, `user-{sessionId}`, `logs-user-{sessionId}`).
50-
3. When the router is unavailable, new cross-replica sessions close with WebSocket code **1013** (`Router unavailable for cross-replica session`).
61+
3. When the router/AMQP backend is unavailable, new cross-replica sessions close with WebSocket code **1013** (`Router unavailable for cross-replica session`).
62+
63+
Plan 18 adds an **8-connection AMQP pool** per replica with overflow recovery — intense log streams must not poison other sessions (no router restart required). **Remote CP** resolves **`router.default.svc.bridge.local`** then default router `host`; **Kubernetes CP** resolves **`router.{namespace}.svc.cluster.local`** then default router `host`. Port from `Routers.messagingPort` (default **5671**).
5164

52-
Same-replica sessions may relay directly without AMQP when both user and agent land on the same pod.
65+
### NATS relay (`nats.enabled=true`)
66+
67+
1. Platform NATS hub must be running with `NatsInstances.isHub=true`.
68+
2. Controller provisions dedicated **`controller`** NATS account/user (not SYS / `admin-hub`) via NATS auth reconcile.
69+
3. Cross-replica exec uses subjects `controller.relay.v1.exec.{sessionId}.agent` / `.user`; logs use `controller.relay.v1.log.{sessionId}.user`. Plain TCP to hub — port from `NatsInstances.serverPort` (default **4222**) for every host in the resolver list.
70+
4. **Remote CP:** Controller resolves **`nats.default.svc.bridge.local`** (Edgelet internal DNS) then hub `host`; both use hub `serverPort`.
71+
5. **Kubernetes CP:** Controller resolves **`nats-server.{namespace}.svc.cluster.local`** then hub `host`.
72+
6. Remote ControlPlane replicas connect to the **hub** NATS only — not local fog NATS leaf.
73+
7. When NATS relay is unavailable, fail-fast semantics match AMQP (close **1013** when configured).
74+
75+
Same-replica sessions may relay directly without AMQP or NATS when both user and agent land on the same pod.
5376

5477
---
5578

@@ -59,7 +82,7 @@ On shutdown, Controller drains WebSocket sessions for up to **`WS_DRAIN_TIMEOUT_
5982

6083
1. Reject new upgrades (`verifyClient` → draining).
6184
2. Close pending users with code **1001** (`Server draining`).
62-
3. Send CLOSE frames, clean exec/log session DB rows, tear down AMQP bridges.
85+
3. Send CLOSE frames, clean exec/log session DB rows, tear down relay bridges (AMQP or NATS).
6386

6487
### Kubernetes manifest example
6588

@@ -113,6 +136,15 @@ node test/load/ws-pairing-load.js --multi-ms 100
113136

114137
The `--multi-ms` mode creates **3 exec sessions per microservice** (100 MS × 3 = 300 pairs) to validate multi-session pairing latency under the same p99 SLO.
115138

139+
**AMQP profile** (`nats.enabled=false`): run the probe above on a dev machine — it exercises in-process `ExecSessionManager` pairing only (no router required). Record p99 from stdout; target **< 5000 ms**.
140+
141+
**NATS profile** (`nats.enabled=true`): the same probe validates session-manager pairing latency (transport-agnostic SLO). For end-to-end NATS relay validation in staging:
142+
143+
1. Deploy Controller with **`NATS_ENABLED=true`** on **2+ replicas** and a platform NATS hub (`NatsInstances.isHub=true`).
144+
2. Confirm **`controller`** NATS account reconcile succeeded (NATS auth logs).
145+
3. Run cross-replica exec/log sessions (user on replica A, agent on replica B) while recording OTEL **`ws_pairing_duration_ms`** p99.
146+
4. Optionally repeat `node test/load/ws-pairing-load.js --pairs 500` against staging API with agent simulators — same **p99 < 5s** SLO applies.
147+
116148
For production validation, repeat against a staging cluster with real agent simulators and record p99 from Controller OTEL histogram `ws_pairing_duration_ms`.
117149

118150
---
@@ -128,7 +160,9 @@ Enable `ENABLE_TELEMETRY=true`. Key metrics (`src/websocket/ws-metrics.js`):
128160
| `ws_pending_pairings` | gauge |
129161
| `ws_pairing_duration_ms` | histogram |
130162
| `ws_amqp_publish_errors` | counter |
131-
| `ws_router_connected` | gauge |
163+
| `ws_amqp_session_saturated` | counter (Plan 18 overflow/backpressure) |
164+
| `ws_router_pool_connections` | gauge (Plan 18 AMQP pool health) |
165+
| `ws_router_pool_unsettled` | gauge (Plan 18 AMQP unsettled deliveries) |
132166

133167
---
134168

docs/swagger.yaml

Lines changed: 84 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -2774,8 +2774,7 @@ paths:
27742774
`{ sessionId, microserviceUuid }` in the `data` field and top-level `sessionId`,
27752775
followed by STDERR "waiting for agent…". Relay frames use `execId` equal to `sessionId`.
27762776
2777-
**HA:** Multi-replica deployments require AMQP router (`WebSocketQueueService`).
2778-
Cross-replica sessions fail fast with close code **1013** when router is unavailable.
2777+
**HA (R112):** Multi-replica deployments require a **relay backend** selected at startup by **`nats.enabled`**: AMQP router queues when `false` (default), NATS Core pub/sub on the platform hub when `true`. Cross-replica sessions fail fast with close code **1013** when the active relay backend is unavailable. See `docs/operations/ws-sessions.md`.
27792778
27802779
See `#/components/schemas/WsExecMessageTypes` and `#/components/schemas/WsCloseCodes`.
27812780
operationId: userMicroserviceExecWebSocket
@@ -6593,11 +6592,16 @@ paths:
65936592
tags:
65946593
- NATS
65956594
summary: Gets NATS creds for specific account user
6595+
description: |
6596+
Returns base64-encoded `.creds` file content for the user.
6597+
6598+
Path `{appName}` is an **application name** for app-linked accounts, or a **NATS account name** for platform/system accounts (`SYS`, leaf accounts, **`controller`**).
65966599
operationId: getNatsUserCreds
65976600
parameters:
65986601
- in: path
65996602
name: appName
66006603
required: true
6604+
description: Application name or NATS account name (e.g. `SYS`, `controller`)
66016605
schema:
66026606
type: string
66036607
- in: path
@@ -7033,7 +7037,7 @@ tags:
70337037
- name: WebSocketSessions
70347038
description: |
70357039
Interactive exec and log streaming WebSocket endpoints (MessagePack binary).
7036-
Multi-replica HA requires AMQP router — see operations/ws-sessions.md.
7040+
Multi-replica HA requires a relay backend per **`nats.enabled`** (AMQP when `false`, NATS Core when `true`) — see operations/ws-sessions.md.
70377041
- name: User
70387042
description: Manage your users
70397043
- name: Secrets
@@ -8276,7 +8280,11 @@ components:
82768280
description: True when microservice is the ControlPlane controller workload (DB column)
82778281
ControllerRegisterRequest:
82788282
type: object
8279-
description: Slim register body for Edgelet ControlPlane controller workload
8283+
description: >-
8284+
Edgelet ControlPlane controller workload register body. Accepts the same
8285+
container/workload fields as user microservice deploy (excluding application,
8286+
serviceAccount, and natsConfig). Ownership fields (iofogUuid, application)
8287+
are derived from the authenticated system fog.
82808288
required:
82818289
- uuid
82828290
- images
@@ -8298,6 +8306,31 @@ components:
82988306
$ref: "#/components/schemas/MicroserviceContainerImage"
82998307
registryId:
83008308
type: integer
8309+
config:
8310+
type: string
8311+
annotations:
8312+
type: string
8313+
hostNetworkMode:
8314+
type: boolean
8315+
isPrivileged:
8316+
type: boolean
8317+
logSize:
8318+
type: integer
8319+
minimum: 0
8320+
runtime:
8321+
type: string
8322+
pidMode:
8323+
type: string
8324+
ipcMode:
8325+
type: string
8326+
runAsUser:
8327+
type: string
8328+
platform:
8329+
type: string
8330+
cpuSetCpus:
8331+
type: string
8332+
memoryLimit:
8333+
type: integer
83018334
ports:
83028335
type: array
83038336
items:
@@ -8306,16 +8339,57 @@ components:
83068339
type: array
83078340
items:
83088341
$ref: "#/components/schemas/VolumeMappingRequest"
8342+
extraHosts:
8343+
type: array
8344+
items:
8345+
type: object
8346+
required:
8347+
- name
8348+
- address
8349+
properties:
8350+
name:
8351+
type: string
8352+
address:
8353+
type: string
83098354
env:
83108355
type: array
83118356
items:
83128357
$ref: "#/components/schemas/AgentEnvRequest"
8313-
config:
8314-
type: string
8315-
hostNetworkMode:
8316-
type: boolean
8317-
runtime:
8318-
type: string
8358+
cmd:
8359+
type: array
8360+
items:
8361+
type: string
8362+
cdiDevices:
8363+
type: array
8364+
items:
8365+
type: string
8366+
capAdd:
8367+
type: array
8368+
items:
8369+
type: string
8370+
capDrop:
8371+
type: array
8372+
items:
8373+
type: string
8374+
healthCheck:
8375+
type: object
8376+
required:
8377+
- test
8378+
properties:
8379+
test:
8380+
type: array
8381+
items:
8382+
type: string
8383+
interval:
8384+
type: integer
8385+
timeout:
8386+
type: integer
8387+
startPeriod:
8388+
type: integer
8389+
startInterval:
8390+
type: integer
8391+
retries:
8392+
type: integer
83198393
schedule:
83208394
type: integer
83218395
enum:

0 commit comments

Comments
 (0)