feat(grpc): host CentralControlService on the central node (T1A.2)

Central now ALSO listens for the seven site→central control messages over
gRPC, alongside the existing ClusterClient path. Nothing flips to gRPC yet —
sites keep CentralTransport=Akka (T1A.3's job); central simply starts also
accepting.

- CentralControlGrpcService (Communication.Grpc): decodes each RPC onto the
  SAME in-process message the ClusterClient path carries, Asks the existing
  CentralCommunicationActor (zero handler-logic changes), encodes the reply via
  the T1A.1 mapper. Readiness-gated like SiteStreamGrpcServer.SetReady —
  Unavailable until AkkaHostedService hands the actor over. Heartbeat stays
  fire-and-forget (Tell, always-OK, never gated on readiness). Ingest reuses the
  shared SiteStreamGrpcServer.AuditIngestAskTimeout constant. Fault→status
  mapping is retry-aware: Unavailable (never dispatched, safe to cross-node
  retry) vs DeadlineExceeded/Internal (it ran, do not re-send elsewhere).

- CentralControlAuthInterceptor (Host): a SEPARATE interceptor class, not a
  variant constructor on ControlPlaneAuthInterceptor. Central's model is per-site
  (verify the Bearer token against the key for the site in the required
  x-scadabridge-site header, via ISitePskProvider) where a site verifies its one
  own-key — a genuinely different model. Fail-closed on every branch: missing or
  blank header, unresolvable key, and mismatched token all → PermissionDenied,
  never pass-through. One public constructor only (the explicit-prefix ctor is
  internal), pinned by a reflection test — a second public ctor makes
  Grpc.AspNetCore's GetFactory() throw per-call and silently disables the gate.

- Explicit Kestrel h2c listener on new option ScadaBridge:Node:CentralGrpcPort
  (default 8083, symmetric with sites), mirroring the Site branch. Additive to
  central's :5000 HTTP/1 surface, which is untouched — gRPC does NOT go through
  Traefik (HTTP/1 only). Registered by type on AddGrpc; service mapped with
  MapGrpcService. Port range-validated by NodeOptionsValidator.

- Rig: publish the central gRPC port 9013:8083 / 9014:8083 on both central nodes
  so a later task can exercise it.

Tests: CentralControlEndToEndTests (Host.Tests, TestServer + real interceptor +
real service over a stub actor) proves auth positives/negatives are
distinguishable and covers unary + the ingest bridge shapes; the interceptor is
registered BY TYPE, never in DI. CentralControlAuthInterceptorTests pins the
per-site gate + one-public-ctor invariant. CentralControlGrpcServiceTests
(Communication.Tests, TestKit) covers the readiness gate, fire-and-forget
heartbeat, and the DeadlineExceeded-vs-Unavailable status mapping. No active
<Protobuf> item. Communication.Tests (356) + Host.Tests (384) green.
This commit is contained in:
Joseph Doherty
2026-07-22 18:53:59 -04:00
parent d7455577a8
commit 780bb9c369
10 changed files with 1245 additions and 0 deletions
@@ -0,0 +1,299 @@
using Akka.Actor;
using Google.Protobuf.WellKnownTypes;
using Grpc.Core;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Audit;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Health;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Notification;
using GrpcStatus = Grpc.Core.Status;
namespace ZB.MOM.WW.ScadaBridge.Communication.Grpc;
/// <summary>
/// Central-hosted gRPC face of the seven site→central control messages
/// (<c>Protos/central_control.proto</c>). Decodes each request onto the SAME in-process
/// message type the Akka <c>ClusterClient</c> path already carries, <c>Ask</c>s
/// <see cref="Actors.CentralCommunicationActor"/>, and encodes the reply back.
/// </summary>
/// <remarks>
/// <para>
/// <b>Direction is inverted from <see cref="SiteStreamGrpcServer"/>.</b> That server runs on a
/// site and central dials in; this one runs on CENTRAL and the site dials in. The two listen on
/// the same port number (8083) on their respective nodes, which is symmetry, not a collision —
/// a node is either central or a site, never both.
/// </para>
/// <para>
/// <b>Zero handler logic lives here.</b> Every RPC lands on a receive
/// <c>CentralCommunicationActor</c> already implements for the ClusterClient path, so the two
/// transports cannot drift in behaviour: the actor is the single implementation, and this class
/// is a codec plus an <c>Ask</c>. That is also why the service takes the actor through
/// <see cref="SetReady"/> rather than resolving anything from DI — the actor is created by the
/// host's Akka bootstrap, not by the container.
/// </para>
/// <para>
/// <b>Fault semantics deliberately differ from <see cref="SiteStreamGrpcServer"/>'s ingest
/// RPCs.</b> That server answers a failed audit ingest with an EMPTY <c>IngestAck</c>; this one
/// fails the call with a non-OK status. Both leave the site's rows <c>Pending</c> for the next
/// drain, so the outcome is the same — but this service replaces the ClusterClient path, whose
/// documented behaviour is to propagate the fault (<c>CentralCommunicationActor</c>'s
/// <c>HandleIngestAuditEvents</c> pipes a <c>Status.Failure</c> back), and preserving that keeps
/// a lost batch visible as a failure rather than as a successful call that acked nothing.
/// </para>
/// <para>
/// <b>Status mapping, and why it is not uniform.</b> A site transport may safely re-send a call
/// to the peer central node only when the call provably never ran. So:
/// </para>
/// <list type="bullet">
/// <item><see cref="StatusCode.Unavailable"/> — this node is not ready; nothing was dispatched,
/// so a cross-node retry is safe and correct.</item>
/// <item><see cref="StatusCode.DeadlineExceeded"/> — the <c>Ask</c> timed out. The message WAS
/// delivered and may have been processed; retrying it on the other node would duplicate work.</item>
/// <item><see cref="StatusCode.Internal"/> — the handler faulted (a piped
/// <see cref="Akka.Actor.Status.Failure"/>, e.g. a database error inside reconcile). Same
/// reasoning: it ran, so do not re-send it elsewhere.</item>
/// </list>
/// </remarks>
public sealed class CentralControlGrpcService : CentralControlService.CentralControlServiceBase
{
private readonly ILogger<CentralControlGrpcService> _logger;
private readonly CommunicationOptions _options;
// Null until the host's Akka bootstrap hands the actor over. Doubles as the readiness
// flag: a call arriving before then cannot be served and is refused with Unavailable.
// Volatile because SetReady runs on the startup thread while calls are served on
// Kestrel's thread pool.
private volatile IActorRef? _central;
/// <summary>
/// Creates the service. <b>This must remain the only public constructor</b> — see
/// <see cref="SetReady"/> for how the actor arrives, and the Host's
/// <c>CentralControlAuthInterceptor</c> for the interceptor-side version of the same rule.
/// </summary>
/// <param name="logger">Logger for readiness and fault diagnostics.</param>
/// <param name="options">Communication options supplying the per-RPC <c>Ask</c> timeouts.</param>
public CentralControlGrpcService(
ILogger<CentralControlGrpcService> logger,
IOptions<CommunicationOptions> options)
{
ArgumentNullException.ThrowIfNull(logger);
ArgumentNullException.ThrowIfNull(options);
_logger = logger;
_options = options.Value;
}
/// <summary>
/// Hands the <c>CentralCommunicationActor</c> to the service and opens the gate. Mirrors
/// <see cref="SiteStreamGrpcServer.SetReady"/>: the gRPC service is a DI singleton created
/// before the actor system exists, so the actor arrives post-construction.
/// </summary>
/// <remarks>
/// The contract is deliberately narrow, exactly as on the site side: it asserts that the
/// actor exists and can receive, NOT that every downstream singleton proxy
/// (<c>notification-outbox</c>, <c>audit-log-ingest</c>) has registered itself yet. Those
/// register moments later in the same startup path, and the actor already answers a call
/// that beats them with the same "not available, retry" reply it gives on the ClusterClient
/// path — so gating readiness on them would add nothing but a longer window in which sites
/// see <see cref="StatusCode.Unavailable"/>.
/// </remarks>
/// <param name="centralCommunicationActor">The central communication actor.</param>
public void SetReady(IActorRef centralCommunicationActor)
{
ArgumentNullException.ThrowIfNull(centralCommunicationActor);
_central = centralCommunicationActor;
}
/// <summary>Exposed for wiring assertions in tests.</summary>
internal bool IsReady => _central is not null;
/// <inheritdoc />
public override async Task<NotificationSubmitAckDto> SubmitNotification(
NotificationSubmitDto request, ServerCallContext context)
{
var central = RequireReady(context);
var ack = await AskAsync<NotificationSubmitAck>(
central,
CentralControlDtoMapper.FromDto(request),
_options.NotificationForwardTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToDto(ack);
}
/// <inheritdoc />
public override async Task<NotificationStatusResponseDto> QueryNotificationStatus(
NotificationStatusQueryDto request, ServerCallContext context)
{
var central = RequireReady(context);
var response = await AskAsync<NotificationStatusResponse>(
central,
CentralControlDtoMapper.FromDto(request),
_options.NotificationForwardTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToDto(response);
}
/// <inheritdoc />
public override async Task<IngestAck> IngestAuditEvents(
AuditEventBatch request, ServerCallContext context)
{
// An empty batch is a no-op the actor need never see; answering it here also means a
// site that drains an empty queue does not fail against a not-yet-ready central.
if (request.Events.Count == 0)
{
return new IngestAck();
}
var central = RequireReady(context);
var reply = await AskAsync<IngestAuditEventsReply>(
central,
CentralControlDtoMapper.FromDto(request),
SiteStreamGrpcServer.AuditIngestAskTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToIngestAck(reply.AcceptedEventIds);
}
/// <inheritdoc />
public override async Task<IngestAck> IngestCachedTelemetry(
CachedTelemetryBatch request, ServerCallContext context)
{
if (request.Packets.Count == 0)
{
return new IngestAck();
}
var central = RequireReady(context);
var reply = await AskAsync<IngestCachedTelemetryReply>(
central,
CentralControlDtoMapper.FromDto(request),
SiteStreamGrpcServer.AuditIngestAskTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToIngestAck(reply.AcceptedEventIds);
}
/// <inheritdoc />
public override async Task<ReconcileSiteResponseDto> ReconcileSite(
ReconcileSiteRequestDto request, ServerCallContext context)
{
var central = RequireReady(context);
var response = await AskAsync<ReconcileSiteResponse>(
central,
CentralControlDtoMapper.FromDto(request),
_options.QueryTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToDto(response);
}
/// <inheritdoc />
public override async Task<SiteHealthReportAckDto> ReportSiteHealth(
SiteHealthReportDto request, ServerCallContext context)
{
var central = RequireReady(context);
var ack = await AskAsync<SiteHealthReportAck>(
central,
CentralControlDtoMapper.FromDto(request),
_options.HealthReportTimeout,
context).ConfigureAwait(false);
return CentralControlDtoMapper.ToDto(ack);
}
/// <summary>
/// Application heartbeat — <b>always answers OK</b>, even when this node is not ready.
/// </summary>
/// <remarks>
/// The heartbeat is fire-and-forget on both sides: nothing on the site consumes the reply,
/// and the site's heartbeat timer must never take a fault (a failing heartbeat that raised
/// an error would be a self-inflicted outage on a purely informational signal). So this is
/// the one RPC that does not go through <see cref="RequireReady"/>: a heartbeat arriving
/// before the actor exists is logged and dropped, exactly as the actor itself drops one
/// that arrives before <c>ICentralHealthAggregator</c> is resolvable. Liveness is still
/// detected — the aggregator's offline timeout fires when the heartbeats stop landing.
/// </remarks>
/// <param name="request">The heartbeat.</param>
/// <param name="context">The gRPC call context.</param>
/// <returns>An empty reply, always.</returns>
public override Task<Empty> Heartbeat(HeartbeatDto request, ServerCallContext context)
{
var central = _central;
if (central is null)
{
_logger.LogDebug(
"Dropped a heartbeat from site {SiteId}: the central communication actor is not "
+ "ready yet. Heartbeats are fire-and-forget, so the call still succeeds.",
request.SiteId);
return Task.FromResult(new Empty());
}
// Tell, never Ask: the actor's HandleHeartbeat sends no reply.
central.Tell(CentralControlDtoMapper.FromDto(request), ActorRefs.NoSender);
return Task.FromResult(new Empty());
}
/// <summary>
/// Returns the central communication actor, or throws <see cref="StatusCode.Unavailable"/>
/// when the host has not finished bringing the actor system up.
/// </summary>
private IActorRef RequireReady(ServerCallContext context)
{
var central = _central;
if (central is not null)
{
return central;
}
_logger.LogWarning(
"Refused a control-plane call to {Method}: the central communication actor is not "
+ "ready yet. Nothing was dispatched, so the caller may retry (including against "
+ "the peer central node).",
context.Method);
throw new RpcException(new GrpcStatus(
StatusCode.Unavailable,
"Central control plane is not ready: the actor system is still starting."));
}
/// <summary>
/// Asks the central actor and maps a fault onto the status code that tells the caller
/// whether a cross-node retry is safe. See the class remarks for the mapping rationale.
/// </summary>
private async Task<TReply> AskAsync<TReply>(
IActorRef central, object message, TimeSpan timeout, ServerCallContext context)
{
try
{
return await central.Ask<TReply>(message, timeout, context.CancellationToken)
.ConfigureAwait(false);
}
catch (AskTimeoutException ex)
{
_logger.LogWarning(ex,
"Control-plane call {Method} timed out after {Timeout} waiting for the central "
+ "communication actor.",
context.Method, timeout);
throw new RpcException(new GrpcStatus(
StatusCode.DeadlineExceeded,
$"Central did not answer within {timeout}."));
}
catch (OperationCanceledException)
{
// The client gave up or its deadline expired; there is no one left to answer.
throw new RpcException(new GrpcStatus(
StatusCode.Cancelled, "The call was cancelled."));
}
catch (Exception ex)
{
_logger.LogError(ex,
"Control-plane call {Method} faulted inside the central communication actor.",
context.Method);
throw new RpcException(new GrpcStatus(
StatusCode.Internal, "Central failed to process the call."));
}
}
}