using Akka.Actor; using Akka.Event; using ZB.MOM.WW.OtOpcUa.Commons.Observability; using ZB.MOM.WW.OtOpcUa.Commons.OpcUa; using ZB.MOM.WW.OtOpcUa.Commons.Types; using ZB.MOM.WW.OtOpcUa.Core.Abstractions; namespace ZB.MOM.WW.OtOpcUa.Runtime.Drivers; // private timer key type — file-scoped so the name stays unique per-file file sealed class HealthPollTick { public static readonly HealthPollTick Instance = new(); private HealthPollTick() { } } /// /// Akka wrapper for a single instance. States: /// /// /// Connecting — calling . /// Connected — initialised; serving Read/Write/Subscribe requests. /// Reconnecting — disconnect observed; periodic retry of Initialize. /// Failed — terminal until parent restarts the actor. /// /// /// Engine wiring (subscriptions → AttributeValueUpdate publishes, ApplyDelta-driven Reinitialize, /// per-tag write Asks) is staged for follow-up F7. This skeleton compiles + has a working /// state machine so the Phase 6 control-plane integration tests can target it. /// public sealed class DriverInstanceActor : ReceiveActor, IWithTimers { public static readonly TimeSpan DefaultReconnectInterval = TimeSpan.FromSeconds(10); public sealed record InitializeRequested(string DriverConfigJson); public sealed record InitializeSucceeded(int Generation); public sealed record InitializeFailed(string Reason, int Generation); public sealed record DisconnectObserved(string Reason); /// #477 — sent to the parent () on every connectivity transition: /// Connected=true on entering the Connected state, false on entering Reconnecting. The host /// annotates this driver's native alarm conditions' source-data Quality from it (comms lost → Bad, /// restored → Good) — independently of alarm transitions, since a comms-lost driver emits no alarm /// events. Fire-and-forget, mirroring . public sealed record ConnectivityChanged(string DriverInstanceId, bool Connected); public sealed record ApplyDelta(string DriverConfigJson, CorrelationId Correlation); public sealed record ApplyResult(bool Success, string? Reason, CorrelationId Correlation); /// /// Sent to the parent () AFTER an has fully /// re-initialised the driver (i.e. completed). The host uses /// it to re-register this driver's dependency-consumer mux adapter: a driver may rebuild its declared /// during reinit (the Calculation driver re-binds its /// calc tags), so the re-register MUST happen after reinit — NOT synchronously alongside the ApplyDelta /// dispatch, which would re-read the STALE ref set. Emitted only on a successful apply. /// /// The instance whose adapter should be re-registered. public sealed record DeltaApplied(string DriverInstanceId); public sealed record WriteAttribute(string TagId, object Value); public sealed record WriteAttributeResult(bool Success, string? Reason); /// /// Sent by when an OPC UA client Acknowledges a NATIVE Part 9 /// condition (resolved from the condition NodeId to this driver via the host's alarm inverse map). /// The actor forwards it to the driver's , carrying the /// authored alarm full-reference (= the driver's ConditionId/AlarmFullReference) and the /// authenticated principal. Mirrors , but the ack is fire-and-forget: /// the driver's returns no per-condition status and the /// OPC UA Part 9 ack already committed the local condition state, so there is no reply to surface. /// /// The authored alarm full-reference the driver correlates the ack on /// (= the equipment tag's FullName/AlarmFullReference). /// Operator-supplied comment forwarded to the upstream alarm system; null when none. /// The authenticated principal performing the acknowledge. public sealed record RouteAlarmAck(string ConditionId, string? Comment, string OperatorUser); public sealed record Subscribe(IReadOnlyList FullReferences, TimeSpan PublishingInterval); /// /// Sets the set of references this driver should keep subscribed for the lifetime of the /// current deployment. Unlike the one-shot , the desired set is /// retained and (re)established automatically every time the actor (re)enters /// Connected — closing the F8b "re-subscribe across reconnects" gap and giving /// a single message to drive the SubscribeBulk pass after an /// apply. Sending an empty set clears the desired subscription. /// /// /// The native-alarm references (alarm-bearing equipment-tag FullNames = the driver's /// ConditionId/AlarmFullReference) this driver should keep an alarm subscription open for. /// An driver suppresses until at /// least one alarm subscription exists, so the actor calls /// with this set to un-gate the native feed. Empty /// (or null) means the driver has no alarm tags. Defaults to null so non-alarm callers are unchanged. /// public sealed record SetDesiredSubscriptions( IReadOnlyList FullReferences, TimeSpan PublishingInterval, IReadOnlyList? AlarmReferences = null); public sealed record SubscriptionEstablished(string DiagnosticId, int ReferenceCount); public sealed record SubscriptionFailed(string Reason); public sealed record Unsubscribe; /// /// Sent by when the AdminUI issues a Reconnect operation. /// Pushes the actor out of Connected into Reconnecting so the transport is /// re-established without fully stopping and respawning the actor. Safe to send in any /// state — a no-op when already Reconnecting or Connecting. /// public sealed record ForceReconnect; /// Published to the actor's parent whenever the subscribed IDriver fires /// . The parent forwards to OpcUaPublishActor. public sealed record AttributeValuePublished(string DriverInstanceId, string FullReference, object? Value, OpcUaQuality Quality, DateTime TimestampUtc); private sealed record DataChangeForward(string FullReference, DataValueSnapshot Snapshot); /// Published to the parent whenever the subscribed driver (an ) fires /// . The parent () projects + routes it /// to the materialised Part 9 condition. Parallels . public sealed record AttributeAlarmPublished(string DriverInstanceId, AlarmEventArgs Args); private sealed record NativeAlarmRaised(AlarmEventArgs Args); /// Self-sent on Connected entry / when alarm refs are (re)pushed, to establish the native-alarm /// subscription that un-gates an driver's feed. Handled async so the /// call is bounded + off the synchronous handlers. private sealed record SubscribeAlarms; /// Self-sent when the wrapped driver raises — /// it observed that the remote's tag set may have changed. Marshals the event off the driver's thread /// onto the actor thread. Handled in EVERY behaviour, including Stubbed and Reconnecting: the raise has no /// connection affinity (a Galaxy redeploy or a TwinCAT symbol-version bump can land while the driver is /// between connects), and dropping it in one state would lose the signal silently. private sealed record RediscoveryRaised(RediscoveryEventArgs Args); /// Self-sent when the wrapped driver raises /// — one of its hosts went Running ↔ Stopped ↔ Faulted. Marshals the event off the driver's probe thread /// onto the actor thread. Carries NO payload: the handler re-pulls /// , so the driver stays the single source of truth /// for the host set and a host that appeared since the last publish is picked up too. Handled in every /// behaviour for the same reason as — a probe tick can land while the /// driver is between connects. private sealed record HostStatusRaised; public sealed class RetryConnect { public static readonly RetryConnect Instance = new(); private RetryConnect() { } } /// Interval between periodic health-poll heartbeats sent to the snapshot store. public static readonly TimeSpan HealthPollInterval = TimeSpan.FromSeconds(30); private readonly IDriver _driver; /// Phase 6.1 resilience seam: every guarded driver-capability call this actor makes is /// routed through this invoker (retry / breaker / in-flight accounting / tracker telemetry). Defaults to /// (a genuine pass-through) when none is supplied — so /// unit tests + the F7 skeleton path behave exactly as a raw driver call. The fused Host builds a /// real per-instance invoker via and injects it at /// spawn. Consumed through the abstraction because Runtime /// is deliberately Polly-free (it does not reference ZB.MOM.WW.OtOpcUa.Core). private readonly IDriverCapabilityInvoker _invoker; private readonly string _driverInstanceId; private readonly string _clusterId; private readonly IDriverHealthPublisher _healthPublisher; private readonly TimeSpan _reconnectInterval; private readonly TimeSpan _healthPollInterval; private readonly ILoggingAdapter _log = Context.GetLogger(); private string? _currentConfigJson; /// Monotonic token tagging each attempt. An init result is /// honoured only when its generation matches the latest; an older result is from a superseded attempt /// (e.g. an adopted a new config mid-(re)connect) and is dropped. Touched only /// on the actor thread, so no lock is needed. private int _initGeneration; /// /// Timestamps of recent Faulted-state transitions; used to compute the 5-minute error count. /// No lock needed — every read/write site runs inside an Akka message handler, which is /// single-threaded per actor instance. /// private readonly Queue _faultTimestamps = new(); /// Active subscription handle (null when not subscribed). Tracks the current live /// subscription; the actor auto-(re)subscribes on (re)connect and on each /// message via / , so callers /// do not need to re-send subscription requests after a reconnect. private ISubscriptionHandle? _subscriptionHandle; private EventHandler? _dataChangeHandler; private EventHandler? _alarmEventHandler; private EventHandler? _rediscoveryHandler; private EventHandler? _hostStatusHandler; /// When the driver last raised , and the reason /// it gave. Null until the first raise. Carried on every subsequent health snapshot so the AdminUI can /// prompt an operator to re-browse the device. /// Deliberately sticky — it is not cleared on reconnect or on a later clean pass. The remote's /// tag set stayed changed; only an operator re-browsing and committing resolves it, and this actor cannot /// observe that happening. A redeploy respawns the child, which clears it naturally. private DateTime? _rediscoveryNeededUtc; private string? _rediscoveryReason; /// The references the host wants kept subscribed (set by ). /// Re-applied on every entry into Connected so values resume after a reconnect or redeploy. private IReadOnlyList _desiredRefs = Array.Empty(); private TimeSpan _desiredInterval = TimeSpan.FromSeconds(1); /// The native-alarm references the host wants kept subscribed (set by /// ). Re-applied on every Connected entry so the alarm /// feed is re-un-gated after a reconnect/redeploy. private IReadOnlyList _desiredAlarmRefs = Array.Empty(); /// The active native-alarm subscription handle for an driver, or /// null when none is established. Reset on so the next Connected entry /// re-subscribes against the freshly re-initialised driver; the null check makes the subscribe /// idempotent across repeated pushes. private IAlarmSubscriptionHandle? _alarmSubscriptionHandle; /// /// Gets or sets the timer scheduler for scheduling reconnection attempts. /// public ITimerScheduler Timers { get; set; } = null!; /// /// Creates a Props object for instantiating a . /// /// The driver instance to wrap. /// Optional interval for reconnection attempts; defaults to 10 seconds. /// If true, the actor starts in stub mode for testing or unavailable platforms. /// Optional health publisher; defaults to so tests and /// stub paths don't need to provide one. /// Optional cluster identifier forwarded in messages; /// defaults to an empty string when not provided (e.g. in unit tests). /// Optional Phase 6.1 resilience invoker wrapping this driver's capability /// calls; defaults to (pass-through) when not supplied. /// Optional period of the health-poll heartbeat, which also drives the /// subscription reconcile (); defaults to /// . Exists so a test can drive the reconcile without waiting 30 s. /// A instance configured to create the wrapped . public static Props Props( IDriver driver, TimeSpan? reconnectInterval = null, bool startStubbed = false, IDriverHealthPublisher? healthPublisher = null, string? clusterId = null, IDriverCapabilityInvoker? invoker = null, TimeSpan? healthPollInterval = null) => Akka.Actor.Props.Create(() => new DriverInstanceActor( driver, reconnectInterval ?? DefaultReconnectInterval, startStubbed, healthPublisher ?? NullDriverHealthPublisher.Instance, clusterId ?? string.Empty, invoker, healthPollInterval)); /// /// Returns true when the driver should boot in DEV-STUB mode based on host platform and /// configured roles. Only the legacy v1 in-process "Galaxy" type stays Windows-only: /// the legacy MXAccess COM proxy (retired in PR 7.2; gated for any leftover DriverInstance /// rows that still reference the old type name). /// The v2 "GalaxyMxGateway" driver talks gRPC to an external mxaccessgw process, /// so it runs on any platform .NET 10 supports — Linux containers included. Not stubbed. /// /// The type identifier of the driver. /// Operational roles configured for this instance. /// True if the driver should start in DEV-STUB mode; otherwise false. public static bool ShouldStub(string driverType, IEnumerable roles) { var isWindowsOnly = driverType is "Galaxy"; if (!OperatingSystem.IsWindows() && isWindowsOnly) return true; if (roles.Contains("dev") && isWindowsOnly) return true; return false; } /// /// Initializes a new instance of the class. /// /// The driver instance to wrap and manage. /// Interval between reconnection attempts. /// If true, start in stub mode for testing or unavailable platforms. /// Sink for health-change notifications; must not be null. /// Cluster identifier forwarded in health snapshots. /// Phase 6.1 resilience invoker wrapping this driver's capability calls; /// defaults to (pass-through) when null. /// Period of the health-poll heartbeat, which also drives the /// subscription reconcile; defaults to when null. public DriverInstanceActor( IDriver driver, TimeSpan reconnectInterval, bool startStubbed = false, IDriverHealthPublisher? healthPublisher = null, string? clusterId = null, IDriverCapabilityInvoker? invoker = null, TimeSpan? healthPollInterval = null) { _driver = driver; _invoker = invoker ?? NullDriverCapabilityInvoker.Instance; _driverInstanceId = driver.DriverInstanceId; _clusterId = clusterId ?? string.Empty; _healthPublisher = healthPublisher ?? NullDriverHealthPublisher.Instance; _reconnectInterval = reconnectInterval; _healthPollInterval = healthPollInterval ?? HealthPollInterval; OtOpcUaTelemetry.DriverInstanceLifecycle.Add(1, new KeyValuePair("event", startStubbed ? "spawn_stub" : "spawn"), new KeyValuePair("driver_type", driver.DriverType)); if (startStubbed) { Context.GetLogger().Info("[DEV-STUB] driver={Name} type={Type}", _driverInstanceId, driver.DriverType); Become(Stubbed); } else { Become(Connecting); } } /// protected override void PreStart() { // Warm up the snapshot store immediately so AdminUI sees current state as soon as the // actor starts, before any state transition fires. Also start the periodic heartbeat so // long-lived Healthy drivers keep their snapshot fresh for newly-joined SignalR clients. // Attach the rediscovery signal before the first publish. Not per-connect: an IRediscoverable raise // has no connection affinity, and a driver can observe a remote change while disconnected. AttachRediscoverySource(); AttachHostStatusSource(); PublishHealthSnapshot(); Timers.StartPeriodicTimer("health-poll", HealthPollTick.Instance, _healthPollInterval); } private void Stubbed() { // Stubbed drivers accept the standard message contracts but return deterministic // success without touching real hardware. Read returns null; Write succeeds. Receive(_ => { /* no-op */ }); Receive(msg => Sender.Tell(new ApplyResult(true, "stubbed", msg.Correlation))); Receive(_ => Sender.Tell(new WriteAttributeResult(true, "stubbed"))); // Stubbed drivers have no upstream alarm system — swallow the ack (it's fire-and-forget, no reply). Receive(_ => { /* stubbed drivers have no alarm backend */ }); Receive(_ => { /* stubbed drivers don't disconnect */ }); Receive(_ => { /* stubbed drivers don't reconnect */ }); Receive(StoreDesiredSubscriptions); // Stubbed drivers never enter Connected, so they never kick discovery; swallow defensively in case a // re-discovery self-tick is ever routed here so it doesn't surface as an Akka Unhandled message. Receive(HandleRediscoveryRaised); Receive(_ => PublishHealthSnapshot()); Receive(_ => PublishHealthSnapshot()); } private void Connecting() { Receive(msg => InitializeAsync(msg.DriverConfigJson)); // Fast-fail writes while still connecting — without this the inbound WriteAttribute dead-letters // and DriverHostActor.HandleRouteNodeWrite waits its full 8s Ask before reporting a generic // "write timeout". Synchronous Receive: Sender.Tell on the actor thread is safe. Receive(_ => Sender.Tell(new WriteAttributeResult(false, "driver not connected"))); // An ack arriving while still connecting can't reach the upstream alarm system; drop it (the ack is // fire-and-forget — no reply to surface — and the OPC UA condition state already committed locally). Receive(_ => _log.Debug("DriverInstance {Id}: alarm ack arrived during connect — dropped (driver not connected)", _driverInstanceId)); Receive(AdoptConfigDuringInit); Receive(msg => { if (msg.Generation != _initGeneration) return; _log.Info("DriverInstance {Id}: connected", _driverInstanceId); Become(Connected); PublishHealthSnapshot(); ResubscribeDesired(); AttachAlarmSource(); SubscribeDesiredAlarms(); }); Receive(msg => { if (msg.Generation != _initGeneration) return; _log.Warning("DriverInstance {Id}: initialize failed: {Reason}", _driverInstanceId, msg.Reason); RecordFault(); Become(Reconnecting); PublishHealthSnapshot(); }); Receive(StoreDesiredSubscriptions); Receive(_ => { /* already connecting — no-op */ }); // ResubscribeDesired self-Tells Subscribe; HandleSubscribeAsync replies SubscriptionEstablished to the // sender, which on the self-resubscribe path is Self. Swallow it (trace only) so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: subscription established ({Count} refs, {Diag})", _driverInstanceId, msg.ReferenceCount, msg.DiagnosticId)); // Symmetric to the SubscriptionEstablished swallow: a failed self-resubscribe replies SubscriptionFailed // to Self (HandleSubscribeAsync already logged the underlying cause). Swallow it so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: resubscribe reported failure: {Reason}", _driverInstanceId, msg.Reason)); // A native alarm transition can race in while still (re)connecting (the driver's feed runs on its // own thread); drop it — the feed re-delivers active alarms once Connected. Trace only. Receive(_ => _log.Debug("DriverInstance {Id}: native alarm arrived during connect — dropped (feed re-delivers)", _driverInstanceId)); // A SubscribeAlarms self-tell (from Connected) can be overtaken by an already-queued disconnect into // this state; swallow it so it doesn't dead-letter — the next Connected entry re-subscribes. Receive(_ => { }); Receive(HandleRediscoveryRaised); Receive(_ => PublishHealthSnapshot()); Receive(_ => PublishHealthSnapshot()); } private void Connected() { // #477 — announce connectivity to the host so it can clear any comms-lost Quality annotation on this // driver's native alarm conditions (Bad → Good). Fire-and-forget; the host defaults conditions to a // non-Good "waiting for initial data" quality at materialise, and this is what confirms them Good. Context.Parent.Tell(new ConnectivityChanged(_driverInstanceId, Connected: true)); ReceiveAsync(HandleApplyDeltaAsync); Receive(msg => { _log.Warning("DriverInstance {Id}: disconnect observed ({Reason}); reconnecting", _driverInstanceId, msg.Reason); DetachSubscription(); RecordFault(); Become(Reconnecting); PublishHealthSnapshot(); }); Receive(_ => { _log.Info("DriverInstance {Id}: ForceReconnect requested by admin; re-entering Reconnecting", _driverInstanceId); DetachSubscription(); Become(Reconnecting); PublishHealthSnapshot(); }); ReceiveAsync(HandleWriteAsync); ReceiveAsync(HandleAcknowledgeAsync); ReceiveAsync(HandleSubscribeAsync); ReceiveAsync(_ => UnsubscribeAsync()); Receive(msg => { StoreDesiredSubscriptions(msg); if (_desiredRefs.Count > 0) Self.Tell(new Subscribe(_desiredRefs, _desiredInterval)); else if (_subscriptionHandle is not null) Self.Tell(new Unsubscribe()); // Native-alarm analogue: un-gate the IAlarmSource feed when alarm tags are (now) present. The // common live path — a deploy delivers SetDesiredSubscriptions while the driver is already // Connected — flows through HERE, so the alarm subscribe must happen on this message, not only // on Connected entry. Unlike the value path above, there is deliberately no unsubscribe-on-empty: // a removed alarm tag's condition node is torn down on the address-space rebuild and // DriverHostActor.ForwardNativeAlarm drops any transition whose ConditionId no longer maps, so a // lingering session-less alarm subscription can never surface a removed alarm. SubscribeDesiredAlarms(); }); ReceiveAsync(HandleSubscribeAlarmsAsync); Receive(OnDataChangeForward); // Native alarm transition marshaled onto the actor thread from the driver's OnAlarmEvent; // project it to the parent the same way DataChangeForward projects AttributeValuePublished. Receive(m => Context.Parent.Tell(new AttributeAlarmPublished(_driverInstanceId, m.Args))); // ResubscribeDesired self-Tells Subscribe; HandleSubscribeAsync replies SubscriptionEstablished to the // sender, which on the self-resubscribe path is Self. Swallow it (trace only) so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: subscription established ({Count} refs, {Diag})", _driverInstanceId, msg.ReferenceCount, msg.DiagnosticId)); // Symmetric to the SubscriptionEstablished swallow: a failed self-resubscribe replies SubscriptionFailed // to Self (HandleSubscribeAsync already logged the underlying cause). Swallow it so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: resubscribe reported failure: {Reason}", _driverInstanceId, msg.Reason)); Receive(HandleRediscoveryRaised); Receive(_ => PublishHealthSnapshot()); Receive(_ => { PublishHealthSnapshot(); ReconcileSubscription(); }); } /// /// Periodic desired-vs-actual subscription reconcile (#488), riding the existing health-poll /// tick. While Connected, a driver that has a desired subscription set but no live handle is /// re-subscribed. /// /// /// /// The desired set is otherwise only (re)applied on a Connected transition /// () or on a deploy's . /// Nothing re-asserted it for a subscription that was lost while the driver stayed /// Connected — a failed subscribe whose retry never came, or a handle dropped down an error /// path. That is the hole this closes. /// /// /// This is a state-machine consistency check, not a probe. It can only see that this /// actor holds no handle; it cannot detect a handle that is present but dead server-side, /// because ISubscriptionHandle is opaque (a DiagnosticId string and nothing /// else). Detecting that needs a liveness member on ISubscribable implemented across /// all eight drivers, each with different subscription mechanics — deliberately deferred /// until a real occurrence shows the handle outliving the subscription. /// /// /// Gated on : a driver that cannot subscribe at all may still have /// been handed a desired set, and re-telling Subscribe to it every tick would fail /// every tick, forever. /// /// /// A redundant Subscribe is possible but harmless: a tick sitting in the mailbox ahead /// of an already-queued Subscribe observes the same "no handle" state and tells a /// second one. is a ReceiveAsync (so no tick can be /// processed mid-subscribe) and drops any prior subscription before establishing the new one, /// so the duplicate costs one extra round-trip and converges. Suppressing it would need a /// pending-subscribe flag threaded through every self-tell site — more state than the race /// is worth. /// /// private void ReconcileSubscription() { if (_driver is not ISubscribable) return; if (_desiredRefs.Count == 0 || _subscriptionHandle is not null) return; _log.Info( "DriverInstance {Id}: connected with {Count} desired subscription ref(s) but no live subscription handle — re-subscribing (periodic reconcile)", _driverInstanceId, _desiredRefs.Count); Self.Tell(new Subscribe(_desiredRefs, _desiredInterval)); } private void Reconnecting() { // #477 — announce comms loss to the host so it annotates this driver's native alarm conditions Bad // (a comms-lost driver emits no alarm events, so this is the ONLY signal that the source is unreachable). Context.Parent.Tell(new ConnectivityChanged(_driverInstanceId, Connected: false)); Receive(_ => InitializeAsync(_currentConfigJson ?? "{}")); // Fast-fail writes while reconnecting (same reason as Connecting — avoids the 8s host Ask // timeout on an inbound write to a transiently-down driver). Synchronous Receive. Receive(_ => Sender.Tell(new WriteAttributeResult(false, "driver not connected"))); // An ack arriving while reconnecting can't reach the upstream alarm system; drop it (fire-and-forget, // no reply — the OPC UA condition state already committed locally on the Part 9 ack). Receive(_ => _log.Debug("DriverInstance {Id}: alarm ack arrived during reconnect — dropped (driver not connected)", _driverInstanceId)); Receive(AdoptConfigDuringInit); Receive(msg => { if (msg.Generation != _initGeneration) return; Timers.Cancel("retry-connect"); _log.Info("DriverInstance {Id}: reconnected", _driverInstanceId); Become(Connected); PublishHealthSnapshot(); ResubscribeDesired(); AttachAlarmSource(); SubscribeDesiredAlarms(); }); // A failure here is a no-op regardless of generation — the retry timer keeps trying the // current config; only a (generation-matched) InitializeSucceeded transitions state. Receive(_ => { /* keep retrying via timer */ }); Receive(StoreDesiredSubscriptions); Receive(_ => { /* already reconnecting — no-op */ }); // ResubscribeDesired self-Tells Subscribe; HandleSubscribeAsync replies SubscriptionEstablished to the // sender, which on the self-resubscribe path is Self. Swallow it (trace only) so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: subscription established ({Count} refs, {Diag})", _driverInstanceId, msg.ReferenceCount, msg.DiagnosticId)); // Symmetric to the SubscriptionEstablished swallow: a failed self-resubscribe replies SubscriptionFailed // to Self (HandleSubscribeAsync already logged the underlying cause). Swallow it so it doesn't dead-letter. Receive(msg => _log.Debug("DriverInstance {Id}: resubscribe reported failure: {Reason}", _driverInstanceId, msg.Reason)); // A native alarm transition can race in while still reconnecting (the driver's feed runs on its // own thread); drop it — the feed re-delivers active alarms once Connected. Trace only. Receive(_ => _log.Debug("DriverInstance {Id}: native alarm arrived during reconnect — dropped (feed re-delivers)", _driverInstanceId)); // A SubscribeAlarms self-tell (from Connected) can be overtaken by an already-queued disconnect into // this state; swallow it so it doesn't dead-letter — the next Connected entry re-subscribes. Receive(_ => { }); Receive(HandleRediscoveryRaised); Receive(_ => PublishHealthSnapshot()); Receive(_ => PublishHealthSnapshot()); Timers.StartPeriodicTimer("retry-connect", RetryConnect.Instance, _reconnectInterval); } private void InitializeAsync(string driverConfigJson) { _currentConfigJson = driverConfigJson; var generation = ++_initGeneration; var self = Self; _ = Task.Run(async () => { try { await _driver.InitializeAsync(driverConfigJson, CancellationToken.None); self.Tell(new InitializeSucceeded(generation)); } catch (Exception ex) { self.Tell(new InitializeFailed(ex.Message, generation)); } }); } /// Adopt a new config while not connected: ApplyDelta in Connecting/Reconnecting re-inits /// immediately with the new config. swaps _currentConfigJson and /// bumps the generation, so the in-flight (old-config) init is superseded and its result is dropped. /// The actor stays in its current state; the new init's result drives the next transition. In /// Reconnecting the retry timer is left running — if this immediate attempt fails it keeps retrying /// the new config (a redundant concurrent attempt is deduped by the generation guard). private void AdoptConfigDuringInit(ApplyDelta msg) { _log.Info("DriverInstance {Id}: ApplyDelta during (re)connect — adopting new config, re-initialising now", _driverInstanceId); InitializeAsync(msg.DriverConfigJson); Sender.Tell(new ApplyResult(true, "config adopted; reinitializing", msg.Correlation)); } private async Task HandleApplyDeltaAsync(ApplyDelta msg) { var replyTo = Sender; try { await _driver.ReinitializeAsync(msg.DriverConfigJson, CancellationToken.None); _currentConfigJson = msg.DriverConfigJson; // The driver may have rebuilt its dependency-consumer interest during reinit (the Calculation // driver re-binds its calc tags + dependency refs). Notify the parent so it re-registers this // driver's mux adapter AFTER reinit completed — a synchronous re-register at ApplyDelta-dispatch // time would re-read the driver's STALE DependencyRefs. Parent Tell (fire-and-forget); the // continuation resumes on the actor context (no ConfigureAwait) so Context.Parent is live. Context.Parent.Tell(new DeltaApplied(_driverInstanceId)); replyTo.Tell(new ApplyResult(true, null, msg.Correlation)); } catch (Exception ex) { replyTo.Tell(new ApplyResult(false, ex.Message, msg.Correlation)); } } private async Task HandleWriteAsync(WriteAttribute msg) { var replyTo = Sender; if (_driver is not IWritable writable) { replyTo.Tell(new WriteAttributeResult(false, "Driver does not implement IWritable")); return; } var request = new[] { new WriteRequest(msg.TagId, msg.Value) }; // Bound the write so a hung backend can't pin this actor forever — retry stays // off by default, but a stalled call still needs an answer. using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); try { // Route through the resilience invoker (Write pipeline: timeout, per-host breaker). Writes // default to NON-idempotent (03/S12): a timed-out write may already have committed at the // device, and replaying a command-shaped write (pulse coil, counter increment, Galaxy // supervisory write) can duplicate an irreversible field action. The invoker's no-retry arm is // therefore authoritative regardless of any configured Write retry. Per-tag opt-in to retries via // WriteIdempotentAttribute is the unshipped forward path (see WriteIdempotentAttribute). The outer // cts is retained as a hard backstop above the (typically tighter) tier timeout. var results = await _invoker.ExecuteWriteAsync( ResolveHost(msg.TagId), isIdempotent: false, async ct => await writable.WriteAsync(request, ct), cts.Token); if (results is { Count: 1 } && IsGoodStatus(results[0].StatusCode)) { replyTo.Tell(new WriteAttributeResult(true, null)); return; } var status = results is { Count: > 0 } ? results[0].StatusCode : 0xFFFFFFFF; replyTo.Tell(new WriteAttributeResult(false, $"StatusCode=0x{status:X8}")); } catch (OperationCanceledException) { replyTo.Tell(new WriteAttributeResult(false, "write timeout")); } catch (Exception ex) { replyTo.Tell(new WriteAttributeResult(false, ex.Message)); } } /// /// Forwards an inbound native-condition acknowledge (routed by from a /// resolved condition NodeId) to the driver's . The driver /// correlates on (= the authored alarm /// full-reference); carries the same reference (the /// driver's ack path keys on ConditionId). Bounded to 5s so a hung backend can't pin this actor. /// Fire-and-forget — the OPC UA Part 9 ack already committed the local condition state and /// returns no per-condition status — so there is no reply; /// a failure is logged and dropped (the local condition stays Acknowledged regardless). /// private async Task HandleAcknowledgeAsync(RouteAlarmAck msg) { if (_driver is not IAlarmSource src) { _log.Warning("DriverInstance {Id}: alarm ack dropped — driver does not implement IAlarmSource", _driverInstanceId); return; } var request = new[] { new AlarmAcknowledgeRequest( SourceNodeId: msg.ConditionId, ConditionId: msg.ConditionId, Comment: msg.Comment, OperatorUser: msg.OperatorUser), }; using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); try { // Ack is a write-shaped operation — the AlarmAcknowledge tier policy never retries. await _invoker.ExecuteAsync( DriverCapability.AlarmAcknowledge, ResolveHost(msg.ConditionId), async ct => await src.AcknowledgeAsync(request, ct), cts.Token); _log.Info("DriverInstance {Id}: acknowledged native condition {Cond} by {User}", _driverInstanceId, msg.ConditionId, msg.OperatorUser); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: native-alarm acknowledge of {Cond} failed", _driverInstanceId, msg.ConditionId); } } private async Task HandleSubscribeAsync(Subscribe msg) { // Capture Sender/Self BEFORE any await. The re-subscribe path below awaits // UnsubscribeAsync, and a real async backend can resume the continuation off Akka's // ActorContext — reading raw Sender/Self/Context past that point throws // NotSupportedException ("no active ActorContext"). Keep ConfigureAwait off the awaits // in this handler so continuations resume on the actor context. var replyTo = Sender; var self = Self; if (_driver is not ISubscribable subscribable) { replyTo.Tell(new SubscriptionFailed("Driver does not implement ISubscribable")); return; } if (_subscriptionHandle is not null) { // Subscribe-twice — drop the prior subscription before establishing the new one. await UnsubscribeAsync(); } try { _dataChangeHandler = (_, args) => self.Tell(new DataChangeForward(args.FullReference, args.Snapshot)); subscribable.OnDataChange += _dataChangeHandler; using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); // A bulk subscribe spans potentially many hosts, so key the pipeline on the driver // instance (single-host fallback) rather than any one reference's host. _subscriptionHandle = await _invoker.ExecuteAsync( DriverCapability.Subscribe, _driverInstanceId, async ct => await subscribable.SubscribeAsync(msg.FullReferences, msg.PublishingInterval, ct), cts.Token); replyTo.Tell(new SubscriptionEstablished(_subscriptionHandle.DiagnosticId, msg.FullReferences.Count)); _log.Info("DriverInstance {Id}: subscribed to {Count} refs ({Diag})", _driverInstanceId, msg.FullReferences.Count, _subscriptionHandle.DiagnosticId); } catch (Exception ex) { DetachSubscription(); _log.Warning(ex, "DriverInstance {Id}: subscribe failed", _driverInstanceId); replyTo.Tell(new SubscriptionFailed(ex.Message)); } } private async Task UnsubscribeAsync() { if (_driver is not ISubscribable subscribable || _subscriptionHandle is null) { DetachSubscription(); return; } try { using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); var handle = _subscriptionHandle; // non-null per the guard above; capture before the async lambda await _invoker.ExecuteAsync( DriverCapability.Subscribe, _driverInstanceId, async ct => await subscribable.UnsubscribeAsync(handle, ct), cts.Token); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: unsubscribe threw (continuing)", _driverInstanceId); } finally { DetachSubscription(); } } /// Tear down the data-change + native-alarm event handlers + null the handle. Called from the /// Unsubscribe path, on PostStop, and on Connected → Reconnecting transitions so a stale handler doesn't /// push data-change / alarm events to an actor that has lost its driver connection. private void DetachSubscription() { if (_driver is ISubscribable subscribable && _dataChangeHandler is not null) { subscribable.OnDataChange -= _dataChangeHandler; } _dataChangeHandler = null; _subscriptionHandle = null; DetachAlarmSource(); } /// Subscribe the driver's native alarm event (if it is an ), /// marshaling each transition to the actor thread. Idempotent; mirrors the OnDataChange attach. private void AttachAlarmSource() { if (_driver is not IAlarmSource src || _alarmEventHandler is not null) return; var self = Self; _alarmEventHandler = (_, e) => self.Tell(new NativeAlarmRaised(e)); src.OnAlarmEvent += _alarmEventHandler; } /// Subscribe the driver's (if it is one), /// marshaling each raise to the actor thread. Idempotent; mirrors . /// Attached once in PreStart rather than per-connect, because the interesting raises /// (a Galaxy redeploy, a TwinCAT symbol-version bump) can happen while the driver is between /// connects, and the event carries no connection affinity. private void AttachRediscoverySource() { if (_driver is not IRediscoverable src || _rediscoveryHandler is not null) return; var self = Self; _rediscoveryHandler = (_, e) => self.Tell(new RediscoveryRaised(e)); src.OnRediscoveryNeeded += _rediscoveryHandler; } /// Symmetric teardown, called from PostStop. Load-bearing: the instance /// can OUTLIVE this actor (the host respawns a child around the same driver object), so a missing /// unsubscribe would accumulate one handler per respawn, each holding a dead Self. private void DetachRediscoverySource() { if (_driver is IRediscoverable src && _rediscoveryHandler is not null) src.OnRediscoveryNeeded -= _rediscoveryHandler; _rediscoveryHandler = null; } /// Subscribe the driver's (if it is /// one), marshaling each transition to the actor thread. Idempotent; mirrors /// , including the PreStart-not-per-connect placement — a probe /// loop can report a host down while the driver itself is between connects. private void AttachHostStatusSource() { if (_driver is not IHostConnectivityProbe src || _hostStatusHandler is not null) return; var self = Self; _hostStatusHandler = (_, _) => self.Tell(new HostStatusRaised()); src.OnHostStatusChanged += _hostStatusHandler; } /// Symmetric teardown, called from PostStop — same leak argument as /// . private void DetachHostStatusSource() { if (_driver is IHostConnectivityProbe src && _hostStatusHandler is not null) src.OnHostStatusChanged -= _hostStatusHandler; _hostStatusHandler = null; } /// Current per-host connectivity, or null when the driver is not an /// . Pulled fresh rather than cached: the driver owns the host set, /// and a stale local copy would be a second source of truth to keep in sync. /// GetHostStatuses() is called UNGUARDED (not through ) on purpose — /// it is a pure in-memory snapshot with no I/O, which is exactly why /// UnwrappedCapabilityCallAnalyzer exempts it. A driver that makes it do I/O breaks that /// contract; the try/catch below keeps a misbehaving one from killing the health publish. private IReadOnlyList? CurrentHostStatuses() { if (_driver is not IHostConnectivityProbe probe) return null; try { return probe.GetHostStatuses(); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: GetHostStatuses threw during health publish; omitting host detail", _driverInstanceId); return null; } } /// Records the driver's rediscovery raise and re-publishes health so the signal reaches the /// AdminUI promptly rather than waiting for the next 30 s heartbeat. /// Advisory only. The served address space is deliberately NOT rebuilt: v3 authors raw tags /// explicitly through the /raw browse-commit flow, so a runtime graft would create nodes no /// operator approved and that no deployment artifact records. This prompts a human to re-browse. private void HandleRediscoveryRaised(RediscoveryRaised msg) { _rediscoveryNeededUtc = DateTime.UtcNow; _rediscoveryReason = msg.Args.Reason; _log.Info( "DriverInstance {Id}: driver reports its tag set may have changed ({Reason}) — surfaced for operator re-browse; the served address space is unchanged", _driverInstanceId, msg.Args.Reason); PublishHealthSnapshot(); } /// Symmetric teardown — called from and PostStop so a stale /// handler never pushes to a disconnected actor. private void DetachAlarmSource() { if (_driver is IAlarmSource src && _alarmEventHandler is not null) src.OnAlarmEvent -= _alarmEventHandler; _alarmEventHandler = null; // Drop our handle so the next Connected entry re-subscribes; the desired alarm refs persist across // reconnects. NOTE: this does NOT tear down the driver-side subscription. For a session-bound // IAlarmSource the old subscription dies with the session (no accumulation). For a session-less feed // (GalaxyDriver's always-on central monitor) it survives an in-place reconnect, so the re-subscribe // is additive — but the driver now collapses to a single live handle on each SubscribeAlarmsAsync // (GalaxyDriver.SubscribeAlarmsAsync clears the set before adding), so handles no longer accumulate // across reconnects. The gate (Count > 0) and the one-shot fan-out are unchanged. _alarmSubscriptionHandle = null; } /// Establish the native-alarm subscription that un-gates an driver's /// feed — the driver suppresses until at least one alarm /// subscription exists. Self-sends the async the Connected behaviour /// handles. Idempotent: a no-op unless the driver is an , alarm refs are /// desired, and no subscription is yet established. private void SubscribeDesiredAlarms() { if (_driver is IAlarmSource && _desiredAlarmRefs.Count > 0 && _alarmSubscriptionHandle is null) Self.Tell(new SubscribeAlarms()); } /// Calls the driver's (bounded) to register the /// alarm subscription that un-gates its native feed, caching the returned handle. Re-checks the guard /// (the desired set may have cleared, or another SubscribeAlarms may have already established a handle, /// while this was queued). Failures are logged and retried on the next Connected entry — the feed simply /// stays gated until then. private async Task HandleSubscribeAlarmsAsync(SubscribeAlarms _) { if (_driver is not IAlarmSource src || _desiredAlarmRefs.Count == 0 || _alarmSubscriptionHandle is not null) return; var refs = _desiredAlarmRefs; using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); try { _alarmSubscriptionHandle = await _invoker.ExecuteAsync( DriverCapability.AlarmSubscribe, _driverInstanceId, async ct => await src.SubscribeAlarmsAsync(refs, ct), cts.Token); _log.Info("DriverInstance {Id}: native-alarm subscription established for {Count} alarm ref(s) ({Diag})", _driverInstanceId, refs.Count, _alarmSubscriptionHandle.DiagnosticId); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: native-alarm subscription failed — feed stays gated until reconnect", _driverInstanceId); } } /// Records the host's desired subscription set without touching the live subscription. /// The set is (re)applied by on the next Connected entry. private void StoreDesiredSubscriptions(SetDesiredSubscriptions msg) { _desiredRefs = msg.FullReferences; _desiredInterval = msg.PublishingInterval; _desiredAlarmRefs = msg.AlarmReferences ?? Array.Empty(); } /// Re-establish the desired subscription after (re)connecting. Self-sends the one-shot /// the Connected behaviour already handles (which drops any prior handle /// first), so values resume streaming after a reconnect or redeploy without host involvement. private void ResubscribeDesired() { if (_desiredRefs.Count > 0) { Self.Tell(new Subscribe(_desiredRefs, _desiredInterval)); } } private void OnDataChangeForward(DataChangeForward msg) { var quality = QualityFromStatus(msg.Snapshot.StatusCode); var ts = msg.Snapshot.SourceTimestampUtc ?? msg.Snapshot.ServerTimestampUtc; Context.Parent.Tell(new AttributeValuePublished(_driverInstanceId, msg.FullReference, msg.Snapshot.Value, quality, ts)); } /// Translate an OPC UA status code to the 3-state projection /// the publish actor consumes. Severity bits (top 2): 00 = Good, 01 = Uncertain, 10/11 = Bad. private static OpcUaQuality QualityFromStatus(uint statusCode) { var severity = statusCode >> 30; return severity switch { 0 => OpcUaQuality.Good, 1 => OpcUaQuality.Uncertain, _ => OpcUaQuality.Bad, }; } private static bool IsGoodStatus(uint statusCode) => (statusCode >> 30) == 0; /// Resolve the resilience-pipeline host key for a single driver-side reference. Multi-host /// drivers (e.g. Modbus across N PLCs) implement so one dead host /// trips only its own breaker; single-host drivers fall back to the driver-instance id. Used for the /// single-reference capability calls (write, alarm ack); bulk/whole-driver calls key on the instance /// id directly. private string ResolveHost(string reference) => _driver is IPerCallHostResolver resolver ? resolver.ResolveHost(reference) : _driverInstanceId; /// Records a transition into a Faulted / error state for the 5-minute sliding counter. private void RecordFault() { _faultTimestamps.Enqueue(DateTime.UtcNow); } /// Returns how many fault transitions occurred in the last 5 minutes. private int ErrorCount5Min() { var cutoff = DateTime.UtcNow.AddMinutes(-5); while (_faultTimestamps.Count > 0 && _faultTimestamps.Peek() < cutoff) _faultTimestamps.Dequeue(); return _faultTimestamps.Count; } /// /// Polls and forwards the snapshot to the health publisher. /// Called on every observable state change and by the periodic /// so the AdminUI snapshot store is warmed up for newly-joined SignalR clients. /// Deduplicates: if the resulting (state, lastSuccess, lastError, errorCount) tuple matches /// the last publish, this call is a no-op. Stops flood-publishing identical Healthy snapshots /// every 30s when nothing has changed. Newly-joined SignalR clients still get the current /// snapshot via DriverStatusHub.JoinDriver which reads the store directly. /// private void PublishHealthSnapshot() { try { var health = _driver.GetHealth(); var errorCount = ErrorCount5Min(); var hostStatuses = CurrentHostStatuses(); // _rediscoveryNeededUtc is PART OF THE FINGERPRINT on purpose. A rediscovery raise on an // otherwise-unchanged Healthy driver leaves (state, lastSuccess, lastError, errorCount) // identical, so without it the dedup below would swallow the very publish that carries the // signal and the operator would never see the prompt. // // The host-status digest is in the fingerprint for EXACTLY the same reason, and the failure // mode is the more likely one of the two: a multi-device driver stays aggregate-Healthy when a // single device drops, so every other fingerprint component is unchanged and the transition — // the whole point of this channel — would be deduped away. It is a flattened STRING because a // tuple holding IReadOnlyList compares by reference, which would never match and so would // defeat the dedup in the opposite direction (re-publishing every 30 s heartbeat). var fingerprint = (health.State, health.LastSuccessfulRead, health.LastError, errorCount, _rediscoveryNeededUtc, HostStatusDigest(hostStatuses)); if (_lastPublishedFingerprint is { } prev && prev.Equals(fingerprint)) return; _lastPublishedFingerprint = fingerprint; _healthPublisher.Publish( _clusterId, _driverInstanceId, health, errorCount, _rediscoveryNeededUtc, _rediscoveryReason, hostStatuses); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: GetHealth threw during health publish; skipping", _driverInstanceId); } } /// Order-insensitive value digest of a host-status list for the publish fingerprint. Null in, /// null out — so "driver is not a probe" and "probe reports no hosts" stay distinguishable rather than /// both collapsing to the empty string. /// Ordered by host name because a driver is free to return its hosts in any order (several build /// the list from a Dictionary), and an order flip would otherwise read as a real transition and /// re-publish forever. private static string? HostStatusDigest(IReadOnlyList? statuses) => statuses is null ? null : string.Join('|', statuses .OrderBy(s => s.HostName, StringComparer.Ordinal) .Select(s => $"{s.HostName}={s.State}@{s.LastChangedUtc:O}")); /// Fingerprint of the last call; null until first publish. private (DriverState State, DateTime? LastSuccess, string? LastError, int ErrorCount, DateTime? RediscoveryNeededUtc, string? HostStatusDigest)? _lastPublishedFingerprint; /// protected override void PostStop() { DetachSubscription(); // MUST happen: the IDriver instance can outlive this actor (the host respawns a child around the // same driver object), so a missing unsubscribe accumulates a handler per respawn holding a dead Self. DetachRediscoverySource(); DetachHostStatusSource(); try { _driver.ShutdownAsync(CancellationToken.None).GetAwaiter().GetResult(); } catch (Exception ex) { _log.Warning(ex, "DriverInstance {Id}: ShutdownAsync threw on PostStop", _driverInstanceId); } OtOpcUaTelemetry.DriverInstanceLifecycle.Add(1, new KeyValuePair("event", "stop"), new KeyValuePair("driver_type", _driver.DriverType)); } }