63c16d6912
Phase 4 of the ClusterClient -> gRPC migration deleted the Akka transports but left the naming behind: `ClusterClientSiteAuditClient` was transport- agnostic and worked unchanged, so it survived the deletion under a name that now describes a transport the repo no longer has. Same for a scatter of doc-comments still framing gRPC as "the new transport" beside an Akka one that is gone. Renames it to `SiteCommunicationAuditClient` (and its test file) and rewrites the stale comments to describe the single transport that exists. Also tightens CLAUDE.md: drops the self-describing directory listing and the 27-component enumeration in favour of the non-obvious parts only. Behaviour-neutral: names and prose only. Recorded as the Phase 4 follow-up in docs/plans/2026-07-22-clusterclient-to-grpc-plan.md.
347 lines
18 KiB
C#
347 lines
18 KiB
C#
using Akka.Actor;
|
|
using Akka.Cluster;
|
|
using Akka.Event;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Artifacts;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Audit;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.DebugView;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Deployment;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Health;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.InboundApi;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Integration;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Lifecycle;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Management;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Notification;
|
|
using ZB.MOM.WW.ScadaBridge.Commons.Messages.RemoteQuery;
|
|
|
|
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
|
|
|
|
/// <summary>
|
|
/// Site-side actor that receives messages from central over the site-hosted
|
|
/// <c>SiteCommandService</c> gRPC plane and routes them to the appropriate local actors.
|
|
/// Also sends heartbeats and health reports to central over its
|
|
/// <see cref="ICentralTransport"/> (<see cref="Grpc.GrpcCentralTransport"/> in production).
|
|
///
|
|
/// Routes all 8 message patterns to local handlers.
|
|
/// </summary>
|
|
public class SiteCommunicationActor : ReceiveActor, IWithTimers
|
|
{
|
|
private readonly ILoggingAdapter _log = Context.GetLogger();
|
|
private readonly string _siteId;
|
|
private readonly CommunicationOptions _options;
|
|
|
|
/// <summary>
|
|
/// Predicate that returns <c>true</c> when this node is
|
|
/// the active member of the local site cluster (used to stamp
|
|
/// <see cref="HeartbeatMessage.IsActive"/>). Production builds default to
|
|
/// the Akka <see cref="Cluster"/> leader check; tests inject a stub so they
|
|
/// do not need a real cluster.
|
|
/// </summary>
|
|
private readonly Func<bool> _isActiveCheck;
|
|
|
|
/// <summary>
|
|
/// Reference to the local Deployment Manager singleton proxy.
|
|
/// </summary>
|
|
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>
|
|
/// The site→central transport. Finalized in <see cref="PreStart"/> to the injected instance, or a
|
|
/// fail-loud <see cref="NoOpCentralTransport"/> when none is supplied — so the seven site→central
|
|
/// sends delegate here rather than owning the wire plumbing inline. Production injects the gRPC
|
|
/// <see cref="Grpc.GrpcCentralTransport"/>.
|
|
/// </summary>
|
|
private ICentralTransport _transport;
|
|
|
|
/// <summary>The transport supplied by the Host (null falls back to <see cref="NoOpCentralTransport"/>).</summary>
|
|
private readonly ICentralTransport? _injectedTransport;
|
|
|
|
/// <summary>Akka timer scheduler injected by the framework via <see cref="IWithTimers"/>.</summary>
|
|
public ITimerScheduler Timers { get; set; } = null!;
|
|
|
|
/// <summary>Initializes the actor, wires all message pattern handlers, and schedules the periodic heartbeat.</summary>
|
|
/// <param name="siteId">The site identifier included in outbound messages.</param>
|
|
/// <param name="options">Communication options including heartbeat interval and transport settings.</param>
|
|
/// <param name="deploymentManagerProxy">Local reference to the Deployment Manager singleton proxy.</param>
|
|
/// <param name="isActiveCheck">
|
|
/// Optional override returning <c>true</c> when this node
|
|
/// is the active member of the site cluster. <c>null</c> uses the shared oldest-Up
|
|
/// evaluator (production wiring passes the Host's singleton-host delegate); tests
|
|
/// pass a stub so they do not need to load Akka.Cluster into the <c>TestKit</c>
|
|
/// ActorSystem.
|
|
/// </param>
|
|
/// <param name="transport">
|
|
/// The site→central transport. Production always passes the Host-built gRPC
|
|
/// <see cref="Grpc.GrpcCentralTransport"/>. <c>null</c> (used by TestKit suites that only
|
|
/// exercise command dispatch) falls back to a fail-loud <see cref="NoOpCentralTransport"/>.
|
|
/// </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(
|
|
string siteId,
|
|
CommunicationOptions options,
|
|
IActorRef deploymentManagerProxy,
|
|
Func<bool>? isActiveCheck = null,
|
|
Func<string, string?>? failOverRole = null,
|
|
ICentralTransport? transport = null,
|
|
SiteCommandDispatcher? dispatcher = null)
|
|
{
|
|
_siteId = siteId;
|
|
_options = options;
|
|
_deploymentManagerProxy = deploymentManagerProxy;
|
|
_isActiveCheck = isActiveCheck ?? DefaultIsActiveCheck;
|
|
_injectedTransport = transport;
|
|
// 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.
|
|
_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);
|
|
|
|
Receive<RegisterLocalHandler>(HandleRegisterLocalHandler);
|
|
|
|
// ── The 27 migrated central→site commands (28th is failover, below) all route
|
|
// through the shared SiteCommandDispatcher — the single routing truth also used by
|
|
// the gRPC SiteCommandGrpcService. The actor's job per command is unchanged: Forward
|
|
// to the resolved target (preserving the central Ask sender so replies route straight
|
|
// 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));
|
|
|
|
// 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
|
|
// 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
|
|
// singleton, and contact rotation picks whichever answers); the target is resolved
|
|
// from cluster state, not from who received the message. The gRPC transport defers the
|
|
// leave to keep ack-before-Leave — see PrepareFailover. (Under the removed ClusterClient
|
|
// path the ack Tell merely enqueued, so the dispatcher resolved-and-left in one step.)
|
|
Receive<TriggerSiteFailover>(msg => Sender.Tell(_dispatcher.HandleFailover(msg)));
|
|
|
|
// The seven site→central sends delegate to the injected transport (GrpcCentralTransport in
|
|
// production). Each handler captures the current Sender as the reply target so central's
|
|
// reply routes straight back to the waiting Ask, not through this actor — preserving the
|
|
// sender-forwarding the original ClusterClient path relied on. The per-message
|
|
// "no transport / not-accepted" fallbacks live inside the transport now.
|
|
|
|
// Notification Outbox: forward a buffered notification (S&F forwarder's Ask → ack back).
|
|
Receive<NotificationSubmit>(msg => _transport.SubmitNotification(msg, Sender));
|
|
|
|
// Notification Outbox: forward a Notify.Status query (Notify helper's Ask → response back).
|
|
Receive<NotificationStatusQuery>(msg => _transport.QueryNotificationStatus(msg, Sender));
|
|
|
|
// Audit Log: forward a batch of site-local audit events (telemetry drain's Ask → reply back).
|
|
Receive<IngestAuditEventsCommand>(msg => _transport.IngestAuditEvents(msg, Sender));
|
|
|
|
// Audit Log: forward a batch of combined cached-call telemetry (telemetry drain's Ask → reply back).
|
|
Receive<IngestCachedTelemetryCommand>(msg => _transport.IngestCachedTelemetry(msg, Sender));
|
|
|
|
// Site startup reconciliation: forward the node's local-inventory request (reconcile Ask → response back).
|
|
Receive<ReconcileSiteRequest>(msg => _transport.ReconcileSite(msg, Sender));
|
|
|
|
// Internal: send heartbeat tick.
|
|
Receive<SendHeartbeat>(_ => SendHeartbeatToCentral());
|
|
|
|
// Internal: forward the periodic health report (health transport's Ask → ack back), so a
|
|
// lost report is observable end-to-end and the sender can restore its per-interval counters.
|
|
Receive<SiteHealthReport>(msg => _transport.ReportSiteHealth(msg, Sender));
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
protected override SupervisorStrategy SupervisorStrategy()
|
|
{
|
|
return new OneForOneStrategy(
|
|
maxNrOfRetries: -1,
|
|
withinTimeRange: Timeout.InfiniteTimeSpan,
|
|
decider: Decider.From(ex =>
|
|
{
|
|
_log.Warning(ex, "Child actor of SiteCommunicationActor faulted, resuming (state preserved)");
|
|
return Directive.Resume;
|
|
}));
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
protected override void PreStart()
|
|
{
|
|
// Finalize the transport now that the actor context (and _log) exist. When the Host injects
|
|
// none, fall back to a fail-loud NoOpCentralTransport (given this actor's logging adapter so
|
|
// "no transport configured" warnings surface). PreStart always runs before any message, so
|
|
// the Receive closures above see a non-null _transport.
|
|
_transport = _injectedTransport ?? new NoOpCentralTransport(_log);
|
|
|
|
_log.Info("SiteCommunicationActor started for site {0}", _siteId);
|
|
|
|
// Schedule periodic heartbeat to central. Uses the application heartbeat
|
|
// cadence — distinct from the Akka.Remote transport failure-detector
|
|
// interval — so the two can be tuned independently.
|
|
Timers.StartPeriodicTimer(
|
|
"heartbeat",
|
|
new SendHeartbeat(),
|
|
TimeSpan.FromSeconds(1), // initial delay
|
|
_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: the site
|
|
/// never acked it under the original ClusterClient path either, 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)
|
|
{
|
|
// 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)
|
|
{
|
|
case LocalHandlerType.EventLog:
|
|
_dispatcher.RegisterEventLogHandler(msg.Handler);
|
|
break;
|
|
case LocalHandlerType.ParkedMessages:
|
|
_dispatcher.RegisterParkedMessageHandler(msg.Handler);
|
|
break;
|
|
case LocalHandlerType.Artifacts:
|
|
_dispatcher.RegisterArtifactHandler(msg.Handler);
|
|
break;
|
|
}
|
|
|
|
_log.Info("Registered local handler for {0}", msg.HandlerType);
|
|
}
|
|
|
|
private void SendHeartbeatToCentral()
|
|
{
|
|
var hostname = Environment.MachineName;
|
|
|
|
// Stamp HeartbeatMessage.IsActive with this node's
|
|
// true active/standby role rather than hard-coding `true`. The field is
|
|
// part of the wire contract (additive-only-evolution) so a future
|
|
// central health dashboard can distinguish "active node down, standby
|
|
// up" from "site fully offline" without a new message type.
|
|
bool isActive;
|
|
try
|
|
{
|
|
isActive = _isActiveCheck();
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
// Defensive: never let a cluster-state read failure abort the
|
|
// heartbeat itself (heartbeats are health signal — their absence is
|
|
// already meaningful). Fall back to the safest non-claiming value:
|
|
// standby. Logged at Debug because this path normally only fires
|
|
// during ActorSystem warm-up.
|
|
_log.Debug(ex,
|
|
"Active-node check threw while sending heartbeat for site {0}; reporting IsActive=false",
|
|
_siteId);
|
|
isActive = false;
|
|
}
|
|
|
|
var heartbeat = new HeartbeatMessage(
|
|
_siteId,
|
|
hostname,
|
|
IsActive: isActive,
|
|
DateTimeOffset.UtcNow);
|
|
|
|
// Fire-and-forget on both transports: a failure here must never fault the heartbeat timer
|
|
// path. Both real transports swallow their own errors; this catch is a belt-and-braces
|
|
// guarantee that no transport (including a future one) can turn a heartbeat into a fault.
|
|
try
|
|
{
|
|
_transport.SendHeartbeat(heartbeat, Self);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_log.Debug(ex, "Heartbeat send for site {0} failed; swallowed (heartbeats are fire-and-forget)", _siteId);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Default active-node check used when no override is supplied: oldest-Up member
|
|
/// semantics via the shared <see cref="ClusterState.ActiveNodeEvaluator"/> — the
|
|
/// same predicate as the S&F delivery gate and the replication resync authority
|
|
/// (review 02 round 2, N1). Unscoped by role: a site cluster's members all carry
|
|
/// the site role, so role scoping is a no-op here; production wiring passes the
|
|
/// Host's role-scoped IClusterNodeProvider delegate anyway. Any other state
|
|
/// (still joining, leaving) reports standby — safe-by-default.
|
|
/// </summary>
|
|
private bool DefaultIsActiveCheck() =>
|
|
ClusterState.ActiveNodeEvaluator.SelfIsOldestUp(Cluster.Get(Context.System));
|
|
|
|
// ── Internal messages ──
|
|
|
|
internal record SendHeartbeat;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Command to register a local actor as a handler for a specific message pattern.
|
|
/// </summary>
|
|
public record RegisterLocalHandler(LocalHandlerType HandlerType, IActorRef Handler);
|
|
|
|
public enum LocalHandlerType
|
|
{
|
|
EventLog,
|
|
ParkedMessages,
|
|
Artifacts
|
|
}
|