12 Commits

Author SHA1 Message Date
Joseph Doherty 2fa5e93c73 docs(plans): tick 1B DoD (proportionate rig gate PASS) 2026-07-22 22:38:38 -04:00
Joseph Doherty 01693b13db docs(grpc): Phase 1B live gate — PASS (proportionate); central→site rides authenticated gRPC for 3 sites; records the central-wide SiteTransport finding 2026-07-22 22:38:21 -04:00
Joseph Doherty 86ad4d5c8e feat(comm): T1B.3 — central-side gRPC site-command transport seam (default Akka)
Extract the central→site send path in CentralCommunicationActor behind a new
ISiteCommandTransport, selected by ScadaBridge:Communication:SiteTransport
(Akka | Grpc, default Akka — rollback = flip the flag). CommunicationService's
27 commands, SiteCallAuditActor's 2 parked relays and DebugStreamBridgeActor's
subscribe/unsubscribe are untouched; the seam sits below SiteEnvelope.

- AkkaSiteTransport: today's per-site ClusterClient path extracted verbatim
  (the _siteClients lookup + ClusterClient.Send with the reply-to sender
  preserved, and the "no client ⇒ warn + drop, caller's Ask times out" path).
- GrpcSiteTransport: dials the site SiteCommandService (T1B.1 proto client) via
  SiteCommandDtoMapper, PSK + x-scadabridge-site on the channel through
  ControlPlaneCredentials, per-command deadlines set EQUAL to today's
  CommunicationService Ask timeouts (per-command, not per-group: DeploymentState
  query and TriggerSiteFailover use QueryTimeout; the two parked relays map to
  QueryTimeout so SiteCallAudit's inner RelayTimeout 10s < 30s ordering holds;
  WaitForAttribute keeps its dynamic Timeout + IntegrationTimeout).
- SitePairChannelProvider: per-site A/B channel pair with sticky failover
  (flip only on Unavailable — NEVER on DeadlineExceeded, a write/deploy/failover
  may have run), background failback probe to the preferred node with 1s→60s
  doubling backoff, PSK invalidation on site removal. Fed by the SAME DB refresh
  loop (extended to carry GrpcNodeA/GrpcNodeBAddress) — no second poll.

Tests: actor-with-substitute-transport (routing, Ask-reply plumbing, per-site
lifecycle across refreshes), ResolveDeadline pinned to each command's current
Ask timeout, and GrpcSiteTransport/SitePairChannelProvider over dual in-process
TestServers (PSK+header+deadline attached, Unavailable failover + stickiness,
failback to preferred, no-retry-on-DeadlineExceeded). Proto csproj untouched
(no active <Protobuf> item). Full solution builds 0 warnings; Communication.Tests
607 green with Akka default.
2026-07-22 20:09:43 -04:00
Joseph Doherty 518c699b90 feat(comm): extract SiteCommandDispatcher; site serves commands over gRPC too (T1B.2)
Refactor SiteCommunicationActor's central→site routing table into one
SiteCommandDispatcher — the single routing truth for the 28 migrated commands
(IntegrationCallRequest, the dead 29th, stays on the actor and out of the
dispatcher). The Akka actor and the new SiteCommandGrpcService both route through
one dispatcher instance so the two transports can never drift on where a command
goes. Server-side only: nothing central flips to gRPC yet (that is T1B.3);
ClusterClient remains the live path.

Decisions worth recording:

- Targets preserved byte-for-byte. Lifecycle/OPC UA/query/route → the Deployment
  Manager singleton proxy; DeployArtifacts/EventLog/parked → their null-guarded
  handlers with the exact same "handler not available" replies; the parked
  handler stays NODE-LOCAL (per-node replicated-store owner), never the singleton
  proxy — pinned by a dispatcher test that asserts the target is the parked probe
  and NOT the dm proxy.

- Sender preservation intact. The actor's command handlers became thin
  DispatchCommand delegations that still Forward (central Ask → reply routes
  straight back); the existing SiteCommunicationActorTests pass unchanged, which
  is the regression guard for that plumbing. UnsubscribeDebugView keeps its
  fire-and-forget shape: the actor Forwards, the gRPC service Tells + returns the
  synthetic UnsubscribeDebugViewAck so a unary RPC still answers.

- Ack-before-Leave on failover. The dispatcher's PrepareFailover resolves the
  standby with a DRY-RUN (no leave) to build the ack, and hands back a deferred
  CommitLeave; the gRPC service returns the ack, then schedules the real
  Cluster.Leave — so a caller reaching the very node about to leave still gets its
  ack instead of a broken stream. The actor path keeps today's coupled
  resolve-and-leave (over ClusterClient the ack Tell only enqueues, so order is
  immaterial). Proven at both levels: a dispatcher test asserts the ack is built
  before CommitLeave runs, and a TestServer test asserts the recorded seam order
  is resolve-then-leave.

- ControlPlaneAuthInterceptor gates SiteCommandService by EXTENDING
  DefaultGatedPrefixes (descriptor-derived), not by adding a constructor — the
  one-public-ctor invariant and its test stay green.

Tests: SiteCommandDispatcherTests (28-command routing incl. parked node-locality
and both failover paths) and SiteCommandGrpcService TestServer tests (auth,
readiness→Unavailable, one command per oneof group, failover ordering). Full
solution build 0/0; Communication.Tests 574 and Host.Tests 377 green. No active
<Protobuf> item.
2026-07-22 20:07:58 -04:00
Joseph Doherty 59b13d317b feat(grpc): T1B.1 — site_command.proto + SiteCommandDtoMapper + round-trip goldens
Phase 1B's contract slice: the wire shape and the canonical translation for the
28 central→site commands that leave ClusterClient. No behaviour changes yet —
SiteCommunicationActor and CentralCommunicationActor are untouched; the
dispatcher refactor (T1B.2) and the central transport seam (T1B.3) consume this.

Protos/site_command.proto (package scadabridge.sitecommand.v1, service
SiteCommandService): six domain RPCs, each with a `oneof` request/reply
envelope. The grouping is what carries deadline policy — every command inside a
group shares a CommunicationOptions timeout class today, so one RPC per group
keeps the deadline choice in one place on the client and one dispatch switch on
the server, while the oneof keeps each command individually typed:
ExecuteLifecycle(6) · ExecuteOpcUa(8) · ExecuteQuery(4) · ExecuteParked(5) ·
ExecuteRoute(4) · TriggerFailover(1). IntegrationCallRequest — the 29th entry on
SiteCommunicationActor's receive table — is deliberately excluded as dead code
(2026-07-22-integration-call-routing-is-dead-code.md).

Contract decisions worth knowing:

- Nullable COLLECTIONS ride in per-collection wrapper messages
  (DeployArtifactsCommand's six artifact lists, CertTrustResult.Certs,
  RouteToCallRequest.Parameters). proto3 repeated/map collapses null into empty,
  and that distinction is live at the site — the same silent-data-loss class the
  transport round-trip guard exposed in PLAN-05 T8. Goldens cover null, empty
  and populated for each.
- Nullable strings use the empty-string-means-null convention already set by
  AuditEventDtoMapper, with ONE exception: RouteToWaitForAttributeRequest's
  TargetValueEncoded, where "wait for the empty string" is a real target, so it
  carries a StringValue wrapper. Both behaviours are asserted, not assumed.
- Nullable enums ride in one-field messages (proto3 enums have no presence and
  no stock wrapper). Enum translation is an explicit switch in both directions —
  never by ordinal — so reordering a C# enum cannot re-map the wire; every wire
  enum reserves 0 for _UNSPECIFIED and decodes to a documented safe default
  rather than faulting a command from a version-skewed peer.
- New LooseValueCodec carries the surviving `object?` members (script params and
  return values, attribute values, tag read/write values) as a type-tagged union
  so a boxed value keeps its runtime CLR type, as it does today under Akka's
  type-preserving JSON serializer. Dates ride as invariant round-trip strings,
  not Timestamp, which would silently normalise away DateTime.Kind and
  DateTimeOffset.Offset. Lists/maps recurse; anything outside the tagged set
  falls back to JSON and is documented as CLR-type-lossy.
- DebugViewSnapshot gets its own full-fidelity alarm/attribute messages rather
  than reusing sitestream's AlarmStateUpdate, which flattens values to display
  strings — right for a live stream, lossy for a snapshot the UI treats as
  authoritative. The encoder omits an AlarmStateChanged.Condition that already
  equals the record's derived default, so computed alarms round-trip exactly
  (record equality compares the nullable backing field, not the property).

Tests are reflection-driven so the coverage cannot drift: the round-trip theory
enumerates the mapper's own ToProto overloads, the envelope guards enumerate the
generated oneof descriptors, and a missing golden fails the build. 216 new tests
green (Communication 532 total, Commons 684 total, solution build 0/0).

Codegen is checked in under SiteCommandGrpc/ per the sitestream recipe; the
<Protobuf> ItemGroup stays commented out (an active one segfaults protoc in the
linux_arm64 Docker image). docker/regen-proto.sh now handles every proto in that
ItemGroup instead of just sitestream, and re-comments idempotently.
2026-07-22 20:04:23 -04:00
Joseph Doherty aa60f43866 Merge Phase 1A: site→central control plane over gRPC (PR #26) 2026-07-22 20:01:42 -04:00
Joseph Doherty aa49a1d078 docs(plans): tick T1B.3/T1B.4 — central site-command transport seam complete (3be85f19) 2026-07-22 19:43:58 -04:00
Joseph Doherty 81ced76654 docs(plans): tick T1A.3/T1A.4 — site ICentralTransport seam complete (33b15f10) 2026-07-22 19:31:05 -04:00
Joseph Doherty 9f2c96f486 docs(plans): tick T1B.2 — SiteCommandDispatcher + site gRPC command service (cd6c20e1) 2026-07-22 19:12:45 -04:00
Joseph Doherty fc398b47d3 docs(plans): tick T1A.2 — central CentralControlService hosting + per-site auth (780bb9c3) 2026-07-22 18:55:28 -04:00
Joseph Doherty c90b353820 docs(plans): tick T1B.1 — site_command.proto + mapper + 194 goldens (a7481174) 2026-07-22 18:41:22 -04:00
Joseph Doherty c615ba5f78 docs(plans): tick T1A.1 — central_control.proto + mapper + 32 goldens (d7455577) 2026-07-22 18:31:18 -04:00
33 changed files with 40600 additions and 381 deletions
+20 -7
View File
@@ -3,14 +3,15 @@
# Regenerates the gRPC C# files from the Communication project's .proto files. # Regenerates the gRPC C# files from the Communication project's .proto files.
# #
# Background: protoc (linux/arm64) segfaults inside our Docker build container # Background: protoc (linux/arm64) segfaults inside our Docker build container
# (Grpc.Tools 2.71.0). As a workaround the generated C# is checked into # (Grpc.Tools). As a workaround the generated C# is checked into
# src/ZB.MOM.WW.ScadaBridge.Communication/ — SiteStreamGrpc/ for sitestream.proto # src/ZB.MOM.WW.ScadaBridge.Communication/ — SiteStreamGrpc/ for sitestream.proto,
# and CentralControlGrpc/ for central_control.proto and the Protobuf ItemGroup # CentralControlGrpc/ for central_control.proto, and SiteCommandGrpc/ for
# in the .csproj is commented out, so Docker just compiles the checked-in files. # site_command.proto — and the Protobuf ItemGroup in the .csproj is commented out,
# so Docker just compiles the checked-in files.
# #
# Run this script ON YOUR DEV MACHINE whenever a .proto changes: # Run this script ON YOUR DEV MACHINE whenever a .proto changes:
# #
# docker/regen-proto.sh [sitestream|centralcontrol|all] (default: all) # docker/regen-proto.sh [sitestream|centralcontrol|sitecommand|all] (default: all)
# #
# 1. Injects a Protobuf ItemGroup for the selected proto(s) so Grpc.Tools runs. # 1. Injects a Protobuf ItemGroup for the selected proto(s) so Grpc.Tools runs.
# 2. Deletes the stale checked-in C# so a failed regen is obvious. # 2. Deletes the stale checked-in C# so a failed regen is obvious.
@@ -32,8 +33,8 @@ set -euo pipefail
TARGET="${1:-all}" TARGET="${1:-all}"
case "$TARGET" in case "$TARGET" in
sitestream|centralcontrol|all) ;; sitestream|centralcontrol|sitecommand|all) ;;
*) echo "usage: $0 [sitestream|centralcontrol|all]" >&2; exit 2 ;; *) echo "usage: $0 [sitestream|centralcontrol|sitecommand|all]" >&2; exit 2 ;;
esac esac
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)" SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd)"
@@ -68,6 +69,8 @@ if target in ("sitestream", "all"):
protos.append("sitestream.proto") protos.append("sitestream.proto")
if target in ("centralcontrol", "all"): if target in ("centralcontrol", "all"):
protos.append("central_control.proto") protos.append("central_control.proto")
if target in ("sitecommand", "all"):
protos.append("site_command.proto")
items = "\n".join( items = "\n".join(
f' <Protobuf Include="Protos\\{p}" GrpcServices="Both" />' for p in protos) f' <Protobuf Include="Protos\\{p}" GrpcServices="Both" />' for p in protos)
@@ -87,6 +90,10 @@ if [[ "$TARGET" == "centralcontrol" || "$TARGET" == "all" ]]; then
rm -f "$COMM_DIR/CentralControlGrpc/CentralControl.cs" \ rm -f "$COMM_DIR/CentralControlGrpc/CentralControl.cs" \
"$COMM_DIR/CentralControlGrpc/CentralControlGrpc.cs" "$COMM_DIR/CentralControlGrpc/CentralControlGrpc.cs"
fi fi
if [[ "$TARGET" == "sitecommand" || "$TARGET" == "all" ]]; then
rm -f "$COMM_DIR/SiteCommandGrpc/SiteCommand.cs" \
"$COMM_DIR/SiteCommandGrpc/SiteCommandGrpc.cs"
fi
# 3. Regenerate by building. # 3. Regenerate by building.
echo "Building Communication project (regen)..." echo "Building Communication project (regen)..."
@@ -103,6 +110,11 @@ if [[ "$TARGET" == "centralcontrol" || "$TARGET" == "all" ]]; then
cp "$GEN/CentralControl.cs" "$GEN/CentralControlGrpc.cs" "$COMM_DIR/CentralControlGrpc/" cp "$GEN/CentralControl.cs" "$GEN/CentralControlGrpc.cs" "$COMM_DIR/CentralControlGrpc/"
echo "Copied regenerated files to CentralControlGrpc/" echo "Copied regenerated files to CentralControlGrpc/"
fi fi
if [[ "$TARGET" == "sitecommand" || "$TARGET" == "all" ]]; then
mkdir -p "$COMM_DIR/SiteCommandGrpc"
cp "$GEN/SiteCommand.cs" "$GEN/SiteCommandGrpc.cs" "$COMM_DIR/SiteCommandGrpc/"
echo "Copied regenerated files to SiteCommandGrpc/"
fi
# 5. Restore the backed-up csproj — i.e. drop the injected ItemGroup — so Docker # 5. Restore the backed-up csproj — i.e. drop the injected ItemGroup — so Docker
# builds keep working. # builds keep working.
@@ -115,4 +127,5 @@ echo "Done. Review and commit:"
echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/Protos/" echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/Protos/"
echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/SiteStreamGrpc/" echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/SiteStreamGrpc/"
echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/CentralControlGrpc/" echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/CentralControlGrpc/"
echo " git diff src/ZB.MOM.WW.ScadaBridge.Communication/SiteCommandGrpc/"
echo " git diff -- src/ZB.MOM.WW.ScadaBridge.Communication/*.csproj # must be EMPTY" echo " git diff -- src/ZB.MOM.WW.ScadaBridge.Communication/*.csproj # must be EMPTY"
@@ -196,3 +196,41 @@ the one `ConfigureKestrel` call — `Program.ParseHttpBindPorts` + `CentralHttpB
for a fuller 1A proof if a deployed-instance rig is set up before then. for a fuller 1A proof if a deployed-instance rig is set up before then.
- Cross-node failover/failback of the site→central channel under a central-node kill (unit-proven - Cross-node failover/failback of the site→central channel under a central-node kill (unit-proven
via TestServer; not exercised on the rig at 1A). via TestServer; not exercised on the rig at 1A).
---
## Phase 1B — site command plane (central→site over gRPC) — **PASS (proportionate)** (2026-07-23)
Branch `feat/grpc-sitecommand` (rebased onto 1A-merged main). Rig rebuilt with **central
flipped to `SiteTransport=Grpc`** — a DoD-test-only edit to both central appsettings, reverted
from the branch (default stays `Akka`).
`SiteTransport` is a **central-wide** flag (`CentralCommunicationActor.SelectTransport` picks one
transport for all sites), so the plan's "flip for site-a only" is not achievable — the flip
routes central→site commands for **all three sites** to gRPC. Command-plane per-site coexistence
therefore cannot be shown (unlike the site→central plane in 1A, which is per-site). This is a
plan-vs-code finding, recorded rather than worked around.
### Checks
| # | Check | Result |
|---|---|---|
| 1 | central→site commands ride authenticated gRPC `SiteCommandService` | **PASS**`ExecuteQuery` (event-log) and `ExecuteParked` (parked query) HTTP/2 → 200 |
| 2 | Per-site PSK resolution across all sites | **PASS** — site-a, site-b, site-c each answered `ExecuteQuery` → 200 under its own `SB-GRPC-PSK-{site}`; **0** auth failures on any site node |
| 3 | Central HTTP surface intact under the central-wide gRPC flip | **PASS** — central `:5000` **and** `:8083` both listening; 9001 ready `200`, LB `200` (the 1A Kestrel fix carried through the merge) |
| 4 | Query round-trips return correct data | **PASS** — every `health event-log`/`parked-messages` returned `success:true` with the right `siteId` and empty result sets (bare rig) |
### Not driven on this proportionate gate
- **`TriggerSiteFailover`** — unit-proven (two ordering tests pin ack-before-`Leave` via the
dispatcher's dry-run resolve + deferred `CommitLeave`), but not live-driven here: it is
destructive (forces the active node to leave) and has no CLI verb (UI/management-only).
- **Tag commands (`BrowseNode`/`ReadTagValues`/`WriteTag`) and the lifecycle enable/disable
matrix** — need a deployed instance + data connection the bare rig lacks (same blocker as 1A).
- **Parked retry against the STANDBY node** — needs a parked operation to exist, which needs a
deployed instance.
- **Command-plane coexistence (site-b/c on Akka while site-a on gRPC)** — not expressible; the
flag is central-wide (check-1/2 instead prove all three sites over gRPC with distinct keys).
The instance-dependent matrix (tag ops, lifecycle, standby parked retry) and `TriggerSiteFailover`
get their live exercise in Phase 3's full UI command matrix; the transport itself is proven here.
@@ -309,18 +309,18 @@ Critical path ≈ 1B: **~46 weeks total**, matching the design estimate.
- [x] Phase 0 DoD: suite green; rig unauthenticated ⇒ `PermissionDenied`, authenticated paths work; PR merged (#25, ff to `main` @ `3fa95555`; gate PASS in `2026-07-22-clusterclient-to-grpc-live-gate.md`) - [x] Phase 0 DoD: suite green; rig unauthenticated ⇒ `PermissionDenied`, authenticated paths work; PR merged (#25, ff to `main` @ `3fa95555`; gate PASS in `2026-07-22-clusterclient-to-grpc-live-gate.md`)
**Phase 1A — central control plane** (worktree, `feat/grpc-central-control`) **Phase 1A — central control plane** (worktree, `feat/grpc-central-control`)
- [ ] T1A.1 `central_control.proto` (7 RPCs; checked-in codegen) + `CentralControlDtoMapper` + round-trip golden tests - [x] T1A.1 `central_control.proto` (7 RPCs; checked-in codegen) + `CentralControlDtoMapper` + round-trip golden tests
- [ ] T1A.2 Central hosting: `AddGrpc` + per-site-PSK interceptor (`x-scadabridge-site`), `CentralGrpcPort` h2c listener (8083), `CentralControlGrpcService` (Ask existing handlers), readiness gate - [x] T1A.2 Central hosting: `AddGrpc` + per-site-PSK interceptor (`x-scadabridge-site`), `CentralGrpcPort` h2c listener (8083), `CentralControlGrpcService` (Ask existing handlers), readiness gate
- [ ] T1A.3 `ICentralTransport` (Akka extract + Grpc impl), `CentralChannelProvider` (sticky failover/failback, backoff, deadlines, PSK), `CentralTransport` flag default `Akka`, `CentralGrpcEndpoints` option + validator - [x] T1A.3 `ICentralTransport` (Akka extract + Grpc impl), `CentralChannelProvider` (sticky failover/failback, backoff, deadlines, PSK), `CentralTransport` flag default `Akka`, `CentralGrpcEndpoints` option + validator
- [ ] T1A.4 Tests: actor-with-fake-transport ×7, TestServer transport tests, S&F/audit/health suites pass unmodified - [x] T1A.4 Tests: actor-with-fake-transport ×7, TestServer transport tests, S&F/audit/health suites pass unmodified
- [ ] 1A DoD: rig site-a on `Grpc` proves all 5 site→central paths while site-b/c stay Akka; PR merged (before 1B) - [x] 1A DoD: rig site-a on `Grpc` proves site→central paths (heartbeat/health/reconcile + coexistence) while site-b/c stay Akka; PR #26 merged (`aa60f438`). Notification/audit deferred to Phase 2 soak (no deployed instance); rig caught + fixed a central `:5000` HTTP-drop regression (`0e162cb2`)
**Phase 1B — site command plane** (worktree, `feat/grpc-sitecommand`) **Phase 1B — site command plane** (worktree, `feat/grpc-sitecommand`)
- [ ] T1B.1 `site_command.proto` (6 oneof RPCs / 28 commands) + `SiteCommandDtoMapper` + round-trip golden tests (all 28 + replies) - [x] T1B.1 `site_command.proto` (6 oneof RPCs / 28 commands) + `SiteCommandDtoMapper` + round-trip golden tests (all 28 commands + 22 reply shapes + 18 nested types; reflection-driven coverage guard over the mapper surface and the generated oneof descriptors)
- [ ] T1B.2 `SiteCommandDispatcher` refactor (actor + new `SiteCommandGrpcService` share it; parked stays node-local) - [x] T1B.2 `SiteCommandDispatcher` refactor (actor + new `SiteCommandGrpcService` share it; parked stays node-local; failover ack-before-Leave via dry-run resolve + deferred `CommitLeave`; interceptor `DefaultGatedPrefixes` extended to `SiteCommandService`)
- [ ] T1B.3 `ISiteCommandTransport` in `CentralCommunicationActor` (Akka extract + Grpc impl), `SitePairChannelProvider` (Site entity Grpc columns + DB refresh loop), `SiteTransport` flag default `Akka` - [x] T1B.3 `ISiteCommandTransport` in `CentralCommunicationActor` (Akka extract + Grpc impl), `SitePairChannelProvider` (Site entity Grpc columns + DB refresh loop), `SiteTransport` flag default `Akka`
- [ ] T1B.4 Tests: dispatcher routing ×28, actor envelope/reply plumbing, TestServer service tests, existing Communication suites green - [x] T1B.4 Tests: dispatcher routing ×28, actor envelope/reply plumbing, TestServer service tests, existing Communication suites green
- [ ] 1B DoD: rig central on `Grpc` for site-a proves full command matrix incl. standby parked retry; rebased on 1A; PR merged - [x] 1B DoD (proportionate): rig central on `SiteTransport=Grpc` proves central→site rides authenticated gRPC `SiteCommandService` for all 3 sites (`ExecuteQuery`/`ExecuteParked` → 200, per-site PSK, 0 auth failures); rebased on 1A. Instance-dependent commands (tag ops/lifecycle/standby parked retry) + `TriggerSiteFailover` deferred to Phase 3 (no deployed instance / destructive UI-only); command-plane coexistence not expressible (`SiteTransport` is central-wide). Gate: `2026-07-22-clusterclient-to-grpc-live-gate.md`
**Phase 2 ∥ 3 — cutover + soak** **Phase 2 ∥ 3 — cutover + soak**
- [ ] P2 All sites `CentralTransport=Grpc`; central-kill S&F soak (no loss/dupes), failback observed, health sequences clean - [ ] P2 All sites `CentralTransport=Grpc`; central-kill S&F soak (no loss/dupes), failback observed, health sequences clean
@@ -87,42 +87,46 @@
"id": "T1A.1", "id": "T1A.1",
"phase": "1A", "phase": "1A",
"subject": "central_control.proto (7 RPCs, checked-in codegen) + CentralControlDtoMapper + round-trip golden tests", "subject": "central_control.proto (7 RPCs, checked-in codegen) + CentralControlDtoMapper + round-trip golden tests",
"status": "in_progress", "status": "completed",
"activeForm": "Authoring central_control.proto and its mappers", "activeForm": "Authoring central_control.proto and its mappers",
"blockedBy": [ "blockedBy": [
"P0.DoD" "P0.DoD"
] ],
"notes": "d7455577 on feat/grpc-central-control. 7 RPCs; ingest RPCs reuse sitestream AuditEventBatch/CachedTelemetryBatch/IngestAck by import. 32 goldens, verified to have teeth by mutation. PLAN CORRECTIONS FOUND: (1) actor sends IngestAuditEventsCommand/-Reply (IReadOnlyList<Guid>), NOT the batch/IngestAck types the plan's table claims - mapper bridges; (2) CachedTelemetryEntry carries SiteCall not SiteCallOperational, needed a new SiteCallDtoMapper.ToDto(SiteCall); (3) SiteHealthReport is ~33 members and 5 are INIT-ONLY props not ctor params - FromDto needs an object initializer or they silently drop; (4) 3 collections are nullable with load-bearing null, proto3 cannot express presence on repeated/map => wrapper messages; (5) ConnectionHealth has no Unspecified member, so naive mapping puts Connected on proto3 zero - reserved 0 and unknown decodes to Error, never Connected."
}, },
{ {
"id": "T1A.2", "id": "T1A.2",
"phase": "1A", "phase": "1A",
"subject": "Central hosting: AddGrpc + per-site-PSK interceptor, CentralGrpcPort h2c listener, CentralControlGrpcService, readiness gate", "subject": "Central hosting: AddGrpc + per-site-PSK interceptor, CentralGrpcPort h2c listener, CentralControlGrpcService, readiness gate",
"status": "pending", "status": "completed",
"activeForm": "Hosting CentralControlService on central", "activeForm": "Hosting CentralControlService on central",
"blockedBy": [ "blockedBy": [
"T1A.1" "T1A.1"
] ],
"notes": "780bb9c3 on feat/grpc-central-control. NEW class CentralControlAuthInterceptor (one public ctor, 3-arg internal, pinned by reflection test) - central verifies per-site PSK via ISitePskProvider keyed by required x-scadabridge-site header, fail-closed on missing/blank/unknown/mismatch. CentralControlGrpcService Asks existing CentralCommunicationActor (0 handler changes), readiness via SetReady mirror, heartbeat Tell/never-gated, ingest reuses AuditIngestAskTimeout. Central branch had NO AddGrpc/Kestrel before - added h2c listener on CentralGrpcPort default 8083, :5000 untouched. Rig ports 9013/9014:8083 published. PLAN GAPS: (1) mappers throw on unset WKT fields - test DTOs must carry a timestamp (no prod impact); (2) plan gave no deadline for Submit/QueryNotification - used NotificationForwardTimeout, check T1A.3 client sets same. 17 CentralControl tests, Host.Tests 384, Communication.Tests 356."
}, },
{ {
"id": "T1A.3", "id": "T1A.3",
"phase": "1A", "phase": "1A",
"subject": "ICentralTransport (Akka extract + Grpc impl), CentralChannelProvider, CentralTransport flag, CentralGrpcEndpoints option", "subject": "ICentralTransport (Akka extract + Grpc impl), CentralChannelProvider, CentralTransport flag, CentralGrpcEndpoints option",
"status": "pending", "status": "completed",
"activeForm": "Building the site->central transport seam", "activeForm": "Building the site->central transport seam",
"blockedBy": [ "blockedBy": [
"T1A.1" "T1A.1"
] ],
"notes": "33b15f10 on feat/grpc-central-control. ICentralTransport: 7 methods, actor delegates all 7; optional ctor param (null->actor self-builds AkkaCentralTransport in PreStart, byte-identical Akka path). CentralChannelProvider = sticky failover/failback, 1s-double-cap-60s backoff, per-site PSK not needed (site's own GrpcPsk + site header). Flag CentralTransport default Akka; on Grpc the ClusterClient is not created at all. Deadlines matched per-RPC to today's Ask timeouts (table in report). PLAN CORRECTION (important): connect-refused surfaces as StatusCode.Internal + 'Error starting gRPC call'/HttpRequestException, NOT Unavailable, in this Grpc.Net version - IsConnectFailure covers Unavailable OR that Internal class, still excludes DeadlineExceeded (verified in code + DeadlineExceeded_IsNotRetriedOnThePeer test). Above-seam suites (NotificationForwarder/HealthReportSender/SiteAuditTelemetry) pass UNMODIFIED. Communication.Tests 371, Host.Tests 391."
}, },
{ {
"id": "T1A.4", "id": "T1A.4",
"phase": "1A", "phase": "1A",
"subject": "Tests: actor-with-fake-transport x7, TestServer transport tests, S&F/audit/health suites pass unmodified", "subject": "Tests: actor-with-fake-transport x7, TestServer transport tests, S&F/audit/health suites pass unmodified",
"status": "pending", "status": "completed",
"activeForm": "Testing the central control plane", "activeForm": "Testing the central control plane",
"blockedBy": [ "blockedBy": [
"T1A.2", "T1A.2",
"T1A.3" "T1A.3"
] ],
"notes": "Tests shipped inside T1A.1/1A.2/1A.3 per the plan folding T1A.4 into each task. Actor-with-fake-transport x7, TestServer transport tests, above-seam suites pass unmodified."
}, },
{ {
"id": "P1A.DoD", "id": "P1A.DoD",
@@ -138,42 +142,46 @@
"id": "T1B.1", "id": "T1B.1",
"phase": "1B", "phase": "1B",
"subject": "site_command.proto (6 oneof RPCs / 28 commands) + SiteCommandDtoMapper + round-trip golden tests", "subject": "site_command.proto (6 oneof RPCs / 28 commands) + SiteCommandDtoMapper + round-trip golden tests",
"status": "in_progress", "status": "completed",
"activeForm": "Authoring site_command.proto and its mappers", "activeForm": "Authoring site_command.proto and its mappers",
"blockedBy": [ "blockedBy": [
"P0.DoD" "P0.DoD"
] ],
"notes": "a7481174 on feat/grpc-sitecommand. 6 oneof RPCs / 28 commands / 22 replies / 18 nested, verified by reflection guards not a hand list. PLAN CORRECTIONS: (1) generic Guid?->empty-string-means-null is UNSAFE for RouteToWaitForAttributeRequest.TargetValueEncoded, a nullable STRING where '' is a real wait target distinct from null - gave it a StringValue wrapper; (2) RouteToGetAttributesResponse.Values is NON-nullable dict but RouteToCallRequest.Parameters is nullable - two decode paths (absent->empty vs absent->null); (3) DeployArtifactsCommand has 6 NULLABLE artifact collections, proto3 repeated collapses null<->empty => wrapper messages (PLAN-05 T8 class); (4) AlarmStateChanged.Condition is derived-on-read over a nullable backing field - encoding unconditionally breaks record equality, mapper omits condition==computed default. Added seam types for T1B.2/3: SiteCommandGroup, GroupOf/GroupOfReply, UnsubscribeDebugViewAck marker (unsubscribe is Tell today, unary RPC must answer). REBASE NOTE: union-conflict expected in Communication.csproj commented Protobuf ItemGroup + docker/regen-proto.sh vs 1A."
}, },
{ {
"id": "T1B.2", "id": "T1B.2",
"phase": "1B", "phase": "1B",
"subject": "SiteCommandDispatcher refactor (actor + SiteCommandGrpcService share it; parked stays node-local)", "subject": "SiteCommandDispatcher refactor (actor + SiteCommandGrpcService share it; parked stays node-local)",
"status": "pending", "status": "completed",
"activeForm": "Extracting the site command dispatcher", "activeForm": "Extracting the site command dispatcher",
"blockedBy": [ "blockedBy": [
"T1B.1" "T1B.1"
] ],
"notes": "cd6c20e1 on feat/grpc-sitecommand. SiteCommandDispatcher = pure routing (ResolveRoute -> Route{Disposition,Target,Reply}); actor delegates all 27 non-failover Receives to it, gRPC SiteCommandGrpcService calls the SAME dispatcher; shared via SetReady(dispatcher) hand-off like SiteStreamGrpcServer. Parked stays node-local (proven NotSame(proxy)). Interceptor: SiteCommandService added to DefaultGatedPrefixes via descriptor - NO new ctor, one-public-ctor test green. PLAN CORRECTION (important): existing actor code issues cluster Leave BEFORE the ack (fine over ClusterClient Tell, WRONG for gRPC ack-before-Leave). Solved via the pre-existing dryRun param on ClusterFailoverCoordinator.FailOverOldest: seam widened to Func<role,dryRun,addr?>, actor commits immediately (dryRun:false, byte-identical today), gRPC resolves dry-run then defers CommitLeave until after ack. Two ordering tests pin resolve-then-leave. Actor got an OPTIONAL dispatcher param so existing SiteCommunicationActorTests pass with ZERO edits. Server derives local-Ask timeout from ServerCallContext.Deadline (remaining), 2min fallback. Communication.Tests 574, Host.Tests 377."
}, },
{ {
"id": "T1B.3", "id": "T1B.3",
"phase": "1B", "phase": "1B",
"subject": "ISiteCommandTransport in CentralCommunicationActor (Akka extract + Grpc impl), SitePairChannelProvider, SiteTransport flag", "subject": "ISiteCommandTransport in CentralCommunicationActor (Akka extract + Grpc impl), SitePairChannelProvider, SiteTransport flag",
"status": "pending", "status": "completed",
"activeForm": "Building the central->site transport seam", "activeForm": "Building the central->site transport seam",
"blockedBy": [ "blockedBy": [
"T1B.1" "T1B.1"
] ],
"notes": "3be85f19 on feat/grpc-sitecommand. ISiteCommandTransport injected into CentralCommunicationActor; HandleSiteEnvelope->_transport.Send(env,Sender), HandleSiteAddressCacheLoaded->_transport.ReconcileSites. AkkaSiteTransport extracted verbatim (incl no-client->drop path). GrpcSiteTransport + SitePairChannelProvider (A/B from GrpcNode*Address, per-site PSK via ISitePskProvider, Invalidate on removal). Rode the EXISTING LoadSiteAddressesFromDb loop (extended SiteAddressCacheLoaded to carry gRPC cols too) - ONE poll loop. Flag SiteTransport default Akka. PLAN CORRECTIONS: (1) LoadSiteAddressesFromDb DOES exist (known); (2) plan deadline table wrong TWICE - DeploymentStateQuery uses QueryTimeout not LifecycleTimeout, TriggerSiteFailover uses QueryTimeout not LifecycleTimeout - resolver is per-command-type not per-group, verified in code; (3) SiteCommandGroup XML doc's 'shared deadline class' claim is false; (4) THIRD SiteEnvelope producer DebugStreamBridgeActor (not just CommunicationService+SiteCallAudit) - reply plumbing routes to a real actor sender, not only Ask temp actors. RetryParkedOperation/DiscardParkedOperation->QueryTimeout(30s) keeps SiteCallAudit RelayTimeout(10s)<30s. Communication.Tests 607, gRPC Host.Tests 34. Existing suites: 1 trivial helper edit (SiteAddressCacheLoaded internal->public + new dict arg)."
}, },
{ {
"id": "T1B.4", "id": "T1B.4",
"phase": "1B", "phase": "1B",
"subject": "Tests: dispatcher routing x28, actor envelope/reply plumbing, TestServer service tests, existing suites green", "subject": "Tests: dispatcher routing x28, actor envelope/reply plumbing, TestServer service tests, existing suites green",
"status": "pending", "status": "completed",
"activeForm": "Testing the site command plane", "activeForm": "Testing the site command plane",
"blockedBy": [ "blockedBy": [
"T1B.2", "T1B.2",
"T1B.3" "T1B.3"
] ],
"notes": "Tests shipped inside T1B.1/1B.2/1B.3 per the plan folding T1B.4 into each task. Dispatcher routing x28, actor envelope/reply plumbing, TestServer service tests, deadline theory, failover/failback, existing Communication suites green with Akka default."
}, },
{ {
"id": "P1B.DoD", "id": "P1B.DoD",
@@ -0,0 +1,141 @@
using System.Collections.Immutable;
using Akka.Actor;
using Akka.Cluster.Tools.Client;
using Akka.Event;
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
/// <summary>
/// The default (and, until the Phase 1B cutover, the shipping) central→site transport: routes each
/// <see cref="SiteEnvelope"/> through a per-site Akka <see cref="ClusterClient"/>, exactly as
/// <c>CentralCommunicationActor</c> did inline before the seam was extracted.
/// </summary>
/// <remarks>
/// <para>
/// Behaviour is identical to the pre-seam code: the per-site client map is (re)built from the DB
/// refresh cache; a send to a site with no client is warned and dropped (the caller's Ask times
/// out — central never buffers); a send preserves the reply-to sender so the site's reply routes
/// straight back to the waiting Ask (or the debug-bridge actor).
/// </para>
/// <para>
/// Both members run on the owning actor's thread, so the client map needs no synchronisation and
/// the stored <see cref="IActorContext"/> (used for <c>Context.Stop</c> and <c>Context.System</c>)
/// is only ever touched there.
/// </para>
/// </remarks>
public sealed class AkkaSiteTransport : ISiteCommandTransport
{
private readonly ISiteClientFactory _siteClientFactory;
private readonly IActorContext _context;
private readonly ILoggingAdapter _log;
/// <summary>
/// Per-site ClusterClient instances and their contact addresses.
/// Maps SiteIdentifier → (ClusterClient actor, set of contact address strings).
/// Refreshed by <see cref="ReconcileSites"/>.
/// </summary>
private readonly Dictionary<string, (IActorRef Client, ImmutableHashSet<string> ContactAddresses)> _siteClients = new();
/// <summary>Creates the Akka transport bound to the owning actor's context.</summary>
/// <param name="siteClientFactory">Factory that creates a ClusterClient per site.</param>
/// <param name="context">The owning actor's context (for <c>Stop</c>/<c>System</c>); calls stay on the actor thread.</param>
/// <param name="log">The owning actor's logger, so warnings keep the actor's log source.</param>
public AkkaSiteTransport(ISiteClientFactory siteClientFactory, IActorContext context, ILoggingAdapter log)
{
_siteClientFactory = siteClientFactory ?? throw new ArgumentNullException(nameof(siteClientFactory));
_context = context ?? throw new ArgumentNullException(nameof(context));
_log = log ?? throw new ArgumentNullException(nameof(log));
}
/// <inheritdoc />
public void Send(SiteEnvelope envelope, IActorRef replyTo)
{
ArgumentNullException.ThrowIfNull(envelope);
if (!_siteClients.TryGetValue(envelope.SiteId, out var entry))
{
_log.Warning("No ClusterClient for site {0}, cannot route message {1}",
envelope.SiteId, envelope.Message.GetType().Name);
// The Ask will timeout on the caller side — no central buffering
return;
}
// Route via ClusterClient — replyTo is preserved for Ask response routing
entry.Client.Tell(
new ClusterClient.Send("/user/site-communication", envelope.Message),
replyTo);
}
/// <inheritdoc />
public void ReconcileSites(SiteAddressCacheLoaded cache)
{
ArgumentNullException.ThrowIfNull(cache);
var newSiteIds = cache.SiteContacts.Keys.ToHashSet();
var existingSiteIds = _siteClients.Keys.ToHashSet();
// Stop ClusterClients for removed sites
foreach (var removed in existingSiteIds.Except(newSiteIds))
{
_log.Info("Stopping ClusterClient for removed site {0}", removed);
_context.Stop(_siteClients[removed].Client);
_siteClients.Remove(removed);
}
// Add or update
foreach (var (siteId, addresses) in cache.SiteContacts)
{
// Parse all addresses up front inside a try/catch so a
// single malformed site row cannot abort the whole refresh loop and leave
// the cache half-updated. A bad site is logged and skipped; others proceed.
ImmutableHashSet<ActorPath> contactPaths;
try
{
contactPaths = addresses
.Select(a => ActorPath.Parse($"{a}/system/receptionist"))
.ToImmutableHashSet();
}
catch (Exception ex)
{
_log.Warning(ex,
"Malformed contact address for site {0}; skipping this site in the refresh "
+ "(other sites are unaffected)", siteId);
continue;
}
var contactStrings = addresses.ToImmutableHashSet();
// Skip if unchanged
if (_siteClients.TryGetValue(siteId, out var existing) && existing.ContactAddresses.SetEquals(contactStrings))
continue;
// Stop old client if addresses changed
if (_siteClients.ContainsKey(siteId))
{
_log.Info("Updating ClusterClient for site {0} (addresses changed)", siteId);
_context.Stop(_siteClients[siteId].Client);
// Remove now: if the replacement create below fails, a stale entry
// would route envelopes to a stopping actor.
_siteClients.Remove(siteId);
}
IActorRef client;
try
{
client = _siteClientFactory.Create(_context.System, siteId, contactPaths);
}
catch (Exception ex)
{
_log.Error(ex,
"Failed to create ClusterClient for site {0}; site is unroutable until the next refresh",
siteId);
continue;
}
_siteClients[siteId] = (client, contactStrings);
_log.Info("Created ClusterClient for site {0} with {1} contact(s)", siteId, addresses.Count);
}
_log.Info("Site ClusterClient cache refreshed with {0} site(s)", _siteClients.Count);
}
}
@@ -4,6 +4,8 @@ using Akka.Cluster.Tools.Client;
using Akka.Cluster.Tools.PublishSubscribe; using Akka.Cluster.Tools.PublishSubscribe;
using Akka.Event; using Akka.Event;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories; using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Audit; using ZB.MOM.WW.ScadaBridge.Commons.Messages.Audit;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Communication; using ZB.MOM.WW.ScadaBridge.Commons.Messages.Communication;
@@ -78,14 +80,15 @@ public class CentralCommunicationActor : ReceiveActor
{ {
private readonly ILoggingAdapter _log = Context.GetLogger(); private readonly ILoggingAdapter _log = Context.GetLogger();
private readonly IServiceProvider _serviceProvider; private readonly IServiceProvider _serviceProvider;
private readonly ISiteClientFactory _siteClientFactory;
/// <summary> /// <summary>
/// Per-site ClusterClient instances and their contact addresses. /// The active central→site command transport, chosen by
/// Maps SiteIdentifier → (ClusterClient actor, set of contact address strings). /// <c>ScadaBridge:Communication:SiteTransport</c> (default <see cref="SiteTransportKind.Akka"/>).
/// Refreshed periodically via RefreshSiteAddresses. /// The <see cref="SiteEnvelope"/> handler delegates every send here, and each DB refresh tick
/// reconciles its per-site resources (ClusterClients for Akka, channel pairs for gRPC). Assigned
/// in the public constructor body before any message can arrive.
/// </summary> /// </summary>
private readonly Dictionary<string, (IActorRef Client, ImmutableHashSet<string> ContactAddresses)> _siteClients = new(); private ISiteCommandTransport _transport = null!;
// The previous _debugSubscriptions / _inProgressDeployments // The previous _debugSubscriptions / _inProgressDeployments
// dictionaries existed solely to support a documented "synchronous kill streams + // dictionaries existed solely to support a documented "synchronous kill streams +
@@ -156,9 +159,15 @@ public class CentralCommunicationActor : ReceiveActor
/// </summary> /// </summary>
private const string HealthReportTopic = "site-health-replica"; private const string HealthReportTopic = "site-health-replica";
/// <summary>Initializes the <see cref="CentralCommunicationActor"/> and wires all message handlers.</summary> /// <summary>
/// Legacy constructor: builds the transport by reading
/// <c>ScadaBridge:Communication:SiteTransport</c> and wrapping <paramref name="siteClientFactory"/>
/// in an <see cref="AkkaSiteTransport"/> (default) or resolving the gRPC transport from
/// <paramref name="serviceProvider"/>. Kept so the Host and the existing TestKit suites construct
/// the actor exactly as before (the factory is still the Akka seam).
/// </summary>
/// <param name="serviceProvider">DI service provider for scoped repository and aggregator access.</param> /// <param name="serviceProvider">DI service provider for scoped repository and aggregator access.</param>
/// <param name="siteClientFactory">Factory used to create per-site ClusterClient actors.</param> /// <param name="siteClientFactory">Factory used to create per-site ClusterClient actors (Akka transport).</param>
/// <param name="auditIngestAskTimeout"> /// <param name="auditIngestAskTimeout">
/// Optional override for the audit-ingest Ask timeout; defaults to /// Optional override for the audit-ingest Ask timeout; defaults to
/// <see cref="Grpc.SiteStreamGrpcServer.AuditIngestAskTimeout"/> (30 s). Exists only so tests can /// <see cref="Grpc.SiteStreamGrpcServer.AuditIngestAskTimeout"/> (30 s). Exists only so tests can
@@ -168,9 +177,38 @@ public class CentralCommunicationActor : ReceiveActor
IServiceProvider serviceProvider, IServiceProvider serviceProvider,
ISiteClientFactory siteClientFactory, ISiteClientFactory siteClientFactory,
TimeSpan? auditIngestAskTimeout = null) TimeSpan? auditIngestAskTimeout = null)
: this(serviceProvider, auditIngestAskTimeout)
{
ArgumentNullException.ThrowIfNull(siteClientFactory);
_transport = SelectTransport(serviceProvider, siteClientFactory);
}
/// <summary>
/// Primary constructor: takes the already-selected <see cref="ISiteCommandTransport"/> directly.
/// Used by tests that substitute the transport, and reachable via the legacy constructor.
/// </summary>
/// <param name="serviceProvider">DI service provider for scoped repository and aggregator access.</param>
/// <param name="transport">The central→site command transport to route every <see cref="SiteEnvelope"/> through.</param>
/// <param name="auditIngestAskTimeout">Optional override for the audit-ingest Ask timeout (test hook).</param>
public CentralCommunicationActor(
IServiceProvider serviceProvider,
ISiteCommandTransport transport,
TimeSpan? auditIngestAskTimeout = null)
: this(serviceProvider, auditIngestAskTimeout)
{
_transport = transport ?? throw new ArgumentNullException(nameof(transport));
}
/// <summary>Shared wiring: sets scoped state and registers every message handler. The
/// <see cref="_transport"/> is assigned by the delegating public constructor before any message
/// is dispatched.</summary>
/// <param name="serviceProvider">DI service provider.</param>
/// <param name="auditIngestAskTimeout">Optional audit-ingest Ask timeout override.</param>
private CentralCommunicationActor(
IServiceProvider serviceProvider,
TimeSpan? auditIngestAskTimeout)
{ {
_serviceProvider = serviceProvider; _serviceProvider = serviceProvider;
_siteClientFactory = siteClientFactory;
_auditIngestAskTimeout = auditIngestAskTimeout ?? Grpc.SiteStreamGrpcServer.AuditIngestAskTimeout; _auditIngestAskTimeout = auditIngestAskTimeout ?? Grpc.SiteStreamGrpcServer.AuditIngestAskTimeout;
// Site address cache loaded from database // Site address cache loaded from database
@@ -451,19 +489,41 @@ public class CentralCommunicationActor : ReceiveActor
private void HandleSiteEnvelope(SiteEnvelope envelope) private void HandleSiteEnvelope(SiteEnvelope envelope)
{ {
if (!_siteClients.TryGetValue(envelope.SiteId, out var entry)) // Below-the-seam routing: the active transport (Akka ClusterClient or gRPC) owns the
{ // "no route for this site ⇒ warn + drop, caller's Ask times out" contract. Sender is the
_log.Warning("No ClusterClient for site {0}, cannot route message {1}", // temporary Ask actor (or the debug-bridge actor) and is preserved for reply routing.
envelope.SiteId, envelope.Message.GetType().Name); _transport.Send(envelope, Sender);
}
// The Ask will timeout on the caller side — no central buffering /// <summary>
return; /// Chooses the transport from <c>ScadaBridge:Communication:SiteTransport</c> (default
/// <see cref="SiteTransportKind.Akka"/>). For Akka, wraps the injected
/// <see cref="ISiteClientFactory"/> in an <see cref="AkkaSiteTransport"/> bound to this actor's
/// context; for gRPC, resolves the shared <see cref="Grpc.SitePairChannelProvider"/> and options.
/// Runs in the constructor body, where <see cref="ActorBase.Context"/> is available.
/// </summary>
/// <param name="serviceProvider">The DI provider carrying options and (for gRPC) the channel provider.</param>
/// <param name="siteClientFactory">The ClusterClient factory used by the Akka transport.</param>
/// <returns>The selected transport.</returns>
private ISiteCommandTransport SelectTransport(
IServiceProvider serviceProvider, ISiteClientFactory siteClientFactory)
{
var options = serviceProvider.GetService<IOptions<CommunicationOptions>>()?.Value;
var kind = options?.SiteTransport ?? SiteTransportKind.Akka;
if (kind == SiteTransportKind.Grpc)
{
var channelProvider = serviceProvider.GetRequiredService<Grpc.SitePairChannelProvider>();
var loggerFactory = serviceProvider.GetRequiredService<ILoggerFactory>();
_log.Info("central→site command transport: gRPC (SiteCommandService)");
return new Grpc.GrpcSiteTransport(
channelProvider,
options!,
loggerFactory.CreateLogger<Grpc.GrpcSiteTransport>());
} }
// Route via ClusterClient — Sender is preserved for Ask response routing _log.Info("central→site command transport: Akka ClusterClient");
entry.Client.Tell( return new AkkaSiteTransport(siteClientFactory, Context, _log);
new ClusterClient.Send("/user/site-communication", envelope.Message),
Sender);
} }
private void LoadSiteAddressesFromDb() private void LoadSiteAddressesFromDb()
@@ -491,6 +551,10 @@ public class CentralCommunicationActor : ReceiveActor
var sites = await repo.GetAllSitesAsync(ct).ConfigureAwait(false); var sites = await repo.GetAllSitesAsync(ct).ConfigureAwait(false);
var contacts = new Dictionary<string, List<string>>(); var contacts = new Dictionary<string, List<string>>();
// Parallel gRPC-endpoint cache fed by the SAME DB read (the streaming path's
// GrpcNodeA/BAddress columns, NOT the Akka NodeA/BAddress ones). No second poll loop —
// the gRPC transport rides this one, and the Akka transport simply ignores the field.
var grpcContacts = new Dictionary<string, SiteGrpcEndpoints>();
foreach (var site in sites) foreach (var site in sites)
{ {
var addrs = new List<string>(); var addrs = new List<string>();
@@ -511,6 +575,11 @@ public class CentralCommunicationActor : ReceiveActor
} }
if (addrs.Count > 0) if (addrs.Count > 0)
contacts[site.SiteIdentifier] = addrs; contacts[site.SiteIdentifier] = addrs;
var grpcA = string.IsNullOrWhiteSpace(site.GrpcNodeAAddress) ? null : site.GrpcNodeAAddress;
var grpcB = string.IsNullOrWhiteSpace(site.GrpcNodeBAddress) ? null : site.GrpcNodeBAddress;
if (grpcA is not null || grpcB is not null)
grpcContacts[site.SiteIdentifier] = new SiteGrpcEndpoints(grpcA, grpcB);
} }
// Freeze the cross-task payload before piping to // Freeze the cross-task payload before piping to
@@ -525,82 +594,20 @@ public class CentralCommunicationActor : ReceiveActor
// address-bearing subset in `frozen`) so the aggregator prunes only // address-bearing subset in `frozen`) so the aggregator prunes only
// genuinely-deleted sites and never an addressless-but-configured one. // genuinely-deleted sites and never an addressless-but-configured one.
var knownSiteIds = sites.Select(s => s.SiteIdentifier).ToList(); var knownSiteIds = sites.Select(s => s.SiteIdentifier).ToList();
return new SiteAddressCacheLoaded(frozen, knownSiteIds); return new SiteAddressCacheLoaded(frozen, knownSiteIds, grpcContacts);
}).PipeTo(self); }).PipeTo(self);
} }
private void HandleSiteAddressCacheLoaded(SiteAddressCacheLoaded msg) private void HandleSiteAddressCacheLoaded(SiteAddressCacheLoaded msg)
{ {
var newSiteIds = msg.SiteContacts.Keys.ToHashSet(); // Per-transport per-site resource reconciliation (create/stop ClusterClients, or
var existingSiteIds = _siteClients.Keys.ToHashSet(); // build/drop gRPC channel pairs). Runs on the actor thread each refresh tick.
_transport.ReconcileSites(msg);
// Stop ClusterClients for removed sites
foreach (var removed in existingSiteIds.Except(newSiteIds))
{
_log.Info("Stopping ClusterClient for removed site {0}", removed);
Context.Stop(_siteClients[removed].Client);
_siteClients.Remove(removed);
}
// Add or update
foreach (var (siteId, addresses) in msg.SiteContacts)
{
// Parse all addresses up front inside a try/catch so a
// single malformed site row cannot abort the whole refresh loop and leave
// the cache half-updated. A bad site is logged and skipped; others proceed.
ImmutableHashSet<ActorPath> contactPaths;
try
{
contactPaths = addresses
.Select(a => ActorPath.Parse($"{a}/system/receptionist"))
.ToImmutableHashSet();
}
catch (Exception ex)
{
_log.Warning(ex,
"Malformed contact address for site {0}; skipping this site in the refresh "
+ "(other sites are unaffected)", siteId);
continue;
}
var contactStrings = addresses.ToImmutableHashSet();
// Skip if unchanged
if (_siteClients.TryGetValue(siteId, out var existing) && existing.ContactAddresses.SetEquals(contactStrings))
continue;
// Stop old client if addresses changed
if (_siteClients.ContainsKey(siteId))
{
_log.Info("Updating ClusterClient for site {0} (addresses changed)", siteId);
Context.Stop(_siteClients[siteId].Client);
// Remove now: if the replacement create below fails, a stale entry
// would route envelopes to a stopping actor.
_siteClients.Remove(siteId);
}
IActorRef client;
try
{
client = _siteClientFactory.Create(Context.System, siteId, contactPaths);
}
catch (Exception ex)
{
_log.Error(ex,
"Failed to create ClusterClient for site {0}; site is unroutable until the next refresh",
siteId);
continue;
}
_siteClients[siteId] = (client, contactStrings);
_log.Info("Created ClusterClient for site {0} with {1} contact(s)", siteId, addresses.Count);
}
_log.Info("Site ClusterClient cache refreshed with {0} site(s)", _siteClients.Count);
// Self-healing eviction: a site deleted from configuration would otherwise // Self-healing eviction: a site deleted from configuration would otherwise
// linger in the aggregator as a permanently-offline tile (and live KPI // linger in the aggregator as a permanently-offline tile (and live KPI
// sample source) forever. Prune on every refresh so it disappears within // sample source) forever. Prune on every refresh so it disappears within
// one refresh interval without needing a dedicated deletion event. // one refresh interval without needing a dedicated deletion event. Transport-agnostic.
_serviceProvider.GetService<ICentralHealthAggregator>()?.PruneUnknownSites(msg.KnownSiteIds); _serviceProvider.GetService<ICentralHealthAggregator>()?.PruneUnknownSites(msg.KnownSiteIds);
} }
@@ -673,8 +680,9 @@ public class CentralCommunicationActor : ReceiveActor
public record RefreshSiteAddresses; public record RefreshSiteAddresses;
/// <summary> /// <summary>
/// Internal message carrying the loaded site contact data from the database. /// Message carrying the loaded site contact data from the database. Per-transport resource
/// ClusterClient creation happens on the actor thread in HandleSiteAddressCacheLoaded. /// reconciliation happens on the actor thread in HandleSiteAddressCacheLoaded, which hands this
/// straight to <see cref="ISiteCommandTransport.ReconcileSites"/>.
/// ///
/// The payload is exposed as <see cref="IReadOnlyDictionary{TKey,TValue}"/> /// The payload is exposed as <see cref="IReadOnlyDictionary{TKey,TValue}"/>
/// of <see cref="IReadOnlyList{T}"/> so the Akka.NET "messages are immutable" /// of <see cref="IReadOnlyList{T}"/> so the Akka.NET "messages are immutable"
@@ -682,9 +690,24 @@ public record RefreshSiteAddresses;
/// discipline. The producer wraps the constructed buckets with /// discipline. The producer wraps the constructed buckets with
/// <c>List&lt;T&gt;.AsReadOnly()</c> before piping to Self. /// <c>List&lt;T&gt;.AsReadOnly()</c> before piping to Self.
/// </summary> /// </summary>
internal record SiteAddressCacheLoaded( /// <param name="SiteContacts">Akka ClusterClient contact addresses per site (from NodeA/NodeBAddress).</param>
/// <param name="KnownSiteIds">Every configured site id, address-bearing or not, for aggregator pruning.</param>
/// <param name="GrpcContacts">
/// gRPC endpoint pairs per site (from GrpcNodeA/GrpcNodeBAddress) — the streaming path's columns,
/// consumed by the gRPC transport and ignored by the Akka one.
/// </param>
public sealed record SiteAddressCacheLoaded(
IReadOnlyDictionary<string, IReadOnlyList<string>> SiteContacts, IReadOnlyDictionary<string, IReadOnlyList<string>> SiteContacts,
IReadOnlyCollection<string> KnownSiteIds); IReadOnlyCollection<string> KnownSiteIds,
IReadOnlyDictionary<string, SiteGrpcEndpoints> GrpcContacts);
/// <summary>
/// A site's gRPC node-pair endpoints, as loaded from <c>Site.GrpcNodeAAddress</c>/
/// <c>GrpcNodeBAddress</c>. Either may be null when only one node has a gRPC address configured.
/// </summary>
/// <param name="NodeA">NodeA gRPC base address (e.g. <c>http://scadabridge-site-a-node-a:8083</c>), or null.</param>
/// <param name="NodeB">NodeB gRPC base address, or null.</param>
public readonly record struct SiteGrpcEndpoints(string? NodeA, string? NodeB);
/// <summary> /// <summary>
/// Peer-replication envelope for a site heartbeat, fanned out over the same /// Peer-replication envelope for a site heartbeat, fanned out over the same
@@ -0,0 +1,43 @@
using Akka.Actor;
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
/// <summary>
/// The central→site command-send seam, injected into <see cref="CentralCommunicationActor"/>
/// below the <see cref="SiteEnvelope"/> handler. Exactly one implementation is active per node,
/// chosen by <c>ScadaBridge:Communication:SiteTransport</c>:
/// <list type="bullet">
/// <item><see cref="AkkaSiteTransport"/> — today's per-site <c>ClusterClient</c> path (default).</item>
/// <item><see cref="ZB.MOM.WW.ScadaBridge.Communication.Grpc.GrpcSiteTransport"/> — the site
/// <c>SiteCommandService</c> gRPC plane.</item>
/// </list>
/// The producers above the seam (<c>CommunicationService</c>'s 27 commands, <c>SiteCallAuditActor</c>'s
/// 2 parked relays, and <c>DebugStreamBridgeActor</c>'s subscribe/unsubscribe) are unchanged — they
/// still <c>Ask</c>/<c>Tell</c> a <see cref="SiteEnvelope"/> to the actor, which delegates here.
/// </summary>
/// <remarks>
/// Both members run on the actor thread (from the <see cref="SiteEnvelope"/> and
/// <c>SiteAddressCacheLoaded</c> handlers), so implementations need no internal synchronisation for
/// their own per-site bookkeeping beyond what a background failback loop requires.
/// </remarks>
public interface ISiteCommandTransport
{
/// <summary>
/// Routes <paramref name="envelope"/>'s message to its site. Any reply the site produces is
/// delivered to <paramref name="replyTo"/> — for an <c>Ask</c> that is the temporary ask actor
/// (completing the caller's task); for a <c>Tell</c>-with-sender (the debug bridge) that is the
/// originating actor. A message with no route (an unknown site) is warned and dropped so the
/// caller's <c>Ask</c> times out, exactly as today's ClusterClient path behaves.
/// </summary>
/// <param name="envelope">The site-addressed command envelope.</param>
/// <param name="replyTo">Where a reply (or a <see cref="Status.Failure"/>) is delivered.</param>
void Send(SiteEnvelope envelope, IActorRef replyTo);
/// <summary>
/// Reconciles per-site transport resources (ClusterClients for Akka, channel pairs for gRPC)
/// against the freshly loaded site set. Called once per DB refresh tick with the same cache
/// message the actor already receives — the ONE DB-poll loop feeds both transports.
/// </summary>
/// <param name="cache">The loaded site address cache (Akka contacts + gRPC endpoints + known ids).</param>
void ReconcileSites(SiteAddressCacheLoaded cache);
}
@@ -0,0 +1,303 @@
using Akka.Actor;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DataConnection;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
/// <summary>
/// The single routing truth for the 28 migrated central→site commands. Given one
/// of those command records it decides which local target answers it —
/// the Deployment Manager singleton proxy, the artifact/event-log/parked-message
/// handlers, or the node-local failover path — preserving EXACTLY the targets and
/// null-guard semantics <see cref="SiteCommunicationActor"/> used when this routing
/// lived inline.
/// </summary>
/// <remarks>
/// <para>
/// Two transports call this one unit: the Akka <see cref="SiteCommunicationActor"/>
/// (ClusterClient) and the new <c>SiteCommandGrpcService</c> (gRPC). Centralising the
/// table means the two can never drift on where a command goes. Each transport keeps
/// its own send mechanics — the actor <c>Forward</c>s (preserving the Ask sender), the
/// gRPC service <c>Ask</c>s and encodes the reply — but both read the same
/// <see cref="Route"/> here.
/// </para>
/// <para>
/// <b>The parked-message handler stays node-local on purpose.</b> A parked retry or
/// discard must run on the node that holds the replicated store row, so parked
/// commands route to the per-node <c>_parkedMessageHandler</c> — never onto the
/// singleton proxy (design §7.3). This is the same target the actor used; the
/// extraction does not "fix" it.
/// </para>
/// <para>
/// <b>Excluded by design:</b> <c>IntegrationCallRequest</c> — the 29th command, dead at
/// both ends. It never enters this dispatcher; the actor keeps its own vestigial handler
/// for it (28 of 29 migrate).
/// </para>
/// </remarks>
public sealed class SiteCommandDispatcher
{
/// <summary>How a resolved command is delivered to its target.</summary>
public enum RouteDisposition
{
/// <summary>Forward/Ask <see cref="Route.Target"/> and route its reply back.</summary>
Forward,
/// <summary>
/// Tell <see cref="Route.Target"/> and answer immediately with
/// <see cref="Route.Reply"/> — the fire-and-forget path (only
/// <see cref="UnsubscribeDebugViewRequest"/>, which the downstream never acks).
/// The actor <c>Forward</c>s it and sends nothing back; the gRPC transport Tells and
/// returns the synthetic ack so a unary RPC still answers.
/// </summary>
TellFireAndForget,
/// <summary>
/// No target is available (a null-guarded handler is unregistered); answer with the
/// synthetic <see cref="Route.Reply"/> — the "handler not available" reply the actor
/// produced inline.
/// </summary>
ImmediateReply
}
/// <summary>A resolved routing decision for one command.</summary>
/// <param name="Disposition">How the command is delivered.</param>
/// <param name="Target">The target actor for <see cref="RouteDisposition.Forward"/>/<see cref="RouteDisposition.TellFireAndForget"/>; <c>null</c> for an immediate reply.</param>
/// <param name="Reply">The synthetic reply for <see cref="RouteDisposition.ImmediateReply"/>/<see cref="RouteDisposition.TellFireAndForget"/>; <c>null</c> for a forward.</param>
public readonly record struct Route(RouteDisposition Disposition, IActorRef? Target, object? Reply);
/// <summary>The outcome of preparing a failover: the ack to send now, plus the leave to run after.</summary>
/// <param name="Ack">The ack the caller must send back before the node leaves.</param>
/// <param name="CommitLeave">
/// When non-<c>null</c> (accepted only), invoking it performs the real graceful
/// <c>Cluster.Leave</c>. It is deliberately deferred so the transport can flush the
/// <see cref="Ack"/> first — a caller reaching the very node that is about to leave still
/// receives the ack rather than a broken stream.
/// </param>
public readonly record struct FailoverOutcome(SiteFailoverAck Ack, Action? CommitLeave);
private readonly string _siteId;
private readonly IActorRef _deploymentManagerProxy;
// (role, dryRun) -> leaving node address, or null when there is no standby.
// dryRun:true resolves the target WITHOUT leaving; dryRun:false performs the Leave.
private readonly Func<string, bool, string?> _resolveFailover;
// Registered at runtime via the actor's RegisterLocalHandler flow. Written on the actor
// thread, read by both the actor and the gRPC service (a Kestrel thread), so volatile.
private volatile IActorRef? _eventLogHandler;
private volatile IActorRef? _parkedMessageHandler;
private volatile IActorRef? _artifactHandler;
/// <summary>Creates the dispatcher.</summary>
/// <param name="siteId">This site's identifier, stamped into synthetic replies and matched by the failover guard.</param>
/// <param name="deploymentManagerProxy">The local Deployment Manager singleton proxy — the target for all lifecycle/OPC UA/query/route commands.</param>
/// <param name="resolveFailover">
/// Resolves (and, when <c>dryRun</c> is false, performs) the graceful leave of the oldest
/// Up member in a role scope. Injected so tests need no real cluster.
/// </param>
public SiteCommandDispatcher(
string siteId,
IActorRef deploymentManagerProxy,
Func<string, bool, string?> resolveFailover)
{
ArgumentNullException.ThrowIfNull(siteId);
ArgumentNullException.ThrowIfNull(deploymentManagerProxy);
ArgumentNullException.ThrowIfNull(resolveFailover);
_siteId = siteId;
_deploymentManagerProxy = deploymentManagerProxy;
_resolveFailover = resolveFailover;
}
/// <summary>Registers the site event-log query handler (a cluster singleton proxy).</summary>
/// <param name="handler">The event-log handler.</param>
public void RegisterEventLogHandler(IActorRef handler) => _eventLogHandler = handler;
/// <summary>Registers the node-local parked-message handler (the replicated-store owner on this node).</summary>
/// <param name="handler">The parked-message handler.</param>
public void RegisterParkedMessageHandler(IActorRef handler) => _parkedMessageHandler = handler;
/// <summary>Registers the artifact-deployment handler.</summary>
/// <param name="handler">The artifact handler.</param>
public void RegisterArtifactHandler(IActorRef handler) => _artifactHandler = handler;
/// <summary>
/// Resolves how one of the 27 non-failover commands is delivered. <see cref="TriggerSiteFailover"/>
/// is handled separately via <see cref="PrepareFailover"/>/<see cref="HandleFailover"/> because its
/// leave is deferred; passing it here throws.
/// </summary>
/// <param name="command">The command to route.</param>
/// <returns>The routing decision.</returns>
/// <exception cref="ArgumentException">The command is not a migrated site command (or is failover).</exception>
public Route ResolveRoute(object command)
{
ArgumentNullException.ThrowIfNull(command);
return command switch
{
// ── Deployment + instance lifecycle → Deployment Manager singleton proxy ──
RefreshDeploymentCommand => ToProxy(),
EnableInstanceCommand => ToProxy(),
DisableInstanceCommand => ToProxy(),
DeleteInstanceCommand => ToProxy(),
DeploymentStateQueryRequest => ToProxy(),
// ── Artifact deployment → artifact handler (null-guarded) ──
DeployArtifactsCommand c => _artifactHandler is { } h
? Forwarded(h)
: Immediate(new ArtifactDeploymentResponse(
c.DeploymentId, _siteId, false, "Artifact handler not available", DateTimeOffset.UtcNow)),
// ── Interactive OPC UA / MxGateway → Deployment Manager singleton proxy ──
// The singleton always lands on the active node, which owns the live sessions.
BrowseNodeCommand => ToProxy(),
SearchAddressSpaceCommand => ToProxy(),
ReadTagValuesCommand => ToProxy(),
VerifyEndpointCommand => ToProxy(),
TrustServerCertCommand => ToProxy(),
ListServerCertsCommand => ToProxy(),
RemoveServerCertCommand => ToProxy(),
WriteTagRequest => ToProxy(),
// ── Remote queries: event log (null-guarded) + debug view (proxy) ──
EventLogQueryRequest r => _eventLogHandler is { } h
? Forwarded(h)
: Immediate(new EventLogQueryResponse(
r.CorrelationId, _siteId, [], null, false, false,
"Event log handler not available", DateTimeOffset.UtcNow)),
DebugSnapshotRequest => ToProxy(),
SubscribeDebugViewRequest => ToProxy(),
// Fire-and-forget: the Deployment Manager never acks an unsubscribe, so the gRPC
// transport Tells it and returns the synthetic ack; the actor just Forwards.
UnsubscribeDebugViewRequest => new Route(
RouteDisposition.TellFireAndForget, _deploymentManagerProxy, UnsubscribeDebugViewAck.Instance),
// ── Parked store-and-forward actions → node-local parked handler (null-guarded) ──
ParkedMessageQueryRequest r => _parkedMessageHandler is { } h
? Forwarded(h)
: Immediate(new ParkedMessageQueryResponse(
r.CorrelationId, _siteId, [], 0, r.PageNumber, r.PageSize, false,
"Parked message handler not available", DateTimeOffset.UtcNow)),
ParkedMessageRetryRequest r => _parkedMessageHandler is { } h
? Forwarded(h)
: Immediate(new ParkedMessageRetryResponse(
r.CorrelationId, false, "Parked message handler not available")),
ParkedMessageDiscardRequest r => _parkedMessageHandler is { } h
? Forwarded(h)
: Immediate(new ParkedMessageDiscardResponse(
r.CorrelationId, false, "Parked message handler not available")),
RetryParkedOperation r => _parkedMessageHandler is { } h
? Forwarded(h)
: Immediate(new ParkedOperationActionAck(
r.CorrelationId, Applied: false, "Parked message handler not available")),
DiscardParkedOperation r => _parkedMessageHandler is { } h
? Forwarded(h)
: Immediate(new ParkedOperationActionAck(
r.CorrelationId, Applied: false, "Parked message handler not available")),
// ── Inbound-API Route.To() relays → Deployment Manager singleton proxy ──
RouteToCallRequest => ToProxy(),
RouteToGetAttributesRequest => ToProxy(),
RouteToSetAttributesRequest => ToProxy(),
RouteToWaitForAttributeRequest => ToProxy(),
TriggerSiteFailover => throw new ArgumentException(
"TriggerSiteFailover is handled by PrepareFailover/HandleFailover, not ResolveRoute.",
nameof(command)),
_ => throw new ArgumentException(
$"'{command.GetType().Name}' is not a migrated site command.", nameof(command))
};
Route ToProxy() => Forwarded(_deploymentManagerProxy);
}
private static Route Forwarded(IActorRef target) => new(RouteDisposition.Forward, target, null);
private static Route Immediate(object reply) => new(RouteDisposition.ImmediateReply, null, reply);
/// <summary>
/// The actor's failover path: resolve the standby AND issue the leave in one step (today's
/// coupled behaviour over ClusterClient, where <c>Tell</c> merely enqueues the ack so leave
/// order is immaterial), then return the ack.
/// </summary>
/// <param name="msg">The failover command.</param>
/// <returns>The ack to send back.</returns>
public SiteFailoverAck HandleFailover(TriggerSiteFailover msg)
=> FailoverCore(msg, commitLeaveImmediately: true).Ack;
/// <summary>
/// The gRPC failover path: resolve the standby WITHOUT leaving, build the ack, and hand back a
/// deferred <see cref="FailoverOutcome.CommitLeave"/>. The caller must send the ack before
/// invoking the leave — otherwise a caller reaching the leaving node sees a broken stream
/// instead of its ack (the ack-before-Leave rule).
/// </summary>
/// <param name="msg">The failover command.</param>
/// <returns>The ack and the deferred leave (leave is <c>null</c> when refused).</returns>
public FailoverOutcome PrepareFailover(TriggerSiteFailover msg)
=> FailoverCore(msg, commitLeaveImmediately: false);
private FailoverOutcome FailoverCore(TriggerSiteFailover msg, bool commitLeaveImmediately)
{
ArgumentNullException.ThrowIfNull(msg);
// A misrouted command must be refused, not silently acted on — acting would fail over a
// site the operator never selected. Checked before resolving so the resolver is untouched.
if (!string.Equals(msg.SiteId, _siteId, StringComparison.Ordinal))
{
return new FailoverOutcome(
new SiteFailoverAck(
msg.CorrelationId, Accepted: false, TargetAddress: null,
ErrorMessage: $"Command addressed to site '{msg.SiteId}' but this node serves '{_siteId}'."),
CommitLeave: null);
}
// Site singletons are scoped to the site-specific role, so failover must target that role.
var role = $"site-{_siteId}";
try
{
if (commitLeaveImmediately)
{
var target = _resolveFailover(role, false);
if (target is null)
{
return new FailoverOutcome(NoPeerAck(msg), CommitLeave: null);
}
return new FailoverOutcome(
new SiteFailoverAck(msg.CorrelationId, Accepted: true, target, ErrorMessage: null),
CommitLeave: null);
}
var resolved = _resolveFailover(role, true);
if (resolved is null)
{
return new FailoverOutcome(NoPeerAck(msg), CommitLeave: null);
}
return new FailoverOutcome(
new SiteFailoverAck(msg.CorrelationId, Accepted: true, resolved, ErrorMessage: null),
CommitLeave: () => _resolveFailover(role, false));
}
catch (Exception ex)
{
// A fault must be reported to the operator, never thrown into supervision (over
// ClusterClient) or surfaced as a broken stream (over gRPC).
return new FailoverOutcome(
new SiteFailoverAck(
msg.CorrelationId, Accepted: false, TargetAddress: null, ErrorMessage: ex.Message),
CommitLeave: null);
}
}
private static SiteFailoverAck NoPeerAck(TriggerSiteFailover msg) => new(
msg.CorrelationId, Accepted: false, TargetAddress: null,
ErrorMessage: "No standby available — failing over a lone node would be an outage.");
}
@@ -36,13 +36,20 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
/// do not need a real cluster. /// do not need a real cluster.
/// </summary> /// </summary>
private readonly Func<bool> _isActiveCheck; private readonly Func<bool> _isActiveCheck;
private readonly Func<string, string?> _failOverRole;
/// <summary> /// <summary>
/// Reference to the local Deployment Manager singleton proxy. /// Reference to the local Deployment Manager singleton proxy.
/// </summary> /// </summary>
private readonly IActorRef _deploymentManagerProxy; private readonly IActorRef _deploymentManagerProxy;
/// <summary>
/// The single routing truth for the 28 migrated central→site commands, shared with
/// the gRPC <c>SiteCommandGrpcService</c> so the two transports cannot drift. In
/// production it is created by the Host and passed in (so the gRPC server holds the
/// same instance); when unset (tests) the actor builds its own from the failover seam.
/// </summary>
private readonly SiteCommandDispatcher _dispatcher;
/// <summary> /// <summary>
/// The site→central transport. Finalized in <see cref="PreStart"/> to the injected instance, /// The site→central transport. Finalized in <see cref="PreStart"/> to the injected instance,
/// or a default <see cref="AkkaCentralTransport"/> (ClusterClient) when none is supplied — so /// or a default <see cref="AkkaCentralTransport"/> (ClusterClient) when none is supplied — so
@@ -54,13 +61,10 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
private readonly ICentralTransport? _injectedTransport; private readonly ICentralTransport? _injectedTransport;
/// <summary> /// <summary>
/// Local actor references for routing specific message patterns. /// Handler for the vestigial <see cref="IntegrationCallRequest"/> — the one command NOT
/// Populated via registration messages. /// migrated to the dispatcher (dead at both ends), so it still routes on the actor.
/// </summary> /// </summary>
private IActorRef? _eventLogHandler;
private IActorRef? _parkedMessageHandler;
private IActorRef? _integrationHandler; private IActorRef? _integrationHandler;
private IActorRef? _artifactHandler;
/// <summary>Akka timer scheduler injected by the framework via <see cref="IWithTimers"/>.</summary> /// <summary>Akka timer scheduler injected by the framework via <see cref="IWithTimers"/>.</summary>
public ITimerScheduler Timers { get; set; } = null!; public ITimerScheduler Timers { get; set; } = null!;
@@ -82,24 +86,41 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
/// Host injects a <see cref="Grpc.GrpcCentralTransport"/> when /// Host injects a <see cref="Grpc.GrpcCentralTransport"/> when
/// <c>ScadaBridge:Communication:CentralTransport</c> is <c>Grpc</c>. /// <c>ScadaBridge:Communication:CentralTransport</c> is <c>Grpc</c>.
/// </param> /// </param>
/// <param name="dispatcher">
/// The shared <see cref="SiteCommandDispatcher"/> (production: created by the Host and also
/// handed to the gRPC command service, so both transports route through one instance).
/// <c>null</c> makes the actor build its own from <paramref name="failOverRole"/> — the shape
/// the existing tests use.
/// </param>
public SiteCommunicationActor( public SiteCommunicationActor(
string siteId, string siteId,
CommunicationOptions options, CommunicationOptions options,
IActorRef deploymentManagerProxy, IActorRef deploymentManagerProxy,
Func<bool>? isActiveCheck = null, Func<bool>? isActiveCheck = null,
Func<string, string?>? failOverRole = null, Func<string, string?>? failOverRole = null,
ICentralTransport? transport = null) ICentralTransport? transport = null,
SiteCommandDispatcher? dispatcher = null)
{ {
_siteId = siteId; _siteId = siteId;
_options = options; _options = options;
_deploymentManagerProxy = deploymentManagerProxy; _deploymentManagerProxy = deploymentManagerProxy;
_isActiveCheck = isActiveCheck ?? DefaultIsActiveCheck; _isActiveCheck = isActiveCheck ?? DefaultIsActiveCheck;
_failOverRole = failOverRole ?? DefaultFailOverRole;
_injectedTransport = transport; _injectedTransport = transport;
// Finalized in PreStart (where _log is usable for the default transport); assigned here // Finalized in PreStart (where _log is usable for the default transport); assigned here
// too so the field is definitely-assigned for the constructor's Receive closures. // too so the field is definitely-assigned for the constructor's Receive closures.
_transport = transport!; _transport = transport!;
// When no shared dispatcher is supplied, build one over the same failover seam the
// actor used before extraction: an injected Func (tests) that both resolves and leaves,
// or the shared ClusterFailoverCoordinator. The actor only ever commits the leave
// immediately (dryRun:false), so an injected coupled Func fits unchanged.
var system = Context.System;
Func<string, bool, string?> resolveFailover = failOverRole is not null
? (role, _) => failOverRole(role)
: (role, dryRun) =>
ClusterState.ClusterFailoverCoordinator.FailOverOldest(system, role, dryRun)?.ToString();
_dispatcher = dispatcher ?? new SiteCommandDispatcher(siteId, deploymentManagerProxy, resolveFailover);
// Registration. Feeding the ClusterClient into the transport is a no-op unless the // Registration. Feeding the ClusterClient into the transport is a no-op unless the
// default/Akka transport is in use — the gRPC transport dials configured endpoints and // default/Akka transport is in use — the gRPC transport dials configured endpoints and
// never receives this message (the Host does not create a ClusterClient for it). // never receives this message (the Host does not create a ClusterClient for it).
@@ -109,37 +130,44 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
}); });
Receive<RegisterLocalHandler>(HandleRegisterLocalHandler); Receive<RegisterLocalHandler>(HandleRegisterLocalHandler);
// Pattern 1: Instance Deployment — forward to Deployment Manager // ── The 27 migrated central→site commands (28th is failover, below) all route
Receive<RefreshDeploymentCommand>(msg => // through the shared SiteCommandDispatcher — the single routing truth also used by
{ // the gRPC SiteCommandGrpcService. The actor's job per command is unchanged: Forward
_log.Debug("Routing RefreshDeploymentCommand for {0} to DeploymentManager", msg.InstanceUniqueName); // to the resolved target (preserving the central Ask sender so replies route straight
_deploymentManagerProxy.Forward(msg); // back to the waiting Ask), or Tell the caller the dispatcher's synthetic reply when a
}); // null-guarded handler is unregistered. See SiteCommandDispatcher for the target of
// each command and why the parked handler stays node-local.
Receive<RefreshDeploymentCommand>(cmd => DispatchCommand(cmd));
Receive<DisableInstanceCommand>(cmd => DispatchCommand(cmd));
Receive<EnableInstanceCommand>(cmd => DispatchCommand(cmd));
Receive<DeleteInstanceCommand>(cmd => DispatchCommand(cmd));
Receive<DeploymentStateQueryRequest>(cmd => DispatchCommand(cmd));
Receive<DeployArtifactsCommand>(cmd => DispatchCommand(cmd));
Receive<SubscribeDebugViewRequest>(cmd => DispatchCommand(cmd));
Receive<UnsubscribeDebugViewRequest>(cmd => DispatchCommand(cmd));
Receive<DebugSnapshotRequest>(cmd => DispatchCommand(cmd));
Receive<RouteToCallRequest>(cmd => DispatchCommand(cmd));
Receive<RouteToGetAttributesRequest>(cmd => DispatchCommand(cmd));
Receive<RouteToSetAttributesRequest>(cmd => DispatchCommand(cmd));
Receive<RouteToWaitForAttributeRequest>(cmd => DispatchCommand(cmd));
Receive<BrowseNodeCommand>(cmd => DispatchCommand(cmd));
Receive<ReadTagValuesCommand>(cmd => DispatchCommand(cmd));
Receive<SearchAddressSpaceCommand>(cmd => DispatchCommand(cmd));
Receive<Commons.Messages.DataConnection.WriteTagRequest>(cmd => DispatchCommand(cmd));
Receive<VerifyEndpointCommand>(cmd => DispatchCommand(cmd));
Receive<TrustServerCertCommand>(cmd => DispatchCommand(cmd));
Receive<ListServerCertsCommand>(cmd => DispatchCommand(cmd));
Receive<RemoveServerCertCommand>(cmd => DispatchCommand(cmd));
Receive<EventLogQueryRequest>(cmd => DispatchCommand(cmd));
Receive<ParkedMessageQueryRequest>(cmd => DispatchCommand(cmd));
Receive<ParkedMessageRetryRequest>(cmd => DispatchCommand(cmd));
Receive<ParkedMessageDiscardRequest>(cmd => DispatchCommand(cmd));
Receive<RetryParkedOperation>(cmd => DispatchCommand(cmd));
Receive<DiscardParkedOperation>(cmd => DispatchCommand(cmd));
// Pattern 2: Lifecycle — forward to Deployment Manager // Integration Routing — the 29th command, NOT migrated to the dispatcher (dead at
Receive<DisableInstanceCommand>(msg => _deploymentManagerProxy.Forward(msg)); // both ends; no production code registers the handler). Kept on the actor so the
Receive<EnableInstanceCommand>(msg => _deploymentManagerProxy.Forward(msg)); // dispatcher's command surface stays the 28 that actually migrate.
Receive<DeleteInstanceCommand>(msg => _deploymentManagerProxy.Forward(msg));
// Query-the-site-before-redeploy — forward to
// the Deployment Manager, which owns the deployed-config store and
// answers with the instance's currently-applied deployment identity.
Receive<DeploymentStateQueryRequest>(msg => _deploymentManagerProxy.Forward(msg));
// Pattern 3: Artifact Deployment — forward to artifact handler if registered
Receive<DeployArtifactsCommand>(msg =>
{
if (_artifactHandler != null)
_artifactHandler.Forward(msg);
else
{
_log.Warning("No artifact handler registered, replying with failure");
Sender.Tell(new ArtifactDeploymentResponse(
msg.DeploymentId, _siteId, false, "Artifact handler not available", DateTimeOffset.UtcNow));
}
});
// Pattern 4: Integration Routing — forward to integration handler
Receive<IntegrationCallRequest>(msg => Receive<IntegrationCallRequest>(msg =>
{ {
if (_integrationHandler != null) if (_integrationHandler != null)
@@ -151,143 +179,15 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
} }
}); });
// Pattern 5: Debug View — forward to Deployment Manager (which routes to Instance Actor)
Receive<SubscribeDebugViewRequest>(msg => _deploymentManagerProxy.Forward(msg));
Receive<UnsubscribeDebugViewRequest>(msg => _deploymentManagerProxy.Forward(msg));
// Pattern 6a: Debug Snapshot (one-shot) — forward to Deployment Manager
Receive<DebugSnapshotRequest>(msg => _deploymentManagerProxy.Forward(msg));
// Inbound API Route.To() — forward to Deployment Manager for instance routing
Receive<RouteToCallRequest>(msg => _deploymentManagerProxy.Forward(msg));
Receive<RouteToGetAttributesRequest>(msg => _deploymentManagerProxy.Forward(msg));
Receive<RouteToSetAttributesRequest>(msg => _deploymentManagerProxy.Forward(msg));
Receive<RouteToWaitForAttributeRequest>(msg => _deploymentManagerProxy.Forward(msg));
// OPC UA Tag Browser (interactive design-time query) — forward to the
// Deployment Manager singleton, which always lands on the active site
// node. Routing to the site-local /user/dcl-manager directly is wrong
// because the standby node has a dcl-manager too, but its
// DataConnectionActor children (which own the live OPC UA sessions)
// only exist on the singleton's node. The singleton then re-forwards
// to its own /user/dcl-manager, which DOES have the connection.
Receive<BrowseNodeCommand>(msg => _deploymentManagerProxy.Forward(msg));
// Test Bindings (interactive design-time read) — same routing rationale
// as BrowseNodeCommand above: the singleton always lands on the
// active site node, which is the node that owns the DataConnectionActor
// children holding the live OPC UA sessions.
Receive<ReadTagValuesCommand>(msg => _deploymentManagerProxy.Forward(msg));
// OPC UA tag-picker address-space search and secured-write execute
// — same singleton routing rationale as BrowseNodeCommand above: the
// DataConnectionActor children that own the live OPC UA sessions exist only
// on the singleton's (active) node, so these must hop through the Deployment
// Manager proxy too. Without these forwards the commands dead-letter and the
// central Ask times out. Forward preserves the central Ask sender so the
// result routes straight back to the waiting Ask.
Receive<SearchAddressSpaceCommand>(msg => _deploymentManagerProxy.Forward(msg));
Receive<Commons.Messages.DataConnection.WriteTagRequest>(msg => _deploymentManagerProxy.Forward(msg));
// OPC UA endpoint Verify — probes a (possibly unsaved) endpoint config
// WITHOUT persisting it. The Deployment Manager singleton's dcl-manager runs
// the probe directly (no existing connection required), so — like the
// commands above — Verify routes through the singleton's active node.
Receive<VerifyEndpointCommand>(msg => _deploymentManagerProxy.Forward(msg));
// OPC UA server-certificate trust management — forward to the
// Deployment Manager singleton, which owns the cross-node trust broadcast.
// The trusted-peer PKI store is node-wide per site node, so a trust/remove
// decision must reach BOTH nodes' CertStoreActor; the singleton broadcasts
// to every site node (list answers from the singleton's own node). The
// singleton always lands on the active node, the same routing rationale as
// BrowseNodeCommand above. Forward preserves the central Ask sender so the
// CertTrustResult routes straight back to the waiting Ask.
Receive<TrustServerCertCommand>(msg => _deploymentManagerProxy.Forward(msg));
Receive<ListServerCertsCommand>(msg => _deploymentManagerProxy.Forward(msg));
Receive<RemoveServerCertCommand>(msg => _deploymentManagerProxy.Forward(msg));
// Pattern 7: Remote Queries
Receive<EventLogQueryRequest>(msg =>
{
if (_eventLogHandler != null)
_eventLogHandler.Forward(msg);
else
{
Sender.Tell(new EventLogQueryResponse(
msg.CorrelationId, _siteId, [], null, false, false,
"Event log handler not available", DateTimeOffset.UtcNow));
}
});
Receive<ParkedMessageQueryRequest>(msg =>
{
if (_parkedMessageHandler != null)
_parkedMessageHandler.Forward(msg);
else
{
Sender.Tell(new ParkedMessageQueryResponse(
msg.CorrelationId, _siteId, [], 0, msg.PageNumber, msg.PageSize, false,
"Parked message handler not available", DateTimeOffset.UtcNow));
}
});
Receive<ParkedMessageRetryRequest>(msg =>
{
if (_parkedMessageHandler != null)
_parkedMessageHandler.Forward(msg);
else
{
Sender.Tell(new ParkedMessageRetryResponse(
msg.CorrelationId, false, "Parked message handler not available"));
}
});
Receive<ParkedMessageDiscardRequest>(msg =>
{
if (_parkedMessageHandler != null)
_parkedMessageHandler.Forward(msg);
else
{
Sender.Tell(new ParkedMessageDiscardResponse(
msg.CorrelationId, false, "Parked message handler not available"));
}
});
// Central→site Retry/Discard relay for parked cached
// operations. SiteCallAuditActor relays these over the command/control
// channel; the parked-message handler executes them against the local
// S&F buffer and replies a ParkedOperationActionAck that routes back to
// the relaying SiteCallAuditActor's Ask.
Receive<RetryParkedOperation>(msg =>
{
if (_parkedMessageHandler != null)
_parkedMessageHandler.Forward(msg);
else
{
Sender.Tell(new ParkedOperationActionAck(
msg.CorrelationId, Applied: false, "Parked message handler not available"));
}
});
Receive<DiscardParkedOperation>(msg =>
{
if (_parkedMessageHandler != null)
_parkedMessageHandler.Forward(msg);
else
{
Sender.Tell(new ParkedOperationActionAck(
msg.CorrelationId, Applied: false, "Parked message handler not available"));
}
});
// Central→site manual failover relay. Central and the site are separate clusters, // Central→site manual failover relay. Central and the site are separate clusters,
// so central can only ask — this node performs the graceful Leave locally, scoped to // so central can only ask — this node performs the graceful Leave locally, scoped to
// the SITE-SPECIFIC role, because that is what site singletons (the Deployment // the SITE-SPECIFIC role, because that is what site singletons (the Deployment
// Manager) are placed on. Either node may receive this (the actor is per-node, not a // Manager) are placed on. Either node may receive this (the actor is per-node, not a
// singleton, and contact rotation picks whichever answers); the target is resolved // singleton, and contact rotation picks whichever answers); the target is resolved
// from cluster state, not from who received the message. // from cluster state, not from who received the message. Over ClusterClient the ack
Receive<TriggerSiteFailover>(HandleTriggerSiteFailover); // Tell merely enqueues, so the dispatcher resolves-and-leaves in one step (the gRPC
// transport defers the leave to keep ack-before-Leave — see PrepareFailover).
Receive<TriggerSiteFailover>(msg => Sender.Tell(_dispatcher.HandleFailover(msg)));
// The seven site→central sends now delegate to the injected transport (ClusterClient by // The seven site→central sends now delegate to the injected transport (ClusterClient by
// default, gRPC when configured). Each handler captures the current Sender as the reply // default, gRPC when configured). Each handler captures the current Sender as the reply
@@ -352,21 +252,47 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
_options.ApplicationHeartbeatInterval); _options.ApplicationHeartbeatInterval);
} }
/// <summary>
/// Executes a resolved command route within the actor: Forward to the target (preserving the
/// central Ask sender so the reply routes straight back to the waiting Ask), or — when a
/// null-guarded handler is unregistered — Tell the caller the dispatcher's synthetic reply.
/// The fire-and-forget disposition (UnsubscribeDebugView) is a plain Forward here, exactly as
/// before: over ClusterClient the site never acked it, so the synthetic ack is a gRPC-only
/// concern.
/// </summary>
/// <param name="command">The migrated central→site command to route.</param>
private void DispatchCommand(object command)
{
var route = _dispatcher.ResolveRoute(command);
switch (route.Disposition)
{
case SiteCommandDispatcher.RouteDisposition.ImmediateReply:
Sender.Tell(route.Reply!);
break;
default:
// Forward and TellFireAndForget both Forward on the actor path.
route.Target!.Forward(command);
break;
}
}
private void HandleRegisterLocalHandler(RegisterLocalHandler msg) private void HandleRegisterLocalHandler(RegisterLocalHandler msg)
{ {
// The migrated handlers live on the shared dispatcher so the gRPC command service sees the
// same registrations. Integration is the one command kept on the actor (see the receive).
switch (msg.HandlerType) switch (msg.HandlerType)
{ {
case LocalHandlerType.EventLog: case LocalHandlerType.EventLog:
_eventLogHandler = msg.Handler; _dispatcher.RegisterEventLogHandler(msg.Handler);
break; break;
case LocalHandlerType.ParkedMessages: case LocalHandlerType.ParkedMessages:
_parkedMessageHandler = msg.Handler; _dispatcher.RegisterParkedMessageHandler(msg.Handler);
break; break;
case LocalHandlerType.Integration: case LocalHandlerType.Integration:
_integrationHandler = msg.Handler; _integrationHandler = msg.Handler;
break; break;
case LocalHandlerType.Artifacts: case LocalHandlerType.Artifacts:
_artifactHandler = msg.Handler; _dispatcher.RegisterArtifactHandler(msg.Handler);
break; break;
} }
@@ -431,69 +357,6 @@ public class SiteCommunicationActor : ReceiveActor, IWithTimers
private bool DefaultIsActiveCheck() => private bool DefaultIsActiveCheck() =>
ClusterState.ActiveNodeEvaluator.SelfIsOldestUp(Cluster.Get(Context.System)); ClusterState.ActiveNodeEvaluator.SelfIsOldestUp(Cluster.Get(Context.System));
/// <summary>
/// Handles a central-initiated site failover. Refuses a command addressed to a different
/// site (a misroute must never silently fail over a site the operator did not select) and
/// refuses when the site pair has no peer to take over. The ack is sent BEFORE the Leave
/// takes effect on the wire, so it still reaches central even when this node is the one
/// leaving.
/// </summary>
private void HandleTriggerSiteFailover(TriggerSiteFailover msg)
{
if (!string.Equals(msg.SiteId, _siteId, StringComparison.Ordinal))
{
_log.Warning(
"Refusing TriggerSiteFailover addressed to site {Requested}; this node serves {Actual}",
msg.SiteId, _siteId);
Sender.Tell(new SiteFailoverAck(
msg.CorrelationId, Accepted: false, TargetAddress: null,
ErrorMessage: $"Command addressed to site '{msg.SiteId}' but this node serves '{_siteId}'."));
return;
}
var role = $"site-{_siteId}";
try
{
var target = _failOverRole(role);
if (target is null)
{
_log.Warning(
"Refusing TriggerSiteFailover for {SiteId}: fewer than 2 Up members in role {Role}, "
+ "so there is no standby to take over", _siteId, role);
Sender.Tell(new SiteFailoverAck(
msg.CorrelationId, Accepted: false, TargetAddress: null,
ErrorMessage: "No standby available — failing over a lone node would be an outage."));
return;
}
_log.Warning(
"Manual failover requested by central for site {SiteId}: {Target} is leaving the "
+ "site cluster gracefully; its singletons hand over to the standby.", _siteId, target);
Sender.Tell(new SiteFailoverAck(msg.CorrelationId, Accepted: true, target, ErrorMessage: null));
}
catch (Exception ex)
{
// A fault here must be reported to the operator, not thrown into supervision —
// restarting the communication actor would drop central's Ask into a timeout and
// lose the reason.
_log.Error(ex, "TriggerSiteFailover for {SiteId} faulted", _siteId);
Sender.Tell(new SiteFailoverAck(
msg.CorrelationId, Accepted: false, TargetAddress: null, ErrorMessage: ex.Message));
}
}
/// <summary>
/// Production failover action: gracefully Leave the oldest Up member carrying
/// <paramref name="role"/>, via the shared <see cref="ClusterState.ClusterFailoverCoordinator"/>
/// so the central and site paths cannot drift. Injected in tests for the same reason
/// <see cref="DefaultIsActiveCheck"/> is — a real Leave needs Akka.Cluster in the
/// ActorSystem, which the TestKit system does not load.
/// </summary>
/// <param name="role">Site-specific role scope.</param>
/// <returns>Address of the leaving node, or null when there is no peer.</returns>
private string? DefaultFailOverRole(string role) =>
ClusterState.ClusterFailoverCoordinator.FailOverOldest(Context.System, role)?.ToString();
// ── Internal messages ── // ── Internal messages ──
internal record SendHeartbeat; internal record SendHeartbeat;
@@ -1,6 +1,5 @@
namespace ZB.MOM.WW.ScadaBridge.Communication; namespace ZB.MOM.WW.ScadaBridge.Communication;
/// <summary>
/// Which transport carries the seven site→central control messages. Selected per node by /// Which transport carries the seven site→central control messages. Selected per node by
/// <c>ScadaBridge:Communication:CentralTransport</c>; the migration ships with /// <c>ScadaBridge:Communication:CentralTransport</c>; the migration ships with
/// <see cref="Akka"/> as the default so nothing flips until a node opts in. /// <see cref="Akka"/> as the default so nothing flips until a node opts in.
@@ -14,6 +13,21 @@ public enum CentralTransportMode
Grpc = 1, Grpc = 1,
} }
/// <summary>
/// Selects the transport the central→site command plane rides on. The Akka
/// per-site <c>ClusterClient</c> path is the default until the gRPC cutover
/// (ClusterClient→gRPC migration, Phase 1B); flipping to <see cref="Grpc"/> is the
/// rollback-by-flag switch.
/// </summary>
public enum SiteTransportKind
{
/// <summary>Route <c>SiteEnvelope</c>s through the per-site Akka <c>ClusterClient</c> (today's default).</summary>
Akka,
/// <summary>Route <c>SiteEnvelope</c>s over the site <c>SiteCommandService</c> gRPC plane.</summary>
Grpc
}
/// <summary> /// <summary>
/// Configuration options for central-site communication, including per-pattern /// Configuration options for central-site communication, including per-pattern
/// timeouts and transport heartbeat settings. /// timeouts and transport heartbeat settings.
@@ -39,6 +53,15 @@ public class CommunicationOptions
/// </remarks> /// </remarks>
public List<string> CentralGrpcEndpoints { get; set; } = new(); public List<string> CentralGrpcEndpoints { get; set; } = new();
/// <summary>
/// Which transport the central→site command plane uses. Default <see cref="SiteTransportKind.Akka"/>
/// (the per-site ClusterClient path) — flipping to <see cref="SiteTransportKind.Grpc"/> moves
/// every <c>SiteEnvelope</c> onto the site <c>SiteCommandService</c> gRPC plane. Selected inside
/// <c>CentralCommunicationActor</c>; <c>CommunicationService</c> and <c>SiteCallAuditActor</c>
/// are unchanged either way. Rollback at any point = flip this back to <c>Akka</c>.
/// </summary>
public SiteTransportKind SiteTransport { get; set; } = SiteTransportKind.Akka;
/// <summary>Timeout for deployment commands (typically longest due to apply logic).</summary> /// <summary>Timeout for deployment commands (typically longest due to apply logic).</summary>
public TimeSpan DeploymentTimeout { get; set; } = TimeSpan.FromMinutes(2); public TimeSpan DeploymentTimeout { get; set; } = TimeSpan.FromMinutes(2);
@@ -0,0 +1,244 @@
using Akka.Actor;
using Grpc.Net.Client;
using Microsoft.Extensions.Logging;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DataConnection;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// The gRPC central→site command transport: encodes each <see cref="SiteEnvelope"/> command with
/// <see cref="SiteCommandDtoMapper"/>, dials the site's <c>SiteCommandService</c> through the sticky
/// <see cref="SitePairChannelProvider"/> (PSK + <c>x-scadabridge-site</c> already on the channel),
/// and routes the decoded reply back to the waiting <c>Ask</c> (or the debug-bridge actor).
/// </summary>
/// <remarks>
/// <para>
/// <b>Per-call deadlines match today's Ask timeouts exactly</b> — see <see cref="ResolveDeadline"/>.
/// Behaviour is otherwise unchanged from the Akka path: an unknown/unconfigured site is warned and
/// dropped (the caller's Ask times out), and a transport fault surfaces to the caller as a
/// <see cref="Status.Failure"/>, which the S&amp;F/audit layers already treat as transient.
/// </para>
/// <para>
/// <b>Cross-node retry is the channel provider's job</b> and happens only on
/// <c>Unavailable</c> — never on <c>DeadlineExceeded</c> (a write/deploy/failover may have run).
/// </para>
/// </remarks>
public sealed class GrpcSiteTransport : ISiteCommandTransport
{
private readonly SitePairChannelProvider _channels;
private readonly CommunicationOptions _options;
private readonly ILogger<GrpcSiteTransport> _logger;
// Sites this transport has pushed into the channel provider — used to diff removals per refresh,
// the gRPC analogue of the Akka transport's _siteClients key set. Touched only on the actor thread.
private readonly HashSet<string> _knownSites = new(StringComparer.Ordinal);
/// <summary>Creates the gRPC transport and ensures the provider's failback loop is running.</summary>
/// <param name="channels">The shared per-site channel-pair provider.</param>
/// <param name="options">Communication options supplying the per-command deadlines.</param>
/// <param name="logger">Logger for drop/fault diagnostics.</param>
public GrpcSiteTransport(
SitePairChannelProvider channels,
CommunicationOptions options,
ILogger<GrpcSiteTransport> logger)
{
_channels = channels ?? throw new ArgumentNullException(nameof(channels));
_options = options ?? throw new ArgumentNullException(nameof(options));
_logger = logger ?? throw new ArgumentNullException(nameof(logger));
_channels.EnsureFailbackLoop();
}
/// <inheritdoc />
public void Send(SiteEnvelope envelope, IActorRef replyTo)
{
ArgumentNullException.ThrowIfNull(envelope);
// UnsubscribeDebugView is the one fire-and-forget command: the site never acks it over Akka
// (the debug bridge Tells and expects nothing), so we run the RPC but do not deliver its
// synthetic ack — mirroring today's no-reply behaviour and avoiding a dead-lettered ack.
var fireAndForget = envelope.Message is UnsubscribeDebugViewRequest;
_ = RunAsync(envelope, replyTo, fireAndForget);
}
private async Task RunAsync(SiteEnvelope envelope, IActorRef replyTo, bool fireAndForget)
{
try
{
var reply = await SendCoreAsync(envelope, CancellationToken.None).ConfigureAwait(false);
if (!fireAndForget && !replyTo.IsNobody())
{
replyTo.Tell(reply, ActorRefs.NoSender);
}
}
catch (SiteChannelUnavailableException)
{
// Parity with the Akka "no ClusterClient for site" path: warn and drop, so the caller's
// Ask times out. Central never buffers.
_logger.LogWarning(
"No gRPC channel for site {SiteId}; dropping {Message} (caller's Ask will time out)",
envelope.SiteId, envelope.Message.GetType().Name);
}
catch (Exception ex)
{
if (!fireAndForget && !replyTo.IsNobody())
{
// A timeout or non-OK status faults the caller's Ask exactly as an Akka Ask timeout
// did — the S&F/audit layers already treat that as transient.
replyTo.Tell(new Status.Failure(ex), ActorRefs.NoSender);
}
else
{
_logger.LogWarning(ex,
"Fire-and-forget {Message} to site {SiteId} faulted",
envelope.Message.GetType().Name, envelope.SiteId);
}
}
}
private Task<object> SendCoreAsync(SiteEnvelope envelope, CancellationToken ct)
{
var message = envelope.Message;
// Compute the absolute deadline ONCE so a cross-node failover retry shares the overall
// budget rather than restarting it.
var deadline = DateTime.UtcNow + ResolveDeadline(message);
var group = SiteCommandDtoMapper.GroupOf(message);
return _channels.ExecuteAsync(
envelope.SiteId,
(channel, callCt) => InvokeAsync(channel, group, message, deadline, callCt),
ct);
}
private static async Task<object> InvokeAsync(
GrpcChannel channel, SiteCommandGroup group, object message, DateTime deadline, CancellationToken ct)
{
var client = new SiteCommandService.SiteCommandServiceClient(channel);
switch (group)
{
case SiteCommandGroup.Lifecycle:
{
var reply = await client.ExecuteLifecycleAsync(
SiteCommandDtoMapper.ToLifecycleRequest(message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromLifecycleReply(reply);
}
case SiteCommandGroup.OpcUa:
{
var reply = await client.ExecuteOpcUaAsync(
SiteCommandDtoMapper.ToOpcUaRequest(message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromOpcUaReply(reply);
}
case SiteCommandGroup.Query:
{
var reply = await client.ExecuteQueryAsync(
SiteCommandDtoMapper.ToQueryRequest(message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromQueryReply(reply);
}
case SiteCommandGroup.Parked:
{
var reply = await client.ExecuteParkedAsync(
SiteCommandDtoMapper.ToParkedRequest(message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromParkedReply(reply);
}
case SiteCommandGroup.Route:
{
var reply = await client.ExecuteRouteAsync(
SiteCommandDtoMapper.ToRouteRequest(message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromRouteReply(reply);
}
case SiteCommandGroup.Failover:
{
var ack = await client.TriggerFailoverAsync(
SiteCommandDtoMapper.ToProto((TriggerSiteFailover)message),
deadline: deadline, cancellationToken: ct);
return SiteCommandDtoMapper.FromProto(ack);
}
default:
throw new ArgumentOutOfRangeException(nameof(group), group, "Unknown site command group.");
}
}
/// <summary>
/// The per-command deadline, set EQUAL to the Ask timeout <c>CommunicationService</c> uses for
/// that command today, so flipping the transport changes nothing about how long a call waits.
/// Note the group is NOT a uniform deadline class: within Lifecycle, <c>DeploymentStateQuery</c>
/// uses <c>QueryTimeout</c> (not <c>LifecycleTimeout</c>), and <c>TriggerSiteFailover</c> uses
/// <c>QueryTimeout</c> (not <c>LifecycleTimeout</c>) — matching the real Ask sites, not the
/// plan's per-group table. The two parked relays (<c>RetryParkedOperation</c>/
/// <c>DiscardParkedOperation</c>) map to <c>QueryTimeout</c> (30s), preserving the
/// <c>SiteCallAuditActor</c> inner <c>RelayTimeout</c> (10s) &lt; 30s ordering. WaitForAttribute
/// keeps its dynamic <c>request.Timeout + IntegrationTimeout</c> budget.
/// </summary>
/// <param name="message">The command being sent.</param>
/// <returns>The deadline duration for that command.</returns>
internal TimeSpan ResolveDeadline(object message) => message switch
{
RefreshDeploymentCommand => _options.DeploymentTimeout,
EnableInstanceCommand or DisableInstanceCommand or DeleteInstanceCommand => _options.LifecycleTimeout,
DeploymentStateQueryRequest => _options.QueryTimeout,
DeployArtifactsCommand => _options.ArtifactDeploymentTimeout,
BrowseNodeCommand or SearchAddressSpaceCommand or ReadTagValuesCommand or VerifyEndpointCommand
or TrustServerCertCommand or ListServerCertsCommand or RemoveServerCertCommand or WriteTagRequest
=> _options.QueryTimeout,
EventLogQueryRequest or DebugSnapshotRequest => _options.QueryTimeout,
SubscribeDebugViewRequest or UnsubscribeDebugViewRequest => _options.DebugViewTimeout,
ParkedMessageQueryRequest or ParkedMessageRetryRequest or ParkedMessageDiscardRequest
or RetryParkedOperation or DiscardParkedOperation
=> _options.QueryTimeout,
RouteToCallRequest or RouteToGetAttributesRequest or RouteToSetAttributesRequest
=> _options.IntegrationTimeout,
RouteToWaitForAttributeRequest r => r.Timeout + _options.IntegrationTimeout,
TriggerSiteFailover => _options.QueryTimeout,
_ => throw new ArgumentException(
$"'{message.GetType().Name}' is not a migrated site command.", nameof(message))
};
/// <inheritdoc />
public void ReconcileSites(SiteAddressCacheLoaded cache)
{
ArgumentNullException.ThrowIfNull(cache);
var desired = cache.GrpcContacts;
// Drop sites that lost their gRPC endpoints (or were deleted).
foreach (var removed in _knownSites.Where(s => !desired.ContainsKey(s)).ToList())
{
_channels.RemoveSite(removed);
_knownSites.Remove(removed);
}
// Create/refresh the rest.
foreach (var (siteId, endpoints) in desired)
{
_channels.UpdateSite(siteId, endpoints.NodeA, endpoints.NodeB);
_knownSites.Add(siteId);
}
}
}
@@ -0,0 +1,219 @@
using System.Collections;
using System.Globalization;
using System.Text.Json;
using Google.Protobuf;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// Codec for the loosely-typed <c>object?</c> members that survive on the site
/// command plane — script parameters and return values, attribute values, and
/// OPC UA tag read/write values — mapping them to and from the type-tagged
/// <see cref="LooseValue"/> proto carrier.
/// </summary>
/// <remarks>
/// <para>
/// <b>Why a tagged union rather than a string or JSON blob.</b> These values
/// reach an operator's screen (Test Bindings, Debug View) and a device write
/// (<c>WriteTag</c>), so collapsing them to text would change behaviour: today
/// the Akka JSON serializer runs with <c>TypeNameHandling</c> on and preserves
/// the boxed CLR type end-to-end. The tags below cover every CLR type these
/// fields actually carry, so those values keep their runtime type across the
/// wire exactly as they do over Akka remoting.
/// </para>
/// <para>
/// <b>The one documented lossy path.</b> Anything outside the tagged set —
/// an exotic numeric (<see cref="byte"/>, <see cref="uint"/>, …), an enum, a
/// POCO — falls back to <see cref="LooseValue.JsonValue"/> and decodes as a
/// <see cref="JsonElement"/> rather than its original CLR type. The value
/// itself is preserved; its CLR identity is not. Collections and string-keyed
/// dictionaries do NOT take that path: they encode recursively (see
/// <see cref="LooseValueList"/>/<see cref="LooseValueMap"/>) and decode to
/// <c>List&lt;object?&gt;</c> / <c>Dictionary&lt;string, object?&gt;</c>, so
/// their ELEMENTS keep their types while the container type widens.
/// </para>
/// <para>
/// <b>Why dates ride as strings.</b> <c>google.protobuf.Timestamp</c>
/// normalises everything to UTC, which silently discards
/// <see cref="DateTime.Kind"/> and <see cref="DateTimeOffset.Offset"/>. For a
/// timestamped tag value that is data loss, not normalisation, so
/// <see cref="DateTime"/>/<see cref="DateTimeOffset"/>/<see cref="TimeSpan"/>
/// use invariant round-trip formats instead. (Message FIELDS that are declared
/// <c>DateTimeOffset</c> in the DTO are a different case and do use
/// <c>Timestamp</c> — see <see cref="SiteCommandDtoMapper"/>.)
/// </para>
/// </remarks>
public static class LooseValueCodec
{
private static readonly JsonSerializerOptions JsonOpts = new() { WriteIndented = false };
/// <summary>
/// Encodes a boxed value for a NULLABLE message field: <c>null</c> returns
/// <c>null</c> so the field is simply left unset on the wire.
/// </summary>
/// <param name="value">The boxed value to encode, or <c>null</c>.</param>
/// <returns>The encoded carrier, or <c>null</c> when <paramref name="value"/> is <c>null</c>.</returns>
public static LooseValue? ToProtoOrNull(object? value) =>
value is null ? null : ToProto(value);
/// <summary>
/// Encodes a boxed value, representing <c>null</c> explicitly as
/// <see cref="LooseNull"/>. Used inside maps and lists, where an absent
/// entry means "no such key/element" rather than "a null value".
/// </summary>
/// <param name="value">The boxed value to encode, or <c>null</c>.</param>
/// <returns>A populated <see cref="LooseValue"/>; never <c>null</c>.</returns>
public static LooseValue ToProto(object? value) => value switch
{
null => new LooseValue { NullValue = new LooseNull() },
bool b => new LooseValue { BoolValue = b },
int i => new LooseValue { Int32Value = i },
long l => new LooseValue { Int64Value = l },
double d => new LooseValue { DoubleValue = d },
float f => new LooseValue { FloatValue = f },
string s => new LooseValue { StringValue = s },
decimal m => new LooseValue { DecimalValue = m.ToString(CultureInfo.InvariantCulture) },
DateTime dt => new LooseValue { DateTimeValue = dt.ToString("O", CultureInfo.InvariantCulture) },
DateTimeOffset dto => new LooseValue { DateTimeOffsetValue = dto.ToString("O", CultureInfo.InvariantCulture) },
Guid g => new LooseValue { GuidValue = g.ToString("D") },
TimeSpan ts => new LooseValue { TimeSpanValue = ts.ToString("c", CultureInfo.InvariantCulture) },
byte[] bytes => new LooseValue { BytesValue = ByteString.CopyFrom(bytes) },
IDictionary dict => new LooseValue { MapValue = ToMap(dict) },
IEnumerable seq => new LooseValue { ListValue = ToList(seq) },
// Escape hatch. Preserves the value, not the CLR type — see the remarks.
_ => new LooseValue { JsonValue = JsonSerializer.Serialize(value, value.GetType(), JsonOpts) }
};
/// <summary>
/// Decodes a carrier back to a boxed value. An unset field
/// (<c>null</c>) and an explicit <see cref="LooseNull"/> both decode to
/// <c>null</c>, so the two encoders above are interchangeable on read.
/// </summary>
/// <param name="value">The carrier to decode, or <c>null</c> when the field was unset.</param>
/// <returns>The decoded boxed value; <c>null</c> for an unset or explicitly-null carrier.</returns>
public static object? FromProto(LooseValue? value)
{
if (value is null)
{
return null;
}
return value.KindCase switch
{
LooseValue.KindOneofCase.None => null,
LooseValue.KindOneofCase.NullValue => null,
LooseValue.KindOneofCase.BoolValue => value.BoolValue,
LooseValue.KindOneofCase.Int32Value => value.Int32Value,
LooseValue.KindOneofCase.Int64Value => value.Int64Value,
LooseValue.KindOneofCase.DoubleValue => value.DoubleValue,
LooseValue.KindOneofCase.FloatValue => value.FloatValue,
LooseValue.KindOneofCase.StringValue => value.StringValue,
LooseValue.KindOneofCase.DecimalValue =>
decimal.Parse(value.DecimalValue, CultureInfo.InvariantCulture),
LooseValue.KindOneofCase.DateTimeValue =>
DateTime.Parse(value.DateTimeValue, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind),
LooseValue.KindOneofCase.DateTimeOffsetValue =>
DateTimeOffset.Parse(value.DateTimeOffsetValue, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind),
LooseValue.KindOneofCase.GuidValue => Guid.Parse(value.GuidValue),
LooseValue.KindOneofCase.TimeSpanValue =>
TimeSpan.ParseExact(value.TimeSpanValue, "c", CultureInfo.InvariantCulture),
LooseValue.KindOneofCase.BytesValue => value.BytesValue.ToByteArray(),
LooseValue.KindOneofCase.ListValue => FromList(value.ListValue),
LooseValue.KindOneofCase.MapValue => FromMap(value.MapValue),
LooseValue.KindOneofCase.JsonValue => JsonSerializer.Deserialize<JsonElement>(value.JsonValue),
_ => null
};
}
/// <summary>
/// Encodes a string-keyed dictionary onto the wire. Null values are carried
/// explicitly (<see cref="LooseNull"/>), so a key present with a null value
/// stays distinct from an absent key.
/// </summary>
/// <param name="values">The dictionary to encode.</param>
/// <returns>A populated <see cref="LooseValueMap"/>.</returns>
public static LooseValueMap ToProtoMap(IReadOnlyDictionary<string, object?> values)
{
ArgumentNullException.ThrowIfNull(values);
var map = new LooseValueMap();
foreach (var (key, value) in values)
{
map.Entries[key] = ToProto(value);
}
return map;
}
/// <summary>
/// Encodes a NULLABLE string-keyed dictionary: <c>null</c> returns
/// <c>null</c> so the field is left unset and the null/empty distinction
/// survives (proto3 <c>map</c> alone cannot express it).
/// </summary>
/// <param name="values">The dictionary to encode, or <c>null</c>.</param>
/// <returns>The encoded map, or <c>null</c> when <paramref name="values"/> is <c>null</c>.</returns>
public static LooseValueMap? ToProtoMapOrNull(IReadOnlyDictionary<string, object?>? values) =>
values is null ? null : ToProtoMap(values);
/// <summary>Decodes a wire map back to a dictionary; an unset field decodes to <c>null</c>.</summary>
/// <param name="map">The wire map, or <c>null</c> when the field was unset.</param>
/// <returns>The decoded dictionary, or <c>null</c>.</returns>
public static IReadOnlyDictionary<string, object?>? FromProtoMapOrNull(LooseValueMap? map) =>
map is null ? null : FromMap(map);
/// <summary>Decodes a wire map back to a dictionary; an unset field decodes to an EMPTY dictionary.</summary>
/// <remarks>For DTO members that are declared non-nullable, so "absent" can only mean "empty".</remarks>
/// <param name="map">The wire map, or <c>null</c> when the field was unset.</param>
/// <returns>The decoded dictionary; empty when <paramref name="map"/> is <c>null</c>.</returns>
public static IReadOnlyDictionary<string, object?> FromProtoMap(LooseValueMap? map) =>
map is null ? new Dictionary<string, object?>() : FromMap(map);
private static LooseValueList ToList(IEnumerable source)
{
var list = new LooseValueList();
foreach (var element in source)
{
list.Items.Add(ToProto(element));
}
return list;
}
private static LooseValueMap ToMap(IDictionary source)
{
var map = new LooseValueMap();
foreach (DictionaryEntry entry in source)
{
// Non-string keys are outside the wire contract; the invariant-culture
// rendering keeps the entry readable rather than dropping it silently.
var key = entry.Key as string
?? Convert.ToString(entry.Key, CultureInfo.InvariantCulture)
?? string.Empty;
map.Entries[key] = ToProto(entry.Value);
}
return map;
}
private static List<object?> FromList(LooseValueList list)
{
var result = new List<object?>(list.Items.Count);
foreach (var item in list.Items)
{
result.Add(FromProto(item));
}
return result;
}
private static Dictionary<string, object?> FromMap(LooseValueMap map)
{
var result = new Dictionary<string, object?>(map.Entries.Count);
foreach (var (key, value) in map.Entries)
{
result[key] = FromProto(value);
}
return result;
}
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,33 @@
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// The six domain groups the 28 migrated central→site commands are partitioned
/// into — one group per <c>SiteCommandService</c> RPC.
/// </summary>
/// <remarks>
/// The partition is not cosmetic: every command inside a group shares a
/// DEADLINE class today (the <c>CommunicationOptions</c> timeout the central
/// <c>Ask</c> uses), so one RPC per group keeps the deadline decision in one
/// place on the client and one dispatch switch on the server, while the
/// <c>oneof</c> envelope keeps each command individually typed.
/// </remarks>
public enum SiteCommandGroup
{
/// <summary>Deployment refresh, instance enable/disable/delete, deployment-state query, artifact deployment.</summary>
Lifecycle,
/// <summary>Interactive OPC UA / MxGateway design-time commands: browse, search, read, verify, cert trust, write tag.</summary>
OpcUa,
/// <summary>Read-only remote queries: site event log and debug view snapshot/subscribe/unsubscribe.</summary>
Query,
/// <summary>Parked store-and-forward message actions and parked cached-operation retry/discard relays.</summary>
Parked,
/// <summary>Inbound-API <c>Route.To()</c> relays: call, get/set attributes, wait for attribute.</summary>
Route,
/// <summary>Operator-initiated manual site-pair failover.</summary>
Failover
}
@@ -0,0 +1,227 @@
using Akka.Actor;
using Grpc.Core;
using Microsoft.Extensions.Logging;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
using GrpcStatus = Grpc.Core.Status;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// gRPC front door for the central→site command plane: decodes a proto command (via
/// <see cref="SiteCommandDtoMapper"/>), routes it through the ONE
/// <see cref="SiteCommandDispatcher"/> the Akka <c>SiteCommunicationActor</c> also uses, and
/// encodes the reply. One routing truth, two transports.
/// </summary>
/// <remarks>
/// <para>
/// Mapped in the site branch next to <c>SiteStreamGrpcServer</c>, behind the same
/// <c>ControlPlaneAuthInterceptor</c> PSK gate and the same readiness convention: calls are
/// rejected with <see cref="StatusCode.Unavailable"/> until <see cref="SetReady"/> is called
/// once the site actor system is up (mirrors <c>SiteStreamGrpcServer.SetReady</c>).
/// </para>
/// <para>
/// <b>Coexistence — server-side only.</b> This makes the site ALSO listen on gRPC for commands;
/// nothing central flips to gRPC here (that is T1B.3). Central still dials sites via ClusterClient.
/// </para>
/// <para>
/// <b>Fire-and-forget.</b> The Akka path never acks <c>UnsubscribeDebugView</c>, but a unary RPC
/// must answer, so the dispatcher marks it <see cref="SiteCommandDispatcher.RouteDisposition.TellFireAndForget"/>
/// and this service Tells the target then returns the synthetic ack.
/// </para>
/// <para>
/// <b>Ack-before-Leave.</b> <c>TriggerFailover</c> resolves the standby WITHOUT leaving
/// (<see cref="SiteCommandDispatcher.PrepareFailover"/>), returns the ack, and only THEN schedules
/// the real <c>Cluster.Leave</c> — so a caller reaching the very node that is about to leave still
/// receives its ack instead of a broken stream.
/// </para>
/// </remarks>
public sealed class SiteCommandGrpcService : SiteCommandService.SiteCommandServiceBase
{
// A local Ask has no client deadline of its own; fall back to a generous ceiling when the
// caller set none (a deadline-less client is a test or an internal caller). When the caller
// DID set a gRPC deadline, honour the remaining time instead.
private static readonly TimeSpan DefaultAskTimeout = TimeSpan.FromMinutes(2);
private readonly ILogger<SiteCommandGrpcService> _logger;
private readonly Action<Action> _leaveScheduler;
// Set once by SetReady after the site actor system and dispatcher exist. Read on Kestrel
// threads, written on the host bring-up thread — volatile.
private volatile SiteCommandDispatcher? _dispatcher;
private volatile bool _ready;
/// <summary>DI constructor.</summary>
/// <param name="logger">Logger for denial/dispatch diagnostics.</param>
public SiteCommandGrpcService(ILogger<SiteCommandGrpcService> logger)
: this(logger, DeferLeaveUntilAfterReply)
{
}
/// <summary>
/// Test constructor letting a test observe when the deferred leave runs. Internal so DI sees a
/// single public constructor.
/// </summary>
/// <param name="logger">Logger.</param>
/// <param name="leaveScheduler">Runs the deferred <c>Cluster.Leave</c> after the ack is returned.</param>
internal SiteCommandGrpcService(ILogger<SiteCommandGrpcService> logger, Action<Action> leaveScheduler)
{
ArgumentNullException.ThrowIfNull(logger);
ArgumentNullException.ThrowIfNull(leaveScheduler);
_logger = logger;
_leaveScheduler = leaveScheduler;
}
/// <summary>
/// Marks the service ready and injects the shared routing table. Called once the site actor
/// system, the Deployment Manager singleton, and the local handlers are up — the same point
/// <c>SiteStreamGrpcServer.SetReady</c> is called.
/// </summary>
/// <param name="dispatcher">The shared dispatcher (the same instance the actor routes through).</param>
public void SetReady(SiteCommandDispatcher dispatcher)
{
ArgumentNullException.ThrowIfNull(dispatcher);
_dispatcher = dispatcher;
_ready = true;
}
/// <summary>Whether the service is accepting commands. Exposed for tests.</summary>
internal bool IsReady => _ready;
/// <inheritdoc />
public override async Task<LifecycleReply> ExecuteLifecycle(LifecycleRequest request, ServerCallContext context)
{
var command = SiteCommandDtoMapper.FromLifecycleRequest(request);
var reply = await DispatchAsync(command, context).ConfigureAwait(false);
return SiteCommandDtoMapper.ToLifecycleReply(reply);
}
/// <inheritdoc />
public override async Task<OpcUaReply> ExecuteOpcUa(OpcUaRequest request, ServerCallContext context)
{
var command = SiteCommandDtoMapper.FromOpcUaRequest(request);
var reply = await DispatchAsync(command, context).ConfigureAwait(false);
return SiteCommandDtoMapper.ToOpcUaReply(reply);
}
/// <inheritdoc />
public override async Task<QueryReply> ExecuteQuery(QueryRequest request, ServerCallContext context)
{
var command = SiteCommandDtoMapper.FromQueryRequest(request);
var reply = await DispatchAsync(command, context).ConfigureAwait(false);
return SiteCommandDtoMapper.ToQueryReply(reply);
}
/// <inheritdoc />
public override async Task<ParkedReply> ExecuteParked(ParkedRequest request, ServerCallContext context)
{
var command = SiteCommandDtoMapper.FromParkedRequest(request);
var reply = await DispatchAsync(command, context).ConfigureAwait(false);
return SiteCommandDtoMapper.ToParkedReply(reply);
}
/// <inheritdoc />
public override async Task<RouteReply> ExecuteRoute(RouteRequest request, ServerCallContext context)
{
var command = SiteCommandDtoMapper.FromRouteRequest(request);
var reply = await DispatchAsync(command, context).ConfigureAwait(false);
return SiteCommandDtoMapper.ToRouteReply(reply);
}
/// <inheritdoc />
public override Task<SiteFailoverAckDto> TriggerFailover(TriggerSiteFailoverDto request, ServerCallContext context)
{
var dispatcher = EnsureReady();
var msg = SiteCommandDtoMapper.FromProto(request);
// Resolve the standby and build the ack WITHOUT leaving. The real leave is deferred so the
// ack is on the wire first (ack-before-Leave).
var outcome = dispatcher.PrepareFailover(msg);
var dto = SiteCommandDtoMapper.ToProto(outcome.Ack);
if (outcome.CommitLeave is not null)
{
_leaveScheduler(outcome.CommitLeave);
}
return Task.FromResult(dto);
}
/// <summary>
/// Routes a decoded command through the shared dispatcher and returns the reply record.
/// Forward → local <c>Ask</c>; fire-and-forget → local <c>Tell</c> + synthetic ack; no target →
/// the dispatcher's synthetic "handler not available" reply.
/// </summary>
private async Task<object> DispatchAsync(object command, ServerCallContext context)
{
var dispatcher = EnsureReady();
var route = dispatcher.ResolveRoute(command);
switch (route.Disposition)
{
case SiteCommandDispatcher.RouteDisposition.ImmediateReply:
return route.Reply!;
case SiteCommandDispatcher.RouteDisposition.TellFireAndForget:
route.Target!.Tell(command);
return route.Reply!;
default:
try
{
return await route.Target!.Ask<object>(
command, AskTimeout(context), context.CancellationToken).ConfigureAwait(false);
}
catch (RpcException)
{
throw;
}
catch (OperationCanceledException)
{
// Client cancelled or the deadline elapsed.
throw new RpcException(new GrpcStatus(
StatusCode.DeadlineExceeded, "Site did not answer within the deadline."));
}
catch (Exception ex)
{
_logger.LogWarning(ex,
"Local dispatch of {Command} faulted", command.GetType().Name);
throw new RpcException(new GrpcStatus(StatusCode.Internal, ex.Message));
}
}
}
private SiteCommandDispatcher EnsureReady()
{
var dispatcher = _dispatcher;
if (!_ready || dispatcher is null)
{
throw new RpcException(new GrpcStatus(
StatusCode.Unavailable, "Site command plane not ready."));
}
return dispatcher;
}
private static TimeSpan AskTimeout(ServerCallContext context)
{
var deadline = context.Deadline;
if (deadline == DateTime.MaxValue)
{
return DefaultAskTimeout;
}
var remaining = deadline - DateTime.UtcNow;
return remaining > TimeSpan.Zero ? remaining : TimeSpan.FromMilliseconds(1);
}
// Default deferred-leave: return control (so the ack serialises) before the graceful Leave
// runs. A yield hands the current continuation back before the leave begins; the leave itself
// is slow-async (member marked Leaving, CoordinatedShutdown over seconds), so the ack is long
// gone by the time Kestrel actually stops.
private static void DeferLeaveUntilAfterReply(Action commitLeave)
=> _ = Task.Run(async () =>
{
await Task.Yield();
commitLeave();
});
}
@@ -0,0 +1,436 @@
using System.Collections.Concurrent;
using Grpc.Core;
using Grpc.Net.Client;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// Raised when a site has no usable gRPC channel — an unknown site, or a site with neither
/// <c>GrpcNodeAAddress</c> nor <c>GrpcNodeBAddress</c> configured. The gRPC transport treats this
/// the way the Akka path treats "no ClusterClient for site": warn and drop, so the caller's Ask
/// times out (central never buffers).
/// </summary>
public sealed class SiteChannelUnavailableException(string siteId)
: Exception($"No gRPC channel is configured for site '{siteId}'.")
{
/// <summary>The site with no configured channel.</summary>
public string SiteId { get; } = siteId;
}
/// <summary>
/// Shared, per-site gRPC channel PAIR provider with sticky failover + failback (design §3.5).
/// Central holds one <see cref="GrpcChannel"/> per site node (NodeA/NodeB) rather than one channel
/// with re-dial logic; NodeA is the preferred node (config order). Every call goes to the current
/// sticky channel; on connect failure / <see cref="StatusCode.Unavailable"/> / readiness rejection
/// it flips to the other node and stays there. A background probe returns to the preferred node
/// once it answers again. Reconnect probing backs off 1 s → doubling → 60 s cap.
/// </summary>
/// <remarks>
/// <para>
/// <b>Retry safety.</b> <see cref="ExecuteAsync{T}"/> retries on the other node ONLY when the first
/// attempt failed with <see cref="StatusCode.Unavailable"/> (provably never reached a server —
/// connect refused or readiness-rejected before response headers). A
/// <see cref="StatusCode.DeadlineExceeded"/>, or any other status, is rethrown WITHOUT a cross-node
/// retry: a <c>WriteTag</c>/<c>DeployArtifacts</c>/<c>TriggerSiteFailover</c> may already have
/// executed, and the layers above tolerate the one-shot ambiguity exactly as ClusterClient's Ask
/// did — the migration must not turn one write into two.
/// </para>
/// <para>
/// <b>Addresses.</b> Fed by <see cref="UpdateSite"/>/<see cref="RemoveSite"/> from the single DB
/// refresh loop (the <c>GrpcNodeAAddress</c>/<c>GrpcNodeBAddress</c> columns). PSK credentials and
/// the <c>x-scadabridge-site</c> header ride every channel via
/// <see cref="ControlPlaneCredentials.WithSiteCredentials"/>; removing a site invalidates its
/// cached key.
/// </para>
/// </remarks>
public sealed class SitePairChannelProvider : IDisposable
{
private static readonly TimeSpan InitialBackoff = TimeSpan.FromSeconds(1);
private static readonly TimeSpan MaxBackoff = TimeSpan.FromSeconds(60);
private static readonly TimeSpan FailbackSweepInterval = TimeSpan.FromSeconds(1);
private readonly ISitePskProvider _pskProvider;
private readonly CommunicationOptions _options;
private readonly ILogger<SitePairChannelProvider> _logger;
private readonly Func<string, HttpMessageHandler>? _handlerFactory;
private readonly Func<GrpcChannel, CancellationToken, Task<bool>> _reachabilityProbe;
private readonly ConcurrentDictionary<string, SitePair> _sites = new(StringComparer.Ordinal);
private readonly CancellationTokenSource _shutdown = new();
private Task? _failbackLoop;
private readonly object _loopGate = new();
private bool _disposed;
/// <summary>Creates the provider.</summary>
/// <param name="pskProvider">Resolves each site's preshared key for channel credentials.</param>
/// <param name="options">Communication options (keepalive settings applied to production channels).</param>
/// <param name="logger">Logger for failover/failback diagnostics.</param>
/// <param name="handlerFactory">
/// Test seam mapping an endpoint to the <see cref="HttpMessageHandler"/> its channel should use
/// (e.g. an in-process <c>TestServer</c> handler). Null in production, where each channel builds
/// a keepalive-configured <see cref="SocketsHttpHandler"/>.
/// </param>
/// <param name="reachabilityProbe">
/// Test seam deciding whether a preferred node is reachable during a failback sweep. Null in
/// production, where <see cref="GrpcChannel.ConnectAsync"/> is used.
/// </param>
public SitePairChannelProvider(
ISitePskProvider pskProvider,
IOptions<CommunicationOptions> options,
ILogger<SitePairChannelProvider> logger,
Func<string, HttpMessageHandler>? handlerFactory = null,
Func<GrpcChannel, CancellationToken, Task<bool>>? reachabilityProbe = null)
{
ArgumentNullException.ThrowIfNull(pskProvider);
ArgumentNullException.ThrowIfNull(options);
ArgumentNullException.ThrowIfNull(logger);
_pskProvider = pskProvider;
_options = options.Value;
_logger = logger;
_handlerFactory = handlerFactory;
_reachabilityProbe = reachabilityProbe ?? DefaultReachabilityProbe;
}
/// <summary>
/// (Re)builds a site's channel pair from its gRPC node endpoints. A no-op when neither endpoint
/// changed. Disposes and rebuilds a channel whose endpoint changed, and resets stickiness to the
/// preferred node when the pair is (re)created. Called on the DB refresh tick.
/// </summary>
/// <param name="siteId">Site identifier.</param>
/// <param name="nodeAEndpoint">NodeA gRPC base address, or null when unconfigured.</param>
/// <param name="nodeBEndpoint">NodeB gRPC base address, or null when unconfigured.</param>
public void UpdateSite(string siteId, string? nodeAEndpoint, string? nodeBEndpoint)
{
ArgumentException.ThrowIfNullOrWhiteSpace(siteId);
var pair = _sites.GetOrAdd(siteId, id => new SitePair(id));
lock (pair.Gate)
{
var changedA = !string.Equals(pair.NodeAEndpoint, nodeAEndpoint, StringComparison.Ordinal);
var changedB = !string.Equals(pair.NodeBEndpoint, nodeBEndpoint, StringComparison.Ordinal);
if (!changedA && !changedB)
{
return;
}
if (changedA)
{
pair.ChannelA?.Dispose();
pair.ChannelA = string.IsNullOrWhiteSpace(nodeAEndpoint) ? null : BuildChannel(nodeAEndpoint, siteId);
pair.NodeAEndpoint = nodeAEndpoint;
}
if (changedB)
{
pair.ChannelB?.Dispose();
pair.ChannelB = string.IsNullOrWhiteSpace(nodeBEndpoint) ? null : BuildChannel(nodeBEndpoint, siteId);
pair.NodeBEndpoint = nodeBEndpoint;
}
// Addresses changed — return to the preferred node and reset failback backoff.
pair.CurrentIsA = pair.ChannelA is not null;
pair.ResetFailback();
_logger.LogInformation(
"gRPC channel pair for site {SiteId} refreshed (nodeA={HasA}, nodeB={HasB})",
siteId, pair.ChannelA is not null, pair.ChannelB is not null);
}
}
/// <summary>Disposes a removed site's channels and invalidates its cached preshared key.</summary>
/// <param name="siteId">Site identifier to remove.</param>
public void RemoveSite(string siteId)
{
if (string.IsNullOrWhiteSpace(siteId))
{
return;
}
if (_sites.TryRemove(siteId, out var pair))
{
lock (pair.Gate)
{
pair.ChannelA?.Dispose();
pair.ChannelB?.Dispose();
pair.ChannelA = null;
pair.ChannelB = null;
}
_pskProvider.Invalidate(siteId);
_logger.LogInformation("gRPC channel pair for site {SiteId} removed", siteId);
}
}
/// <summary>
/// Runs <paramref name="call"/> against the site's current sticky channel, failing over to the
/// other node once on <see cref="StatusCode.Unavailable"/> (never on
/// <see cref="StatusCode.DeadlineExceeded"/>). Ensures the background failback loop is running.
/// </summary>
/// <typeparam name="T">The RPC reply type.</typeparam>
/// <param name="siteId">Target site.</param>
/// <param name="call">The RPC to run against a given channel.</param>
/// <param name="ct">Cancellation token.</param>
/// <returns>The RPC reply.</returns>
/// <exception cref="SiteChannelUnavailableException">The site has no configured channel.</exception>
public async Task<T> ExecuteAsync<T>(
string siteId, Func<GrpcChannel, CancellationToken, Task<T>> call, CancellationToken ct)
{
ArgumentException.ThrowIfNullOrWhiteSpace(siteId);
ArgumentNullException.ThrowIfNull(call);
EnsureFailbackLoop();
if (!_sites.TryGetValue(siteId, out var pair))
{
throw new SiteChannelUnavailableException(siteId);
}
GrpcChannel primary;
GrpcChannel? secondary;
bool primaryIsA;
lock (pair.Gate)
{
primaryIsA = pair.CurrentIsA;
primary = (primaryIsA ? pair.ChannelA : pair.ChannelB)
?? pair.ChannelA ?? pair.ChannelB
?? throw new SiteChannelUnavailableException(siteId);
// Recompute which node `primary` actually is (the current side may be null).
primaryIsA = ReferenceEquals(primary, pair.ChannelA);
secondary = primaryIsA ? pair.ChannelB : pair.ChannelA;
}
try
{
return await call(primary, ct).ConfigureAwait(false);
}
catch (RpcException ex) when (ex.StatusCode == StatusCode.Unavailable && secondary is not null)
{
// Provably never reached a server → safe to try the other node, even for a write.
_logger.LogWarning(
"Site {SiteId} node {FromNode} unavailable; failing over to node {ToNode}",
siteId, primaryIsA ? "A" : "B", primaryIsA ? "B" : "A");
FlipTo(pair, toIsA: !primaryIsA);
return await call(secondary, ct).ConfigureAwait(false);
}
}
/// <summary>
/// Probes a site's preferred node (NodeA) and, if reachable, returns stickiness to it. The
/// background loop calls this per due site; tests call it to drive failback deterministically.
/// </summary>
/// <param name="siteId">Site identifier.</param>
/// <param name="ct">Cancellation token.</param>
/// <returns>True when the site is now (or already) pointed at its preferred node.</returns>
internal async Task<bool> TryFailbackAsync(string siteId, CancellationToken ct)
{
if (!_sites.TryGetValue(siteId, out var pair))
{
return false;
}
GrpcChannel? preferred;
lock (pair.Gate)
{
if (pair.CurrentIsA || pair.ChannelA is null)
{
return true; // already on preferred (or no preferred channel to fail back to)
}
preferred = pair.ChannelA;
}
var reachable = await _reachabilityProbe(preferred, ct).ConfigureAwait(false);
lock (pair.Gate)
{
if (reachable)
{
pair.CurrentIsA = true;
pair.ResetFailback();
_logger.LogInformation("Site {SiteId} preferred node A reachable again; failing back", siteId);
return true;
}
pair.BumpFailbackBackoff();
return false;
}
}
/// <summary>Test/diagnostic accessor: true when the site currently targets its preferred node (A).</summary>
/// <param name="siteId">Site identifier.</param>
/// <returns>True when on NodeA (or when the site is unknown).</returns>
internal bool IsOnPreferredNode(string siteId)
=> !_sites.TryGetValue(siteId, out var pair) || pair.CurrentIsA;
/// <summary>Starts the background failback loop if not already running. Idempotent.</summary>
public void EnsureFailbackLoop()
{
if (_failbackLoop is not null || _disposed)
{
return;
}
lock (_loopGate)
{
if (_failbackLoop is null && !_disposed)
{
_failbackLoop = Task.Run(() => FailbackLoopAsync(_shutdown.Token));
}
}
}
private async Task FailbackLoopAsync(CancellationToken ct)
{
using var timer = new PeriodicTimer(FailbackSweepInterval);
try
{
while (await timer.WaitForNextTickAsync(ct).ConfigureAwait(false))
{
var now = DateTime.UtcNow;
foreach (var kvp in _sites)
{
var pair = kvp.Value;
bool due;
lock (pair.Gate)
{
due = !pair.CurrentIsA && pair.ChannelA is not null && now >= pair.NextProbeUtc;
}
if (due)
{
try
{
await TryFailbackAsync(kvp.Key, ct).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception ex)
{
_logger.LogDebug(ex, "Failback probe for site {SiteId} faulted", kvp.Key);
}
}
}
}
}
catch (OperationCanceledException)
{
// Shutting down.
}
}
private void FlipTo(SitePair pair, bool toIsA)
{
lock (pair.Gate)
{
pair.CurrentIsA = toIsA;
if (!toIsA)
{
// Failed over off the preferred node → schedule the first failback probe.
pair.ScheduleFailback(InitialBackoff);
}
else
{
pair.ResetFailback();
}
}
}
private GrpcChannel BuildChannel(string endpoint, string siteId)
{
var channelOptions = new GrpcChannelOptions
{
HttpHandler = _handlerFactory is not null
? _handlerFactory(endpoint)
: new SocketsHttpHandler
{
KeepAlivePingDelay = _options.GrpcKeepAlivePingDelay,
KeepAlivePingTimeout = _options.GrpcKeepAlivePingTimeout,
KeepAlivePingPolicy = HttpKeepAlivePingPolicy.Always,
EnableMultipleHttp2Connections = true
}
}.WithSiteCredentials(_pskProvider, siteId);
return GrpcChannel.ForAddress(endpoint, channelOptions);
}
private static async Task<bool> DefaultReachabilityProbe(GrpcChannel channel, CancellationToken ct)
{
try
{
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(ct);
timeout.CancelAfter(TimeSpan.FromSeconds(5));
await channel.ConnectAsync(timeout.Token).ConfigureAwait(false);
return true;
}
catch
{
return false;
}
}
/// <inheritdoc />
public void Dispose()
{
if (_disposed)
{
return;
}
_disposed = true;
_shutdown.Cancel();
try
{
_failbackLoop?.Wait(TimeSpan.FromSeconds(2));
}
catch
{
// Best-effort loop drain on shutdown.
}
_shutdown.Dispose();
foreach (var pair in _sites.Values)
{
lock (pair.Gate)
{
pair.ChannelA?.Dispose();
pair.ChannelB?.Dispose();
}
}
_sites.Clear();
}
/// <summary>Per-site mutable channel state, guarded by <see cref="Gate"/>.</summary>
private sealed class SitePair(string siteId)
{
public string SiteId { get; } = siteId;
public readonly object Gate = new();
public string? NodeAEndpoint;
public string? NodeBEndpoint;
public GrpcChannel? ChannelA;
public GrpcChannel? ChannelB;
/// <summary>Sticky current node; NodeA is preferred. True = pointing at NodeA.</summary>
public bool CurrentIsA = true;
/// <summary>Next time a failback probe is due (while failed over). MaxValue = not scheduled.</summary>
public DateTime NextProbeUtc = DateTime.MaxValue;
private TimeSpan _backoff = InitialBackoff;
public void ScheduleFailback(TimeSpan delay)
{
_backoff = delay;
NextProbeUtc = DateTime.UtcNow + _backoff;
}
public void BumpFailbackBackoff()
{
var next = TimeSpan.FromTicks(Math.Min(_backoff.Ticks * 2, MaxBackoff.Ticks));
_backoff = next;
NextProbeUtc = DateTime.UtcNow + _backoff;
}
public void ResetFailback()
{
_backoff = InitialBackoff;
NextProbeUtc = DateTime.MaxValue;
}
}
}
@@ -0,0 +1,857 @@
syntax = "proto3";
option csharp_namespace = "ZB.MOM.WW.ScadaBridge.Communication.Grpc";
package scadabridge.sitecommand.v1;
import "google/protobuf/timestamp.proto";
import "google/protobuf/duration.proto";
import "google/protobuf/wrappers.proto";
//
// Site command plane (central site) the gRPC replacement for the per-site
// ClusterClient command/control channel handled by SiteCommunicationActor.
//
// The 28 migrated commands are grouped into six domain RPCs rather than 28
// individual RPCs, because the grouping is what carries the DEADLINE policy:
// every command inside one group shares a timeout class today (see
// CommunicationOptions), so one RPC per group keeps the deadline choice in one
// place on the client and one dispatch switch on the server. A `oneof` request
// / reply envelope preserves per-command typing inside the group.
//
// IntegrationCallRequest (the 29th command) is deliberately NOT here it is
// dead code at both ends; see
// docs/known-issues/2026-07-22-integration-call-routing-is-dead-code.md.
//
// EVOLUTION RULE (same as sitestream.proto): field numbers are never reused and
// changes are additive only. New commands take the next free oneof tag; an
// older peer simply sees an unset oneof and answers with a clean error.
//
// The generated C# is CHECKED IN under Communication/SiteCommandGrpc/ because
// protoc segfaults inside the linux_arm64 Docker build image. Regenerate with
// docker/regen-proto.sh (or by hand per the recipe in the .csproj) never by
// leaving an active <Protobuf> item in the project file.
//
service SiteCommandService {
// Deployment + instance lifecycle (6 commands).
rpc ExecuteLifecycle(LifecycleRequest) returns (LifecycleReply);
// Interactive OPC UA / MxGateway design-time commands (8 commands).
rpc ExecuteOpcUa(OpcUaRequest) returns (OpcUaReply);
// Remote read-only queries: event log + debug view (4 commands).
rpc ExecuteQuery(QueryRequest) returns (QueryReply);
// Parked store-and-forward message + cached-operation actions (5 commands).
rpc ExecuteParked(ParkedRequest) returns (ParkedReply);
// Inbound-API Route.To() relays (4 commands).
rpc ExecuteRoute(RouteRequest) returns (RouteReply);
// Operator-initiated manual site-pair failover (1 command).
rpc TriggerFailover(TriggerSiteFailoverDto) returns (SiteFailoverAckDto);
}
//
// Shared value carrier
//
// Explicit null marker. Needed because a map<string, LooseValue> entry always
// has a value message present, so "the dictionary holds a null for this key"
// cannot be expressed by absence the way a nullable message FIELD can.
message LooseNull {}
message LooseValueList { repeated LooseValue items = 1; }
message LooseValueMap { map<string, LooseValue> entries = 1; }
// Type-tagged carrier for the `object?` members the command plane still uses
// (script parameters and return values, attribute values, tag read/write
// values). The tags cover every CLR type these fields actually carry, so those
// values round-trip with their runtime type intact; anything else falls back to
// `json_value` (see the mapper's documented lossiness note).
//
// Date/time and decimal ride as invariant round-trip STRINGS rather than
// google.protobuf.Timestamp: Timestamp normalises everything to UTC and would
// silently drop DateTime.Kind / DateTimeOffset.Offset, which for an operator's
// tag value is data loss, not normalisation.
message LooseValue {
oneof kind {
LooseNull null_value = 1;
bool bool_value = 2;
int32 int32_value = 3;
int64 int64_value = 4;
double double_value = 5;
float float_value = 6;
string string_value = 7;
string decimal_value = 8; // invariant-culture round-trip
string date_time_value = 9; // DateTime, "O" (preserves Kind)
string date_time_offset_value = 10; // DateTimeOffset, "O" (preserves Offset)
string guid_value = 11; // "D"
string time_span_value = 12; // "c" (constant/round-trip)
bytes bytes_value = 13;
LooseValueList list_value = 14;
LooseValueMap map_value = 15;
string json_value = 16; // fallback; decodes to JsonElement
}
}
//
// Enums. Every enum reserves 0 for _UNSPECIFIED so a value that is absent on
// the wire is never mistaken for the first CLR member. The mapper translates
// explicitly in both directions (never by ordinal), so reordering a C# enum
// cannot silently re-map the wire.
//
enum DeploymentStatusDto {
DEPLOYMENT_STATUS_DTO_UNSPECIFIED = 0;
DEPLOYMENT_STATUS_DTO_PENDING = 1;
DEPLOYMENT_STATUS_DTO_IN_PROGRESS = 2;
DEPLOYMENT_STATUS_DTO_SUCCESS = 3;
DEPLOYMENT_STATUS_DTO_FAILED = 4;
}
enum BrowseNodeClassDto {
BROWSE_NODE_CLASS_DTO_UNSPECIFIED = 0;
BROWSE_NODE_CLASS_DTO_OBJECT = 1;
BROWSE_NODE_CLASS_DTO_VARIABLE = 2;
BROWSE_NODE_CLASS_DTO_METHOD = 3;
BROWSE_NODE_CLASS_DTO_OTHER = 4;
}
enum BrowseFailureKindDto {
BROWSE_FAILURE_KIND_DTO_UNSPECIFIED = 0;
BROWSE_FAILURE_KIND_DTO_CONNECTION_NOT_FOUND = 1;
BROWSE_FAILURE_KIND_DTO_CONNECTION_NOT_CONNECTED = 2;
BROWSE_FAILURE_KIND_DTO_NOT_BROWSABLE = 3;
BROWSE_FAILURE_KIND_DTO_TIMEOUT = 4;
BROWSE_FAILURE_KIND_DTO_SERVER_ERROR = 5;
}
enum ReadTagValuesFailureKindDto {
READ_TAG_VALUES_FAILURE_KIND_DTO_UNSPECIFIED = 0;
READ_TAG_VALUES_FAILURE_KIND_DTO_CONNECTION_NOT_FOUND = 1;
READ_TAG_VALUES_FAILURE_KIND_DTO_CONNECTION_NOT_CONNECTED = 2;
READ_TAG_VALUES_FAILURE_KIND_DTO_TIMEOUT = 3;
READ_TAG_VALUES_FAILURE_KIND_DTO_SERVER_ERROR = 4;
}
enum VerifyFailureKindDto {
VERIFY_FAILURE_KIND_DTO_UNSPECIFIED = 0;
VERIFY_FAILURE_KIND_DTO_UNREACHABLE = 1;
VERIFY_FAILURE_KIND_DTO_AUTH_FAILED = 2;
VERIFY_FAILURE_KIND_DTO_UNTRUSTED_CERTIFICATE = 3;
VERIFY_FAILURE_KIND_DTO_TIMEOUT = 4;
VERIFY_FAILURE_KIND_DTO_SERVER_ERROR = 5;
}
enum StoreAndForwardCategoryDto {
STORE_AND_FORWARD_CATEGORY_DTO_UNSPECIFIED = 0;
STORE_AND_FORWARD_CATEGORY_DTO_EXTERNAL_SYSTEM = 1;
STORE_AND_FORWARD_CATEGORY_DTO_NOTIFICATION = 2;
STORE_AND_FORWARD_CATEGORY_DTO_CACHED_DB_WRITE = 3;
}
enum AlarmStateDto {
ALARM_STATE_DTO_UNSPECIFIED = 0;
ALARM_STATE_DTO_ACTIVE = 1;
ALARM_STATE_DTO_NORMAL = 2;
}
enum AlarmLevelDto {
ALARM_LEVEL_DTO_UNSPECIFIED = 0;
ALARM_LEVEL_DTO_NONE = 1;
ALARM_LEVEL_DTO_LOW = 2;
ALARM_LEVEL_DTO_LOW_LOW = 3;
ALARM_LEVEL_DTO_HIGH = 4;
ALARM_LEVEL_DTO_HIGH_HIGH = 5;
}
enum AlarmKindDto {
ALARM_KIND_DTO_UNSPECIFIED = 0;
ALARM_KIND_DTO_COMPUTED = 1;
ALARM_KIND_DTO_NATIVE_OPC_UA = 2;
ALARM_KIND_DTO_NATIVE_MX_ACCESS = 3;
}
enum AlarmShelveStateDto {
ALARM_SHELVE_STATE_DTO_UNSPECIFIED = 0;
ALARM_SHELVE_STATE_DTO_UNSHELVED = 1;
ALARM_SHELVE_STATE_DTO_ONE_SHOT_SHELVED = 2;
ALARM_SHELVE_STATE_DTO_TIMED_SHELVED = 3;
ALARM_SHELVE_STATE_DTO_PERMANENT_SHELVED = 4;
}
//
// RPC 1 ExecuteLifecycle
//
message RefreshDeploymentCommandDto {
string deployment_id = 1;
string instance_unique_name = 2;
string revision_hash = 3;
string deployed_by = 4;
google.protobuf.Timestamp timestamp = 5;
string central_fetch_base_url = 6;
string fetch_token = 7;
}
message EnableInstanceCommandDto {
string command_id = 1;
string instance_unique_name = 2;
google.protobuf.Timestamp timestamp = 3;
}
message DisableInstanceCommandDto {
string command_id = 1;
string instance_unique_name = 2;
google.protobuf.Timestamp timestamp = 3;
}
message DeleteInstanceCommandDto {
string command_id = 1;
string instance_unique_name = 2;
google.protobuf.Timestamp timestamp = 3;
}
message DeploymentStateQueryRequestDto {
string correlation_id = 1;
string instance_unique_name = 2;
google.protobuf.Timestamp timestamp = 3;
}
message SharedScriptArtifactDto {
string name = 1;
string code = 2;
string parameter_definitions = 3; // empty string represents null
string return_definition = 4; // empty string represents null
}
message ExternalSystemArtifactDto {
string name = 1;
string endpoint_url = 2;
string auth_type = 3;
string auth_configuration = 4; // empty string represents null
string method_definitions_json = 5; // empty string represents null
int32 timeout_seconds = 6;
}
message DatabaseConnectionArtifactDto {
string name = 1;
string connection_string = 2;
int32 max_retries = 3;
google.protobuf.Duration retry_delay = 4;
}
message NotificationListArtifactDto {
string name = 1;
repeated string recipient_emails = 2;
}
message DataConnectionArtifactDto {
string name = 1;
string protocol = 2;
string primary_configuration_json = 3; // empty string represents null
string backup_configuration_json = 4; // empty string represents null
int32 failover_retry_count = 5;
}
message SmtpConfigurationArtifactDto {
string name = 1;
string server = 2;
int32 port = 3;
string auth_mode = 4;
string from_address = 5;
string username = 6; // empty string represents null
string password = 7; // empty string represents null
string oauth_config = 8; // empty string represents null
}
// Each artifact collection on DeployArtifactsCommand is a NULLABLE list, and
// proto3 `repeated` cannot distinguish null from empty. Wrapping each in its
// own message makes the null/empty distinction a message-presence question,
// which proto3 does model the same silent-data-loss class the transport
// round-trip guard caught in PLAN-05 T8.
message SharedScriptArtifactListDto { repeated SharedScriptArtifactDto items = 1; }
message ExternalSystemArtifactListDto { repeated ExternalSystemArtifactDto items = 1; }
message DatabaseConnectionArtifactListDto { repeated DatabaseConnectionArtifactDto items = 1; }
message NotificationListArtifactListDto { repeated NotificationListArtifactDto items = 1; }
message DataConnectionArtifactListDto { repeated DataConnectionArtifactDto items = 1; }
message SmtpConfigurationArtifactListDto { repeated SmtpConfigurationArtifactDto items = 1; }
message DeployArtifactsCommandDto {
string deployment_id = 1;
SharedScriptArtifactListDto shared_scripts = 2; // absent => null
ExternalSystemArtifactListDto external_systems = 3; // absent => null
DatabaseConnectionArtifactListDto database_connections = 4; // absent => null
NotificationListArtifactListDto notification_lists = 5; // absent => null
DataConnectionArtifactListDto data_connections = 6; // absent => null
SmtpConfigurationArtifactListDto smtp_configurations = 7; // absent => null
google.protobuf.Timestamp timestamp = 8;
}
message DeploymentStatusResponseDto {
string deployment_id = 1;
string instance_unique_name = 2;
DeploymentStatusDto status = 3;
string error_message = 4; // empty string represents null
google.protobuf.Timestamp timestamp = 5;
}
message InstanceLifecycleResponseDto {
string command_id = 1;
string instance_unique_name = 2;
bool success = 3;
string error_message = 4; // empty string represents null
google.protobuf.Timestamp timestamp = 5;
}
message DeploymentStateQueryResponseDto {
string correlation_id = 1;
string instance_unique_name = 2;
bool is_deployed = 3;
string applied_deployment_id = 4; // empty string represents null
string applied_revision_hash = 5; // empty string represents null
google.protobuf.Timestamp timestamp = 6;
}
message ArtifactDeploymentResponseDto {
string deployment_id = 1;
string site_id = 2;
bool success = 3;
string error_message = 4; // empty string represents null
google.protobuf.Timestamp timestamp = 5;
}
message LifecycleRequest {
oneof command {
RefreshDeploymentCommandDto refresh_deployment = 1;
EnableInstanceCommandDto enable_instance = 2;
DisableInstanceCommandDto disable_instance = 3;
DeleteInstanceCommandDto delete_instance = 4;
DeploymentStateQueryRequestDto deployment_state_query = 5;
DeployArtifactsCommandDto deploy_artifacts = 6;
}
}
message LifecycleReply {
oneof reply {
DeploymentStatusResponseDto deployment_status = 1;
InstanceLifecycleResponseDto instance_lifecycle = 2;
DeploymentStateQueryResponseDto deployment_state_query = 3;
ArtifactDeploymentResponseDto artifact_deployment = 4;
}
}
//
// RPC 2 ExecuteOpcUa
//
message BrowseNodeCommandDto {
string connection_name = 1;
string parent_node_id = 2; // empty string represents null (browse root)
string continuation_token = 3; // empty string represents null (first page)
string site_identifier = 4; // empty string represents null
}
message BrowseNodeDto {
string node_id = 1;
string display_name = 2;
BrowseNodeClassDto node_class = 3;
bool has_children = 4;
string data_type = 5; // empty string represents null
google.protobuf.Int32Value value_rank = 6; // absent => null
google.protobuf.BoolValue writable = 7; // absent => null
}
message BrowseFailureDto {
BrowseFailureKindDto kind = 1;
string message = 2;
}
message BrowseNodeResultDto {
repeated BrowseNodeDto children = 1;
bool truncated = 2;
BrowseFailureDto failure = 3; // absent => null (success)
string continuation_token = 4; // empty string represents null (final page)
}
message SearchAddressSpaceCommandDto {
string connection_name = 1;
string query = 2;
int32 max_depth = 3;
int32 max_results = 4;
string site_identifier = 5; // empty string represents null
}
message AddressSpaceMatchDto {
BrowseNodeDto node = 1;
string path = 2;
}
message SearchAddressSpaceResultDto {
repeated AddressSpaceMatchDto matches = 1;
bool cap_reached = 2;
BrowseFailureDto failure = 3; // absent => null (success)
}
message ReadTagValuesCommandDto {
string connection_name = 1;
repeated string tag_paths = 2;
}
message TagReadOutcomeDto {
string tag_path = 1;
bool success = 2;
LooseValue value = 3; // absent => null
string quality = 4;
google.protobuf.Timestamp timestamp = 5;
string error_message = 6; // empty string represents null
}
message ReadTagValuesFailureDto {
ReadTagValuesFailureKindDto kind = 1;
string message = 2;
}
message ReadTagValuesResultDto {
repeated TagReadOutcomeDto outcomes = 1;
ReadTagValuesFailureDto failure = 2; // absent => null (success)
}
message VerifyEndpointCommandDto {
string connection_name = 1;
string protocol = 2;
string config_json = 3;
string site_identifier = 4; // empty string represents null
}
message ServerCertInfoDto {
string thumbprint = 1;
string subject = 2;
string issuer = 3;
google.protobuf.Timestamp not_before_utc = 4;
google.protobuf.Timestamp not_after_utc = 5;
string der_base64 = 6;
}
// failure_kind is a NULLABLE enum on the CLR side, and proto3 enums have no
// presence, so it rides inside a one-field message exactly like Int32Value.
message VerifyFailureKindValue { VerifyFailureKindDto value = 1; }
message VerifyEndpointResultDto {
bool success = 1;
VerifyFailureKindValue failure_kind = 2; // absent => null (success)
string error = 3; // empty string represents null
ServerCertInfoDto cert = 4; // absent => null
}
message TrustServerCertCommandDto {
string connection_name = 1;
string der_base64 = 2;
string thumbprint = 3;
string site_identifier = 4; // empty string represents null
}
message ListServerCertsCommandDto {
string site_identifier = 1; // empty string represents null
}
message RemoveServerCertCommandDto {
string thumbprint = 1;
string site_identifier = 2; // empty string represents null
}
message TrustedCertInfoDto {
string thumbprint = 1;
string subject = 2;
string issuer = 3;
google.protobuf.Timestamp not_before_utc = 4;
google.protobuf.Timestamp not_after_utc = 5;
bool rejected = 6;
}
// CertTrustResult.Certs is a nullable list (null for trust/remove, populated
// for list) same null-vs-empty problem as the artifact collections.
message TrustedCertInfoListDto { repeated TrustedCertInfoDto items = 1; }
message CertTrustResultDto {
bool success = 1;
string error = 2; // empty string represents null
TrustedCertInfoListDto certs = 3; // absent => null
}
message WriteTagRequestDto {
string correlation_id = 1;
string connection_name = 2;
string tag_path = 3;
LooseValue value = 4; // absent => null
google.protobuf.Timestamp timestamp = 5;
}
message WriteTagResponseDto {
string correlation_id = 1;
bool success = 2;
string error_message = 3; // empty string represents null
google.protobuf.Timestamp timestamp = 4;
}
message OpcUaRequest {
oneof command {
BrowseNodeCommandDto browse_node = 1;
SearchAddressSpaceCommandDto search_address_space = 2;
ReadTagValuesCommandDto read_tag_values = 3;
VerifyEndpointCommandDto verify_endpoint = 4;
TrustServerCertCommandDto trust_server_cert = 5;
ListServerCertsCommandDto list_server_certs = 6;
RemoveServerCertCommandDto remove_server_cert = 7;
WriteTagRequestDto write_tag = 8;
}
}
message OpcUaReply {
oneof reply {
BrowseNodeResultDto browse_node = 1;
SearchAddressSpaceResultDto search_address_space = 2;
ReadTagValuesResultDto read_tag_values = 3;
VerifyEndpointResultDto verify_endpoint = 4;
CertTrustResultDto cert_trust = 5;
WriteTagResponseDto write_tag = 6;
}
}
//
// RPC 3 ExecuteQuery
//
message EventLogQueryRequestDto {
string correlation_id = 1;
string site_id = 2;
google.protobuf.Timestamp from = 3; // absent => null
google.protobuf.Timestamp to = 4; // absent => null
string event_type = 5; // empty string represents null
string severity = 6; // empty string represents null
string instance_id = 7; // empty string represents null
string keyword_filter = 8; // empty string represents null
string continuation_token = 9; // empty string represents null
int32 page_size = 10;
google.protobuf.Timestamp timestamp = 11;
}
message EventLogEntryDto {
string id = 1;
google.protobuf.Timestamp timestamp = 2;
string event_type = 3;
string severity = 4;
string instance_id = 5; // empty string represents null
string source = 6;
string message = 7;
string details = 8; // empty string represents null
}
message EventLogQueryResponseDto {
string correlation_id = 1;
string site_id = 2;
repeated EventLogEntryDto entries = 3;
string continuation_token = 4; // empty string represents null
bool has_more = 5;
bool success = 6;
string error_message = 7; // empty string represents null
google.protobuf.Timestamp timestamp = 8;
}
message DebugSnapshotRequestDto {
string instance_unique_name = 1;
string correlation_id = 2;
}
message SubscribeDebugViewRequestDto {
string instance_unique_name = 1;
string correlation_id = 2;
}
message UnsubscribeDebugViewRequestDto {
string instance_unique_name = 1;
string correlation_id = 2;
}
// UnsubscribeDebugView is a Tell today (CommunicationService.UnsubscribeDebugView
// fires and forgets). A unary RPC must still answer something, so the site sends
// this empty ack; the central transport ignores it to keep the caller-visible
// fire-and-forget semantics identical.
message UnsubscribeDebugViewAckDto {}
message AlarmConditionStateDto {
bool active = 1;
bool acknowledged = 2;
google.protobuf.BoolValue confirmed = 3; // absent => null (not confirmable)
AlarmShelveStateDto shelve = 4;
bool suppressed = 5;
int32 severity = 6;
}
message DebugAttributeValueDto {
string instance_unique_name = 1;
string attribute_path = 2;
string attribute_name = 3;
LooseValue value = 4; // absent => null
string quality = 5;
google.protobuf.Timestamp timestamp = 6;
}
// Full-fidelity projection of Commons AlarmStateChanged. Deliberately NOT the
// sitestream AlarmStateUpdate: that one flattens the value to a display string
// and the shelve state to free text, which is right for a live stream but would
// lose data on a snapshot the central UI renders as authoritative state.
message DebugAlarmStateDto {
string instance_unique_name = 1;
string alarm_name = 2;
AlarmStateDto state = 3;
int32 priority = 4;
google.protobuf.Timestamp timestamp = 5;
AlarmLevelDto level = 6;
string message = 7;
AlarmKindDto kind = 8;
// Absent means "not explicitly set" the CLR record then derives the
// computed default from state + priority. See the mapper's note on why the
// encoder omits a condition that already equals that derived default.
AlarmConditionStateDto condition = 9;
string source_reference = 10;
string alarm_type_name = 11;
string category = 12;
string operator_user = 13;
string operator_comment = 14;
google.protobuf.Timestamp original_raise_time = 15; // absent => null
string current_value = 16;
string limit_value = 17;
string native_source_canonical_name = 18;
bool is_configured_placeholder = 19;
}
message DebugViewSnapshotDto {
string instance_unique_name = 1;
repeated DebugAttributeValueDto attribute_values = 2;
repeated DebugAlarmStateDto alarm_states = 3;
google.protobuf.Timestamp snapshot_timestamp = 4;
bool instance_not_found = 5;
}
message QueryRequest {
oneof command {
EventLogQueryRequestDto event_log_query = 1;
DebugSnapshotRequestDto debug_snapshot = 2;
SubscribeDebugViewRequestDto subscribe_debug_view = 3;
UnsubscribeDebugViewRequestDto unsubscribe_debug_view = 4;
}
}
message QueryReply {
oneof reply {
EventLogQueryResponseDto event_log_query = 1;
DebugViewSnapshotDto debug_view_snapshot = 2;
UnsubscribeDebugViewAckDto unsubscribe_debug_view = 3;
}
}
//
// RPC 4 ExecuteParked
//
message ParkedMessageQueryRequestDto {
string correlation_id = 1;
string site_id = 2;
int32 page_number = 3;
int32 page_size = 4;
google.protobuf.Timestamp timestamp = 5;
}
message ParkedMessageEntryDto {
string message_id = 1;
string target_system = 2;
string method_name = 3;
string error_message = 4;
int32 attempt_count = 5;
google.protobuf.Timestamp original_timestamp = 6;
google.protobuf.Timestamp last_attempt_timestamp = 7;
int32 max_attempts = 8;
StoreAndForwardCategoryDto category = 9;
string origin_instance = 10; // empty string represents null
}
message ParkedMessageQueryResponseDto {
string correlation_id = 1;
string site_id = 2;
repeated ParkedMessageEntryDto messages = 3;
int32 total_count = 4;
int32 page_number = 5;
int32 page_size = 6;
bool success = 7;
string error_message = 8; // empty string represents null
google.protobuf.Timestamp timestamp = 9;
}
message ParkedMessageRetryRequestDto {
string correlation_id = 1;
string site_id = 2;
string message_id = 3;
google.protobuf.Timestamp timestamp = 4;
}
message ParkedMessageRetryResponseDto {
string correlation_id = 1;
bool success = 2;
string error_message = 3; // empty string represents null
}
message ParkedMessageDiscardRequestDto {
string correlation_id = 1;
string site_id = 2;
string message_id = 3;
google.protobuf.Timestamp timestamp = 4;
}
message ParkedMessageDiscardResponseDto {
string correlation_id = 1;
bool success = 2;
string error_message = 3; // empty string represents null
}
message RetryParkedOperationDto {
string correlation_id = 1;
string tracked_operation_id = 2; // TrackedOperationId GUID, "D" format
}
message DiscardParkedOperationDto {
string correlation_id = 1;
string tracked_operation_id = 2; // TrackedOperationId GUID, "D" format
}
message ParkedOperationActionAckDto {
string correlation_id = 1;
bool applied = 2;
string error_message = 3; // empty string represents null
}
message ParkedRequest {
oneof command {
ParkedMessageQueryRequestDto parked_message_query = 1;
ParkedMessageRetryRequestDto parked_message_retry = 2;
ParkedMessageDiscardRequestDto parked_message_discard = 3;
RetryParkedOperationDto retry_parked_operation = 4;
DiscardParkedOperationDto discard_parked_operation = 5;
}
}
message ParkedReply {
oneof reply {
ParkedMessageQueryResponseDto parked_message_query = 1;
ParkedMessageRetryResponseDto parked_message_retry = 2;
ParkedMessageDiscardResponseDto parked_message_discard = 3;
ParkedOperationActionAckDto parked_operation_action = 4;
}
}
//
// RPC 5 ExecuteRoute
//
message RouteToCallRequestDto {
string correlation_id = 1;
string instance_unique_name = 2;
string script_name = 3;
LooseValueMap parameters = 4; // absent => null (distinct from empty)
google.protobuf.Timestamp timestamp = 5;
string parent_execution_id = 6; // empty string represents null
}
message RouteToCallResponseDto {
string correlation_id = 1;
bool success = 2;
LooseValue return_value = 3; // absent => null
string error_message = 4; // empty string represents null
google.protobuf.Timestamp timestamp = 5;
}
message RouteToGetAttributesRequestDto {
string correlation_id = 1;
string instance_unique_name = 2;
repeated string attribute_names = 3;
google.protobuf.Timestamp timestamp = 4;
string parent_execution_id = 5; // empty string represents null
}
message RouteToGetAttributesResponseDto {
string correlation_id = 1;
LooseValueMap values = 2; // non-nullable on the CLR side; absent => empty
bool success = 3;
string error_message = 4; // empty string represents null
google.protobuf.Timestamp timestamp = 5;
}
message RouteToSetAttributesRequestDto {
string correlation_id = 1;
string instance_unique_name = 2;
map<string, string> attribute_values = 3;
google.protobuf.Timestamp timestamp = 4;
string parent_execution_id = 5; // empty string represents null
}
message RouteToSetAttributesResponseDto {
string correlation_id = 1;
bool success = 2;
string error_message = 3; // empty string represents null
google.protobuf.Timestamp timestamp = 4;
}
message RouteToWaitForAttributeRequestDto {
string correlation_id = 1;
string instance_unique_name = 2;
string attribute_name = 3;
// NOT the empty-string-means-null convention: "wait for this attribute to
// become the empty string" is a legitimate target that must stay distinct
// from "no target supplied", so this one nullable string rides a wrapper.
google.protobuf.StringValue target_value_encoded = 4;
google.protobuf.Duration timeout = 5;
google.protobuf.Timestamp timestamp = 6;
string parent_execution_id = 7; // empty string represents null
bool require_good_quality = 8;
}
message RouteToWaitForAttributeResponseDto {
string correlation_id = 1;
bool matched = 2;
LooseValue value = 3; // absent => null
string quality = 4; // empty string represents null
bool timed_out = 5;
bool success = 6;
string error_message = 7; // empty string represents null
google.protobuf.Timestamp timestamp = 8;
}
message RouteRequest {
oneof command {
RouteToCallRequestDto route_to_call = 1;
RouteToGetAttributesRequestDto route_to_get_attributes = 2;
RouteToSetAttributesRequestDto route_to_set_attributes = 3;
RouteToWaitForAttributeRequestDto route_to_wait_for_attribute = 4;
}
}
message RouteReply {
oneof reply {
RouteToCallResponseDto route_to_call = 1;
RouteToGetAttributesResponseDto route_to_get_attributes = 2;
RouteToSetAttributesResponseDto route_to_set_attributes = 3;
RouteToWaitForAttributeResponseDto route_to_wait_for_attribute = 4;
}
}
//
// RPC 6 TriggerFailover
//
message TriggerSiteFailoverDto {
string correlation_id = 1;
string site_id = 2;
}
message SiteFailoverAckDto {
string correlation_id = 1;
bool accepted = 2;
string target_address = 3; // empty string represents null
string error_message = 4; // empty string represents null
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,559 @@
// <auto-generated>
// Generated by the protocol buffer compiler. DO NOT EDIT!
// source: Protos/site_command.proto
// </auto-generated>
#pragma warning disable 0414, 1591, 8981, 0612
#region Designer generated code
using grpc = global::Grpc.Core;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc {
public static partial class SiteCommandService
{
static readonly string __ServiceName = "scadabridge.sitecommand.v1.SiteCommandService";
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static void __Helper_SerializeMessage(global::Google.Protobuf.IMessage message, grpc::SerializationContext context)
{
#if !GRPC_DISABLE_PROTOBUF_BUFFER_SERIALIZATION
if (message is global::Google.Protobuf.IBufferMessage)
{
context.SetPayloadLength(message.CalculateSize());
global::Google.Protobuf.MessageExtensions.WriteTo(message, context.GetBufferWriter());
context.Complete();
return;
}
#endif
context.Complete(global::Google.Protobuf.MessageExtensions.ToByteArray(message));
}
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static class __Helper_MessageCache<T>
{
public static readonly bool IsBufferMessage = global::System.Reflection.IntrospectionExtensions.GetTypeInfo(typeof(global::Google.Protobuf.IBufferMessage)).IsAssignableFrom(typeof(T));
}
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static T __Helper_DeserializeMessage<T>(grpc::DeserializationContext context, global::Google.Protobuf.MessageParser<T> parser) where T : global::Google.Protobuf.IMessage<T>
{
#if !GRPC_DISABLE_PROTOBUF_BUFFER_SERIALIZATION
if (__Helper_MessageCache<T>.IsBufferMessage)
{
return parser.ParseFrom(context.PayloadAsReadOnlySequence());
}
#endif
return parser.ParseFrom(context.PayloadAsNewBuffer());
}
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest> __Marshaller_scadabridge_sitecommand_v1_LifecycleRequest = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply> __Marshaller_scadabridge_sitecommand_v1_LifecycleReply = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest> __Marshaller_scadabridge_sitecommand_v1_OpcUaRequest = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply> __Marshaller_scadabridge_sitecommand_v1_OpcUaReply = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest> __Marshaller_scadabridge_sitecommand_v1_QueryRequest = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply> __Marshaller_scadabridge_sitecommand_v1_QueryReply = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest> __Marshaller_scadabridge_sitecommand_v1_ParkedRequest = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply> __Marshaller_scadabridge_sitecommand_v1_ParkedReply = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest> __Marshaller_scadabridge_sitecommand_v1_RouteRequest = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply> __Marshaller_scadabridge_sitecommand_v1_RouteReply = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto> __Marshaller_scadabridge_sitecommand_v1_TriggerSiteFailoverDto = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Marshaller<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto> __Marshaller_scadabridge_sitecommand_v1_SiteFailoverAckDto = grpc::Marshallers.Create(__Helper_SerializeMessage, context => __Helper_DeserializeMessage(context, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto.Parser));
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply> __Method_ExecuteLifecycle = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply>(
grpc::MethodType.Unary,
__ServiceName,
"ExecuteLifecycle",
__Marshaller_scadabridge_sitecommand_v1_LifecycleRequest,
__Marshaller_scadabridge_sitecommand_v1_LifecycleReply);
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply> __Method_ExecuteOpcUa = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply>(
grpc::MethodType.Unary,
__ServiceName,
"ExecuteOpcUa",
__Marshaller_scadabridge_sitecommand_v1_OpcUaRequest,
__Marshaller_scadabridge_sitecommand_v1_OpcUaReply);
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply> __Method_ExecuteQuery = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply>(
grpc::MethodType.Unary,
__ServiceName,
"ExecuteQuery",
__Marshaller_scadabridge_sitecommand_v1_QueryRequest,
__Marshaller_scadabridge_sitecommand_v1_QueryReply);
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply> __Method_ExecuteParked = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply>(
grpc::MethodType.Unary,
__ServiceName,
"ExecuteParked",
__Marshaller_scadabridge_sitecommand_v1_ParkedRequest,
__Marshaller_scadabridge_sitecommand_v1_ParkedReply);
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply> __Method_ExecuteRoute = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply>(
grpc::MethodType.Unary,
__ServiceName,
"ExecuteRoute",
__Marshaller_scadabridge_sitecommand_v1_RouteRequest,
__Marshaller_scadabridge_sitecommand_v1_RouteReply);
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
static readonly grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto> __Method_TriggerFailover = new grpc::Method<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto>(
grpc::MethodType.Unary,
__ServiceName,
"TriggerFailover",
__Marshaller_scadabridge_sitecommand_v1_TriggerSiteFailoverDto,
__Marshaller_scadabridge_sitecommand_v1_SiteFailoverAckDto);
/// <summary>Service descriptor</summary>
public static global::Google.Protobuf.Reflection.ServiceDescriptor Descriptor
{
get { return global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteCommandReflection.Descriptor.Services[0]; }
}
/// <summary>Base class for server-side implementations of SiteCommandService</summary>
[grpc::BindServiceMethod(typeof(SiteCommandService), "BindService")]
public abstract partial class SiteCommandServiceBase
{
/// <summary>
/// Deployment + instance lifecycle (6 commands).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply> ExecuteLifecycle(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
/// <summary>
/// Interactive OPC UA / MxGateway design-time commands (8 commands).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply> ExecuteOpcUa(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
/// <summary>
/// Remote read-only queries: event log + debug view (4 commands).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply> ExecuteQuery(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
/// <summary>
/// Parked store-and-forward message + cached-operation actions (5 commands).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply> ExecuteParked(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
/// <summary>
/// Inbound-API Route.To() relays (4 commands).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply> ExecuteRoute(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
/// <summary>
/// Operator-initiated manual site-pair failover (1 command).
/// </summary>
/// <param name="request">The request received from the client.</param>
/// <param name="context">The context of the server-side call handler being invoked.</param>
/// <returns>The response to send back to the client (wrapped by a task).</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::System.Threading.Tasks.Task<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto> TriggerFailover(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto request, grpc::ServerCallContext context)
{
throw new grpc::RpcException(new grpc::Status(grpc::StatusCode.Unimplemented, ""));
}
}
/// <summary>Client for SiteCommandService</summary>
public partial class SiteCommandServiceClient : grpc::ClientBase<SiteCommandServiceClient>
{
/// <summary>Creates a new client for SiteCommandService</summary>
/// <param name="channel">The channel to use to make remote calls.</param>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public SiteCommandServiceClient(grpc::ChannelBase channel) : base(channel)
{
}
/// <summary>Creates a new client for SiteCommandService that uses a custom <c>CallInvoker</c>.</summary>
/// <param name="callInvoker">The callInvoker to use to make remote calls.</param>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public SiteCommandServiceClient(grpc::CallInvoker callInvoker) : base(callInvoker)
{
}
/// <summary>Protected parameterless constructor to allow creation of test doubles.</summary>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
protected SiteCommandServiceClient() : base()
{
}
/// <summary>Protected constructor to allow creation of configured clients.</summary>
/// <param name="configuration">The client configuration.</param>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
protected SiteCommandServiceClient(ClientBaseConfiguration configuration) : base(configuration)
{
}
/// <summary>
/// Deployment + instance lifecycle (6 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply ExecuteLifecycle(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteLifecycle(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Deployment + instance lifecycle (6 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply ExecuteLifecycle(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_ExecuteLifecycle, null, options, request);
}
/// <summary>
/// Deployment + instance lifecycle (6 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply> ExecuteLifecycleAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteLifecycleAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Deployment + instance lifecycle (6 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply> ExecuteLifecycleAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_ExecuteLifecycle, null, options, request);
}
/// <summary>
/// Interactive OPC UA / MxGateway design-time commands (8 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply ExecuteOpcUa(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteOpcUa(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Interactive OPC UA / MxGateway design-time commands (8 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply ExecuteOpcUa(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_ExecuteOpcUa, null, options, request);
}
/// <summary>
/// Interactive OPC UA / MxGateway design-time commands (8 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply> ExecuteOpcUaAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteOpcUaAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Interactive OPC UA / MxGateway design-time commands (8 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply> ExecuteOpcUaAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_ExecuteOpcUa, null, options, request);
}
/// <summary>
/// Remote read-only queries: event log + debug view (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply ExecuteQuery(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteQuery(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Remote read-only queries: event log + debug view (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply ExecuteQuery(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_ExecuteQuery, null, options, request);
}
/// <summary>
/// Remote read-only queries: event log + debug view (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply> ExecuteQueryAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteQueryAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Remote read-only queries: event log + debug view (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply> ExecuteQueryAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_ExecuteQuery, null, options, request);
}
/// <summary>
/// Parked store-and-forward message + cached-operation actions (5 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply ExecuteParked(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteParked(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Parked store-and-forward message + cached-operation actions (5 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply ExecuteParked(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_ExecuteParked, null, options, request);
}
/// <summary>
/// Parked store-and-forward message + cached-operation actions (5 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply> ExecuteParkedAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteParkedAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Parked store-and-forward message + cached-operation actions (5 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply> ExecuteParkedAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_ExecuteParked, null, options, request);
}
/// <summary>
/// Inbound-API Route.To() relays (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply ExecuteRoute(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteRoute(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Inbound-API Route.To() relays (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply ExecuteRoute(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_ExecuteRoute, null, options, request);
}
/// <summary>
/// Inbound-API Route.To() relays (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply> ExecuteRouteAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return ExecuteRouteAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Inbound-API Route.To() relays (4 commands).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply> ExecuteRouteAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_ExecuteRoute, null, options, request);
}
/// <summary>
/// Operator-initiated manual site-pair failover (1 command).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto TriggerFailover(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return TriggerFailover(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Operator-initiated manual site-pair failover (1 command).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The response received from the server.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto TriggerFailover(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto request, grpc::CallOptions options)
{
return CallInvoker.BlockingUnaryCall(__Method_TriggerFailover, null, options, request);
}
/// <summary>
/// Operator-initiated manual site-pair failover (1 command).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="headers">The initial metadata to send with the call. This parameter is optional.</param>
/// <param name="deadline">An optional deadline for the call. The call will be cancelled if deadline is hit.</param>
/// <param name="cancellationToken">An optional token for canceling the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto> TriggerFailoverAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto request, grpc::Metadata headers = null, global::System.DateTime? deadline = null, global::System.Threading.CancellationToken cancellationToken = default(global::System.Threading.CancellationToken))
{
return TriggerFailoverAsync(request, new grpc::CallOptions(headers, deadline, cancellationToken));
}
/// <summary>
/// Operator-initiated manual site-pair failover (1 command).
/// </summary>
/// <param name="request">The request to send to the server.</param>
/// <param name="options">The options for the call.</param>
/// <returns>The call object.</returns>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public virtual grpc::AsyncUnaryCall<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto> TriggerFailoverAsync(global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto request, grpc::CallOptions options)
{
return CallInvoker.AsyncUnaryCall(__Method_TriggerFailover, null, options, request);
}
/// <summary>Creates a new instance of client from given <c>ClientBaseConfiguration</c>.</summary>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
protected override SiteCommandServiceClient NewInstance(ClientBaseConfiguration configuration)
{
return new SiteCommandServiceClient(configuration);
}
}
/// <summary>Creates service definition that can be registered with a server</summary>
/// <param name="serviceImpl">An object implementing the server-side handling logic.</param>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public static grpc::ServerServiceDefinition BindService(SiteCommandServiceBase serviceImpl)
{
return grpc::ServerServiceDefinition.CreateBuilder()
.AddMethod(__Method_ExecuteLifecycle, serviceImpl.ExecuteLifecycle)
.AddMethod(__Method_ExecuteOpcUa, serviceImpl.ExecuteOpcUa)
.AddMethod(__Method_ExecuteQuery, serviceImpl.ExecuteQuery)
.AddMethod(__Method_ExecuteParked, serviceImpl.ExecuteParked)
.AddMethod(__Method_ExecuteRoute, serviceImpl.ExecuteRoute)
.AddMethod(__Method_TriggerFailover, serviceImpl.TriggerFailover).Build();
}
/// <summary>Register service method with a service binder with or without implementation. Useful when customizing the service binding logic.
/// Note: this method is part of an experimental API that can change or be removed without any prior notice.</summary>
/// <param name="serviceBinder">Service methods will be bound by calling <c>AddMethod</c> on this object.</param>
/// <param name="serviceImpl">An object implementing the server-side handling logic.</param>
[global::System.CodeDom.Compiler.GeneratedCode("grpc_csharp_plugin", null)]
public static void BindService(grpc::ServiceBinderBase serviceBinder, SiteCommandServiceBase serviceImpl)
{
serviceBinder.AddMethod(__Method_ExecuteLifecycle, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.LifecycleReply>(serviceImpl.ExecuteLifecycle));
serviceBinder.AddMethod(__Method_ExecuteOpcUa, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.OpcUaReply>(serviceImpl.ExecuteOpcUa));
serviceBinder.AddMethod(__Method_ExecuteQuery, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.QueryReply>(serviceImpl.ExecuteQuery));
serviceBinder.AddMethod(__Method_ExecuteParked, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.ParkedReply>(serviceImpl.ExecuteParked));
serviceBinder.AddMethod(__Method_ExecuteRoute, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteRequest, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.RouteReply>(serviceImpl.ExecuteRoute));
serviceBinder.AddMethod(__Method_TriggerFailover, serviceImpl == null ? null : new grpc::UnaryServerMethod<global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.TriggerSiteFailoverDto, global::ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteFailoverAckDto>(serviceImpl.TriggerFailover));
}
}
}
#endregion
@@ -10,6 +10,9 @@
<ItemGroup> <ItemGroup>
<InternalsVisibleTo Include="ZB.MOM.WW.ScadaBridge.Communication.Tests" /> <InternalsVisibleTo Include="ZB.MOM.WW.ScadaBridge.Communication.Tests" />
<InternalsVisibleTo Include="ZB.MOM.WW.ScadaBridge.IntegrationTests" /> <InternalsVisibleTo Include="ZB.MOM.WW.ScadaBridge.IntegrationTests" />
<!-- GrpcSiteTransport failover/failback TestServer coverage (T1B.3) lives in Host.Tests,
which owns the ASP.NET TestHost stack; it drives the provider's internal failback seams. -->
<InternalsVisibleTo Include="ZB.MOM.WW.ScadaBridge.Host.Tests" />
</ItemGroup> </ItemGroup>
<ItemGroup> <ItemGroup>
@@ -32,30 +35,31 @@
<ProjectReference Include="../ZB.MOM.WW.ScadaBridge.HealthMonitoring/ZB.MOM.WW.ScadaBridge.HealthMonitoring.csproj" /> <ProjectReference Include="../ZB.MOM.WW.ScadaBridge.HealthMonitoring/ZB.MOM.WW.ScadaBridge.HealthMonitoring.csproj" />
</ItemGroup> </ItemGroup>
<!-- gRPC proto generation. The compiled C# is checked in — SiteStreamGrpc/ <!-- gRPC proto generation. The compiled C# is checked in — under SiteStreamGrpc/
(Sitestream.cs + SitestreamGrpc.cs) for Protos/sitestream.proto, and (Sitestream.cs + SitestreamGrpc.cs) for Protos/sitestream.proto,
CentralControlGrpc/ (CentralControl.cs + CentralControlGrpc.cs) for CentralControlGrpc/ (CentralControl.cs + CentralControlGrpc.cs) for
Protos/central_control.proto — because protoc segfaults inside our Protos/central_control.proto, and SiteCommandGrpc/ (SiteCommand.cs +
linux_arm64 Docker build image. To regenerate after schema changes run SiteCommandGrpc.cs) for Protos/site_command.proto — because protoc segfaults
`docker/regen-proto.sh [sitestream|centralcontrol|all]`, which does all inside our linux_arm64 Docker build image. To regenerate after schema changes
of the following and always leaves this file as it found it: run `docker/regen-proto.sh [sitestream|centralcontrol|sitecommand|all]`, which
1. Temporarily uncomment the Protobuf ItemGroup below (just the line does all of the following and always leaves this file as it found it:
for the proto you changed — the other file's checked-in C# is 1. Temporarily uncomment the Protobuf ItemGroup below (just the line for
already compiled, so enabling both at once duplicates types). the proto you changed — the other files' checked-in C# is already
compiled, so enabling several at once duplicates types).
2. Delete the matching checked-in *.cs. 2. Delete the matching checked-in *.cs.
3. `dotnet build` (on macOS) — Grpc.Tools writes fresh files to obj/. 3. `dotnet build` (on macOS) — Grpc.Tools writes fresh files to obj/.
4. Copy obj/Debug/net10.0/Protos/*.cs into the matching folder. 4. Copy obj/Debug/net10.0/Protos/*.cs into the matching folder.
5. Re-comment the ItemGroup. 5. Re-comment the ItemGroup.
central_control.proto imports sitestream.proto, so protoc resolves it central_control.proto imports sitestream.proto, so protoc resolves it from the
from the project-relative path without sitestream.proto needing its own project-relative path without sitestream.proto needing its own Protobuf item.
Protobuf item. An ACTIVE Protobuf item must never be committed — it breaks the Docker image
An ACTIVE Protobuf item must never be committed — it breaks the Docker build. Eventually we should switch the Docker build image to one with a working
image build. Eventually we should switch the Docker build image to one protoc on arm64. -->
with a working protoc on arm64. -->
<!-- <!--
<ItemGroup> <ItemGroup>
<Protobuf Include="Protos\sitestream.proto" GrpcServices="Both" /> <Protobuf Include="Protos\sitestream.proto" GrpcServices="Both" />
<Protobuf Include="Protos\central_control.proto" GrpcServices="Both" /> <Protobuf Include="Protos\central_control.proto" GrpcServices="Both" />
<Protobuf Include="Protos\site_command.proto" GrpcServices="Both" />
</ItemGroup> </ItemGroup>
--> -->
@@ -863,7 +863,20 @@ akka {{
_nodeOptions.SiteId); _nodeOptions.SiteId);
} }
// Create SiteCommunicationActor for receiving messages from central // The ONE routing table for central→site commands, shared by the Akka
// SiteCommunicationActor (below) and the gRPC SiteCommandGrpcService (SetReady at the end
// of this method). Its failover seam resolves (dryRun) or performs the graceful Leave via
// the shared ClusterFailoverCoordinator so the central/site paths cannot drift.
var siteCommandDispatcher = new SiteCommandDispatcher(
_nodeOptions.SiteId!,
dmProxy,
(role, dryRun) => ZB.MOM.WW.ScadaBridge.Communication.ClusterState.ClusterFailoverCoordinator
.FailOverOldest(_actorSystem!, role, dryRun)?.ToString());
// Create SiteCommunicationActor for receiving messages from central. It routes commands
// through the shared dispatcher; RegisterLocalHandler (below) registers the handlers INTO
// that dispatcher, so the gRPC command service sees the same registrations. The site→central
// transport is selected above (default Akka); both are passed in.
var siteCommActor = _actorSystem.ActorOf( var siteCommActor = _actorSystem.ActorOf(
Props.Create(() => new SiteCommunicationActor( Props.Create(() => new SiteCommunicationActor(
_nodeOptions.SiteId!, _nodeOptions.SiteId!,
@@ -871,7 +884,8 @@ akka {{
dmProxy, dmProxy,
activeNodeCheck, activeNodeCheck,
null, null,
centralTransport)), centralTransport,
siteCommandDispatcher)),
"site-communication"); "site-communication");
// Register local handlers with SiteCommunicationActor // Register local handlers with SiteCommunicationActor
@@ -1141,5 +1155,12 @@ akka {{
grpcServer?.SetOperationTrackingStore(siteTrackingStore); grpcServer?.SetOperationTrackingStore(siteTrackingStore);
} }
grpcServer?.SetReady(_actorSystem!); grpcServer?.SetReady(_actorSystem!);
// Site command plane: hand the shared dispatcher to the gRPC command service and flip it
// ready at the same point as the streaming server — the actor graph exists and the local
// handlers have been registered into the dispatcher, so it can route commands now.
var commandService = _serviceProvider
.GetService<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteCommandGrpcService>();
commandService?.SetReady(siteCommandDispatcher);
} }
} }
@@ -51,12 +51,18 @@ namespace ZB.MOM.WW.ScadaBridge.Host;
public sealed class ControlPlaneAuthInterceptor : Interceptor public sealed class ControlPlaneAuthInterceptor : Interceptor
{ {
/// <summary> /// <summary>
/// Service prefixes gated by default. Read from the generated <c>sitestream.proto</c> /// Service prefixes gated by default. Read from the generated service descriptors —
/// package/service names — <c>package sitestream; service SiteStreamService</c>. /// <c>package sitestream; service SiteStreamService</c> (real-time data + audit pull) and
/// Later phases append their own services here. /// <c>package scadabridge.sitecommand.v1; service SiteCommandService</c> (the T1B command
/// plane). Later phases append their own services here rather than adding a second
/// interceptor or constructor (see the public constructor's remarks).
/// </summary> /// </summary>
public static readonly IReadOnlyList<string> DefaultGatedPrefixes = public static readonly IReadOnlyList<string> DefaultGatedPrefixes =
new[] { "/sitestream.SiteStreamService/" }; new[]
{
$"/{SiteStreamService.Descriptor.FullName}/",
$"/{SiteCommandService.Descriptor.FullName}/",
};
private readonly IReadOnlyList<string> _gatedPrefixes; private readonly IReadOnlyList<string> _gatedPrefixes;
private readonly IOptions<CommunicationOptions> _options; private readonly IOptions<CommunicationOptions> _options;
+15
View File
@@ -167,6 +167,14 @@ try
}); });
builder.Services.AddSingleton< builder.Services.AddSingleton<
ZB.MOM.WW.ScadaBridge.Communication.Grpc.CentralControlGrpcService>(); ZB.MOM.WW.ScadaBridge.Communication.Grpc.CentralControlGrpcService>();
// Shared per-site gRPC channel-pair provider (sticky failover/failback) backing the
// GrpcSiteTransport when ScadaBridge:Communication:SiteTransport=Grpc. Central-only, and a
// no-op until CentralCommunicationActor builds the gRPC transport (its failback loop starts
// on first use), so registering it unconditionally is harmless under the default Akka path.
builder.Services.AddSingleton<
ZB.MOM.WW.ScadaBridge.Communication.Grpc.SitePairChannelProvider>();
builder.Services.AddHealthMonitoring(); builder.Services.AddHealthMonitoring();
builder.Services.AddCentralHealthAggregation(); builder.Services.AddCentralHealthAggregation();
builder.Services.AddExternalSystemGateway(); builder.Services.AddExternalSystemGateway();
@@ -610,6 +618,10 @@ try
options.Interceptors.Add<ControlPlaneAuthInterceptor>(); options.Interceptors.Add<ControlPlaneAuthInterceptor>();
}); });
builder.Services.AddSingleton<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteStreamGrpcServer>(); builder.Services.AddSingleton<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteStreamGrpcServer>();
// Site command plane (central→site) — the gRPC peer of SiteStreamGrpcServer. Both share
// the h2c listener and the ControlPlaneAuthInterceptor PSK gate; both are readiness-gated
// (SetReady, below). Central still dials over ClusterClient until T1B.3.
builder.Services.AddSingleton<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteCommandGrpcService>();
// Existing site service registrations (this is also where LocalDb and its // Existing site service registrations (this is also where LocalDb and its
// replication engine are registered — see SiteServiceRegistration) // replication engine are registered — see SiteServiceRegistration)
@@ -631,6 +643,9 @@ try
// Map gRPC service — resolves the singleton SiteStreamGrpcServer from DI // Map gRPC service — resolves the singleton SiteStreamGrpcServer from DI
app.MapGrpcService<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteStreamGrpcServer>(); app.MapGrpcService<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteStreamGrpcServer>();
// Map the site command plane (central→site) alongside it, on the same gated listener.
app.MapGrpcService<ZB.MOM.WW.ScadaBridge.Communication.Grpc.SiteCommandGrpcService>();
// The passive half of LocalDb replication: the peer node dials THIS endpoint. // The passive half of LocalDb replication: the peer node dials THIS endpoint.
// It shares the Kestrel h2c listener the site gRPC server already uses, so no // It shares the Kestrel h2c listener the site gRPC server already uses, so no
// listener or port changes are needed. Mapping it is harmless with no peer // listener or port changes are needed. Mapping it is harmless with no peer
@@ -39,7 +39,8 @@ public class CentralCommunicationActorClientLifecycleTests : TestKit
private static SiteAddressCacheLoaded Load(string siteId, params string[] addrs) => private static SiteAddressCacheLoaded Load(string siteId, params string[] addrs) =>
new(new Dictionary<string, IReadOnlyList<string>> new(new Dictionary<string, IReadOnlyList<string>>
{ [siteId] = addrs.ToList().AsReadOnly() }, { [siteId] = addrs.ToList().AsReadOnly() },
new[] { siteId }); new[] { siteId },
new Dictionary<string, SiteGrpcEndpoints>());
[Fact] [Fact]
public void PeriodicRefresh_PrunesDeletedSites_FromHealthAggregator() public void PeriodicRefresh_PrunesDeletedSites_FromHealthAggregator()
@@ -0,0 +1,136 @@
using Akka.Actor;
using Akka.TestKit.Xunit2;
using Microsoft.Extensions.DependencyInjection;
using NSubstitute;
using ZB.MOM.WW.ScadaBridge.Commons.Entities.Sites;
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// T1B.3 seam tests: <see cref="CentralCommunicationActor"/> routing a <see cref="SiteEnvelope"/>
/// through an injected <see cref="ISiteCommandTransport"/> (transport-agnostic), proving the
/// envelope reaches the transport, the reply routes back to the waiting Ask, and the DB refresh
/// tick reconciles per-site transport resources.
/// </summary>
public class CentralCommunicationActorTransportTests : TestKit
{
public CentralCommunicationActorTransportTests() : base(@"akka.loglevel = WARNING") { }
private static Site GrpcSite(string id, string? grpcA = "http://a:8083", string? grpcB = "http://b:8083") =>
new("Test " + id, id)
{
NodeAAddress = $"akka.tcp://scadabridge@{id}-a:8081",
GrpcNodeAAddress = grpcA,
GrpcNodeBAddress = grpcB
};
private (IActorRef actor, ISiteCommandTransport transport, ISiteRepository repo) CreateActor(
IEnumerable<Site>? sites = null)
{
var repo = Substitute.For<ISiteRepository>();
repo.GetAllSitesAsync(Arg.Any<CancellationToken>()).Returns(sites?.ToList() ?? new List<Site>());
var services = new ServiceCollection();
services.AddScoped(_ => repo);
var sp = services.BuildServiceProvider();
var transport = Substitute.For<ISiteCommandTransport>();
var actor = Sys.ActorOf(Props.Create(() => new CentralCommunicationActor(sp, transport, (TimeSpan?)null)));
return (actor, transport, repo);
}
[Fact]
public void SiteEnvelope_IsRoutedToTheTransport_WithTheSenderAsReplyTo()
{
var (actor, transport, _) = CreateActor();
var probe = CreateTestProbe();
var cmd = new EnableInstanceCommand("cmd-1", "Site1.Pump1", DateTimeOffset.UtcNow);
actor.Tell(new SiteEnvelope("site-a", cmd), probe.Ref);
AwaitAssert(() => transport.Received(1).Send(
Arg.Is<SiteEnvelope>(e => e.SiteId == "site-a" && ReferenceEquals(e.Message, cmd)),
probe.Ref));
}
[Fact]
public async Task AskReply_FromTransport_RoutesBackToTheWaitingAsk()
{
var (actor, transport, _) = CreateActor();
// The transport substitute stands in for the site: when handed the envelope it delivers a
// reply to the captured replyTo, exactly as a real transport pipes the site's reply back.
var reply = new DeploymentStatusResponse(
"dep-1", "Site1.Pump1", DeploymentStatus.Success, null, DateTimeOffset.UtcNow);
transport
.When(t => t.Send(Arg.Any<SiteEnvelope>(), Arg.Any<IActorRef>()))
.Do(ci => ci.Arg<IActorRef>().Tell(reply));
var cmd = new RefreshDeploymentCommand(
"dep-1", "Site1.Pump1", "hash", "multi-role", DateTimeOffset.UtcNow, "http://c:5000", "tok");
var result = await actor.Ask<DeploymentStatusResponse>(
new SiteEnvelope("site-a", cmd), TimeSpan.FromSeconds(3));
Assert.Equal("dep-1", result.DeploymentId);
Assert.Equal(DeploymentStatus.Success, result.Status);
}
[Fact]
public void DbRefresh_ReconcilesTheTransport_WithTheLoadedGrpcEndpoints()
{
var (_, transport, _) = CreateActor(new[] { GrpcSite("site-a") });
// PreStart fires the refresh at Zero; the loaded cache is handed to the transport.
AwaitAssert(() => transport.Received().ReconcileSites(
Arg.Is<SiteAddressCacheLoaded>(c =>
c.KnownSiteIds.Contains("site-a")
&& c.GrpcContacts.ContainsKey("site-a")
&& c.GrpcContacts["site-a"].NodeA == "http://a:8083"
&& c.GrpcContacts["site-a"].NodeB == "http://b:8083")),
TimeSpan.FromSeconds(3));
}
[Fact]
public void SiteAddedAndRemoved_AcrossRefreshes_FlowsThroughToTheTransport()
{
var (actor, transport, repo) = CreateActor(new[] { GrpcSite("site-a") });
AwaitAssert(() => transport.Received().ReconcileSites(
Arg.Is<SiteAddressCacheLoaded>(c => c.GrpcContacts.ContainsKey("site-a"))),
TimeSpan.FromSeconds(3));
// site-a removed, site-b added.
repo.GetAllSitesAsync(Arg.Any<CancellationToken>()).Returns(new List<Site> { GrpcSite("site-b") });
actor.Tell(new RefreshSiteAddresses());
AwaitAssert(() => transport.Received().ReconcileSites(
Arg.Is<SiteAddressCacheLoaded>(c =>
c.GrpcContacts.ContainsKey("site-b")
&& !c.GrpcContacts.ContainsKey("site-a")
&& c.KnownSiteIds.Contains("site-b")
&& !c.KnownSiteIds.Contains("site-a"))),
TimeSpan.FromSeconds(3));
}
[Fact]
public void SiteWithoutGrpcAddresses_IsAbsentFromGrpcContacts_ButStillKnown()
{
// A site with only Akka addresses (no gRPC columns) is a known site but carries no gRPC
// endpoint — the gRPC transport must not try to dial it, the Akka one still can.
var akkaOnly = new Site("Akka only", "site-x")
{
NodeAAddress = "akka.tcp://scadabridge@site-x-a:8081"
};
var (_, transport, _) = CreateActor(new[] { akkaOnly });
AwaitAssert(() => transport.Received().ReconcileSites(
Arg.Is<SiteAddressCacheLoaded>(c =>
c.KnownSiteIds.Contains("site-x") && !c.GrpcContacts.ContainsKey("site-x"))),
TimeSpan.FromSeconds(3));
}
}
@@ -0,0 +1,107 @@
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DataConnection;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Commons.Types;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// Pins <see cref="GrpcSiteTransport.ResolveDeadline"/> to the EXACT Ask timeout
/// <see cref="CommunicationService"/> uses for each command today, so flipping the transport cannot
/// change how long any call waits. Every timeout is given a distinct value so a wrong mapping is
/// caught, not masked by two defaults that happen to be equal (QueryTimeout == IntegrationTimeout in
/// production).
/// </summary>
public class GrpcSiteTransportDeadlineTests
{
private static readonly CommunicationOptions Opts = new()
{
DeploymentTimeout = TimeSpan.FromSeconds(120),
LifecycleTimeout = TimeSpan.FromSeconds(30),
ArtifactDeploymentTimeout = TimeSpan.FromSeconds(60),
QueryTimeout = TimeSpan.FromSeconds(31),
IntegrationTimeout = TimeSpan.FromSeconds(32),
DebugViewTimeout = TimeSpan.FromSeconds(10)
};
private static readonly GrpcSiteTransport Transport = new(
new SitePairChannelProvider(
new NoKeyProvider(), Options.Create(Opts), NullLogger<SitePairChannelProvider>.Instance),
Opts,
NullLogger<GrpcSiteTransport>.Instance);
private static readonly DateTimeOffset T = DateTimeOffset.UtcNow;
public static IEnumerable<object[]> Cases()
{
// Lifecycle group — note it is NOT one deadline class.
yield return [new RefreshDeploymentCommand("d", "i", "h", "by", T, "u", "t"), Opts.DeploymentTimeout];
yield return [new EnableInstanceCommand("c", "i", T), Opts.LifecycleTimeout];
yield return [new DisableInstanceCommand("c", "i", T), Opts.LifecycleTimeout];
yield return [new DeleteInstanceCommand("c", "i", T), Opts.LifecycleTimeout];
// DeploymentStateQuery uses QueryTimeout in CommunicationService, NOT LifecycleTimeout.
yield return [new DeploymentStateQueryRequest("c", "i", T), Opts.QueryTimeout];
yield return [new DeployArtifactsCommand("d", null, null, null, null, null, null, T), Opts.ArtifactDeploymentTimeout];
// OPC UA group — all QueryTimeout.
yield return [new BrowseNodeCommand("conn", null, null, null), Opts.QueryTimeout];
yield return [new SearchAddressSpaceCommand("conn", "q", 1, 1, null), Opts.QueryTimeout];
yield return [new ReadTagValuesCommand("conn", []), Opts.QueryTimeout];
yield return [new VerifyEndpointCommand("conn", "OpcUa", "{}", null), Opts.QueryTimeout];
yield return [new TrustServerCertCommand("conn", "ZGVy", "AA", null), Opts.QueryTimeout];
yield return [new ListServerCertsCommand(null), Opts.QueryTimeout];
yield return [new RemoveServerCertCommand("AA", null), Opts.QueryTimeout];
yield return [new WriteTagRequest("c", "conn", "tag", 1, T), Opts.QueryTimeout];
// Query group.
yield return [new EventLogQueryRequest("c", "s", null, null, null, null, null, null, null, 10, T), Opts.QueryTimeout];
yield return [new DebugSnapshotRequest("i", "c"), Opts.QueryTimeout];
yield return [new SubscribeDebugViewRequest("i", "c"), Opts.DebugViewTimeout];
yield return [new UnsubscribeDebugViewRequest("i", "c"), Opts.DebugViewTimeout];
// Parked group — the two relays keep QueryTimeout (30s) so SiteCallAudit's inner
// RelayTimeout (10s) still expires first.
yield return [new ParkedMessageQueryRequest("c", "s", 1, 10, T), Opts.QueryTimeout];
yield return [new ParkedMessageRetryRequest("c", "s", "m", T), Opts.QueryTimeout];
yield return [new ParkedMessageDiscardRequest("c", "s", "m", T), Opts.QueryTimeout];
yield return [new RetryParkedOperation("c", new TrackedOperationId(Guid.NewGuid())), Opts.QueryTimeout];
yield return [new DiscardParkedOperation("c", new TrackedOperationId(Guid.NewGuid())), Opts.QueryTimeout];
// Route group.
yield return [new RouteToCallRequest("c", "i", "m", null, T), Opts.IntegrationTimeout];
yield return [new RouteToGetAttributesRequest("c", "i", [], T), Opts.IntegrationTimeout];
yield return [new RouteToSetAttributesRequest("c", "i", new Dictionary<string, string>(), T), Opts.IntegrationTimeout];
// Failover — QueryTimeout in CommunicationService, NOT LifecycleTimeout.
yield return [new TriggerSiteFailover("c", "s"), Opts.QueryTimeout];
}
[Theory]
[MemberData(nameof(Cases))]
public void ResolveDeadline_MatchesTodaysAskTimeout(object command, TimeSpan expected)
{
Assert.Equal(expected, Transport.ResolveDeadline(command));
}
[Fact]
public void WaitForAttribute_UsesItsDynamicTimeoutPlusIntegrationSlack()
{
// CommunicationService.RouteToWaitForAttributeAsync uses request.Timeout + IntegrationTimeout.
var wait = new RouteToWaitForAttributeRequest("c", "i", "attr", "10", TimeSpan.FromSeconds(45), T);
Assert.Equal(TimeSpan.FromSeconds(45) + Opts.IntegrationTimeout, Transport.ResolveDeadline(wait));
}
private sealed class NoKeyProvider : ISitePskProvider
{
public ValueTask<string> GetAsync(string siteId, CancellationToken ct) => new("k");
public void Invalidate(string siteId) { }
}
}
@@ -0,0 +1,153 @@
using System.Text.Json;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// Goldens for <see cref="LooseValueCodec"/> — the type-tagged carrier for the
/// <c>object?</c> members of the site command plane (script parameters and
/// return values, attribute values, tag reads/writes).
/// </summary>
/// <remarks>
/// The contract these tests pin down is that a boxed value keeps its RUNTIME
/// CLR TYPE across the wire, not merely its printed form. Today's Akka JSON
/// serializer preserves it, and an operator-facing tag value that silently
/// turns from <c>int</c> into <c>long</c> (or from <c>DateTime</c> into a
/// UTC-normalised copy) is a behaviour change, not a refactor.
/// </remarks>
public class LooseValueCodecTests
{
public static IEnumerable<object[]> TaggedScalars() =>
[
[true],
[false],
[42],
[-1],
[long.MaxValue],
[3.5d],
[1.5f],
["a string"],
[string.Empty],
[12.3456789m],
[new DateTime(2026, 7, 22, 1, 2, 3, DateTimeKind.Utc)],
[new DateTimeOffset(2026, 7, 22, 1, 2, 3, TimeSpan.FromHours(-5))],
[Guid.Parse("11111111-2222-3333-4444-555555555555")],
[TimeSpan.FromMinutes(90)],
[new byte[] { 1, 2, 3 }]
];
[Theory]
[MemberData(nameof(TaggedScalars))]
public void TaggedValues_KeepTheirClrType(object value)
{
var restored = LooseValueCodec.FromProto(LooseValueCodec.ToProto(value));
Assert.NotNull(restored);
Assert.Equal(value.GetType(), restored.GetType());
StructuralEquality.AssertDeepEqual(value, restored, value.GetType().Name);
}
/// <summary>
/// A local <see cref="DateTime"/> must come back local, and an offset
/// <see cref="DateTimeOffset"/> must come back with the same offset — the
/// reason these ride as round-trip strings instead of proto timestamps.
/// </summary>
[Fact]
public void DateValues_KeepKindAndOffset()
{
var local = new DateTime(2026, 7, 22, 1, 2, 3, DateTimeKind.Local);
var offset = new DateTimeOffset(2026, 7, 22, 1, 2, 3, TimeSpan.FromHours(5.5));
var restoredLocal = Assert.IsType<DateTime>(LooseValueCodec.FromProto(LooseValueCodec.ToProto(local)));
var restoredOffset =
Assert.IsType<DateTimeOffset>(LooseValueCodec.FromProto(LooseValueCodec.ToProto(offset)));
Assert.Equal(DateTimeKind.Local, restoredLocal.Kind);
Assert.Equal(local, restoredLocal);
Assert.Equal(TimeSpan.FromHours(5.5), restoredOffset.Offset);
}
/// <summary>An unset carrier and an explicit null marker both decode to null.</summary>
[Fact]
public void NullIsCarriedTwoWays()
{
Assert.Null(LooseValueCodec.FromProto(null));
Assert.Null(LooseValueCodec.FromProto(LooseValueCodec.ToProto(null)));
Assert.Null(LooseValueCodec.ToProtoOrNull(null));
}
/// <summary>
/// A dictionary entry whose value is null stays distinct from an absent key —
/// the reason <c>LooseNull</c> exists at all.
/// </summary>
[Fact]
public void NullMapValue_SurvivesAsAPresentKey()
{
var original = new Dictionary<string, object?> { ["present"] = null };
var restored = LooseValueCodec.FromProtoMap(LooseValueCodec.ToProtoMap(original));
Assert.True(restored.ContainsKey("present"));
Assert.Null(restored["present"]);
}
/// <summary>A null dictionary stays null; an empty one stays empty.</summary>
[Fact]
public void NullMap_StaysDistinctFromEmptyMap()
{
Assert.Null(LooseValueCodec.ToProtoMapOrNull(null));
Assert.Null(LooseValueCodec.FromProtoMapOrNull(null));
var empty = LooseValueCodec.FromProtoMapOrNull(
LooseValueCodec.ToProtoMapOrNull(new Dictionary<string, object?>()));
Assert.NotNull(empty);
Assert.Empty(empty);
}
/// <summary>Nested lists and maps recurse, so element values keep their types too.</summary>
[Fact]
public void NestedCollections_KeepElementTypes()
{
object original = new List<object?>
{
1,
"two",
null,
new Dictionary<string, object?> { ["deep"] = 3.5d }
};
var restored = LooseValueCodec.FromProto(LooseValueCodec.ToProto(original));
StructuralEquality.AssertDeepEqual(original, restored, "list");
}
/// <summary>
/// The documented widening: a typed collection keeps its ELEMENTS' types but
/// the container comes back as <c>List&lt;object?&gt;</c>. Recorded here so
/// the behaviour is a decision rather than a surprise.
/// </summary>
[Fact]
public void TypedList_WidensToObjectList()
{
var restored = LooseValueCodec.FromProto(LooseValueCodec.ToProto(new List<int> { 1, 2, 3 }));
var list = Assert.IsType<List<object?>>(restored);
Assert.Equal([1, 2, 3], list.Cast<int>());
}
/// <summary>
/// The documented lossy escape hatch: a value outside the tagged set keeps its
/// DATA but not its CLR identity — it decodes as a <see cref="JsonElement"/>.
/// Nothing on the command plane carries such a value today; the fallback
/// exists so an unexpected one cannot fault a command.
/// </summary>
[Fact]
public void UntaggedValue_FallsBackToJsonAndLosesItsClrType()
{
var restored = LooseValueCodec.FromProto(LooseValueCodec.ToProto((byte)7));
var element = Assert.IsType<JsonElement>(restored);
Assert.Equal("7", element.GetRawText());
}
}
@@ -0,0 +1,360 @@
using Akka.Actor;
using Akka.TestKit.Xunit2;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Integration;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Commons.Types;
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// The single routing truth for the 28 migrated central→site commands. Proves each command
/// resolves to EXACTLY the target the actor used inline before the extraction — including the
/// node-local parked handler (never the singleton proxy) and the local failover path — so the
/// Akka actor and the gRPC service can share one table without drift.
/// </summary>
public class SiteCommandDispatcherTests : TestKit
{
private const string SiteId = "site1";
private SiteCommandDispatcher Build(IActorRef dmProxy, Func<string, bool, string?>? failover = null)
=> new(SiteId, dmProxy, failover ?? ((_, _) => null));
// ── The 19 commands that forward to the Deployment Manager singleton proxy ──
/// <summary>The proxy-routed commands (lifecycle, OPC UA, debug snapshot/subscribe, route).</summary>
public static IEnumerable<object[]> ProxyCommandTypes() => new[]
{
typeof(Commons.Messages.Deployment.RefreshDeploymentCommand),
typeof(Commons.Messages.Lifecycle.EnableInstanceCommand),
typeof(Commons.Messages.Lifecycle.DisableInstanceCommand),
typeof(Commons.Messages.Lifecycle.DeleteInstanceCommand),
typeof(Commons.Messages.Deployment.DeploymentStateQueryRequest),
typeof(Commons.Messages.Management.BrowseNodeCommand),
typeof(Commons.Messages.Management.SearchAddressSpaceCommand),
typeof(Commons.Messages.Management.ReadTagValuesCommand),
typeof(Commons.Messages.Management.VerifyEndpointCommand),
typeof(Commons.Messages.Management.TrustServerCertCommand),
typeof(Commons.Messages.Management.ListServerCertsCommand),
typeof(Commons.Messages.Management.RemoveServerCertCommand),
typeof(Commons.Messages.DataConnection.WriteTagRequest),
typeof(Commons.Messages.DebugView.DebugSnapshotRequest),
typeof(Commons.Messages.DebugView.SubscribeDebugViewRequest),
typeof(Commons.Messages.InboundApi.RouteToCallRequest),
typeof(Commons.Messages.InboundApi.RouteToGetAttributesRequest),
typeof(Commons.Messages.InboundApi.RouteToSetAttributesRequest),
typeof(Commons.Messages.InboundApi.RouteToWaitForAttributeRequest),
}.Select(t => new object[] { t });
[Theory]
[MemberData(nameof(ProxyCommandTypes))]
public void ProxyCommands_ForwardToDeploymentManager(Type commandType)
{
var dm = CreateTestProbe();
var dispatcher = Build(dm.Ref);
var command = SiteCommandSamples.All[commandType][0];
var route = dispatcher.ResolveRoute(command);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition);
Assert.Same(dm.Ref, route.Target);
}
[Fact]
public void UnsubscribeDebugView_IsFireAndForget_ToDeploymentManager_WithSyntheticAck()
{
// Fire-and-forget: the Deployment Manager never acks an unsubscribe. The actor Forwards it;
// the gRPC transport Tells it and returns this synthetic ack so a unary RPC still answers.
var dm = CreateTestProbe();
var dispatcher = Build(dm.Ref);
var command = SiteCommandSamples.All[typeof(Commons.Messages.DebugView.UnsubscribeDebugViewRequest)][0];
var route = dispatcher.ResolveRoute(command);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.TellFireAndForget, route.Disposition);
Assert.Same(dm.Ref, route.Target);
Assert.Same(UnsubscribeDebugViewAck.Instance, route.Reply);
}
// ── Artifact handler (null-guarded) ──
[Fact]
public void DeployArtifacts_WithHandler_ForwardsToArtifactHandler_NotTheProxy()
{
var dm = CreateTestProbe();
var artifact = CreateTestProbe();
var dispatcher = Build(dm.Ref);
dispatcher.RegisterArtifactHandler(artifact.Ref);
var route = dispatcher.ResolveRoute(SiteCommandSamples.All[typeof(DeployArtifactsCommand)][0]);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition);
Assert.Same(artifact.Ref, route.Target);
Assert.NotSame(dm.Ref, route.Target);
}
[Fact]
public void DeployArtifacts_WithoutHandler_RepliesHandlerNotAvailable()
{
var dm = CreateTestProbe();
var dispatcher = Build(dm.Ref);
var command = (DeployArtifactsCommand)SiteCommandSamples.All[typeof(DeployArtifactsCommand)][0];
var route = dispatcher.ResolveRoute(command);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.ImmediateReply, route.Disposition);
var reply = Assert.IsType<ArtifactDeploymentResponse>(route.Reply);
Assert.False(reply.Success);
Assert.Equal("Artifact handler not available", reply.ErrorMessage);
Assert.Equal(command.DeploymentId, reply.DeploymentId);
Assert.Equal(SiteId, reply.SiteId);
}
// ── Event-log handler (null-guarded) ──
[Fact]
public void EventLogQuery_WithHandler_ForwardsToEventLogHandler()
{
var dm = CreateTestProbe();
var eventLog = CreateTestProbe();
var dispatcher = Build(dm.Ref);
dispatcher.RegisterEventLogHandler(eventLog.Ref);
var route = dispatcher.ResolveRoute(SiteCommandSamples.All[typeof(EventLogQueryRequest)][0]);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition);
Assert.Same(eventLog.Ref, route.Target);
}
[Fact]
public void EventLogQuery_WithoutHandler_RepliesHandlerNotAvailable()
{
var dm = CreateTestProbe();
var dispatcher = Build(dm.Ref);
var command = (EventLogQueryRequest)SiteCommandSamples.All[typeof(EventLogQueryRequest)][0];
var route = dispatcher.ResolveRoute(command);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.ImmediateReply, route.Disposition);
var reply = Assert.IsType<EventLogQueryResponse>(route.Reply);
Assert.False(reply.Success);
Assert.Equal(command.CorrelationId, reply.CorrelationId);
}
// ── Parked handler (null-guarded) — stays NODE-LOCAL on purpose (replicated store) ──
/// <summary>The five parked commands, each of which routes to the per-node parked handler.</summary>
public static IEnumerable<object[]> ParkedCommandTypes() => new[]
{
typeof(ParkedMessageQueryRequest),
typeof(ParkedMessageRetryRequest),
typeof(ParkedMessageDiscardRequest),
typeof(RetryParkedOperation),
typeof(DiscardParkedOperation),
}.Select(t => new object[] { t });
[Theory]
[MemberData(nameof(ParkedCommandTypes))]
public void ParkedCommands_WithHandler_RouteToNodeLocalParkedHandler_NeverTheProxy(Type commandType)
{
// Node-locality proof: a parked retry/discard must run on the node holding the replicated
// store row, so it goes to the per-node parked handler — NEVER the singleton proxy. This is
// the constraint the extraction must not "fix".
var dm = CreateTestProbe();
var parked = CreateTestProbe();
var dispatcher = Build(dm.Ref);
dispatcher.RegisterParkedMessageHandler(parked.Ref);
var route = dispatcher.ResolveRoute(SiteCommandSamples.All[commandType][0]);
Assert.Equal(SiteCommandDispatcher.RouteDisposition.Forward, route.Disposition);
Assert.Same(parked.Ref, route.Target);
Assert.NotSame(dm.Ref, route.Target);
}
[Fact]
public void ParkedMessageQuery_WithoutHandler_RepliesHandlerNotAvailable()
{
var dispatcher = Build(CreateTestProbe().Ref);
var command = (ParkedMessageQueryRequest)SiteCommandSamples.All[typeof(ParkedMessageQueryRequest)][0];
var route = dispatcher.ResolveRoute(command);
var reply = Assert.IsType<ParkedMessageQueryResponse>(route.Reply);
Assert.False(reply.Success);
Assert.Equal(command.CorrelationId, reply.CorrelationId);
Assert.Equal(command.PageNumber, reply.PageNumber);
}
[Fact]
public void ParkedMessageRetry_WithoutHandler_RepliesHandlerNotAvailable()
{
var dispatcher = Build(CreateTestProbe().Ref);
var command = (ParkedMessageRetryRequest)SiteCommandSamples.All[typeof(ParkedMessageRetryRequest)][0];
var reply = Assert.IsType<ParkedMessageRetryResponse>(dispatcher.ResolveRoute(command).Reply);
Assert.False(reply.Success);
Assert.Equal(command.CorrelationId, reply.CorrelationId);
}
[Fact]
public void ParkedMessageDiscard_WithoutHandler_RepliesHandlerNotAvailable()
{
var dispatcher = Build(CreateTestProbe().Ref);
var command = (ParkedMessageDiscardRequest)SiteCommandSamples.All[typeof(ParkedMessageDiscardRequest)][0];
var reply = Assert.IsType<ParkedMessageDiscardResponse>(dispatcher.ResolveRoute(command).Reply);
Assert.False(reply.Success);
Assert.Equal(command.CorrelationId, reply.CorrelationId);
}
[Fact]
public void RetryParkedOperation_WithoutHandler_RepliesNotAppliedAck()
{
var dispatcher = Build(CreateTestProbe().Ref);
var command = new RetryParkedOperation("corr-x", TrackedOperationId.New());
var reply = Assert.IsType<ParkedOperationActionAck>(dispatcher.ResolveRoute(command).Reply);
Assert.False(reply.Applied);
Assert.Equal("corr-x", reply.CorrelationId);
Assert.NotNull(reply.ErrorMessage);
}
[Fact]
public void DiscardParkedOperation_WithoutHandler_RepliesNotAppliedAck()
{
var dispatcher = Build(CreateTestProbe().Ref);
var command = new DiscardParkedOperation("corr-y", TrackedOperationId.New());
var reply = Assert.IsType<ParkedOperationActionAck>(dispatcher.ResolveRoute(command).Reply);
Assert.False(reply.Applied);
Assert.Equal("corr-y", reply.CorrelationId);
}
// ── Commands that must NOT enter ResolveRoute ──
[Fact]
public void ResolveRoute_RejectsFailover_ItGoesThroughPrepareFailover()
{
var dispatcher = Build(CreateTestProbe().Ref);
Assert.Throws<ArgumentException>(() =>
dispatcher.ResolveRoute(new TriggerSiteFailover("c", SiteId)));
}
[Fact]
public void ResolveRoute_RejectsTheExcludedIntegrationCommand()
{
// IntegrationCallRequest is the 29th command, dead at both ends and deliberately excluded
// (28 of 29 migrate). It never enters the dispatcher.
var dispatcher = Build(CreateTestProbe().Ref);
var command = new IntegrationCallRequest(
"c", SiteId, "inst", "es", "m", new Dictionary<string, object?>(), DateTimeOffset.UtcNow);
Assert.Throws<ArgumentException>(() => dispatcher.ResolveRoute(command));
}
// ── Failover (local path) ──
[Fact]
public void HandleFailover_ResolvesAndLeaves_ThenAcksWithTheTarget()
{
string? roleAsked = null;
var leaveIssued = false;
Func<string, bool, string?> resolve = (role, dryRun) =>
{
roleAsked = role;
leaveIssued = !dryRun;
return "akka.tcp://scadabridge@site1-a:8082";
};
var dispatcher = Build(CreateTestProbe().Ref, resolve);
var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-1", SiteId));
Assert.True(ack.Accepted);
Assert.Equal("corr-1", ack.CorrelationId);
Assert.Equal("akka.tcp://scadabridge@site1-a:8082", ack.TargetAddress);
Assert.Null(ack.ErrorMessage);
// Site singletons are scoped to the site-specific role.
Assert.Equal("site-site1", roleAsked);
// The actor path leaves in one step (dryRun:false).
Assert.True(leaveIssued);
}
[Fact]
public void HandleFailover_RefusesWhenThereIsNoPeer()
{
var dispatcher = Build(CreateTestProbe().Ref, (_, _) => null);
var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-2", SiteId));
Assert.False(ack.Accepted);
Assert.Null(ack.TargetAddress);
Assert.NotNull(ack.ErrorMessage);
}
[Fact]
public void HandleFailover_RefusesACommandAddressedToAnotherSite_WithoutTouchingTheResolver()
{
var invoked = false;
var dispatcher = Build(CreateTestProbe().Ref, (_, _) => { invoked = true; return "addr"; });
var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-3", "site2"));
Assert.False(ack.Accepted);
Assert.Contains("site2", ack.ErrorMessage);
Assert.False(invoked);
}
[Fact]
public void HandleFailover_FaultInTheLeave_IsReportedNotThrown()
{
var dispatcher = Build(CreateTestProbe().Ref,
(_, _) => throw new InvalidOperationException("cluster unavailable"));
var ack = dispatcher.HandleFailover(new TriggerSiteFailover("corr-4", SiteId));
Assert.False(ack.Accepted);
Assert.Contains("cluster unavailable", ack.ErrorMessage);
}
[Fact]
public void PrepareFailover_BuildsTheAckFromADryRun_WithoutLeaving_UntilCommitLeaveRuns()
{
// The gRPC ack-before-Leave proof at the routing-truth level: PrepareFailover resolves the
// standby with a DRY-RUN only (so the ack is built without leaving), and hands back a
// deferred CommitLeave. Only invoking CommitLeave — what the gRPC service does AFTER the
// ack is on the wire — performs the real leave.
var events = new List<string>();
Func<string, bool, string?> resolve = (_, dryRun) =>
{
events.Add(dryRun ? "resolve" : "leave");
return "akka.tcp://scadabridge@site1-a:8082";
};
var dispatcher = Build(CreateTestProbe().Ref, resolve);
var outcome = dispatcher.PrepareFailover(new TriggerSiteFailover("corr-5", SiteId));
Assert.True(outcome.Ack.Accepted);
Assert.Equal("akka.tcp://scadabridge@site1-a:8082", outcome.Ack.TargetAddress);
Assert.NotNull(outcome.CommitLeave);
// Building the ack did NOT leave.
Assert.Equal(new[] { "resolve" }, events);
outcome.CommitLeave!();
// The leave runs strictly after — the caller sends the ack first.
Assert.Equal(new[] { "resolve", "leave" }, events);
}
[Fact]
public void PrepareFailover_WhenRefused_HasNoDeferredLeave()
{
var dispatcher = Build(CreateTestProbe().Ref, (_, _) => null);
var outcome = dispatcher.PrepareFailover(new TriggerSiteFailover("corr-6", SiteId));
Assert.False(outcome.Ack.Accepted);
Assert.Null(outcome.CommitLeave);
}
}
@@ -0,0 +1,569 @@
using System.Reflection;
using Google.Protobuf;
using Google.Protobuf.Reflection;
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Protocol;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Streaming;
using ZB.MOM.WW.ScadaBridge.Commons.Types.Alarms;
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// Round-trip goldens for <see cref="SiteCommandDtoMapper"/> — the 28 migrated
/// central→site commands, all 22 reply shapes, and every nested type they
/// carry, each proven to survive record → proto → record unchanged.
/// </summary>
/// <remarks>
/// <para>
/// The suite is REFLECTION-DRIVEN on purpose. Both the per-type theory and the
/// coverage guards enumerate the mapper's real surface (its <c>ToProto</c>
/// overloads) and the real proto contract (the generated <c>oneof</c>
/// descriptors), so adding a command without a golden fails the build rather
/// than quietly shipping an untested field. A hand-maintained list across 50
/// types would drift on the first change.
/// </para>
/// <para>
/// Where a round trip is not bit-exact, the mapper is the thing that gets fixed
/// — never the fixture. The two deliberate normalisations
/// (empty-string-means-null, and the implicit alarm <c>Condition</c>) are
/// asserted explicitly below rather than papered over.
/// </para>
/// </remarks>
public class SiteCommandDtoMapperGoldenTests
{
private static readonly IReadOnlyDictionary<string, Type> TypesByName =
SiteCommandSamples.All.Keys.ToDictionary(t => t.FullName!, t => t);
/// <summary>One theory case per (mapped type, golden sample).</summary>
/// <returns>Type name and sample index pairs covering every golden fixture.</returns>
public static IEnumerable<object[]> AllSamples() =>
SiteCommandSamples.All
.OrderBy(kv => kv.Key.FullName, StringComparer.Ordinal)
.SelectMany(kv => Enumerable.Range(0, kv.Value.Length)
.Select(i => new object[] { kv.Key.FullName!, i }));
[Theory]
[MemberData(nameof(AllSamples))]
public void RoundTrips_Unchanged(string typeName, int sampleIndex)
{
var type = TypesByName[typeName];
var original = SiteCommandSamples.All[type][sampleIndex];
var restored = RoundTrip(original, type);
StructuralEquality.AssertDeepEqual(original, restored, type.Name);
}
// ─────────────────────────────────────────────────────────────────────
// Coverage guards
// ─────────────────────────────────────────────────────────────────────
/// <summary>
/// Every non-enum type the mapper can project must have goldens. This is the
/// guard that makes the suite self-maintaining.
/// </summary>
[Fact]
public void EveryMappedType_HasGoldenSamples()
{
var missing = MappedRecordTypes()
.Where(t => !SiteCommandSamples.All.ContainsKey(t))
.Select(t => t.Name)
.OrderBy(n => n, StringComparer.Ordinal)
.ToList();
Assert.True(
missing.Count == 0,
"SiteCommandDtoMapper projects types with no golden sample: " + string.Join(", ", missing));
}
/// <summary>Every golden sample type must have at least two samples (populated + minimal).</summary>
[Fact]
public void EveryGoldenType_HasAtLeastOneSample()
{
var empty = SiteCommandSamples.All
.Where(kv => kv.Value.Length == 0)
.Select(kv => kv.Key.Name)
.ToList();
Assert.True(empty.Count == 0, "Golden types with no samples: " + string.Join(", ", empty));
}
/// <summary>
/// The contract carries exactly 28 commands — 29 on
/// <c>SiteCommunicationActor</c>'s receive table minus the dead
/// <c>IntegrationCallRequest</c>. If the site gains a command, this count
/// moves deliberately, not silently.
/// </summary>
[Fact]
public void CommandInventory_Is28_AcrossSixGroups()
{
var commands = CommandSampleTypes().ToList();
Assert.Equal(28, commands.Count);
Assert.Equal(
[6, 8, 4, 5, 4, 1],
new[]
{
SiteCommandGroup.Lifecycle, SiteCommandGroup.OpcUa, SiteCommandGroup.Query,
SiteCommandGroup.Parked, SiteCommandGroup.Route, SiteCommandGroup.Failover
}.Select(g => commands.Count(t => SiteCommandDtoMapper.GroupOf(Sample(t)) == g)).ToArray());
}
/// <summary>The contract carries 22 reply shapes across the same six groups.</summary>
[Fact]
public void ReplyInventory_Is22_AcrossSixGroups()
{
var replies = ReplyTypes().ToList();
Assert.Equal(22, replies.Count);
Assert.Equal(
[4, 6, 3, 4, 4, 1],
new[]
{
SiteCommandGroup.Lifecycle, SiteCommandGroup.OpcUa, SiteCommandGroup.Query,
SiteCommandGroup.Parked, SiteCommandGroup.Route, SiteCommandGroup.Failover
}.Select(g => replies.Count(t => SiteCommandDtoMapper.GroupOfReply(SampleReply(t)) == g)).ToArray());
}
/// <summary>
/// Packing every command golden must exercise EVERY <c>oneof</c> case
/// declared in the four multi-command request envelopes. A new proto case
/// with no producing command fails here.
/// </summary>
[Theory]
[MemberData(nameof(RequestEnvelopes))]
public void EveryRequestOneofCase_IsProducedByACommand(string envelopeName)
{
var produced = CommandSampleTypes()
.Select(t => PackCommand(Sample(t)))
.OfType<IMessage>()
.Where(m => m.Descriptor.Name == envelopeName)
.Select(OneofCaseName)
.ToHashSet(StringComparer.Ordinal);
var declared = DescriptorFor(envelopeName).Oneofs[0].Fields
.Select(f => f.PropertyName)
.ToHashSet(StringComparer.Ordinal);
Assert.Equal(declared.OrderBy(n => n, StringComparer.Ordinal), produced.OrderBy(n => n, StringComparer.Ordinal));
}
/// <summary>The reply-envelope mirror of <see cref="EveryRequestOneofCase_IsProducedByACommand"/>.</summary>
[Theory]
[MemberData(nameof(ReplyEnvelopes))]
public void EveryReplyOneofCase_IsProducedByAReply(string envelopeName)
{
var produced = ReplyTypes()
.Select(t => PackReply(SampleReply(t)))
.OfType<IMessage>()
.Where(m => m.Descriptor.Name == envelopeName)
.Select(OneofCaseName)
.ToHashSet(StringComparer.Ordinal);
var declared = DescriptorFor(envelopeName).Oneofs[0].Fields
.Select(f => f.PropertyName)
.ToHashSet(StringComparer.Ordinal);
Assert.Equal(declared.OrderBy(n => n, StringComparer.Ordinal), produced.OrderBy(n => n, StringComparer.Ordinal));
}
/// <summary>Names of the four multi-command request envelopes.</summary>
/// <returns>Envelope message names.</returns>
public static IEnumerable<object[]> RequestEnvelopes() =>
[["LifecycleRequest"], ["OpcUaRequest"], ["QueryRequest"], ["ParkedRequest"], ["RouteRequest"]];
/// <summary>Names of the four multi-reply reply envelopes.</summary>
/// <returns>Envelope message names.</returns>
public static IEnumerable<object[]> ReplyEnvelopes() =>
[["LifecycleReply"], ["OpcUaReply"], ["QueryReply"], ["ParkedReply"], ["RouteReply"]];
/// <summary>Every command golden survives the full envelope pack/unpack, not just its own message.</summary>
[Fact]
public void EveryCommandGolden_SurvivesItsEnvelope()
{
foreach (var (type, samples) in SiteCommandSamples.All)
{
if (!IsCommand(type))
{
continue;
}
foreach (var sample in samples)
{
var restored = UnpackCommand(PackCommand(sample));
StructuralEquality.AssertDeepEqual(sample, restored, $"{type.Name}(envelope)");
}
}
}
/// <summary>Every reply golden survives the full envelope pack/unpack.</summary>
[Fact]
public void EveryReplyGolden_SurvivesItsEnvelope()
{
foreach (var (type, samples) in SiteCommandSamples.All)
{
if (!IsReply(type))
{
continue;
}
foreach (var sample in samples)
{
var restored = UnpackReply(PackReply(sample));
StructuralEquality.AssertDeepEqual(sample, restored, $"{type.Name}(envelope)");
}
}
}
/// <summary>The fire-and-forget unsubscribe ack is a real envelope case, not a special path.</summary>
[Fact]
public void UnsubscribeDebugViewAck_RoundTripsThroughQueryReply()
{
var envelope = SiteCommandDtoMapper.ToQueryReply(UnsubscribeDebugViewAck.Instance);
Assert.Equal(QueryReply.ReplyOneofCase.UnsubscribeDebugView, envelope.ReplyCase);
Assert.Same(UnsubscribeDebugViewAck.Instance, SiteCommandDtoMapper.FromQueryReply(envelope));
}
// ─────────────────────────────────────────────────────────────────────
// Enum exhaustiveness
// ─────────────────────────────────────────────────────────────────────
[Theory]
[InlineData(DeploymentStatus.Pending)]
[InlineData(DeploymentStatus.InProgress)]
[InlineData(DeploymentStatus.Success)]
[InlineData(DeploymentStatus.Failed)]
public void DeploymentStatus_RoundTrips(DeploymentStatus value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(BrowseNodeClass.Object)]
[InlineData(BrowseNodeClass.Variable)]
[InlineData(BrowseNodeClass.Method)]
[InlineData(BrowseNodeClass.Other)]
public void BrowseNodeClass_RoundTrips(BrowseNodeClass value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(BrowseFailureKind.ConnectionNotFound)]
[InlineData(BrowseFailureKind.ConnectionNotConnected)]
[InlineData(BrowseFailureKind.NotBrowsable)]
[InlineData(BrowseFailureKind.Timeout)]
[InlineData(BrowseFailureKind.ServerError)]
public void BrowseFailureKind_RoundTrips(BrowseFailureKind value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(ReadTagValuesFailureKind.ConnectionNotFound)]
[InlineData(ReadTagValuesFailureKind.ConnectionNotConnected)]
[InlineData(ReadTagValuesFailureKind.Timeout)]
[InlineData(ReadTagValuesFailureKind.ServerError)]
public void ReadTagValuesFailureKind_RoundTrips(ReadTagValuesFailureKind value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(VerifyFailureKind.Unreachable)]
[InlineData(VerifyFailureKind.AuthFailed)]
[InlineData(VerifyFailureKind.UntrustedCertificate)]
[InlineData(VerifyFailureKind.Timeout)]
[InlineData(VerifyFailureKind.ServerError)]
public void VerifyFailureKind_RoundTrips(VerifyFailureKind value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(StoreAndForwardCategory.ExternalSystem)]
[InlineData(StoreAndForwardCategory.Notification)]
[InlineData(StoreAndForwardCategory.CachedDbWrite)]
public void StoreAndForwardCategory_RoundTrips(StoreAndForwardCategory value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(AlarmState.Active)]
[InlineData(AlarmState.Normal)]
public void AlarmState_RoundTrips(AlarmState value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(AlarmLevel.None)]
[InlineData(AlarmLevel.Low)]
[InlineData(AlarmLevel.LowLow)]
[InlineData(AlarmLevel.High)]
[InlineData(AlarmLevel.HighHigh)]
public void AlarmLevel_RoundTrips(AlarmLevel value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(AlarmKind.Computed)]
[InlineData(AlarmKind.NativeOpcUa)]
[InlineData(AlarmKind.NativeMxAccess)]
public void AlarmKind_RoundTrips(AlarmKind value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
[Theory]
[InlineData(AlarmShelveState.Unshelved)]
[InlineData(AlarmShelveState.OneShotShelved)]
[InlineData(AlarmShelveState.TimedShelved)]
[InlineData(AlarmShelveState.PermanentShelved)]
public void AlarmShelveState_RoundTrips(AlarmShelveState value) =>
Assert.Equal(value, SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(value)));
/// <summary>
/// An unspecified wire enum — what a peer on an older contract sends — must
/// decode to the documented safe default instead of faulting the command.
/// </summary>
[Fact]
public void UnspecifiedWireEnums_DecodeToSafeDefaults()
{
Assert.Equal(DeploymentStatus.Failed, SiteCommandDtoMapper.FromProto(DeploymentStatusDto.Unspecified));
Assert.Equal(BrowseNodeClass.Other, SiteCommandDtoMapper.FromProto(BrowseNodeClassDto.Unspecified));
Assert.Equal(BrowseFailureKind.ServerError, SiteCommandDtoMapper.FromProto(BrowseFailureKindDto.Unspecified));
Assert.Equal(
ReadTagValuesFailureKind.ServerError,
SiteCommandDtoMapper.FromProto(ReadTagValuesFailureKindDto.Unspecified));
Assert.Equal(VerifyFailureKind.ServerError, SiteCommandDtoMapper.FromProto(VerifyFailureKindDto.Unspecified));
Assert.Equal(
StoreAndForwardCategory.ExternalSystem,
SiteCommandDtoMapper.FromProto(StoreAndForwardCategoryDto.Unspecified));
Assert.Equal(AlarmState.Normal, SiteCommandDtoMapper.FromProto(AlarmStateDto.Unspecified));
Assert.Equal(AlarmLevel.None, SiteCommandDtoMapper.FromProto(AlarmLevelDto.Unspecified));
Assert.Equal(AlarmKind.Computed, SiteCommandDtoMapper.FromProto(AlarmKindDto.Unspecified));
Assert.Equal(AlarmShelveState.Unshelved, SiteCommandDtoMapper.FromProto(AlarmShelveStateDto.Unspecified));
}
// ─────────────────────────────────────────────────────────────────────
// The two deliberate normalisations, asserted rather than hidden
// ─────────────────────────────────────────────────────────────────────
/// <summary>
/// Nullable strings ride as plain proto3 strings, so an empty one comes back
/// as null. Documented in the mapper; asserted here so it stays a choice.
/// </summary>
[Fact]
public void EmptyNullableString_NormalisesToNull()
{
var original = new SiteFailoverAck("corr", Accepted: false, TargetAddress: "", ErrorMessage: "");
var restored = SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(original));
Assert.Null(restored.TargetAddress);
Assert.Null(restored.ErrorMessage);
}
/// <summary>
/// …but NOT for a wait target, where the empty string is a real value. This
/// is why that one field carries a <c>StringValue</c> wrapper.
/// </summary>
[Fact]
public void EmptyWaitTarget_StaysDistinctFromNull()
{
var empty = SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(
new RouteToWaitForAttributeRequest("c", "i", "a", "", TimeSpan.FromSeconds(1), DateTimeOffset.UtcNow)));
var missing = SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(
new RouteToWaitForAttributeRequest("c", "i", "a", null, TimeSpan.FromSeconds(1), DateTimeOffset.UtcNow)));
Assert.Equal(string.Empty, empty.TargetValueEncoded);
Assert.Null(missing.TargetValueEncoded);
}
/// <summary>
/// A computed alarm leaves <see cref="AlarmStateChanged.Condition"/> implicit
/// (the record derives it from State + Priority). The encoder omits a
/// condition that already equals that derived value, so the record comes back
/// byte-for-byte equal — including under record equality, which compares the
/// nullable backing field, not the property.
/// </summary>
[Fact]
public void ComputedAlarm_KeepsItsImplicitCondition()
{
var original = new AlarmStateChanged("inst", "alarm", AlarmState.Active, 500, DateTimeOffset.UtcNow);
var dto = SiteCommandDtoMapper.ToProto(original);
var restored = SiteCommandDtoMapper.FromProto(dto);
Assert.Null(dto.Condition);
Assert.Equal(original, restored);
}
/// <summary>
/// The flip side of the normalisation: a condition set EXPLICITLY to the value
/// the record would have derived comes back implicit. The
/// <see cref="AlarmStateChanged.Condition"/> value is identical — only the
/// record's private "was it set?" bit differs, which no consumer can observe.
/// </summary>
[Fact]
public void ExplicitConditionEqualToTheDerivedDefault_NormalisesToImplicit()
{
var original = new AlarmStateChanged("inst", "alarm", AlarmState.Normal, 100, DateTimeOffset.UtcNow)
{
Condition = AlarmConditionStateFactory.ForComputed(AlarmState.Normal, 100)
};
var restored = SiteCommandDtoMapper.FromProto(SiteCommandDtoMapper.ToProto(original));
Assert.Equal(original.Condition, restored.Condition);
}
/// <summary>A native alarm's explicit, non-derived condition is carried verbatim.</summary>
[Fact]
public void NativeAlarm_KeepsItsExplicitCondition()
{
var condition = new AlarmConditionState(true, false, true, AlarmShelveState.PermanentShelved, true, 999);
var original = new AlarmStateChanged("inst", "alarm", AlarmState.Active, 500, DateTimeOffset.UtcNow)
{
Kind = AlarmKind.NativeMxAccess,
Condition = condition
};
var dto = SiteCommandDtoMapper.ToProto(original);
var restored = SiteCommandDtoMapper.FromProto(dto);
Assert.NotNull(dto.Condition);
Assert.Equal(condition, restored.Condition);
}
// ─────────────────────────────────────────────────────────────────────
// Group classification + rejection
// ─────────────────────────────────────────────────────────────────────
/// <summary><c>IntegrationCallRequest</c> is excluded by design and must be rejected, not silently dropped.</summary>
[Fact]
public void IntegrationCallRequest_IsRejected()
{
var dead = new Commons.Messages.Integration.IntegrationCallRequest(
"corr", "site-a", "Instance", "System", "Method",
new Dictionary<string, object?>(), DateTimeOffset.UtcNow);
Assert.Throws<ArgumentException>(() => SiteCommandDtoMapper.GroupOf(dead));
}
/// <summary>Packing a command into the wrong group's envelope is a hard error, not a silent no-op.</summary>
[Fact]
public void PackingIntoTheWrongGroup_Throws() =>
Assert.Throws<ArgumentException>(() =>
SiteCommandDtoMapper.ToLifecycleRequest(new DebugSnapshotRequest("inst", "corr")));
/// <summary>An envelope with no oneof set (a newer peer's unknown case) surfaces as a clear failure.</summary>
[Fact]
public void UnsetOneof_ThrowsNotSupported() =>
Assert.Throws<NotSupportedException>(() => SiteCommandDtoMapper.FromLifecycleRequest(new LifecycleRequest()));
// ─────────────────────────────────────────────────────────────────────
// Reflection plumbing
// ─────────────────────────────────────────────────────────────────────
private static object RoundTrip(object original, Type type)
{
var toProto = typeof(SiteCommandDtoMapper)
.GetMethods(BindingFlags.Public | BindingFlags.Static)
.Single(m => m.Name == "ToProto" && m.GetParameters() is [{ } p] && p.ParameterType == type);
var wire = toProto.Invoke(null, [original])!;
var fromProto = typeof(SiteCommandDtoMapper)
.GetMethods(BindingFlags.Public | BindingFlags.Static)
.Single(m => m.Name == "FromProto"
&& m.GetParameters() is [{ } p]
&& p.ParameterType == wire.GetType());
return fromProto.Invoke(null, [wire])!;
}
private static IEnumerable<Type> MappedRecordTypes() =>
typeof(SiteCommandDtoMapper)
.GetMethods(BindingFlags.Public | BindingFlags.Static)
.Where(m => m.Name == "ToProto" && m.GetParameters().Length == 1)
.Select(m => m.GetParameters()[0].ParameterType)
.Where(t => !t.IsEnum)
.Distinct();
private static IEnumerable<Type> CommandSampleTypes() =>
SiteCommandSamples.All.Keys.Where(IsCommand);
/// <summary>
/// Reply shapes = the mapper's own classification, plus the synthetic
/// unsubscribe ack, which has no <c>ToProto</c> overload of its own.
/// </summary>
private static IEnumerable<Type> ReplyTypes() =>
SiteCommandSamples.All.Keys.Where(IsReply).Append(typeof(UnsubscribeDebugViewAck));
private static bool IsCommand(Type type) => Classifies(type, SiteCommandDtoMapper.GroupOf);
private static bool IsReply(Type type) => Classifies(type, SiteCommandDtoMapper.GroupOfReply);
private static bool Classifies(Type type, Func<object, SiteCommandGroup> classify)
{
try
{
classify(SiteCommandSamples.All[type][0]);
return true;
}
catch (ArgumentException)
{
return false;
}
}
private static object Sample(Type type) => SiteCommandSamples.All[type][0];
private static object SampleReply(Type type) =>
type == typeof(UnsubscribeDebugViewAck) ? UnsubscribeDebugViewAck.Instance : Sample(type);
private static object PackCommand(object command) => SiteCommandDtoMapper.GroupOf(command) switch
{
SiteCommandGroup.Lifecycle => SiteCommandDtoMapper.ToLifecycleRequest(command),
SiteCommandGroup.OpcUa => SiteCommandDtoMapper.ToOpcUaRequest(command),
SiteCommandGroup.Query => SiteCommandDtoMapper.ToQueryRequest(command),
SiteCommandGroup.Parked => SiteCommandDtoMapper.ToParkedRequest(command),
SiteCommandGroup.Route => SiteCommandDtoMapper.ToRouteRequest(command),
_ => SiteCommandDtoMapper.ToProto((TriggerSiteFailover)command)
};
private static object UnpackCommand(object envelope) => envelope switch
{
LifecycleRequest r => SiteCommandDtoMapper.FromLifecycleRequest(r),
OpcUaRequest r => SiteCommandDtoMapper.FromOpcUaRequest(r),
QueryRequest r => SiteCommandDtoMapper.FromQueryRequest(r),
ParkedRequest r => SiteCommandDtoMapper.FromParkedRequest(r),
RouteRequest r => SiteCommandDtoMapper.FromRouteRequest(r),
TriggerSiteFailoverDto d => SiteCommandDtoMapper.FromProto(d),
_ => throw new InvalidOperationException($"Unknown request envelope {envelope.GetType().Name}.")
};
private static object PackReply(object reply) => SiteCommandDtoMapper.GroupOfReply(reply) switch
{
SiteCommandGroup.Lifecycle => SiteCommandDtoMapper.ToLifecycleReply(reply),
SiteCommandGroup.OpcUa => SiteCommandDtoMapper.ToOpcUaReply(reply),
SiteCommandGroup.Query => SiteCommandDtoMapper.ToQueryReply(reply),
SiteCommandGroup.Parked => SiteCommandDtoMapper.ToParkedReply(reply),
SiteCommandGroup.Route => SiteCommandDtoMapper.ToRouteReply(reply),
_ => SiteCommandDtoMapper.ToProto((SiteFailoverAck)reply)
};
private static object UnpackReply(object envelope) => envelope switch
{
LifecycleReply r => SiteCommandDtoMapper.FromLifecycleReply(r),
OpcUaReply r => SiteCommandDtoMapper.FromOpcUaReply(r),
QueryReply r => SiteCommandDtoMapper.FromQueryReply(r),
ParkedReply r => SiteCommandDtoMapper.FromParkedReply(r),
RouteReply r => SiteCommandDtoMapper.FromRouteReply(r),
SiteFailoverAckDto d => SiteCommandDtoMapper.FromProto(d),
_ => throw new InvalidOperationException($"Unknown reply envelope {envelope.GetType().Name}.")
};
private static MessageDescriptor DescriptorFor(string name) =>
SiteCommandReflection.Descriptor.MessageTypes.Single(m => m.Name == name);
private static string OneofCaseName(IMessage envelope)
{
var oneof = envelope.Descriptor.Oneofs[0];
var field = oneof.Accessor.GetCaseFieldDescriptor(envelope);
Assert.NotNull(field);
return field.PropertyName;
}
}
@@ -0,0 +1,304 @@
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Protocol;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DataConnection;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Streaming;
using ZB.MOM.WW.ScadaBridge.Commons.Types;
using ZB.MOM.WW.ScadaBridge.Commons.Types.Alarms;
using ZB.MOM.WW.ScadaBridge.Commons.Types.Enums;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// Golden fixtures for the site-command round-trip suite: for every command,
/// reply and nested type the mapper handles, at least one MAXIMALLY populated
/// sample and one MINIMAL sample (every nullable null, every optional
/// collection empty or absent).
/// </summary>
/// <remarks>
/// <para>
/// The two-sample rule is the point. A single "typical" sample proves only that
/// the happy path survives; the pairs are what catch a dropped optional field or
/// a null/empty collapse — the class of silent data loss the transport
/// round-trip guard exposed in PLAN-05 T8, which per-field unit tests had all
/// missed.
/// </para>
/// <para>
/// <see cref="SiteCommandDtoMapperGoldenTests"/> drives this table by
/// reflection and FAILS when the mapper gains a type that has no sample here, so
/// the coverage cannot silently drift as commands are added.
/// </para>
/// </remarks>
internal static class SiteCommandSamples
{
private static readonly DateTimeOffset T1 = new(2026, 7, 22, 13, 45, 12, 345, TimeSpan.Zero);
private static readonly DateTimeOffset T2 = new(2026, 1, 2, 3, 4, 5, TimeSpan.Zero);
private static readonly DateTime D1 = new(2026, 3, 4, 5, 6, 7, DateTimeKind.Utc);
private static readonly DateTime D2 = new(2027, 3, 4, 5, 6, 7, DateTimeKind.Utc);
private static readonly Guid G1 = Guid.Parse("11111111-2222-3333-4444-555555555555");
/// <summary>
/// Every CLR type the mapper round-trips, mapped to its golden samples.
/// Enums are excluded — they get their own exhaustive test.
/// </summary>
public static IReadOnlyDictionary<Type, object[]> All { get; } = Build();
private static Dictionary<Type, object[]> Build()
{
var samples = new Dictionary<Type, object[]>();
void Add<T>(params T[] values) where T : notnull =>
samples[typeof(T)] = values.Cast<object>().ToArray();
// ── Lifecycle commands ──
Add(new RefreshDeploymentCommand(
"dep-1", "Site1.Pump1", "hash-abc", "multi-role", T1, "http://central:5000", "tok-1"));
Add(new EnableInstanceCommand("cmd-1", "Site1.Pump1", T1));
Add(new DisableInstanceCommand("cmd-2", "Site1.Pump1", T2));
Add(new DeleteInstanceCommand("cmd-3", "Site1.Pump1", T1));
Add(new DeploymentStateQueryRequest("corr-1", "Site1.Pump1", T1));
Add(
// Full: every artifact collection populated.
new DeployArtifactsCommand(
"dep-2",
[SharedScript(), new SharedScriptArtifact("s2", "return 2;", null, null)],
[ExternalSystem()],
[new DatabaseConnectionArtifact("db", "Server=x;", 3, TimeSpan.FromSeconds(7.5))],
[new NotificationListArtifact("ops", ["a@x.com", "b@x.com"])],
[DataConnection()],
[Smtp()],
T1),
// Minimal: every collection NULL — must not come back as empty lists.
new DeployArtifactsCommand("dep-3", null, null, null, null, null, null, T2),
// Boundary: every collection present but EMPTY — must not come back null.
new DeployArtifactsCommand("dep-4", [], [], [], [], [], [], T1));
Add(SharedScript(), new SharedScriptArtifact("bare", "", null, null));
Add(ExternalSystem(), new ExternalSystemArtifact("bare", "http://x", "None", null, null, 0));
Add(new DatabaseConnectionArtifact("db", "Server=x;", 3, TimeSpan.FromMilliseconds(1500)),
new DatabaseConnectionArtifact("db2", "", 0, TimeSpan.Zero));
Add(new NotificationListArtifact("ops", ["a@x.com"]), new NotificationListArtifact("empty", []));
Add(DataConnection(), new DataConnectionArtifact("bare", "OpcUa", null, null, 0));
Add(Smtp(), new SmtpConfigurationArtifact("bare", "smtp", 25, "None", "f@x.com", null, null, null));
// ── Lifecycle replies ──
Add(new DeploymentStatusResponse("dep-1", "Site1.Pump1", DeploymentStatus.Success, null, T1),
new DeploymentStatusResponse("dep-1", "Site1.Pump1", DeploymentStatus.Failed, "boom", T2));
Add(new InstanceLifecycleResponse("cmd-1", "Site1.Pump1", true, null, T1),
new InstanceLifecycleResponse("cmd-1", "Site1.Pump1", false, "nope", T2));
Add(new DeploymentStateQueryResponse("corr-1", "Site1.Pump1", true, "dep-1", "hash-abc", T1),
new DeploymentStateQueryResponse("corr-1", "Site1.Pump1", false, null, null, T2));
Add(new ArtifactDeploymentResponse("dep-2", "site-a", true, null, T1),
new ArtifactDeploymentResponse("dep-2", "site-a", false, "handler missing", T2));
// ── OPC UA commands ──
Add(new BrowseNodeCommand("conn", "ns=2;s=Root", "cursor-1", "site-a"),
new BrowseNodeCommand("conn", null, null, null));
Add(new SearchAddressSpaceCommand("conn", "pump", 4, 200, "site-a"),
new SearchAddressSpaceCommand("conn", "", 0, 0, null));
Add(new ReadTagValuesCommand("conn", ["a", "b"]), new ReadTagValuesCommand("conn", []));
Add(new VerifyEndpointCommand("conn", "OpcUa", "{\"url\":\"opc.tcp://x\"}", "site-a"),
new VerifyEndpointCommand("conn", "OpcUa", "{}", null));
Add(new TrustServerCertCommand("conn", "ZGVy", "AABB", "site-a"),
new TrustServerCertCommand("conn", "ZGVy", "AABB", null));
Add(new ListServerCertsCommand("site-a"), new ListServerCertsCommand(null));
Add(new RemoveServerCertCommand("AABB", "site-a"), new RemoveServerCertCommand("AABB", null));
Add(new WriteTagRequest("corr-2", "conn", "ns=2;s=Speed", 42.5d, T1),
new WriteTagRequest("corr-2", "conn", "ns=2;s=Speed", null, T2),
new WriteTagRequest("corr-2", "conn", "ns=2;s=Flag", true, T1),
new WriteTagRequest("corr-2", "conn", "ns=2;s=Name", "manual", T1));
// ── OPC UA replies ──
Add(new BrowseNodeResult([Node(), NodeBare()], true, null, "cursor-2"),
new BrowseNodeResult([], false, new BrowseFailure(BrowseFailureKind.Timeout, "timed out"), null));
Add(Node(), NodeBare());
Add(new BrowseFailure(BrowseFailureKind.ConnectionNotConnected, "not connected"),
new BrowseFailure(BrowseFailureKind.ServerError, ""));
Add(new SearchAddressSpaceResult([new AddressSpaceMatch(Node(), "/Root/Pump1")], true, null),
new SearchAddressSpaceResult([], false, new BrowseFailure(BrowseFailureKind.NotBrowsable, "no")));
Add(new AddressSpaceMatch(Node(), "/Root/Pump1"), new AddressSpaceMatch(NodeBare(), ""));
Add(new ReadTagValuesResult([Outcome(), OutcomeFailed()], null),
new ReadTagValuesResult([], new ReadTagValuesFailure(ReadTagValuesFailureKind.ConnectionNotFound, "gone")));
Add(Outcome(), OutcomeFailed());
Add(new ReadTagValuesFailure(ReadTagValuesFailureKind.Timeout, "slow"),
new ReadTagValuesFailure(ReadTagValuesFailureKind.ServerError, ""));
Add(new VerifyEndpointResult(true, null, null, null),
new VerifyEndpointResult(false, VerifyFailureKind.UntrustedCertificate, "untrusted", Cert()),
new VerifyEndpointResult(false, VerifyFailureKind.Unreachable, "refused", null));
Add(Cert());
Add(new CertTrustResult(true, null, [Trusted(), TrustedRejected()]),
new CertTrustResult(true, null, null),
new CertTrustResult(false, "partial failure", []));
Add(Trusted(), TrustedRejected());
Add(new WriteTagResponse("corr-2", true, null, T1),
new WriteTagResponse("corr-2", false, "denied", T2));
// ── Query commands ──
Add(new EventLogQueryRequest(
"corr-3", "site-a", T1, T2, "Lifecycle", "Warning", "Site1.Pump1", "restart", "cursor-3", 50, T1),
new EventLogQueryRequest("corr-3", "site-a", null, null, null, null, null, null, null, 25, T2));
Add(new DebugSnapshotRequest("Site1.Pump1", "corr-4"));
Add(new SubscribeDebugViewRequest("Site1.Pump1", "corr-5"));
Add(new UnsubscribeDebugViewRequest("Site1.Pump1", "corr-6"));
// ── Query replies ──
Add(new EventLogQueryResponse("corr-3", "site-a", [Entry(), EntryBare()], "cursor-4", true, true, null, T1),
new EventLogQueryResponse("corr-3", "site-a", [], null, false, false, "handler missing", T2));
Add(Entry(), EntryBare());
Add(new DebugViewSnapshot("Site1.Pump1", [AttrValue(), AttrValueNull()], [ComputedAlarm(), NativeAlarm()], T1),
new DebugViewSnapshot("Site1.Pump1", [], [], T2, InstanceNotFound: true));
Add(AttrValue(), AttrValueNull(), AttrValueList());
Add(ComputedAlarm(), NativeAlarm());
Add(new AlarmConditionState(true, false, true, AlarmShelveState.TimedShelved, true, 900),
new AlarmConditionState(false, true, null, AlarmShelveState.Unshelved, false, 0));
// ── Parked commands ──
Add(new ParkedMessageQueryRequest("corr-7", "site-a", 2, 25, T1));
Add(new ParkedMessageRetryRequest("corr-8", "site-a", "msg-1", T1));
Add(new ParkedMessageDiscardRequest("corr-9", "site-a", "msg-1", T2));
Add(new RetryParkedOperation("corr-10", new TrackedOperationId(G1)),
new RetryParkedOperation("corr-10", default));
Add(new DiscardParkedOperation("corr-11", new TrackedOperationId(G1)),
new DiscardParkedOperation("corr-11", default));
// ── Parked replies ──
Add(new ParkedMessageQueryResponse("corr-7", "site-a", [Parked(), ParkedBare()], 2, 1, 25, true, null, T1),
new ParkedMessageQueryResponse("corr-7", "site-a", [], 0, 1, 25, false, "handler missing", T2));
Add(Parked(), ParkedBare());
Add(new ParkedMessageRetryResponse("corr-8", true),
new ParkedMessageRetryResponse("corr-8", false, "not parked"));
Add(new ParkedMessageDiscardResponse("corr-9", true),
new ParkedMessageDiscardResponse("corr-9", false, "not parked"));
Add(new ParkedOperationActionAck("corr-10", true),
new ParkedOperationActionAck("corr-10", false, "handler missing"));
// ── Route commands ──
Add(new RouteToCallRequest("corr-12", "Site1.Pump1", "Start", Parameters(), T1, G1),
new RouteToCallRequest("corr-12", "Site1.Pump1", "Start", null, T2),
new RouteToCallRequest("corr-12", "Site1.Pump1", "Start", new Dictionary<string, object?>(), T1));
Add(new RouteToGetAttributesRequest("corr-13", "Site1.Pump1", ["Speed", "Flow"], T1, G1),
new RouteToGetAttributesRequest("corr-13", "Site1.Pump1", [], T2));
Add(new RouteToSetAttributesRequest(
"corr-14", "Site1.Pump1", new Dictionary<string, string> { ["Speed"] = "10" }, T1, G1),
new RouteToSetAttributesRequest("corr-14", "Site1.Pump1", new Dictionary<string, string>(), T2));
Add(new RouteToWaitForAttributeRequest(
"corr-15", "Site1.Pump1", "Speed", "10", TimeSpan.FromSeconds(30), T1, G1, true),
new RouteToWaitForAttributeRequest(
"corr-15", "Site1.Pump1", "Speed", null, TimeSpan.Zero, T2),
// "" is a legitimate wait target and must NOT collapse to null.
new RouteToWaitForAttributeRequest(
"corr-15", "Site1.Pump1", "Speed", "", TimeSpan.FromMinutes(1), T1));
// ── Route replies ──
Add(new RouteToCallResponse("corr-12", true, 17L, null, T1),
new RouteToCallResponse("corr-12", false, null, "script faulted", T2));
Add(new RouteToGetAttributesResponse("corr-13", Parameters(), true, null, T1),
new RouteToGetAttributesResponse("corr-13", new Dictionary<string, object?>(), false, "no instance", T2));
Add(new RouteToSetAttributesResponse("corr-14", true, null, T1),
new RouteToSetAttributesResponse("corr-14", false, "locked", T2));
Add(new RouteToWaitForAttributeResponse("corr-15", true, 10, "Good", false, true, null, T1),
new RouteToWaitForAttributeResponse("corr-15", false, null, null, true, true, null, T2));
// ── Failover ──
Add(new TriggerSiteFailover("corr-16", "site-a"));
Add(new SiteFailoverAck("corr-16", true, "akka.tcp://scadabridge@node-b:8081", null),
new SiteFailoverAck("corr-16", false, null, "no standby available"));
return samples;
}
private static SharedScriptArtifact SharedScript() =>
new("Calc", "return 1;", "{\"p\":\"int\"}", "{\"r\":\"int\"}");
private static ExternalSystemArtifact ExternalSystem() =>
new("Mes", "https://mes/api", "ApiKey", "{\"key\":\"x\"}", "[{\"name\":\"Post\"}]", 45);
private static DataConnectionArtifact DataConnection() =>
new("Plc1", "OpcUa", "{\"url\":\"opc.tcp://a\"}", "{\"url\":\"opc.tcp://b\"}", 5);
private static SmtpConfigurationArtifact Smtp() =>
new("Default", "smtp.host", 587, "Basic", "from@x.com", "user", "secret", "{\"tenant\":\"t\"}");
private static BrowseNode Node() => new("ns=2;s=Pump1.Speed", "Speed", BrowseNodeClass.Variable, false, "Double", -1, true);
private static BrowseNode NodeBare() => new("ns=2;s=Root", "Root", BrowseNodeClass.Object, true);
private static TagReadOutcome Outcome() => new("ns=2;s=Speed", true, 12.5d, "Good", T1, null);
private static TagReadOutcome OutcomeFailed() => new("ns=2;s=Bad", false, null, "Bad", T2, "no such node");
private static ServerCertInfo Cert() => new("AABB", "CN=server", "CN=issuer", D1, D2, "ZGVy");
private static TrustedCertInfo Trusted() => new("AABB", "CN=server", "CN=issuer", D1, D2, false);
private static TrustedCertInfo TrustedRejected() => new("CCDD", "CN=other", "CN=issuer", D1, D2, true);
private static EventLogEntry Entry() =>
new(G1.ToString("D"), T1, "Lifecycle", "Warning", "Site1.Pump1", "InstanceActor", "restarted", "{\"n\":1}");
private static EventLogEntry EntryBare() =>
new(Guid.Empty.ToString("D"), T2, "System", "Info", null, "Host", "started", null);
private static ParkedMessageEntry Parked() =>
new("msg-1", "Mes", "Post", "500 from server", 7, T1, T2, 10, StoreAndForwardCategory.CachedDbWrite, "Site1.Pump1");
private static ParkedMessageEntry ParkedBare() =>
new("msg-2", "Mes", "Post", "", 0, T2, T2);
private static AttributeValueChanged AttrValue() =>
new("Site1.Pump1", "Pump1.Speed", "Speed", 12.5d, "Good", T1);
private static AttributeValueChanged AttrValueNull() =>
new("Site1.Pump1", "Pump1.Speed", "Speed", null, "Bad", T2);
private static AttributeValueChanged AttrValueList() =>
new("Site1.Pump1", "Pump1.Trend", "Trend", new List<object?> { 1, 2.5d, "x", null }, "Good", T1);
/// <summary>A computed alarm: <c>Condition</c> left implicit, which is the common case.</summary>
private static AlarmStateChanged ComputedAlarm() =>
new("Site1.Pump1", "HighSpeed", AlarmState.Active, 700, T1)
{
Level = AlarmLevel.HighHigh,
Message = "Speed critically high"
};
/// <summary>A mirrored native alarm: every native enrichment field populated, <c>Condition</c> explicit.</summary>
private static AlarmStateChanged NativeAlarm() =>
new("Site1.Pump1", "Tank01.Level.HiHi", AlarmState.Active, 850, T2)
{
Level = AlarmLevel.None,
Message = "Level high",
Kind = AlarmKind.NativeOpcUa,
Condition = new AlarmConditionState(true, false, false, AlarmShelveState.OneShotShelved, true, 850),
SourceReference = "Tank01.Level.HiHi",
AlarmTypeName = "AnalogLimitAlarm.HiHi",
Category = "Process",
OperatorUser = "multi-role",
OperatorComment = "ack'd at panel",
OriginalRaiseTime = T1,
CurrentValue = "91.2",
LimitValue = "90.0",
NativeSourceCanonicalName = "Tank01.TankAlarms",
IsConfiguredPlaceholder = true
};
private static Dictionary<string, object?> Parameters() => new()
{
["count"] = 3,
["ratio"] = 1.25d,
["name"] = "batch-1",
["enabled"] = true,
["missing"] = null,
["when"] = T1,
["id"] = G1,
["items"] = new List<object?> { 1L, "two", null },
["nested"] = new Dictionary<string, object?> { ["inner"] = 5f }
};
}
@@ -0,0 +1,157 @@
using System.Collections;
using System.Reflection;
using System.Text.Json;
namespace ZB.MOM.WW.ScadaBridge.Communication.Tests;
/// <summary>
/// Structural deep-equality assertion for the site-command round-trip goldens.
/// </summary>
/// <remarks>
/// <para>
/// Record equality is NOT usable here. A C# record's generated <c>Equals</c>
/// compares members with <c>EqualityComparer&lt;T&gt;.Default</c>, which for
/// <c>IReadOnlyList&lt;T&gt;</c> / <c>IReadOnlyDictionary&lt;,&gt;</c> members
/// degrades to reference equality — so <c>Assert.Equal(dto, roundTripped)</c>
/// would fail on every collection-bearing message even when the mapper is
/// perfect, and (worse) would pass vacuously nowhere useful. This walker
/// compares by shape instead: dictionaries by key, sequences element-wise, and
/// everything else property-by-property down to leaf values.
/// </para>
/// <para>
/// The failure message carries the full property path, so a dropped field
/// reports as e.g. <c>DeployArtifactsCommand.ExternalSystems[0].TimeoutSeconds</c>
/// rather than an opaque "objects differ".
/// </para>
/// </remarks>
internal static class StructuralEquality
{
/// <summary>Asserts two object graphs are structurally identical.</summary>
/// <param name="expected">The original value.</param>
/// <param name="actual">The value that came back through the mapper.</param>
/// <param name="rootName">Name used as the root of the reported property path.</param>
public static void AssertDeepEqual(object? expected, object? actual, string rootName)
=> Compare(expected, actual, rootName);
private static void Compare(object? expected, object? actual, string path)
{
if (expected is null && actual is null)
{
return;
}
if (expected is null || actual is null)
{
Assert.Fail($"{path}: expected {Describe(expected)} but got {Describe(actual)}.");
return;
}
// JsonElement is the documented lossy escape hatch of LooseValueCodec; it
// has no useful Equals, so compare the canonical JSON text.
if (expected is JsonElement expectedJson && actual is JsonElement actualJson)
{
Assert.Equal(expectedJson.GetRawText(), actualJson.GetRawText());
return;
}
// Collections are compared by CONTENT, not by concrete container type.
// Members are declared IReadOnlyList<T>/IReadOnlyDictionary<,>, so the
// backing type (a collection-expression array here, a List<T> out of the
// mapper) is not part of the contract. Element and value types below are
// still compared strictly.
if (expected is not string && actual is not string)
{
if (expected is IDictionary expectedMap && actual is IDictionary actualMap)
{
CompareDictionaries(expectedMap, actualMap, path);
return;
}
if (expected is IEnumerable expectedSeq && actual is IEnumerable actualSeq)
{
CompareSequences(expectedSeq, actualSeq, path);
return;
}
}
var expectedType = expected.GetType();
var actualType = actual.GetType();
if (expectedType != actualType)
{
Assert.Fail(
$"{path}: type changed across the round trip — expected {expectedType.Name}, got {actualType.Name}.");
return;
}
if (IsLeaf(expectedType))
{
Assert.True(
Equals(expected, actual),
$"{path}: expected '{expected}' but got '{actual}'.");
return;
}
CompareProperties(expected, actual, expectedType, path);
}
private static void CompareDictionaries(IDictionary expected, IDictionary actual, string path)
{
Assert.True(
expected.Count == actual.Count,
$"{path}: entry count changed — expected {expected.Count}, got {actual.Count}.");
foreach (DictionaryEntry entry in expected)
{
Assert.True(
actual.Contains(entry.Key),
$"{path}: key '{entry.Key}' is missing after the round trip.");
Compare(entry.Value, actual[entry.Key], $"{path}['{entry.Key}']");
}
}
private static void CompareSequences(IEnumerable expected, IEnumerable actual, string path)
{
var expectedItems = expected.Cast<object?>().ToList();
var actualItems = actual.Cast<object?>().ToList();
Assert.True(
expectedItems.Count == actualItems.Count,
$"{path}: element count changed — expected {expectedItems.Count}, got {actualItems.Count}.");
for (var i = 0; i < expectedItems.Count; i++)
{
Compare(expectedItems[i], actualItems[i], $"{path}[{i}]");
}
}
private static void CompareProperties(object expected, object actual, Type type, string path)
{
var properties = type
.GetProperties(BindingFlags.Public | BindingFlags.Instance)
.Where(p => p.CanRead && p.GetIndexParameters().Length == 0)
.ToList();
Assert.True(properties.Count > 0, $"{path}: {type.Name} exposes no readable properties to compare.");
foreach (var property in properties)
{
Compare(property.GetValue(expected), property.GetValue(actual), $"{path}.{property.Name}");
}
}
private static bool IsLeaf(Type type) =>
type.IsPrimitive
|| type.IsEnum
|| type == typeof(string)
|| type == typeof(decimal)
|| type == typeof(DateTime)
|| type == typeof(DateTimeOffset)
|| type == typeof(TimeSpan)
|| type == typeof(Guid)
// Value types with no collection members (e.g. TrackedOperationId) have a
// correct structural Equals of their own.
|| (type.IsValueType && !type.IsGenericType);
private static string Describe(object? value) => value is null ? "<null>" : $"{value.GetType().Name}('{value}')";
}
@@ -0,0 +1,218 @@
using Grpc.Core;
using Grpc.Net.Client;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.TestHost;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Communication;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Host.Tests;
/// <summary>
/// T1B.3: <see cref="SitePairChannelProvider"/> failover/failback + credential/deadline proof over
/// two in-process gRPC <see cref="TestServer"/>s (NodeA preferred, NodeB standby). Exercises the
/// real client stack — PSK call credentials, per-call deadline, sticky failover, no-retry on
/// <see cref="StatusCode.DeadlineExceeded"/>, and failback to the preferred node.
/// </summary>
public sealed class GrpcSiteTransportFailoverTests : IAsyncLifetime
{
private const string SiteId = "site-1";
private const string SiteKey = "the-site-1-key";
private const string EndpointA = "http://node-a";
private const string EndpointB = "http://node-b";
private IHost _hostA = null!;
private IHost _hostB = null!;
private StubSiteCommandService _stubA = null!;
private StubSiteCommandService _stubB = null!;
private SitePairChannelProvider _provider = null!;
private volatile bool _preferredReachable;
/// <inheritdoc />
public async Task InitializeAsync()
{
(_hostA, _stubA) = await StartServerAsync();
(_hostB, _stubB) = await StartServerAsync();
var handlerA = _hostA.GetTestServer().CreateHandler();
var handlerB = _hostB.GetTestServer().CreateHandler();
_provider = new SitePairChannelProvider(
new FixedPskProvider(SiteKey),
Options.Create(new CommunicationOptions()),
NullLogger<SitePairChannelProvider>.Instance,
handlerFactory: endpoint => endpoint == EndpointA ? handlerA : handlerB,
reachabilityProbe: (_, _) => Task.FromResult(_preferredReachable));
_provider.UpdateSite(SiteId, EndpointA, EndpointB);
}
/// <inheritdoc />
public async Task DisposeAsync()
{
_provider.Dispose();
await _hostA.StopAsync();
await _hostB.StopAsync();
_hostA.Dispose();
_hostB.Dispose();
}
private static async Task<(IHost, StubSiteCommandService)> StartServerAsync()
{
var stub = new StubSiteCommandService();
var host = await new HostBuilder()
.ConfigureWebHost(web => web
.UseTestServer()
.ConfigureServices(services =>
{
services.AddGrpc();
services.AddSingleton(stub);
})
.Configure(app =>
{
app.UseRouting();
app.UseEndpoints(e => e.MapGrpcService<StubSiteCommandService>());
}))
.StartAsync();
return (host, stub);
}
private static LifecycleRequest EnableReq() =>
SiteCommandDtoMapper.ToLifecycleRequest(
new EnableInstanceCommand("cmd-1", "Site1.Pump1", DateTimeOffset.UtcNow));
private Task<LifecycleReply> CallAsync(TimeSpan? deadline = null) =>
_provider.ExecuteAsync(
SiteId,
(channel, ct) =>
{
var client = new SiteCommandService.SiteCommandServiceClient(channel);
return client.ExecuteLifecycleAsync(
EnableReq(),
deadline: DateTime.UtcNow + (deadline ?? TimeSpan.FromSeconds(10)),
cancellationToken: ct).ResponseAsync;
},
CancellationToken.None);
[Fact]
public async Task HealthyCall_HitsPreferredNodeA_WithPskAndSiteHeaderAndDeadline()
{
var reply = await CallAsync();
Assert.Equal(LifecycleReply.ReplyOneofCase.InstanceLifecycle, reply.ReplyCase);
Assert.Equal(1, _stubA.LifecycleCalls);
Assert.Equal(0, _stubB.LifecycleCalls);
Assert.Equal($"Bearer {SiteKey}", _stubA.LastAuthHeader);
Assert.Equal(SiteId, _stubA.LastSiteHeader);
Assert.True(_stubA.DeadlineWasSet, "the client must set a per-call deadline");
Assert.True(_provider.IsOnPreferredNode(SiteId));
}
[Fact]
public async Task Unavailable_OnNodeA_FailsOverToNodeB_AndStaysSticky()
{
_stubA.ThrowUnavailable = true;
var reply = await CallAsync();
// Failed over: A was tried (and threw), B answered.
Assert.Equal(LifecycleReply.ReplyOneofCase.InstanceLifecycle, reply.ReplyCase);
Assert.Equal(1, _stubA.LifecycleCalls);
Assert.Equal(1, _stubB.LifecycleCalls);
Assert.False(_provider.IsOnPreferredNode(SiteId));
// Sticky: the next call goes straight to B without re-touching A.
await CallAsync();
Assert.Equal(1, _stubA.LifecycleCalls);
Assert.Equal(2, _stubB.LifecycleCalls);
}
[Fact]
public async Task Failback_ReturnsToPreferredNodeA_OncePreferredIsReachableAgain()
{
_stubA.ThrowUnavailable = true;
await CallAsync(); // flips to B
Assert.False(_provider.IsOnPreferredNode(SiteId));
// NodeA recovers; the failback probe reports it reachable.
_stubA.ThrowUnavailable = false;
_preferredReachable = true;
var back = await _provider.TryFailbackAsync(SiteId, CancellationToken.None);
Assert.True(back);
Assert.True(_provider.IsOnPreferredNode(SiteId));
var beforeA = _stubA.LifecycleCalls;
await CallAsync();
Assert.Equal(beforeA + 1, _stubA.LifecycleCalls); // next call back on A
}
[Fact]
public async Task DeadlineExceeded_OnNodeA_IsNotRetriedOnNodeB()
{
// A DeadlineExceeded is ambiguous (a WriteTag/Deploy/Failover may already have executed), so
// — unlike Unavailable — it must NOT fail over to B. Modelled by A returning the status
// directly, isolating the retry-decision from TestServer's own timeout mechanics.
_stubA.StatusToThrow = StatusCode.DeadlineExceeded;
var ex = await Assert.ThrowsAsync<RpcException>(() => CallAsync());
Assert.Equal(StatusCode.DeadlineExceeded, ex.StatusCode);
Assert.Equal(1, _stubA.LifecycleCalls);
Assert.Equal(0, _stubB.LifecycleCalls); // B was never tried
Assert.True(_provider.IsOnPreferredNode(SiteId), "a deadline must not flip stickiness");
}
[Fact]
public async Task UnknownSite_Throws_SiteChannelUnavailable()
{
await Assert.ThrowsAsync<SiteChannelUnavailableException>(
() => _provider.ExecuteAsync<LifecycleReply>(
"not-configured", (_, _) => Task.FromResult(new LifecycleReply()), CancellationToken.None));
}
private sealed class FixedPskProvider(string key) : ISitePskProvider
{
public ValueTask<string> GetAsync(string siteId, CancellationToken ct) => new(key);
public void Invalidate(string siteId) { }
}
/// <summary>Stub SiteCommandService recording what the client sent and letting a test steer faults.</summary>
private sealed class StubSiteCommandService : SiteCommandService.SiteCommandServiceBase
{
private int _lifecycleCalls;
public int LifecycleCalls => Volatile.Read(ref _lifecycleCalls);
public string? LastAuthHeader { get; private set; }
public string? LastSiteHeader { get; private set; }
public bool DeadlineWasSet { get; private set; }
public volatile bool ThrowUnavailable;
public StatusCode? StatusToThrow;
public override Task<LifecycleReply> ExecuteLifecycle(LifecycleRequest request, ServerCallContext context)
{
Interlocked.Increment(ref _lifecycleCalls);
LastAuthHeader = context.RequestHeaders
.FirstOrDefault(h => h.Key == ControlPlaneCredentials.AuthorizationHeader)?.Value;
LastSiteHeader = context.RequestHeaders
.FirstOrDefault(h => h.Key == ControlPlaneCredentials.SiteHeader)?.Value;
DeadlineWasSet = context.Deadline != DateTime.MaxValue;
if (ThrowUnavailable)
{
throw new RpcException(new Status(StatusCode.Unavailable, "node down"));
}
if (StatusToThrow is { } status)
{
throw new RpcException(new Status(status, "modelled fault"));
}
return Task.FromResult(SiteCommandDtoMapper.ToLifecycleReply(
new InstanceLifecycleResponse("cmd-1", "Site1.Pump1", true, null, DateTimeOffset.UtcNow)));
}
}
}
@@ -0,0 +1,285 @@
using Akka.Actor;
using Grpc.Core;
using Grpc.Net.Client;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.TestHost;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
using ZB.MOM.WW.ScadaBridge.Communication;
using ZB.MOM.WW.ScadaBridge.Communication.Actors;
using ZB.MOM.WW.ScadaBridge.Communication.Grpc;
namespace ZB.MOM.WW.ScadaBridge.Host.Tests;
/// <summary>
/// The site command plane's gRPC front door (T1B.2): the same PSK gate and readiness convention
/// as <c>SiteStreamGrpcServer</c>, and a decode→dispatch→encode round trip through the ONE
/// <see cref="SiteCommandDispatcher"/> for one command in each of the six oneof groups. Runs
/// in-process over <see cref="TestServer"/> — no ports, no containers — with the interceptor
/// registered BY TYPE on <c>AddGrpc</c> (never pre-registered in DI), the shape that keeps
/// <c>Grpc.AspNetCore</c>'s own activation in the picture. See
/// <see cref="ControlPlaneAuthEndToEndTests"/> for why that matters.
/// </summary>
public sealed class SiteCommandGrpcServiceTests : IDisposable
{
private const string SiteKey = "the-site-a-command-key";
private const string SiteId = "site-a";
private readonly ActorSystem _system = ActorSystem.Create("sitecmd-tests");
/// <inheritdoc />
public void Dispose() => _system.Dispose();
/// <summary>A dispatcher whose Deployment Manager proxy is a canned-reply responder.</summary>
private SiteCommandDispatcher DispatcherWithResponder(out IActorRef responder)
{
responder = _system.ActorOf(Props.Create(() => new Responder()));
var dispatcher = new SiteCommandDispatcher(SiteId, responder, (_, _) => null);
// The parked handler is node-local; point it at the responder too so the Parked group can
// be exercised end-to-end (the null-guard path is covered by the dispatcher unit tests).
dispatcher.RegisterParkedMessageHandler(responder);
return dispatcher;
}
private SiteCommandGrpcService ReadyService()
{
var service = new SiteCommandGrpcService(NullLogger<SiteCommandGrpcService>.Instance);
service.SetReady(DispatcherWithResponder(out _));
return service;
}
private static async Task<IHost> StartHost(SiteCommandGrpcService service)
=> await new HostBuilder()
.ConfigureWebHost(web => web
.UseTestServer()
.ConfigureServices(services =>
{
// BY TYPE on AddGrpc — never AddSingleton the interceptor (see the sibling test).
services.AddGrpc(o => o.Interceptors.Add<ControlPlaneAuthInterceptor>());
services.AddSingleton(Options.Create(new CommunicationOptions { GrpcPsk = SiteKey }));
services.AddSingleton(service);
})
.Configure(app =>
{
app.UseRouting();
app.UseEndpoints(e => e.MapGrpcService<SiteCommandGrpcService>());
}))
.StartAsync();
private static SiteCommandService.SiteCommandServiceClient Client(IHost host, string? key)
{
var server = host.GetTestServer();
var options = new GrpcChannelOptions { HttpHandler = server.CreateHandler() };
if (key is not null)
{
options.WithSiteCredentials(new FixedPskProvider(key), SiteId);
}
var channel = GrpcChannel.ForAddress(server.BaseAddress, options);
return new SiteCommandService.SiteCommandServiceClient(channel);
}
private sealed class FixedPskProvider(string key) : ISitePskProvider
{
public ValueTask<string> GetAsync(string siteId, CancellationToken ct) => new(key);
public void Invalidate(string siteId) { }
}
// ── Auth gating (delegated to ControlPlaneAuthInterceptor, proven wired here) ──
[Fact]
public async Task NoCredentials_IsRejected_WithPermissionDenied()
{
using var host = await StartHost(ReadyService());
var client = Client(host, key: null);
var ex = await Assert.ThrowsAsync<RpcException>(() => client.ExecuteQueryAsync(
new QueryRequest { UnsubscribeDebugView = new UnsubscribeDebugViewRequestDto() }).ResponseAsync);
Assert.Equal(StatusCode.PermissionDenied, ex.StatusCode);
}
[Fact]
public async Task WrongKey_IsRejected_WithPermissionDenied()
{
using var host = await StartHost(ReadyService());
var client = Client(host, "some-other-key");
var ex = await Assert.ThrowsAsync<RpcException>(() => client.ExecuteQueryAsync(
new QueryRequest { UnsubscribeDebugView = new UnsubscribeDebugViewRequestDto() }).ResponseAsync);
Assert.Equal(StatusCode.PermissionDenied, ex.StatusCode);
}
// ── Readiness ──
[Fact]
public async Task BeforeSetReady_IsRejected_WithUnavailable_EvenWithACorrectKey()
{
// Auth passes (correct key); the readiness gate then rejects until the site actor graph is up.
var unready = new SiteCommandGrpcService(NullLogger<SiteCommandGrpcService>.Instance);
using var host = await StartHost(unready);
var client = Client(host, SiteKey);
var ex = await Assert.ThrowsAsync<RpcException>(() => client.ExecuteLifecycleAsync(
new LifecycleRequest { EnableInstance = new EnableInstanceCommandDto { CommandId = "c" } }).ResponseAsync);
Assert.Equal(StatusCode.Unavailable, ex.StatusCode);
}
// ── decode → dispatch → encode, one command per oneof group ──
[Fact]
public async Task ExecuteLifecycle_RoundTripsThroughTheDispatcher()
{
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var command = new EnableInstanceCommand("cmd-1", "Site1.Pump1", DateTimeOffset.UtcNow);
var reply = await client.ExecuteLifecycleAsync(
new LifecycleRequest { EnableInstance = SiteCommandDtoMapper.ToProto(command) });
var decoded = Assert.IsType<InstanceLifecycleResponse>(SiteCommandDtoMapper.FromLifecycleReply(reply));
Assert.True(decoded.Success);
Assert.Equal("cmd-1", decoded.CommandId);
}
[Fact]
public async Task ExecuteOpcUa_RoundTripsThroughTheDispatcher()
{
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var command = new BrowseNodeCommand("conn", null, null, null);
var reply = await client.ExecuteOpcUaAsync(
new OpcUaRequest { BrowseNode = SiteCommandDtoMapper.ToProto(command) });
Assert.IsType<BrowseNodeResult>(SiteCommandDtoMapper.FromOpcUaReply(reply));
}
[Fact]
public async Task ExecuteQuery_DebugSnapshot_RoundTripsThroughTheDispatcher()
{
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var command = new DebugSnapshotRequest("Site1.Pump1", "corr-4");
var reply = await client.ExecuteQueryAsync(
new QueryRequest { DebugSnapshot = SiteCommandDtoMapper.ToProto(command) });
var decoded = Assert.IsType<DebugViewSnapshot>(SiteCommandDtoMapper.FromQueryReply(reply));
Assert.Equal("Site1.Pump1", decoded.InstanceUniqueName);
}
[Fact]
public async Task ExecuteQuery_UnsubscribeDebugView_ReturnsTheFireAndForgetAck()
{
// Fire-and-forget: the service Tells the target and returns the synthetic ack so a unary
// RPC still answers, keeping the caller-visible fire-and-forget semantics.
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var reply = await client.ExecuteQueryAsync(new QueryRequest
{
UnsubscribeDebugView = SiteCommandDtoMapper.ToProto(
new UnsubscribeDebugViewRequest("Site1.Pump1", "corr-6")),
});
Assert.Equal(QueryReply.ReplyOneofCase.UnsubscribeDebugView, reply.ReplyCase);
}
[Fact]
public async Task ExecuteParked_RoundTripsThroughTheNodeLocalHandler()
{
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var command = new ParkedMessageQueryRequest("corr-7", SiteId, 2, 25, DateTimeOffset.UtcNow);
var reply = await client.ExecuteParkedAsync(
new ParkedRequest { ParkedMessageQuery = SiteCommandDtoMapper.ToProto(command) });
var decoded = Assert.IsType<ParkedMessageQueryResponse>(SiteCommandDtoMapper.FromParkedReply(reply));
Assert.True(decoded.Success);
Assert.Equal("corr-7", decoded.CorrelationId);
}
[Fact]
public async Task ExecuteRoute_RoundTripsThroughTheDispatcher()
{
using var host = await StartHost(ReadyService());
var client = Client(host, SiteKey);
var command = new RouteToCallRequest(
"corr-12", "Site1.Pump1", "Start", new Dictionary<string, object?>(), DateTimeOffset.UtcNow, null);
var reply = await client.ExecuteRouteAsync(
new RouteRequest { RouteToCall = SiteCommandDtoMapper.ToProto(command) });
var decoded = Assert.IsType<RouteToCallResponse>(SiteCommandDtoMapper.FromRouteReply(reply));
Assert.True(decoded.Success);
Assert.Equal("corr-12", decoded.CorrelationId);
}
// ── Failover: the reply completes before the leave is initiated ──
[Fact]
public async Task TriggerFailover_ReturnsTheAck_BeforeTheLeaveIsInitiated()
{
// The resolver records the order of its calls: PrepareFailover does a DRY-RUN resolve to
// build the ack (recorded synchronously, before the RPC returns), and the real leave runs
// only on the deferred CommitLeave — so the recorded order is always resolve-then-leave,
// i.e. the ack is on the wire before the node begins leaving.
var events = new List<string>();
var leaveHappened = new ManualResetEventSlim(false);
Func<string, bool, string?> resolve = (_, dryRun) =>
{
lock (events) { events.Add(dryRun ? "resolve" : "leave"); }
if (!dryRun) leaveHappened.Set();
return "akka.tcp://scadabridge@site-a-node-a:8082";
};
var service = new SiteCommandGrpcService(NullLogger<SiteCommandGrpcService>.Instance);
var dispatcher = new SiteCommandDispatcher(SiteId, _system.ActorOf(Props.Create(() => new Responder())), resolve);
service.SetReady(dispatcher);
using var host = await StartHost(service);
var client = Client(host, SiteKey);
var reply = await client.TriggerFailoverAsync(
new TriggerSiteFailoverDto { CorrelationId = "corr-16", SiteId = SiteId });
Assert.True(reply.Accepted);
Assert.Equal("akka.tcp://scadabridge@site-a-node-a:8082", reply.TargetAddress);
Assert.True(leaveHappened.Wait(TimeSpan.FromSeconds(5)), "the deferred leave never ran");
lock (events)
{
Assert.Equal(new[] { "resolve", "leave" }, events);
}
}
/// <summary>Deployment Manager stand-in: replies each command with a canned reply of the right group.</summary>
private sealed class Responder : ReceiveActor
{
public Responder() => ReceiveAny(msg => Sender.Tell(ReplyFor(msg)));
private static object ReplyFor(object m) => m switch
{
EnableInstanceCommand e =>
new InstanceLifecycleResponse(e.CommandId, e.InstanceUniqueName, true, null, DateTimeOffset.UtcNow),
BrowseNodeCommand => new BrowseNodeResult([], false, null, null),
DebugSnapshotRequest d => new DebugViewSnapshot(d.InstanceUniqueName, [], [], DateTimeOffset.UtcNow),
RouteToCallRequest r =>
new RouteToCallResponse(r.CorrelationId, true, null, null, DateTimeOffset.UtcNow),
ParkedMessageQueryRequest p => new ParkedMessageQueryResponse(
p.CorrelationId, p.SiteId, [], 0, p.PageNumber, p.PageSize, true, null, DateTimeOffset.UtcNow),
_ => new Akka.Actor.Status.Failure(new InvalidOperationException($"no canned reply for {m.GetType().Name}")),
};
}
}