diff --git a/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md b/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md new file mode 100644 index 00000000..0390a1f7 --- /dev/null +++ b/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md @@ -0,0 +1,650 @@ +# Per-Cluster Mesh Phase 5 — gRPC live-telemetry stream Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: execute this with superpowers-extended-cc:subagent-driven-development +> (chosen for the program's prior phases). Fresh subagent per task; classification-driven review chain. + +**Goal:** Replace the DPS fan-out of the live-telemetry observability channels with one gRPC +server-streaming contract — **each driver node hosts the telemetry gRPC server (Kestrel h2c), +central dials in** — so AdminUI live panels keep working once the meshes split (Phase 6) and no +longer depend on a shared Akka gossip ring for fleet observability. Ships as a per-node **dark +switch** (`Telemetry:Mode` = `Dps` default | `Grpc`), authenticated fail-closed from day one. + +**Architecture:** Mirror ScadaBridge's `SiteStreamService` shape (recon'd 2026-07-23), adapted to +OtOpcUa's substrate. Node side: the four telemetry publish seams also emit into a **node-local +in-process hub** (`ITelemetryLocalHub`) — DPS publishing is left untouched, so the switch is pure — +and a streaming gRPC service fans the hub to connected clients through per-subscriber bounded +`DropOldest` channels. Central side: a supervisor actor discovers driver-node gRPC endpoints from +`ClusterNode` rows (`Host` + `GrpcPort`, added in Phase 1), keeps one reconnecting dialer per node, +converts streamed events back to the domain records, and feeds **the exact same in-process sinks the +DPS SignalR bridges feed today** (`IInProcessBroadcaster`, +`IInProcessBroadcaster`, `IDriverStatusSnapshotStore`, +`IDriverResilienceStatusStore`). The AdminUI components are untouched: only the bridge's *upstream* +swaps. + +**Tech Stack:** .NET 10, `Grpc.AspNetCore` / `Grpc.Net.Client` (already in-repo via Phase 3), +`Grpc.Tools` codegen, Akka.NET, `System.Threading.Channels`, xUnit + Shouldly. + +--- + +## Scope (settled with the user 2026-07-23) + +**In scope — migrate these four node→central observability channels** to the stream as four `oneof` +event kinds: + +| DPS topic today | Message record | Central sink fed today | +|---|---|---| +| `alerts` | `AlarmTransitionEvent` | `IInProcessBroadcaster` (+ `AlertHub`) | +| `script-logs` | `ScriptLogEntry` | `IInProcessBroadcaster` (+ `ScriptLogHub`) | +| `driver-health` | `DriverHealthChanged` | `IDriverStatusSnapshotStore` (+ `DriverStatusHub`) | +| `driver-resilience-status` | `DriverResilienceStatusChanged` | `IDriverResilienceStatusStore` (no hub) | + +**Explicitly deferred, with rationale (do NOT migrate in Phase 5):** + +- **`redundancy-state`** — bidirectional, built from `Cluster.State`, **pair-local control-plane** + that drives ServiceLevel + the Primary gate (consumed by `OpcUaPublishActor`, + `ScriptedAlarmHostActor`, `DriverHostActor`, `HistorianAdapterActor`). It stays on DPS in both + MeshTransport modes today, and under Phase 6 it is pair-local and works in-mesh. Central's *display* + of each pair's redundancy is a Phase 6 cross-mesh concern (possibly a later added event kind), not a + Phase 5 observability panel. +- **`fleet-status`** — **central-internal**: `FleetStatusBroadcaster` (admin singleton) builds it + from the admin node's own cluster membership/reachability/leader events; `Fleet.razor` **polls the + ConfigDB** and ignores the feed entirely. It is not a node→central stream, and its live UI path is + already DB-polled. Revisit in Phase 6 when central loses gossip visibility of site nodes. +- **`deployment-acks`** — already rides the Phase 2 ClusterClient transport when + `MeshTransport:Mode=ClusterClient`; it is a command-plane reply, not telemetry. + +This narrowing mirrors how Phases 1 and 3 honestly scoped down from the program sketch. The program +doc's Phase 5 line lists all seven; this plan records the four that are genuinely live node→central +observability and defers the rest with reasons above. **Update the program doc + design §6.3 in +Task 10.** + +## Direction & the dark switch (read before any task) + +- **Node = server, central = client** (the load-bearing ScadaBridge inversion). Telemetry + originates on driver nodes; central/admin consumes it. A fused `admin,driver` node both hosts (as + driver) and dials (as admin) — it dials itself plus its pair peer, same as central dials site nodes. +- **`Telemetry:Mode` is read at startup; both code paths are compiled into every binary.** Flipping + the flag is an appsettings/env change + restart — NOT a rebuild — exactly the Phase 2/3 dark-switch + discipline (`OTOPCUA_CONFIG_MODE` on the rig). The node **always** hosts the gRPC server when + `Telemetry:GrpcListenPort > 0` and **always** emits into the local hub AND publishes DPS, in both + modes — so central can ingest from either side without a node redeploy. Only **central's ingest + source** switches: `Dps` → today's four DPS SignalR bridges subscribe and feed the sinks; `Grpc` → + those four bridges are NOT spawned and the dial supervisor feeds the identical sinks instead. +- **Auth from day one.** Reuse the fail-closed `FixedTimeEquals` bearer interceptor pattern + (`ConfigServeAuthInterceptor` / `LocalDbSyncAuthInterceptor`). Shared node key + `Telemetry:ApiKey` (serve side) == `TelemetryDial:ApiKey` (central side). This supersedes design + §6.3's "unauthenticated for now" — ScadaBridge itself closed that gap with this same pattern. + +## Reuse map (from recon 2026-07-23 — exact sites) + +- **Proto + codegen:** add `telemetry.proto` beside `src/Core/ZB.MOM.WW.OtOpcUa.Commons/Protos/deployment_artifact.proto`; register a `` in `ZB.MOM.WW.OtOpcUa.Commons.csproj:29` (same block as the existing item). Generated types are shared by node (Runtime/Host) and central (AdminUI) from the one Commons reference. +- **Server hosting + Kestrel h2c:** the dedicated-listener block `Program.cs:434-533` and the `MapGrpcService` gate `Program.cs:577-580`. Add `telemetryListenPort` alongside `syncListenPort`/`configServeGrpcPort`. +- **Auth interceptor:** copy `src/Server/ZB.MOM.WW.OtOpcUa.Host/Configuration/ConfigServeAuthInterceptor.cs`; add to the shared `AddGrpc` pipeline at `Program.cs:425-431`. +- **Client dialing:** `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Deployment/GrpcDeploymentArtifactFetcher.cs` (channel-cache, Bearer metadata, h2c `GrpcChannel.ForAddress`, linked-CTS deadline). +- **Options + validator:** `ConfigSourceOptions`/`ConfigSourceOptionsValidator` in `src/Core/ZB.MOM.WW.OtOpcUa.Cluster/`; registered via `AddValidatedOptions` in `ServiceCollectionExtensions.AddOtOpcUaCluster` (`:32`). +- **Node discovery:** `CentralCommunicationActor.cs:198-219` already reads `ClusterNodes` (enabled, non-maintenance) selecting `NodeId, Host, AkkaPort` — extend the same query shape to `GrpcPort` for telemetry dial targets. `ClusterNode.GrpcPort` (nullable) already exists (`Entities/ClusterNode.cs:46`) explicitly for "the Phase 5 telemetry stream." +- **Central sinks (the untouched seam):** registered in `HubServiceCollectionExtensions.AddOtOpcUaDriverStatusServices` (`:29-35`); DPS bridges spawned in `WithOtOpcUaSignalRBridges` (`:53-83`). +- **ScadaBridge reference to mirror:** `SiteStreamGrpcServer.cs` (relay-actor + bounded `DropOldest` channel + lifecycle/cleanup + concurrency cap + max-stream-lifetime), `SiteStreamGrpcClient.cs` / `SiteStreamGrpcClientFactory.cs` (channel cache, keepalive), `SiteAlarmAggregatorActor.cs` (generation-stamped, budget-limited, self-healing reconnect), `ProtoContractTests.cs` (reflection-over-oneof contract lock). + +--- + +## Tasks + +### Task 0: Telemetry proto contract + codegen + contract-lock test + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** Task 1, Task 2 + +**Files:** +- Create: `src/Core/ZB.MOM.WW.OtOpcUa.Commons/Protos/telemetry.proto` +- Modify: `src/Core/ZB.MOM.WW.OtOpcUa.Commons/ZB.MOM.WW.OtOpcUa.Commons.csproj` (add the `` item next to line 29) +- Test: `tests/Core/ZB.MOM.WW.OtOpcUa.Commons.Tests/Protos/TelemetryProtoContractTests.cs` (create; match the existing Commons.Tests project layout — if there is no Commons.Tests project, add the test to the nearest existing Core test project that already references generated Commons types and note it) + +**Step 1: Write `telemetry.proto`.** Package `telemetry.v1`, `csharp_namespace = ZB.MOM.WW.OtOpcUa.Commons.Protos.Telemetry.V1`. One server-streaming RPC + a `oneof` envelope with the four event kinds. Every enum carries a `*_UNSPECIFIED = 0` zero value; every field comment marks additive intent. + +```proto +syntax = "proto3"; +package telemetry.v1; +option csharp_namespace = "ZB.MOM.WW.OtOpcUa.Commons.Protos.Telemetry.V1"; + +import "google/protobuf/timestamp.proto"; + +// Central dials each driver node and opens Subscribe; the node streams its own live telemetry. +// Additive-only field evolution: never renumber/reuse a tag; a pre-field peer must read a new +// field's proto3 default correctly. Locked by TelemetryProtoContractTests. +service TelemetryStreamService { + rpc Subscribe(TelemetryStreamRequest) returns (stream TelemetryEvent); +} + +message TelemetryStreamRequest { + string correlation_id = 1; // safe Akka path element; becomes the relay actor name on the node +} + +message TelemetryEvent { + string correlation_id = 1; + oneof event { + AlarmTransition alarm_transition = 2; // <- DPS topic "alerts" + ScriptLog script_log = 3; // <- DPS topic "script-logs" + DriverHealth driver_health = 4; // <- DPS topic "driver-health" + DriverResilienceStatus driver_resilience = 5; // <- DPS topic "driver-resilience-status" + } +} + +// The four event bodies below mirror the C# records field-for-field. Encode the domain records' +// enums as proto enums (UNSPECIFIED=0), and DateTimes as google.protobuf.Timestamp. +message AlarmTransition { /* fields mirroring Commons/Messages/Alerts/AlarmTransitionEvent.cs */ } +message ScriptLog { /* fields mirroring Commons/Messages/Logging/ScriptLogEntry.cs */ } +message DriverHealth { /* fields mirroring Commons/Messages/Drivers/DriverHealthChanged.cs */ } +message DriverResilienceStatus{ /* fields mirroring Commons/Messages/Drivers/DriverResilienceStatusChanged.cs */ } +``` + +The implementer MUST open each of the four domain records (paths in the scope table above) and +transcribe every field into the corresponding message with a stable tag number, choosing +`Timestamp`/`Int32`/`string`/`bool`/enum as the type dictates. Where a record enum exists, define a +matching proto enum with `_UNSPECIFIED = 0`. This is the contract — get it complete, because the +mapping tasks (7) and the lock test below depend on it. + +**Step 2: Register codegen.** In the `.csproj`, next to the existing `deployment_artifact.proto` item: +```xml + +``` + +**Step 3: Write the contract-lock test (fails first).** Mirror ScadaBridge `ProtoContractTests.AllOneofVariants_HaveConversionHandlers`: reflect over `Enum.GetValues()` minus `None`, assert count- and membership-equality against a hand-maintained `HandledCases` array that Task 7's converter will also key off. This fails until the proto compiles and the array is filled. + +```csharp +[Fact] +public void EveryOneofVariant_IsAccountedFor() +{ + var variants = Enum.GetValues() + .Where(c => c != TelemetryEvent.EventOneofCase.None).ToArray(); + variants.ShouldBe(TelemetryProtoContract.HandledCases, ignoreOrder: true); +} +``` +Add a tiny `TelemetryProtoContract.HandledCases` constant array in Commons (the single source both +this test and the Task-7 converter reference) listing the four cases. + +**Step 4: `dotnet build ZB.MOM.WW.OtOpcUa.slnx`** — expect the generated types to appear; test goes green. + +**Step 5: Commit.** `feat(mesh-phase5): telemetry.proto contract + oneof + contract-lock test` + +--- + +### Task 1: Telemetry options (serve + dial) + validator + registration + +**Classification:** small +**Estimated implement time:** ~4 min +**Parallelizable with:** Task 0, Task 2 + +**Files:** +- Create: `src/Core/ZB.MOM.WW.OtOpcUa.Cluster/TelemetryOptions.cs` +- Modify: `src/Core/ZB.MOM.WW.OtOpcUa.Cluster/ServiceCollectionExtensions.cs` (register in `AddOtOpcUaCluster`, ~lines 40-49 block) +- Test: `tests/Core/ZB.MOM.WW.OtOpcUa.Cluster.Tests/TelemetryOptionsValidatorTests.cs` + +**Step 1:** Define two options classes and one validator, mirroring `ConfigSourceOptions`/`ConfigServeOptions` + `ConfigSourceOptionsValidator` (which reads roles from `IConfiguration`, NOT `IClusterRoleInfo`). + +```csharp +public sealed class TelemetryOptions // section "Telemetry" (serve side, node) +{ + public const string SectionName = "Telemetry"; + public string Mode { get; set; } = "Dps"; // Dps | Grpc (central-ingest selector; harmless on node) + public int GrpcListenPort { get; set; } // 0 = disabled; driver node's telemetry h2c port + public string ApiKey { get; set; } = ""; // shared node key; supply via ${secret:}/env +} + +public sealed class TelemetryDialOptions // section "TelemetryDial" (central) +{ + public const string SectionName = "TelemetryDial"; + public string Mode { get; set; } = "Dps"; // Dps | Grpc + public string ApiKey { get; set; } = ""; // must equal the nodes' Telemetry:ApiKey + public int ContactRefreshSeconds { get; set; } = 60; + public int CallTimeoutSeconds { get; set; } = 30; +} +``` + +**Step 2: Validator rules** (`TelemetryOptionsValidator : IValidateOptions`, ctor-inject `IConfiguration`, read `Cluster:Roles` like `ConfigSourceOptionsValidator:98`): +- `Mode` must be `Dps` or `Grpc` (case-insensitive) — else Fail. +- A **driver-role** node with `Mode=Grpc` must have `GrpcListenPort > 0` — else Fail ("nothing to serve"). +- `Mode=Grpc` with a non-empty role set must have a non-empty `ApiKey` (fail-closed: refuse to host an un-keyed telemetry surface) — else Fail. +- Add a sibling `TelemetryDialOptionsValidator`: `Mode` in {Dps,Grpc}; `Mode=Grpc` ⇒ `ApiKey` non-empty. + +**Step 3: Register** via `AddValidatedOptions(configuration, TelemetryOptions.SectionName)` and the dial equivalent, in `AddOtOpcUaCluster`. + +**Step 4:** Tests — Grpc+driver+port0 fails; Grpc+empty-key fails; Dps passes with defaults; unknown Mode fails. Run `dotnet test tests/Core/ZB.MOM.WW.OtOpcUa.Cluster.Tests`. + +**Step 5: Commit.** `feat(mesh-phase5): Telemetry/TelemetryDial options + fail-closed validators` + +--- + +### Task 2: Node-local in-process telemetry hub + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** Task 0, Task 1 + +**Files:** +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Telemetry/ITelemetryLocalHub.cs` +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Telemetry/TelemetryLocalHub.cs` +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.Runtime.Tests/Telemetry/TelemetryLocalHubTests.cs` + +**Design:** a process-wide singleton carrying **this node's own** telemetry (never cluster-wide — that +is what keeps it correct on both the current single mesh and the Phase-6 pair mesh; DPS would leak +peers' events). It holds a **domain-typed** union (not proto) so the publish seams stay decoupled from +generated types. Two responsibilities: + +1. `Emit(TelemetryItem)` — fan the item to every currently-subscribed writer through **bounded + per-subscriber channels with `FullMode = DropOldest`** (live view is lossy under backpressure by + design; a slow central never blocks the node). +2. **Snapshot replay for last-value channels.** Keep a last-value cache for `DriverHealth` (keyed by + driver instance) and `DriverResilienceStatus` (keyed by `(instance, host)`); `alerts`/`script-logs` + are append logs and are NOT cached. On `Subscribe`, first drain the cached snapshots to the new + writer, then attach it for live deltas — the simplified equivalent of ScadaBridge's seed-then-stream, + so a central reconnect immediately re-primes the driver-health/resilience stores. + +```csharp +public abstract record TelemetryItem +{ + public sealed record Alarm(AlarmTransitionEvent E) : TelemetryItem; + public sealed record Script(ScriptLogEntry E) : TelemetryItem; + public sealed record Health(DriverHealthChanged E) : TelemetryItem; + public sealed record Resilience(DriverResilienceStatusChanged E) : TelemetryItem; +} + +public interface ITelemetryLocalHub +{ + void Emit(TelemetryItem item); + // Returns a reader that first yields cached snapshots, then live deltas; disposing the + // subscription detaches and completes the channel. + ITelemetrySubscription Subscribe(int boundedCapacity); +} +``` + +`TelemetryLocalHub` uses a `ConcurrentDictionary>` of subscribers + +two `ConcurrentDictionary` snapshot caches. `Emit` updates the snapshot cache (for Health/Resilience) +then `TryWrite`s to each subscriber (drop-oldest handled by the bounded channel). `Subscribe` snapshots +the caches into the new channel under a brief lock ordering guarantee (cache-then-attach) so no delta +is lost across the attach boundary. + +**Register** as a singleton in the driver-role DI branch (Task 3 wires the producers; do the +`AddSingleton()` here in the Runtime `ServiceCollectionExtensions` +`hasDriver` path, or wherever driver-role Runtime services register — locate and match). + +**Tests:** Emit-before-Subscribe is not seen except via snapshot (Health/Resilience last-value IS seen; +Alarm/Script are not); two subscribers each get their own copy; a full channel drops oldest not newest; +dispose detaches. Run the Runtime.Tests subset. + +**Step 5: Commit.** `feat(mesh-phase5): node-local telemetry hub (snapshot-replay + drop-oldest fan-out)` + +--- + +### Task 3: Tap the four publish seams into the hub + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** none (needs Task 2) + +**Files (all Modify) — add a hub `Emit` beside the existing DPS `Publish`, leaving DPS intact:** +- `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/AkkaDriverHealthPublisher.cs:36` (driver-health) +- `src/Server/ZB.MOM.WW.OtOpcUa.Host/Drivers/DriverResilienceStatusPublisherService.cs:80` (resilience) +- `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/ScriptedAlarms/ScriptedAlarmHostActor.cs:379` + `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Drivers/DriverHostActor.cs:1405` (alerts — two producers) +- `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/VirtualTags/VirtualTagActor.cs:261` + `src/Server/ZB.MOM.WW.OtOpcUa.Runtime/Scripting/DpsScriptLogPublisher.cs:46` (script-logs — two producers) +- Test: extend the nearest existing publisher tests, or add `tests/.../Telemetry/PublishSeamEmitsToHubTests.cs` + +**Approach:** inject `ITelemetryLocalHub` into each producer (constructor for the services/publishers; +for actors, pass via `Props` — match how each already receives dependencies). At each seam, immediately +after the existing `Mediator.Tell(new Publish(...))`, add `_telemetryHub.Emit(new TelemetryItem.Xxx(msg))`. +**Do NOT remove or gate the DPS publish** — it is the `Dps`-mode path and must remain unconditional so +the switch stays pure. The hub is a no-op sink until a client subscribes, so this is safe on every node +regardless of mode. + +For actors constructed where a hub isn't readily resolvable, resolve the singleton once at spawn and +thread it through `Props` (do not `DependencyResolver` inside the actor per-message). Where a producer +is only present on driver-role nodes, the hub singleton is guaranteed registered (Task 2). + +**Tests:** each producer, when driven, results in exactly one `hub.Emit` of the correct `TelemetryItem` +subtype carrying the same payload it published to DPS. Use a fake `ITelemetryLocalHub` capturing emits. + +**Step 5: Commit.** `feat(mesh-phase5): tap the 4 telemetry publish seams into the local hub (DPS intact)` + +--- + +### Task 4: Node-side gRPC streaming service + +**Classification:** high-risk +**Estimated implement time:** ~5 min (mirror ScadaBridge closely) +**Parallelizable with:** none (needs Task 0, Task 2) + +**Files:** +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.Host/Grpc/TelemetryStreamGrpcService.cs` (place beside `DeploymentArtifactService.cs`; alias the generated base to avoid the name collision, as that file does at its lines 9-10) +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.Host/Grpc/TelemetryProtoMap.Node.cs` (domain→proto mapping; the reverse of Task 7) +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.Host.Tests/Grpc/TelemetryStreamGrpcServiceTests.cs` + +**Design — port `SiteStreamGrpcServer.RunSubscriptionStreamAsync` (`.cs:253-385`) shape:** + +```csharp +public sealed class TelemetryStreamGrpcService : GeneratedTelemetryBase +{ + // ctor: ITelemetryLocalHub hub, ILogger, IOptions (for a GrpcMaxConcurrentStreams knob) + public override async Task Subscribe(TelemetryStreamRequest request, + IServerStreamWriter responseStream, ServerCallContext context) + { + // 1. Validate correlation_id is a safe id (it labels the subscription; reject empty/oversized). + // 2. Concurrency cap (default 100) -> throw RpcException(ResourceExhausted) when exceeded. + // 3. sub = _hub.Subscribe(boundedCapacity: 1000); // snapshot-then-live, DropOldest inside the hub + // 4. MaxStreamLifetime linked-CTS CancelAfter (default 4h) ORed with context.CancellationToken, + // so a zombie stream terminates even if h2c keepalive misses it. + // 5. await foreach (item in sub.Reader.ReadAllAsync(ct)) + // await responseStream.WriteAsync(TelemetryProtoMap.ToProto(item, request.CorrelationId), ct); + // 6. finally: sub.Dispose(); balance any opened/closed gauge; swallow OperationCanceled as normal. + } +} +``` + +Key correctness points (from the ScadaBridge recon — do not skip): +- The **bounded DropOldest channel lives in the hub** (Task 2), so the service just pumps its reader — + a slow/blocked central cannot back-pressure the node's actor threads. +- Wrap the whole body so a client disconnect (`RpcException`/`OperationCanceledException`) exits cleanly + and disposes the subscription; never let it bubble as a fault. +- No readiness race: the hub is a plain singleton available at host build, so unlike ScadaBridge's + `SetReady(ActorSystem)` gate there is nothing to defer — but if the driver actor system isn't up yet + the hub simply has no snapshots and no deltas, which is fine. + +**`TelemetryProtoMap.ToProto`**: switch on `TelemetryItem` subtype → build the matching proto message, +wrap in `TelemetryEvent { CorrelationId, = ... }`. Transcribe every field (the reverse of the +`.proto` transcription in Task 0). DateTimes → `Timestamp.FromDateTime(utc)`. + +**Tests:** in-memory — subscribe, emit each of the four item kinds into the hub, assert the service +writes the correct proto event with fields intact; a cancelled `ServerCallContext` ends the stream and +disposes; the concurrency cap throws `ResourceExhausted` past the limit. (Full two-host wire test is +covered by the Task-12 live gate + optionally a Task-8 boundary test.) + +**Step 5: Commit.** `feat(mesh-phase5): node-side TelemetryStreamService (hub -> server-streaming)` + +--- + +### Task 5: Telemetry stream auth interceptor + +**Classification:** small +**Estimated implement time:** ~3 min +**Parallelizable with:** none (needs Task 1) + +**Files:** +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.Host/Configuration/TelemetryStreamAuthInterceptor.cs` (copy `ConfigServeAuthInterceptor.cs` verbatim, change the gated prefix + options source) +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.Host.Tests/Configuration/TelemetryStreamAuthInterceptorTests.cs` + +**Steps:** +- `ServicePrefix = "/telemetry.v1.TelemetryStreamService/"`; read `IOptions.ApiKey`. +- Override all four handler kinds (especially **ServerStreaming** — that is the RPC shape here) so the + gate holds; non-matching paths pass through; empty key ⇒ throw (fail-closed); `FixedTimeEquals` on the + Bearer token; reject with `PermissionDenied`. +- **Keep exactly one public constructor** (the ScadaBridge recon flagged that `Grpc.AspNetCore` silently + stops authorizing an interceptor with >1 ctor — pin it with a reflection test). + +**Tests:** right key passes; wrong key `PermissionDenied`; empty configured key rejects all; a call to a +different service path passes through. + +**Step 5: Commit.** `feat(mesh-phase5): fail-closed bearer interceptor for the telemetry stream` + +--- + +### Task 6: Kestrel wiring + map the node telemetry server + +**Classification:** standard +**Estimated implement time:** ~4 min +**Parallelizable with:** none (needs Task 1, Task 4, Task 5) + +**Files:** +- Modify: `src/Server/ZB.MOM.WW.OtOpcUa.Host/Program.cs` (the `AddGrpc` block 425-431, the dedicated-listener block 434-533, and the map block 577-580) +- Test: none new (covered by Task 12 live gate); build must stay green + +**Steps — extend, do not rewrite, the existing block:** + +1. **Interceptor registration** (`Program.cs:425-431`): inside the `AddGrpc(o => …)`, add + `if (hasDriver) o.Interceptors.Add();` (alongside the LocalDbSync one — + both are driver-side, both path-scoped, harmless when their service is unmapped). + +2. **Port resolution** (near `:453-454`): + ```csharp + var telemetryListenPort = hasDriver ? builder.Configuration.GetValue("Telemetry:GrpcListenPort") : 0; + ``` + Add `|| telemetryListenPort > 0` to the block guard at `:455`. + +3. **Bind it** inside the `ConfigureKestrel` closure (`:517-526`): capture `var telemetryPortToBind = telemetryListenPort;` beside the other two, and add + `if (telemetryPortToBind > 0) kestrel.ListenAnyIP(telemetryPortToBind, o => o.Protocols = HttpProtocols.Http2);`. + +4. **HTTPS-refuse branch** (`:500-511`): set `telemetryListenPort = 0;` too, and add it to the log + message (all dedicated listeners disable together when the host serves HTTPS). + +5. **Map the service** (after `:580`): + ```csharp + if (hasDriver && telemetryListenPort > 0) + app.MapGrpcService(); + ``` + +6. Update the block's header comment (434-452) to mention the third dedicated port. + +**Verify:** `dotnet build`; a fused node with all three ports set binds all three exactly once (the +"re-bind existing surface once" invariant already holds — the new port is just one more +`ListenAnyIP`). A driver-only node with only `Telemetry:GrpcListenPort` set binds only that. + +**Step 5: Commit.** `feat(mesh-phase5): host the telemetry gRPC server on driver nodes (dedicated h2c port)` + +--- + +### Task 7: Central per-node dialer client + proto→domain mapping + +**Classification:** high-risk +**Estimated implement time:** ~5 min +**Parallelizable with:** none (needs Task 0) + +**Files:** +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.ControlPlane/Telemetry/TelemetryStreamClient.cs` +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.ControlPlane/Telemetry/TelemetryProtoMap.Central.cs` (proto→domain; reverse of Task 4; keyed off `TelemetryProtoContract.HandledCases`) +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.ControlPlane.Tests/Telemetry/TelemetryProtoMapTests.cs` + +**Design — mirror `SiteStreamGrpcClient.cs` + `GrpcDeploymentArtifactFetcher` channel handling:** + +```csharp +public sealed class TelemetryStreamClient : IDisposable +{ + // ctor: string endpoint (http://host:port), string apiKey, IClientFactory seam for tests + // - GrpcChannel.ForAddress(endpoint) over http:// => prior-knowledge h2c; cache one channel per endpoint. + // - HTTP/2 keepalive: PingDelay 15s / Timeout 10s / Always (SocketsHttpHandler on the channel). + public async Task RunAsync(string correlationId, Action onEvent, + Action onError, CancellationToken ct) + { + // headers: authorization: Bearer {apiKey} + // using call = client.Subscribe(new TelemetryStreamRequest{CorrelationId=correlationId}, headers, ct) + // await foreach (evt in call.ResponseStream.ReadAllAsync(ct)) onEvent(TelemetryProtoMap.ToDomain(evt)); + // RpcException(Cancelled) on shutdown => normal; anything else => onError (the reconnect trigger). + } +} +``` + +`TelemetryProtoMap.ToDomain(TelemetryEvent)` switches on `evt.EventCase` and reconstructs the domain +record (`AlarmTransitionEvent` / `ScriptLogEntry` / `DriverHealthChanged` / `DriverResilienceStatusChanged`). +It MUST cover exactly `TelemetryProtoContract.HandledCases` — a `default: throw` on an unknown case makes +"new event added but not mapped" a loud runtime failure, and the Task-0 lock test makes it a compile-time-ish +guard. + +**Tests:** round-trip each domain record → `TelemetryProtoMap.ToProto` (Task 4) → `ToDomain` → assert +field-equality across all four kinds and all enum values (wire-fidelity, like ScadaBridge +`ProtoRoundtripTests`); an unmapped `EventCase` throws. + +**Step 5: Commit.** `feat(mesh-phase5): central telemetry dialer client + proto<->domain round-trip` + +--- + +### Task 8: Central dial supervisor actor (discovery + reconnect + feed sinks) + +**Classification:** high-risk +**Estimated implement time:** ~5 min +**Parallelizable with:** none (needs Task 1, Task 7) + +**Files:** +- Create: `src/Server/ZB.MOM.WW.OtOpcUa.ControlPlane/Telemetry/TelemetryDialSupervisor.cs` +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.ControlPlane.Tests/Telemetry/TelemetryDialSupervisorTests.cs` + +**Design — an admin-role actor holding one dialer per driver node, feeding the existing central sinks:** + +- **Discovery:** on `PreStart` and every `TelemetryDialOptions.ContactRefreshSeconds` (+ on an + admin-change signal if one is already wired for `CentralCommunicationActor`), read enabled, + non-maintenance `ClusterNode` rows selecting `NodeId, Host, GrpcPort` (extend the + `CentralCommunicationActor.cs:198-219` query). A row with `GrpcPort == null` is skipped with a Warning + (the node exposes no telemetry port — honest, not an error). Build `http://{Host}:{GrpcPort}` dial + targets; add dialers for new nodes, stop dialers for removed ones. +- **Per-node dialer loop** (mirror `SiteAlarmAggregatorActor.OpenGrpcStream`/`HandleGrpcError`, + simplified — one node per dialer, no NodeA/NodeB flip since each `ClusterNode` row is its own node): + - Open `TelemetryStreamClient.RunAsync`; `onEvent` posts each `TelemetryItem` to the supervisor which + routes to the sink (below); `onError` schedules a reconnect. + - **Reconnect:** immediate first retry, then fixed `_reconnectDelay` (5s) backoff. A **monotonic + generation stamp** per node makes late errors from a superseded stream ignorable. The dialer does + **not** die on repeated failure — it keeps retrying (observability, not data plane); log the first + failure + every Nth. + - A `correlationId` per (node, generation) — a safe id string, e.g. `central-{NodeId}-{gen}`. +- **Feed the sinks (the untouched seam):** inject the four singletons and route by item type: + - `Alarm` → `IInProcessBroadcaster.Publish(e)` + - `Script` → `IInProcessBroadcaster.Publish(e)` + - `Health` → `IDriverStatusSnapshotStore.Upsert(e)` + - `Resilience` → `IDriverResilienceStatusStore.Upsert(e)` +- **Connection indicator:** drive the two broadcasters' `SetConnected(...)` off aggregate stream health + — connected when **≥1** node stream is up, disconnected when all are down — matching today's + DPS `SubscribeAck`/`PostStop` pill semantics on `/alerts` and `/script-log`. (The two store-backed + panels have no connection flag today; leave their per-row staleness as-is — noted for Task 10 docs.) + +**Registration** is done in Task 9 (mode-gated). This task only builds + unit-tests the actor with a +fake `TelemetryStreamClient` factory and fake sinks. + +**Tests:** discovery adds/removes dialers as the (fake) `ClusterNode` set changes; a null `GrpcPort` +row is skipped; an emitted item of each kind lands on the correct sink; a stream error triggers a +scheduled reconnect with an incremented generation; a late event from a superseded generation is +dropped; broadcaster `IsConnected` flips true on first up / false on all-down. + +**Step 5: Commit.** `feat(mesh-phase5): central telemetry dial supervisor (discover + reconnect + feed sinks)` + +--- + +### Task 9: Wire the dark switch (central ingest source) + +**Classification:** standard +**Estimated implement time:** ~4 min +**Parallelizable with:** none (needs Task 8) + +**Files:** +- Modify: `src/Server/ZB.MOM.WW.OtOpcUa.AdminUI/Hubs/HubServiceCollectionExtensions.cs` (`WithOtOpcUaSignalRBridges`, lines 53-83) +- Modify: the AdminUI Akka configurator call site that invokes `WithOtOpcUaSignalRBridges` (locate; it's inside the `hasAdmin` branch per the method's own doc example) — pass or resolve the mode +- Test: `tests/Server/ZB.MOM.WW.OtOpcUa.AdminUI.Tests/Hubs/TelemetryModeWiringTests.cs` + +**Steps:** +- Inside `WithActors`, resolve `IOptions` and read `Mode`. +- **`Mode == Dps` (default):** spawn all five bridges exactly as today (unchanged behaviour). +- **`Mode == Grpc`:** spawn the **fleet-status bridge only** (out of Phase 5 scope, stays on DPS), and + do **NOT** spawn the four telemetry DPS bridges (alert, script-log, driver-status, resilience). + Instead spawn `TelemetryDialSupervisor` (Task 8), resolving the four sink singletons + the + `IDbContextFactory` (admin nodes have ConfigDb) + `IOptions`. +- The sinks (`AddOtOpcUaDriverStatusServices`, `:29-35`) are registered identically in both modes — no + change there; that is the whole point of the swap. + +**Guard rails:** +- Keep the fleet-status bridge on DPS in both modes (deferred). +- A `Grpc`-mode admin node still needs the sinks registered (it does — same call), so the AdminUI + components resolve them regardless. + +**Tests:** with `Mode=Dps`, the registry has the four telemetry bridge keys and no supervisor; with +`Mode=Grpc`, it has the supervisor and the fleet-status bridge but none of the four telemetry bridge +keys. (Use the Akka.Hosting `ActorRegistry` + a TestKit or the existing bridge-registration test +harness if one exists.) + +**Step 5: Commit.** `feat(mesh-phase5): dark-switch central telemetry ingest (Dps bridges | Grpc dialer)` + +--- + +### Task 10: Docs — Telemetry section, supersede §6.3, program status + +**Classification:** small +**Estimated implement time:** ~4 min +**Parallelizable with:** Task 11 (disjoint files) + +**Files (Modify/Create):** +- Create: `docs/Telemetry.md` — the stream architecture (direction, dark switch, auth, the four channels, deferred three + why, reconnect story, per-panel connection indicator status). +- Modify: `docs/Configuration.md` — add `Telemetry` + `TelemetryDial` sections (keys table, the `Mode` dark switch, port + shared-key `${secret:}` guidance). +- Modify: `docs/Redundancy.md` — a short "Command/observability transport" cross-ref noting redundancy-state stays DPS (pair-local) and telemetry moved to gRPC. +- Modify: `docs/plans/2026-07-21-per-cluster-mesh-design.md` §6.3 — supersede note: "telemetry stream authenticated from day one (fail-closed bearer), superseding the earlier 'unauthenticated for now'; ScadaBridge closed the same gap identically." +- Modify: `docs/plans/2026-07-22-per-cluster-mesh-program.md` — Phase 5 section + Tracking row: record the four-channel scope, the three deferrals + rationale, the dark switch, auth-from-day-one. +- Modify: `CLAUDE.md` — a "Live telemetry transport (`Telemetry`/`TelemetryDial`)" section mirroring the MeshTransport/ConfigSource sections (default `Dps`, node hosts server, central dials, four channels, deferred three, fail-closed key). + +**Step 5: Commit.** `docs(mesh-phase5): telemetry stream — Telemetry.md, config, supersede §6.3, program status` + +--- + +### Task 11: Rig config — telemetry ports + keys + flip env + +**Classification:** small +**Estimated implement time:** ~4 min +**Parallelizable with:** Task 10 (disjoint files); needs Task 6 + Task 9 landed for the flip to mean anything + +**Files:** +- Modify: `docker-dev/docker-compose.yml` +- Modify: any committed `appsettings.json` defaults that must carry the (default-OFF) keys — verify none are needed since `Mode` defaults to `Dps` and port defaults to 0. + +**Steps:** +- Give every **driver-role** node (central-1/2 fused + the four site nodes) `Telemetry__GrpcListenPort` + (pick a non-colliding port — verify against 4053 Akka / 4055 ConfigServe / the LocalDb sync port; use + e.g. **4056**) and `Telemetry__ApiKey: "telemetry-docker-dev-key"` (committed dev-secret exception, like + `configserve-docker-dev-key`). +- Give the **central (admin)** nodes `TelemetryDial__ApiKey: "telemetry-docker-dev-key"` (matching). +- Add each driver node's `GrpcPort` to its `ClusterNode` seed row (the rig's SQL seed / migration seed) so + central discovers the telemetry endpoint — mirror how `AkkaPort` was seeded in Phase 1. +- Leave `Telemetry__Mode`/`TelemetryDial__Mode` unset (⇒ `Dps` default). Document the flip in the compose + header: set `Telemetry__Mode=Grpc` on the driver nodes **and** `TelemetryDial__Mode=Grpc` on the central + nodes at `docker compose up`, then recreate. +- `lmxopcua-fix`/rig note: this is the local `docker-dev` rig, not the shared driver-fixture host. + +**Step 5: Commit.** `chore(mesh-phase5): rig telemetry ports + shared key + ClusterNode.GrpcPort seed` + +--- + +### Task 12: Live gate + +**Classification:** high-risk (gate, not a code change — but the phase is not done until it passes) +**Estimated implement time:** live rig run (not subagent wall-time) +**Parallelizable with:** none (needs everything) + +**Files:** +- Create: `docs/plans/2026-07-23-mesh-phase5-live-gate.md` (record) + +**Procedure (on the local `docker-dev` rig, matching prior phases' gates):** +1. Build the Phase-5 image (`--platform linux/amd64` build stage; `--progress=plain`) and recreate the rig in **default `Dps` mode** — confirm all six nodes up, `/alerts`, `/script-log`, `/hosts` panels live exactly as before (regression check: DPS path unchanged). +2. Flip: set `Telemetry__Mode=Grpc` on all driver nodes and `TelemetryDial__Mode=Grpc` on the central nodes; `docker compose up` (recreate). +3. **Exit gate leg A — panels green over gRPC:** drive a scripted-alarm transition and a script-log emission and a driver health/resilience change; confirm `/alerts`, `/script-log`, and the `/hosts` driver table update live — **with the four DPS telemetry bridges NOT spawned** (grep the central logs to confirm the Grpc branch ran and the four bridge keys are absent). Confirm the `/alerts` + `/script-log` connection pill reads "live". +4. **Exit gate leg B — kill-and-reconnect:** `docker kill` a site driver node (or drop its telemetry port), watch the central dialer log the error + retries and the pill flip to disconnected; restart the node and confirm the dialer reconnects, the driver-health/resilience stores re-prime from the node hub's snapshot replay, and the pill returns to live — **every stream recovers**. +5. **Optionally** confirm the auth gate: a dial with a wrong `TelemetryDial:ApiKey` gets `PermissionDenied` and no panel data (a deliberate misconfig probe). +6. Record all legs + evidence in the gate doc; note any deviations (as Phases 1-4 gates did). + +**Exit gate (program):** all AdminUI live panels green against the rig with the telemetry DPS bridges +off; kill-and-reconnect of the central dialer recovers every stream. + +**Step: Commit** the gate record; then finish per superpowers-extended-cc:finishing-a-development-branch +(merge `feat/mesh-phase5` to master locally, push, update the scadaproj umbrella index, update memory). + +--- + +## Risks / watch-items carried into execution + +- **Single-mesh correctness of the node-local hub.** The hub must carry only *this node's* events + (Task 2) — never subscribe it to cluster DPS, or on the current single mesh every node would stream + every other node's events and central would double-count. Snapshot replay + drop-oldest come from the + hub, not the service. +- **Do not gate/remove the DPS publish** at the four seams (Task 3) — that is the `Dps`-mode path and + the pure-dark-switch guarantee. The hub is a no-op until a client connects. +- **Interceptor single-ctor pin** (Task 5) — `Grpc.AspNetCore` silently disables an interceptor with + >1 public ctor. +- **Kestrel "re-bind once" invariant** (Task 6) — the third dedicated port is just one more + `ListenAnyIP`; the existing-surface re-bind still happens exactly once. Don't duplicate it. +- **Frame size is a non-issue** — server-streaming telemetry events are small and never cross Akka + remoting; the 128 KB ClusterClient frame trap does not apply to this path. +- **Additive-only proto evolution** — new event kinds/fields only; the Task-0 lock test + the + Task-7 `default: throw` are the guard. Never renumber a tag. +- **Deferred-channel honesty** — redundancy-state / fleet-status / deployment-acks are out by design, + documented in Task 10; a future reviewer must not read "seven topics" in the program doc and think + Phase 5 missed three. diff --git a/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md.tasks.json b/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md.tasks.json new file mode 100644 index 00000000..92fd68f0 --- /dev/null +++ b/docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md.tasks.json @@ -0,0 +1,22 @@ +{ + "planPath": "docs/plans/2026-07-23-mesh-phase5-grpc-telemetry-stream.md", + "program": "per-cluster-mesh", + "phase": 5, + "branch": "feat/mesh-phase5", + "tasks": [ + {"id": 0, "subject": "Task 0: telemetry.proto contract + codegen + contract-lock test", "classification": "standard", "status": "pending", "parallelizableWith": [1, 2]}, + {"id": 1, "subject": "Task 1: Telemetry/TelemetryDial options + fail-closed validators", "classification": "small", "status": "pending", "parallelizableWith": [0, 2]}, + {"id": 2, "subject": "Task 2: node-local telemetry hub (snapshot-replay + drop-oldest)", "classification": "standard", "status": "pending", "parallelizableWith": [0, 1]}, + {"id": 3, "subject": "Task 3: tap the 4 publish seams into the hub (DPS intact)", "classification": "standard", "status": "pending", "blockedBy": [2]}, + {"id": 4, "subject": "Task 4: node-side TelemetryStreamService (hub -> server-streaming)","classification": "high-risk", "status": "pending", "blockedBy": [0, 2]}, + {"id": 5, "subject": "Task 5: fail-closed bearer interceptor for the telemetry stream", "classification": "small", "status": "pending", "blockedBy": [1]}, + {"id": 6, "subject": "Task 6: host the telemetry gRPC server on driver nodes (Kestrel)", "classification": "standard", "status": "pending", "blockedBy": [1, 4, 5]}, + {"id": 7, "subject": "Task 7: central telemetry dialer client + proto<->domain round-trip","classification": "high-risk","status": "pending", "blockedBy": [0]}, + {"id": 8, "subject": "Task 8: central telemetry dial supervisor (discover+reconnect+sinks)","classification": "high-risk","status": "pending", "blockedBy": [1, 7]}, + {"id": 9, "subject": "Task 9: dark-switch central telemetry ingest (Dps bridges | Grpc)", "classification": "standard", "status": "pending", "blockedBy": [8]}, + {"id": 10, "subject": "Task 10: docs — Telemetry.md, config, supersede §6.3, program status","classification": "small", "status": "pending", "parallelizableWith": [11], "blockedBy": [9]}, + {"id": 11, "subject": "Task 11: rig telemetry ports + shared key + ClusterNode.GrpcPort seed","classification": "small", "status": "pending", "parallelizableWith": [10], "blockedBy": [6, 9]}, + {"id": 12, "subject": "Task 12: live gate (panels green over gRPC + kill-and-reconnect)", "classification": "high-risk","status": "pending", "blockedBy": [10, 11]} + ], + "lastUpdated": "2026-07-23T00:00:00Z" +}