using Akka.Actor;
using Akka.Cluster.Tools.PublishSubscribe;
using Akka.Event;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using ZB.MOM.WW.ScadaBridge.Commons.Interfaces.Repositories;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Audit;
using ZB.MOM.WW.ScadaBridge.Commons.Messages.Communication;
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 ZB.MOM.WW.ScadaBridge.HealthMonitoring;
namespace ZB.MOM.WW.ScadaBridge.Communication.Actors;
///
/// Central-side actor that routes messages from central to site clusters over the gRPC
/// SiteCommandService command plane (). Resolves site
/// addresses from the database on a periodic refresh cycle and reconciles the transport's per-site
/// channel pairs.
///
/// All 8 message patterns routed through this actor.
/// Ask timeout on connection drop (no central buffering). Debug streams killed on interruption.
///
public class CentralCommunicationActor : ReceiveActor
{
private readonly ILoggingAdapter _log = Context.GetLogger();
private readonly IServiceProvider _serviceProvider;
///
/// The central→site command transport (the gRPC , built by
/// the Host and injected). The handler delegates every send here, and
/// each DB refresh tick reconciles its per-site channel pairs. Assigned in the constructor body
/// before any message can arrive.
///
private ISiteCommandTransport _transport = null!;
// The previous _debugSubscriptions / _inProgressDeployments
// dictionaries existed solely to support a documented "synchronous kill streams +
// mark deployments failed on site disconnect" workflow triggered by
// ConnectionStateChanged. No production code ever emitted that message — only
// the unit test did — so the workflow was dead from end to end. Disconnect
// detection is owned by the underlying transports: the gRPC keepalive PING
// signals stream interruption in ~25s (handled by DebugStreamBridgeActor's own
// reconnection logic), and an Ask round-trip for a deploy times out at the
// CommunicationService layer (caller sees failure). The tracking dicts +
// ConnectionStateChanged record + HandleConnectionStateChanged handler are
// removed; see docs/requirements/Component-Communication.md "Connection
// Failure Behavior" for the keepalive-based contract that survives.
private ICancelable? _refreshSchedule;
///
/// Per-actor lifecycle CTS threaded into the periodic
/// repository call so a hung MS SQL
/// connection is bounded by actor shutdown rather than holding piped tasks
/// open indefinitely. Cancelled in ; never reset.
///
private readonly CancellationTokenSource _lifecycleCts = new();
///
/// Proxy for the central NotificationOutboxActor cluster singleton.
/// Set via — the Host creates the singleton proxy
/// after this actor and registers it (mirrors how the site-side actor receives its
/// runtime s). Null until registration completes; a notification
/// arriving before then is rejected with a non-accepted ack so the site retries.
///
private IActorRef? _notificationOutboxProxy;
///
/// Proxy for the central AuditLogIngestActor cluster
/// singleton. Set via — the Host creates the
/// singleton proxy after this actor and registers it (mirrors
/// ). Null until registration completes;
/// an audit ingest command arriving before then is answered with an empty
/// reply so the site keeps its rows Pending and retries.
///
/// Once registered, the handler Asks this proxy and pipes the reply straight
/// back to the caller. On an Ask timeout or a faulted reply, PipeTo forwards a
/// to the caller — the fault propagates rather
/// than being swallowed. This differs from the gRPC handler
/// (SiteStreamGrpcServer), which catches the exception and returns an
/// empty ack; here the faulted Ask is the transient signal the site relies on
/// (see ).
///
private IActorRef? _auditIngestProxy;
///
/// Default Ask timeout for routing audit ingest commands to the
/// Effective Ask timeout for audit ingest routing. Defaults to
/// (30 s) — the two
/// audit-ingest entry points (the site stream server and the control plane) share one source of truth
/// for the timeout. Overridable via the constructor so tests can exercise the
/// timeout/fault path without waiting 30 s. When the window is exceeded the Ask
/// faults and that fault is piped back to the caller as a
/// (see ).
///
private readonly TimeSpan _auditIngestAskTimeout;
///
/// DistributedPubSub topic used to fan health reports out to the peer
/// central node so both per-node aggregators stay in sync. See
/// for the protocol rationale.
///
private const string HealthReportTopic = "site-health-replica";
///
/// Constructs the actor over the given central→site command transport (production injects the
/// gRPC , built by the Host from the DI
/// ; tests substitute a fake).
///
/// DI service provider for scoped repository and aggregator access.
/// The central→site command transport to route every through.
/// Optional override for the audit-ingest Ask timeout (test hook).
public CentralCommunicationActor(
IServiceProvider serviceProvider,
ISiteCommandTransport transport,
TimeSpan? auditIngestAskTimeout = null)
: this(serviceProvider, auditIngestAskTimeout)
{
_transport = transport ?? throw new ArgumentNullException(nameof(transport));
}
/// Shared wiring: sets scoped state and registers every message handler. The
/// is assigned by the delegating public constructor before any message
/// is dispatched.
/// DI service provider.
/// Optional audit-ingest Ask timeout override.
private CentralCommunicationActor(
IServiceProvider serviceProvider,
TimeSpan? auditIngestAskTimeout)
{
_serviceProvider = serviceProvider;
_auditIngestAskTimeout = auditIngestAskTimeout ?? Grpc.SiteStreamGrpcServer.AuditIngestAskTimeout;
// Site address cache loaded from database
Receive(HandleSiteAddressCacheLoaded);
// Periodic refresh trigger
Receive(_ => LoadSiteAddressesFromDb());
// A faulted LoadSiteAddressesFromDb task is piped here as a
// Status.Failure. Without this handler the failure was an unhandled message
// (debug-level only) and the refresh failed silently — operators could not
// distinguish "no sites configured" from "database is down". Log at Warning.
Receive(failure =>
_log.Warning(failure.Cause,
"Failed to load site addresses from the database; the per-site gRPC channel "
+ "cache was not refreshed and may be stale or empty"));
// Health monitoring: heartbeats and health reports from sites
Receive(HandleHeartbeat);
Receive(HandleSiteHealthReport);
Receive(r => ProcessLocally(r.Report));
Receive(r => MarkHeartbeatLocally(r.Heartbeat));
Receive(_ => { /* DistributedPubSub subscribe confirmation */ });
// Route enveloped messages to sites
Receive(HandleSiteEnvelope);
// Notification Outbox: the Host registers the outbox singleton proxy after this
// actor is created (the proxy cannot exist before this actor's construction).
Receive(msg =>
{
_notificationOutboxProxy = msg.OutboxProxy;
_log.Info("Registered notification outbox proxy");
});
// Notification Outbox ingest: a site forwards a buffered NotificationSubmit to the
// central cluster. Forward to the outbox proxy so the original
// Sender (the site's transport path) is preserved and the NotificationSubmitAck
// routes straight back to the site.
Receive(HandleNotificationSubmit);
// Notification Outbox status query: forward to the outbox proxy, preserving Sender
// so the NotificationStatusResponse routes back to the querying site.
Receive(HandleNotificationStatusQuery);
// Audit Log: the Host registers the AuditLogIngestActor singleton
// proxy after this actor is created (the proxy cannot exist before this
// actor's construction).
Receive(msg =>
{
_auditIngestProxy = msg.AuditIngestActor;
_log.Info("Registered audit ingest proxy");
});
// Audit Log site→central ingest: a site forwards a batch of audit
// events to the central cluster. Ask the ingest proxy
// and pipe the IngestAuditEventsReply back to the original Sender (the
// site's transport path) so the site can flip its rows to Forwarded.
Receive(HandleIngestAuditEvents);
// Audit Log combined-telemetry ingest: routes to the same proxy
// the same way; the proxy replies with an IngestCachedTelemetryReply.
Receive(HandleIngestCachedTelemetry);
// Startup reconciliation: a site node forwards its local deployed inventory on
// startup. Resolve the scoped ReconcileService, diff the
// inventory against central's expected set, and pipe the ReconcileSiteResponse
// (gap fetch tokens + orphans) straight back to the site node.
Receive(HandleReconcileSiteRequest);
}
private void HandleNotificationSubmit(NotificationSubmit msg)
{
if (_notificationOutboxProxy == null)
{
// No outbox proxy registered yet. A non-accepted ack makes the site's
// Store-and-Forward forwarder treat this as transient and retry later.
_log.Warning(
"Cannot route NotificationSubmit {0} — notification outbox not available",
msg.NotificationId);
Sender.Tell(new NotificationSubmitAck(
msg.NotificationId, Accepted: false, Error: "notification outbox not available"));
return;
}
_log.Debug("Routing NotificationSubmit {0} to the notification outbox", msg.NotificationId);
_notificationOutboxProxy.Forward(msg);
}
private void HandleNotificationStatusQuery(NotificationStatusQuery msg)
{
if (_notificationOutboxProxy == null)
{
// No outbox proxy registered yet. Reply Found: false so the querying site
// falls back to its local Store-and-Forward buffer to resolve the status.
_log.Warning(
"Cannot route NotificationStatusQuery {0} — notification outbox not available",
msg.NotificationId);
Sender.Tell(new NotificationStatusResponse(
msg.CorrelationId, Found: false, Status: "Unknown",
RetryCount: 0, LastError: null, DeliveredAt: null));
return;
}
_log.Debug("Routing NotificationStatusQuery {0} to the notification outbox", msg.NotificationId);
_notificationOutboxProxy.Forward(msg);
}
private void HandleIngestAuditEvents(IngestAuditEventsCommand msg)
{
if (_auditIngestProxy == null)
{
// No ingest proxy registered yet (host startup race). Reply with an
// empty IngestAuditEventsReply so the site keeps its rows Pending and
// retries — the same behaviour as the gRPC handler's wiring-race path.
_log.Warning(
"Cannot route IngestAuditEventsCommand ({0} events) — audit ingest not available",
msg.Events.Count);
Sender.Tell(new IngestAuditEventsReply(Array.Empty()));
return;
}
// Capture Sender before the async/PipeTo — Akka resets Sender between
// dispatches. The reply is piped straight back to the calling site node.
// On an Ask timeout or a faulted reply, PipeTo delivers a Status.Failure to
// replyTo: the fault propagates to the caller rather than being swallowed.
// The site's own Ask through this path then faults, and the site drain loop
// treats that as a transient failure — rows stay Pending and are retried on
// the next tick. (The gRPC handler instead returns an empty ack on fault;
// propagating the fault here is the cleaner transient signal.)
var replyTo = Sender;
_log.Debug("Routing IngestAuditEventsCommand ({0} events) to the audit ingest actor", msg.Events.Count);
_auditIngestProxy.Ask(msg, _auditIngestAskTimeout)
.PipeTo(replyTo);
}
private void HandleIngestCachedTelemetry(IngestCachedTelemetryCommand msg)
{
if (_auditIngestProxy == null)
{
_log.Warning(
"Cannot route IngestCachedTelemetryCommand ({0} entries) — audit ingest not available",
msg.Entries.Count);
Sender.Tell(new IngestCachedTelemetryReply(Array.Empty()));
return;
}
var replyTo = Sender;
_log.Debug("Routing IngestCachedTelemetryCommand ({0} entries) to the audit ingest actor", msg.Entries.Count);
_auditIngestProxy.Ask(msg, _auditIngestAskTimeout)
.PipeTo(replyTo);
}
///
/// Startup reconciliation (site→central): resolve the scoped
/// in a DI scope, diff the node's reported inventory
/// against central's expected set, and pipe the
/// back to the site node's transport path. The actor stays thin — all the diff
/// and staging logic lives in the service. Mirrors the DB-access pattern used by
/// (Task.Run + CreateScope + PipeTo) and the
/// Sender-preservation pattern of .
///
/// On a faulted task PipeTo delivers a to the node; its
/// Ask faults and it simply retries reconcile on the next startup — reconcile is
/// best-effort, so the fault is allowed to propagate rather than being swallowed.
///
private void HandleReconcileSiteRequest(ReconcileSiteRequest msg)
{
// Capture Sender before the async/PipeTo — Akka resets Sender between dispatches.
var replyTo = Sender;
// Bound the DB work by the actor lifecycle. The CTS may have
// been disposed by PostStop on a racing late message; treat that as "actor gone".
CancellationToken ct;
try
{
ct = _lifecycleCts.Token;
}
catch (ObjectDisposedException)
{
return;
}
_log.Debug(
"Handling ReconcileSiteRequest from site {0} node {1} ({2} local instance(s))",
msg.SiteIdentifier, msg.NodeId, msg.LocalNameToRevisionHash.Count);
Task.Run(async () =>
{
using var scope = _serviceProvider.CreateScope();
var service = scope.ServiceProvider.GetRequiredService();
return await service.ReconcileAsync(msg, ct).ConfigureAwait(false);
}).PipeTo(replyTo);
}
private void HandleHeartbeat(HeartbeatMessage heartbeat)
{
MarkHeartbeatLocally(heartbeat);
// Fan the heartbeat out to the peer central node so BOTH aggregators mark
// it, regardless of which central node the site delivered
// to. Without this, a heartbeat that only ever reaches one node leaves the
// other node's aggregator blind to that site's liveness after a failover
// (arch review 02, Low). MarkHeartbeat is idempotent (timestamp overwrite),
// so self-delivery via pub-sub is harmless.
try
{
DistributedPubSub.Get(Context.System).Mediator.Tell(
new Publish(HealthReportTopic, new SiteHeartbeatReplica(heartbeat)));
}
catch
{
// No-op in non-clustered hosts (TestKit).
}
}
///
/// Marks a site heartbeat on the local aggregator without re-broadcasting.
/// Used for both site-originated heartbeats and peer-replicated ones.
///
private void MarkHeartbeatLocally(HeartbeatMessage heartbeat)
{
var aggregator = _serviceProvider.GetService();
aggregator?.MarkHeartbeat(heartbeat.SiteId, heartbeat.Timestamp);
}
///
/// Handles a report delivered directly from a site:
/// process locally, then fan out to the peer central node so its
/// aggregator stays in sync.
///
private void HandleSiteHealthReport(SiteHealthReport report)
{
ProcessLocally(report);
try
{
DistributedPubSub.Get(Context.System).Mediator.Tell(
new Publish(HealthReportTopic, new SiteHealthReportReplica(report)));
}
catch
{
// No-op in non-clustered hosts (TestKit).
}
// Ack the site so its AkkaHealthReportTransport Ask completes and the
// report-loss counter-restore path can observe delivery (review 01
// [Medium]). Guarded so the peer-replica path (SiteHealthReportReplica,
// which arrives without an Ask sender) never dead-letters an ack.
if (!Sender.IsNobody())
{
Sender.Tell(new SiteHealthReportAck(report.SiteId, report.SequenceNumber, Accepted: true));
}
}
///
/// Applies a report to the local aggregator without re-broadcasting.
/// Used for both site-originated reports and peer-replicated ones — the
/// aggregator is idempotent via sequence-number comparison.
///
private void ProcessLocally(SiteHealthReport report)
{
var aggregator = _serviceProvider.GetService();
if (aggregator != null)
{
aggregator.ProcessReport(report);
}
else
{
_log.Warning("ICentralHealthAggregator not available, dropping health report from site {0}", report.SiteId);
}
}
// HandleConnectionStateChanged removed — no production
// caller emitted ConnectionStateChanged, so the workflow ran only in tests.
// Disconnect detection is owned by the transport layers (gRPC keepalive +
// Ask timeout).
private void HandleSiteEnvelope(SiteEnvelope envelope)
{
// Below-the-seam routing: the gRPC transport owns the "no route for this site ⇒ warn +
// drop, caller's Ask times out" contract. Sender is the temporary Ask actor (or the
// debug-bridge actor) and is preserved for reply routing.
_transport.Send(envelope, Sender);
}
private void LoadSiteAddressesFromDb()
{
var self = Self;
// Pass the actor's lifecycle CT into the repository
// call so a hung database query is cancelled when the actor stops
// rather than leaving the piped task to accumulate. Captured locally
// because the lifecycle CTS may have been disposed by PostStop on a
// racing late tick; treat that as "actor gone, give up".
CancellationToken ct;
try
{
ct = _lifecycleCts.Token;
}
catch (ObjectDisposedException)
{
return;
}
Task.Run(async () =>
{
using var scope = _serviceProvider.CreateScope();
var repo = scope.ServiceProvider.GetRequiredService();
var sites = await repo.GetAllSitesAsync(ct).ConfigureAwait(false);
var contacts = new Dictionary>();
// 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();
foreach (var site in sites)
{
var addrs = new List();
if (!string.IsNullOrWhiteSpace(site.NodeAAddress))
{
var addr = site.NodeAAddress;
// Strip actor path suffix if present (legacy format)
var idx = addr.IndexOf("/user/");
if (idx > 0) addr = addr.Substring(0, idx);
addrs.Add(addr);
}
if (!string.IsNullOrWhiteSpace(site.NodeBAddress))
{
var addr = site.NodeBAddress;
var idx = addr.IndexOf("/user/");
if (idx > 0) addr = addr.Substring(0, idx);
addrs.Add(addr);
}
if (addrs.Count > 0)
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
// Self. The message record exposes read-only types (
// IReadOnlyDictionary / IReadOnlyList) so the Akka.NET message-
// immutability convention is enforced by type, not just convention.
var frozen = contacts.ToDictionary(
kvp => kvp.Key,
kvp => (IReadOnlyList)kvp.Value.AsReadOnly());
// Carry the full configured-site identifier set (not just the
// address-bearing subset in `frozen`) so the aggregator prunes only
// genuinely-deleted sites and never an addressless-but-configured one.
var knownSiteIds = sites.Select(s => s.SiteIdentifier).ToList();
return new SiteAddressCacheLoaded(frozen, knownSiteIds, grpcContacts);
}).PipeTo(self);
}
private void HandleSiteAddressCacheLoaded(SiteAddressCacheLoaded msg)
{
// Per-site transport resource reconciliation (build/drop gRPC channel
// pairs). Runs on the actor thread each refresh tick.
_transport.ReconcileSites(msg);
// Self-healing eviction: a site deleted from configuration would otherwise
// linger in the aggregator as a permanently-offline tile (and live KPI
// sample source) forever. Prune on every refresh so it disappears within
// one refresh interval without needing a dedicated deletion event. Transport-agnostic.
_serviceProvider.GetService()?.PruneUnknownSites(msg.KnownSiteIds);
}
// TrackMessageForCleanup removed — the dicts it fed
// existed solely to support the dead ConnectionStateChanged workflow.
///
protected override SupervisorStrategy SupervisorStrategy()
{
return new OneForOneStrategy(
maxNrOfRetries: -1,
withinTimeRange: Timeout.InfiniteTimeSpan,
decider: Decider.From(ex =>
{
_log.Warning(ex, "Child actor of CentralCommunicationActor faulted, resuming (state preserved)");
return Directive.Resume;
}));
}
///
protected override void PreStart()
{
_log.Info("CentralCommunicationActor started");
// Subscribe to the peer-replication topic so we receive health reports
// delivered to the other central node and keep our local aggregator
// in sync (a site may deliver its report to either central node).
// Tolerant of non-clustered hosts (TestKit) where the extension is absent.
try
{
DistributedPubSub.Get(Context.System).Mediator.Tell(
new Subscribe(HealthReportTopic, Self));
}
catch (Exception ex)
{
_log.Debug("DistributedPubSub not available — peer health replication disabled: {0}", ex.Message);
}
// Schedule periodic refresh of site addresses from the database
_refreshSchedule = Context.System.Scheduler.ScheduleTellRepeatedlyCancelable(
TimeSpan.Zero,
TimeSpan.FromSeconds(60),
Self,
new RefreshSiteAddresses(),
ActorRefs.NoSender);
}
///
protected override void PostStop()
{
_log.Info("CentralCommunicationActor stopped");
_refreshSchedule?.Cancel();
// Cancel any in-flight LoadSiteAddressesFromDb so a
// hung MS SQL query does not outlive the actor.
try
{
_lifecycleCts.Cancel();
}
catch (ObjectDisposedException)
{
// Double-stop is benign.
}
_lifecycleCts.Dispose();
}
}
///
/// Command to trigger a refresh of site addresses from the database.
///
public record RefreshSiteAddresses;
///
/// Message carrying the loaded site contact data from the database. Per-transport resource
/// reconciliation happens on the actor thread in HandleSiteAddressCacheLoaded, which hands this
/// straight to .
///
/// The payload is exposed as
/// of so the Akka.NET "messages are immutable"
/// convention is enforced at the type level rather than relying on producer
/// discipline. The producer wraps the constructed buckets with
/// List<T>.AsReadOnly() before piping to Self.
///
/// Akka remote addresses per site (from NodeA/NodeBAddress).
/// Every configured site id, address-bearing or not, for aggregator pruning.
///
/// gRPC endpoint pairs per site (from GrpcNodeA/GrpcNodeBAddress) — the streaming path's columns,
/// consumed by the gRPC transport and ignored by the Akka one.
///
public sealed record SiteAddressCacheLoaded(
IReadOnlyDictionary> SiteContacts,
IReadOnlyCollection KnownSiteIds,
IReadOnlyDictionary GrpcContacts);
///
/// A site's gRPC node-pair endpoints, as loaded from Site.GrpcNodeAAddress/
/// GrpcNodeBAddress. Either may be null when only one node has a gRPC address configured.
///
/// NodeA gRPC base address (e.g. http://scadabridge-site-a-node-a:8083), or null.
/// NodeB gRPC base address, or null.
public readonly record struct SiteGrpcEndpoints(string? NodeA, string? NodeB);
///
/// Peer-replication envelope for a site heartbeat, fanned out over the same
/// DistributedPubSub topic as so both
/// central aggregators mark heartbeats regardless of which node the site
/// delivered to. Only ever travels central↔central (same assembly
/// version, rolled together), so the default reflective serializer is safe.
///
public sealed record SiteHeartbeatReplica(HeartbeatMessage Heartbeat);
///
/// Notification sent to debug view subscribers when the stream is terminated
/// due to site disconnection.
///
public record DebugStreamTerminated(string SiteId, string CorrelationId);
///
/// Registers the central NotificationOutboxActor singleton proxy with the
/// so site-forwarded
/// and messages can be routed to it. Sent by the Host
/// after the outbox singleton proxy is created.
///
public record RegisterNotificationOutbox(IActorRef OutboxProxy);
///
/// Registers the central AuditLogIngestActor singleton proxy with the
/// so site-forwarded
/// and
/// messages can be routed to it. Sent by the Host after the audit-ingest
/// singleton proxy is created. Lives here (not in Commons) because
/// ZB.MOM.WW.ScadaBridge.Commons has no Akka package reference and cannot hold an
/// field.
///
public sealed record RegisterAuditIngest(IActorRef AuditIngestActor);