using System.Collections.Frozen; using System.Globalization; using System.Text.Json; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; using ZB.MOM.WW.OtOpcUa.Core.Abstractions; namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect; /// /// MTConnect Agent implementation of — the lifecycle half. Everything the /// driver serves sits behind one , so the whole class is /// exercised against canned XML with no socket. /// /// /// /// The constructor opens nothing and builds nothing. It does not even call /// agentClientFactory — the Wave-0 universal discovery browser constructs a throwaway /// instance purely to ask whether the driver can browse, and a constructor that materialized /// a client would open a connection pool per browse probe against a possibly-unreachable /// Agent. Every connection is owned by . /// /// /// The driverConfigJson argument is honoured, not decorative. The runtime hands /// a config change to a live instance — DriverInstanceActor's ApplyDelta /// calls with the new document and never rebuilds the object /// — so a driver that read only its constructor options would turn every MTConnect config /// edit into a silent no-op that still seals the deployment green. This class therefore /// re-parses the document () whenever it carries a real body, /// matching S7/AbCip/TwinCAT/Galaxy rather than the Modbus/FOCAS/OpcUaClient drivers that /// ignore the argument. A blank / {} / [] document means "no config supplied" /// and keeps the constructor options, which is what unit tests and the factory path rely on. /// /// /// Failure is never dressed as success (#485). A failed initialize disposes whatever /// client it built, drops the probe cache and the observation index, publishes /// with the error text, and rethrows — the runtime's /// retry/backoff loop is the recovery path. In particular a /probe that succeeds while /// the priming /current fails is Faulted, not "Healthy with no data": the index IS the /// read surface, so serving that state would render every tag /// BadWaitingForInitialData indefinitely behind a green driver, and the /// /sample cursor (Task 11) only exists because /current reported a /// nextSequence. /// /// /// Scope. This type implements , , /// , , /// and — the last two /// built on state this class already caches (the Agent instanceId and the probe /// model). There is deliberately no IWritable: MTConnect is a read-only protocol, /// which is also why every discovered leaf is /// . /// /// /// The read path never writes the shared observation index — see /// . That index has exactly one writer: the /sample pump /// (). The same single-writer discipline covers the two /// Task-13 fields: the connectivity probe loop is the sole writer of /// , and it publishes nothing into /// (see ). /// /// public sealed class MTConnectDriver : IDriver, IReadable, ISubscribable, ITagDiscovery, IHostConnectivityProbe, IRediscoverable { /// /// Coarse per-entry byte weights behind . The contract asks /// for an approximate driver-attributable footprint that Core polls every 30 s to /// decide whether to ask for a flush, so an order-of-magnitude estimate from the two things /// that actually scale with a plant (declared data items, live observations) is the right /// precision. Measuring real retained size would need a heap walk per poll. /// private const int ProbeModelBytesPerDataItem = 512; /// Estimated retained bytes per live dataItemId → snapshot entry. private const int IndexBytesPerEntry = 256; /// Estimated retained bytes per authored . private const int TagBytesPerEntry = 256; /// /// The Agent could not be reached, answered unusably, or blew the per-call deadline. Distinct /// from the index's BadNoCommunication, which means the Agent answered and said the /// machine has no value — the two failures need different operator responses. /// private const uint BadCommunicationError = 0x80050000u; /// /// There is no Agent client at all: the driver has not been initialized, its initialize /// faulted, or it has been shut down. Deliberately not one of the index's codes — none of /// them would say "this driver is not running", which is the only true statement available. /// private const uint BadNotConnected = 0x808A0000u; /// /// Floor the reconnect backoff grows from, in milliseconds. Load-bearing: the shipped /// default is 0 (an immediate /// first retry, which is what an operator wants after a one-off blip) and /// 0 × multiplier is still 0 — so without a floor, the "geometric" backoff would be /// an unbounded hot loop against an Agent that stays down. Mirrors /// ModbusTcpTransport.ConnectWithBackoffAsync. /// private const int BackoffGrowthFloorMs = 100; /// Growth factor used when the authored multiplier could not grow the delay (≤ 1). private const double DefaultBackoffMultiplier = 2.0; /// /// Backoff cap used when is authored /// non-positive. Treated as "unset" rather than as "no cap" for the same reason a /// multiplier ≤ 1 falls back to : a zero cap clamps /// EVERY delay to zero, which turns the reconnect loop into a spin against a refused /// connection — one Warning per iteration, as fast as the socket can fail. Mirrors the /// shipped default. /// private const int DefaultMaxBackoffMs = 30_000; /// /// How long teardown waits for a background loop (the /sample pump, the connectivity /// probe) to unwind before abandoning the wait and carrying on. /// /// /// An unbounded wait here is not merely slow — it wedges an actor-system thread. /// DriverInstanceActor.PostStop calls with /// GetAwaiter().GetResult(), i.e. a synchronous block on an Akka dispatcher thread, /// while this driver holds . So a subscriber callback that never /// returns would take the dispatcher thread AND every subsequent lifecycle call on this /// driver with it, permanently. Five seconds matches the budget /// DriverInstanceActor already puts around . /// private static readonly TimeSpan TeardownWait = TimeSpan.FromSeconds(5); /// /// Config-JSON reader options, mirroring the sibling driver factories. Note there is /// deliberately no JsonStringEnumConverter: enum-carrying DTO fields stay /// string? and go through , which is the project-wide /// guard against the numerically-serialized-enum trap that faults a driver at deploy time. /// private static readonly JsonSerializerOptions JsonOptions = new() { PropertyNameCaseInsensitive = true, ReadCommentHandling = JsonCommentHandling.Skip, AllowTrailingCommas = true, }; /// /// Reader options for the emptiness probe in . Kept in step with /// so a document the real parse would accept is never judged /// malformed here (a stray trailing comma must not change which branch a config takes). /// private static readonly JsonDocumentOptions DocumentOptions = new() { CommentHandling = JsonCommentHandling.Skip, AllowTrailingCommas = true, }; private readonly string _driverInstanceId; private readonly ILogger _logger; private readonly Func _agentClientFactory; /// /// How the connectivity probe loop waits between ticks. Injected so the cadence is a seam /// rather than a wall clock: the loop is a timer, and a suite that drove it by sleeping would /// be flaky by construction. Defaults to . /// private readonly Func _probeScheduler; /// /// Serializes the three lifecycle entry points against each other. Initialize / Reinitialize /// / Shutdown all swap the client and the served caches; overlapping them would let a /// shutdown dispose a client an in-flight initialize is about to publish. /// /// /// /// ⚠️ NON-REENTRANT — never call a public lifecycle method from code that already /// holds this. has no owner tracking, so a re-entrant /// call does not fail: it deadlocks the driver permanently, with no exception and /// no log line to diagnose it from. /// /// /// This is aimed squarely at the /sample pump (Task 11). Its ring-buffer /// re-baseline ("on a sequence gap, re-/current and resume") is naturally written /// as "just call the existing re-init logic" — and /// / both take this /// semaphore, so from inside the pump loop that is a hang, not a re-prime. Re-baseline /// must instead call directly on the /// client the pump already holds (through , to keep the /// deadline) and Apply the result to . Any future /// shared step belongs in a private, non-locking helper that both the pump and the /// lock-holding lifecycle methods can call — never in the public entry points. /// /// /// The inverse direction is already safe: the lifecycle methods call /// while holding this, so that method — and whatever /// Task 11 fills it with — must never re-enter the lifecycle either. /// /// private readonly SemaphoreSlim _lifecycle = new(1, 1); /// /// Guards the subscription registry and the pump's control block (, /// , ). /// /// /// /// Lock order is → this, never the reverse. The lifecycle /// methods take this lock (to start/stop the pump) while holding the semaphore; /// therefore releases this lock before awaiting /// , and nothing under this lock ever waits on /// . /// /// /// Nothing is awaited while it is held, and no subscriber callback is raised under /// it. A handler is caller code: it may block, throw, or re-enter this driver, and /// holding a lock across it would turn a slow subscriber into a stalled pump. /// /// private readonly Lock _subscriptionLock = new(); /// Live subscriptions by id. Guarded by . private readonly Dictionary _subscriptions = []; /// /// Guards the two host-connectivity fields, which the probe loop writes and /// reads from an arbitrary caller thread. Held for a field /// read/write only — never across an await, and never while a status handler runs. /// private readonly Lock _probeLock = new(); /// /// The Agent's last observed reachability. Written only by /// , so the state and the event stream can never disagree: /// every value this field has ever held was announced by exactly one /// . That is why a successful /// does not seed it — raising the event from there would run /// caller code while the non-reentrant semaphore is held (the same /// hazard that makes the pump, not , republish to subscribers), /// and seeding it silently would leave a consumer that tracks the event stream disagreeing /// with . until the first tick /// is the honest answer. Guarded by . /// private HostState _hostState = HostState.Unknown; /// When last changed. Guarded by . private DateTime _hostStateChangedUtc = DateTime.UtcNow; /// /// Cancels the connectivity probe loop. Written only under , which /// serializes every start and stop of the loop. /// private CancellationTokenSource? _probeCts; /// The connectivity probe loop's task, or null when none runs. See . private Task? _probeTask; /// /// The Agent instanceId the last was raised for — /// seeded at the commit point of with the id the priming /// /current reported. This is what makes rediscovery fire once per change /// rather than once per chunk: every chunk after a restart carries the new id, and a /// re-baseline that fails leaves the pump still holding the old one. /// private long _announcedInstanceId; private MTConnectDriverOptions _options; /// /// The tag tables derived from — the RawPath ↔ dataItemId /// translation the whole data plane runs through, plus the definition set the observation /// index coerces against. Rebuilt (never mutated) whenever options are installed, and always /// published in the same breath as them, so a lock-free reader can never pair one options /// document's tags with another's RawPaths. /// private TagBinding _binding; private MTConnectObservationIndex _index; private DriverHealth _health = new(DriverState.Unknown, null, null); private long _nextSubscriptionId; /// Cancels the shared /sample pump. Guarded by . private CancellationTokenSource? _pumpCts; /// The shared /sample pump task. Guarded by . private Task? _pumpTask; /// /// The Agent answered /sample with something that cannot be streamed at all, so the /// pump gave up. Latched — reconnecting reproduces the identical answer forever, and a /// later is not new information. Cleared only by a lifecycle /// start, which is the one event that means an operator changed something. Guarded by /// . /// private bool _streamUnsupported; /// /// Whether the /sample pump is currently failing. Kept separate from /// because /sample and /current are different /// Agent requests that fail independently — see . /// private volatile bool _streamDegraded; /// Whether the /current read path is currently failing. See . private volatile bool _readDegraded; /// /// Consecutive /sample reconnect attempts that have not yet been vindicated by /// a delivered chunk — the input to . Reset by any stream that /// actually delivered something, so an agent that drops a connection once an hour is never /// treated as one that has been failing all day. /// private volatile int _reconnectAttempts; /// /// Latch behind the once-per-stream-generation subscriber-fault Warning. A consumer that /// throws on every value would otherwise emit one Warning per observation per handle /// — thousands a second on a busy Agent, which buries the very log it is trying to raise. /// private int _subscriberFaultLogged; /// /// The live agent client plus the /probe cache derived from it, held as one /// immutable snapshot behind a single reference and only ever replaced by a /// of a whole new instance — the same discipline /// uses, and for the same reason. /// /// /// /// Why one field instead of three. These three values are read outside /// — by , by /// (Task 12's browse), and by /// — while every write happens under it. As three /// separate fields they could be observed mid-update: a reader could see a dropped /// ProbeModel beside a stale non-zero data-item count and report a footprint for /// a cache that no longer exists, or pair one initialize's client with the previous /// one's model. Packing them makes every observation internally consistent by /// construction, with no lock on the read path. /// /// /// Why not simply take the lock in those readers. Both readers await the /// network, so they would hold the lifecycle lock across an HTTP round trip — a browse /// could stall a for a full RequestTimeoutMs — and /// every added lock site widens the non-reentrancy hazard documented on /// . It would also make the read path inconsistent with /// , which is deliberately lock-free. /// /// /// What this does NOT fix. A reader can still capture a session moments before a /// concurrent teardown disposes its client, and then call it — no lock-free scheme can /// prevent that, and holding the lock across the call is the trade rejected above. /// What it guarantees is that the failure is named: the resulting /// is caught at the seam and reported as the same /// "driver is not connected" error a caller already has to handle, rather than escaping /// raw out of a browse. /// /// private AgentSession? _session; /// Creates a driver instance. Opens no connection and builds no agent client. /// /// The typed configuration this instance starts from. A config document supplied to /// / supersedes it. /// /// Stable logical id of this driver instance, from the config DB. /// /// Builds the Agent client from the effective options. Defaults to the production /// ; tests inject a canned-XML fake. Called from /// , never from this constructor. /// /// Optional logger; defaults to the null logger. /// /// How the connectivity probe loop waits between ticks. Defaults to /// ; tests inject a hand-driven ticker so /// the probe's cadence is asserted deterministically instead of slept through. /// public MTConnectDriver( MTConnectDriverOptions options, string driverInstanceId, Func? agentClientFactory = null, ILogger? logger = null, Func? probeScheduler = null) { ArgumentNullException.ThrowIfNull(options); ArgumentException.ThrowIfNullOrWhiteSpace(driverInstanceId); _driverInstanceId = driverInstanceId; _logger = logger ?? NullLogger.Instance; _agentClientFactory = agentClientFactory ?? (o => new MTConnectAgentClient(o, logger)); _probeScheduler = probeScheduler ?? Task.Delay; // Sets _options AND the tag tables derived from it. No /probe model is available here (the // constructor dials nothing), so a raw tag whose blob declares no type falls back to String // until InitializeAsync rebuilds the binding against the Agent's own declaration. _options = options; _binding = BuildBinding(options, probeModel: null); // Never null, so every consumer (Task 10's ReadAsync included) can ask it for a snapshot // without a null check: an un-primed index answers BadWaitingForInitialData for an authored // tag and BadNodeIdUnknown otherwise, which is exactly the truth before Initialize runs. _index = new MTConnectObservationIndex(_binding.IndexTags); } /// public string DriverInstanceId => _driverInstanceId; /// /// /// Sourced from the constant rather than a literal — the /// repo-wide fix for the historical TwinCat/Focas drift, where a driver's own /// spelling diverged from the constant every dispatch map keys off. /// public string DriverType => DriverTypeNames.MTConnect; /// /// The live dataItemId → snapshot map that backs the subscription surface: /// primed by and thereafter written only by the /sample /// pump (Task 11). Exposed to tests as the observable proof that initialize actually primed /// it — and, in MTConnectReadTests, that leaves it alone. /// /// /// Read does not consume this. indexes its own /current /// into a per-call snapshot instead, so the pump keeps a single writer — see the failure /// analysis in 's remarks. /// internal MTConnectObservationIndex ObservationIndex => Volatile.Read(ref _index); /// /// The cached /probe device model, or null when it has not been fetched or was /// dropped by . Prefer /// , which cannot observe the flushed hole. /// internal MTConnectProbeModel? CachedProbeModel => Volatile.Read(ref _session)?.ProbeModel; /// /// The Agent's instanceId as of the last successful /current — including the /// re-baseline the /sample pump runs when it sees the id change mid-stream. Task 13 /// watches this for change to raise rediscovery (an Agent restart invalidates every cached /// id). /// internal long? AgentInstanceId => Volatile.Read(ref _session)?.InstanceId; /// /// The sequence the next /sample stream will open at — the priming /current's /// nextSequence, advanced by every re-baseline and by the pump as it stops. Paired /// with in one session record, so the two are never observed /// out of step. /// internal long? AgentNextSequence => Volatile.Read(ref _session)?.NextSequence; /// /// Consecutive /sample reconnect attempts not yet vindicated by a delivered chunk — /// the input to , exposed so the reset policy can be asserted /// without measuring how long anything slept. /// internal int ReconnectAttempts => _reconnectAttempts; /// /// The shared /sample pump's task while one is running, else null. Exposed to /// tests as the deterministic teardown barrier: "the pump has stopped" is a task completion, /// never a sleep. Guarded, so a test never observes a half-installed pump. /// internal Task? SampleStreamTask { get { lock (_subscriptionLock) { return _pumpTask; } } } /// /// The connectivity probe loop's task while one is running, else null. Exposed to /// tests as the deterministic teardown barrier — "the probe loop has stopped" is a task /// completion, never a sleep — and as the proof that Probe.Enabled = false starts /// nothing at all. /// internal Task? ProbeLoopTask => Volatile.Read(ref _probeTask); /// The options currently in force — the constructor's, or the last document applied. internal MTConnectDriverOptions EffectiveOptions => _options; /// /// /// Builds the Agent client from the effective options, runs one deadline-bounded /// /probe (caching the device model and its data-item count), then one /// deadline-bounded /current to prime the observation index and record the Agent's /// instanceId. Publishes on success; on any failure /// disposes the client, clears the caches, publishes with /// the error text, and rethrows. /// public async Task InitializeAsync(string driverConfigJson, CancellationToken cancellationToken) { await _lifecycle.WaitAsync(cancellationToken).ConfigureAwait(false); try { // Parse BEFORE tearing down — see ResolveIncomingOptions. Initialize is reachable on a // live instance (DriverInstanceActor re-initialises in Connecting/Reconnecting), and // destroying a working client before discovering the new document is unreadable would // be the exact outcome ReinitializeAsync is written to avoid. var options = ResolveIncomingOptions(driverConfigJson); // A re-Initialize on a live instance must not orphan the previous client — nor the pump // that is enumerating it. Stopping the stream FIRST is the same ordering // ReinitializeAsync and ShutdownAsync use, and for the same reason: a pump left // enumerating a disposed client reads the disposal as a dropped connection and // reconnects against an object that no longer exists. await StopSampleStreamAsync(cancellationToken).ConfigureAwait(false); await StopProbeLoopAsync(cancellationToken).ConfigureAwait(false); await TeardownCoreAsync().ConfigureAwait(false); await StartCoreAsync(options, existingClient: null, cancellationToken).ConfigureAwait(false); } finally { _lifecycle.Release(); } } /// /// /// /// Stops the /sample stream, then takes one of two paths depending on whether the /// new document changes anything the reads at /// construction: /// /// /// Endpoint/transport changed (AgentUri, DeviceName, /// RequestTimeoutMs, HeartbeatMs, SampleIntervalMs, /// SampleCount) ⇒ full teardown and rebuild. The plan names only the first two, /// but the other four are baked into the client's deadlines, its watchdog window, and its /// /sample query string, so keeping the old client because the URI happened not to /// change would silently discard the operator's edit — the same class of defect as /// ignoring the config document altogether. /// /// /// Otherwise ("config-only": authored tags, probe cadence, reconnect backoff) ⇒ /// the live client and its connection pool are kept, and everything derived from config /// is rebuilt: the probe model is re-fetched, the observation index is recreated against /// the new tag types, and one /current re-primes it. Nothing config-derived /// survives, so a retyped tag cannot keep coercing to its old type. /// /// public async Task ReinitializeAsync(string driverConfigJson, CancellationToken cancellationToken) { await _lifecycle.WaitAsync(cancellationToken).ConfigureAwait(false); try { var incoming = ResolveIncomingOptions(driverConfigJson); var existing = Volatile.Read(ref _session)?.Client; await StopSampleStreamAsync(cancellationToken).ConfigureAwait(false); await StopProbeLoopAsync(cancellationToken).ConfigureAwait(false); if (existing is null || RequiresClientRebuild(_options, incoming)) { await TeardownCoreAsync().ConfigureAwait(false); await StartCoreAsync(incoming, existingClient: null, cancellationToken).ConfigureAwait(false); return; } await StartCoreAsync(incoming, existing, cancellationToken).ConfigureAwait(false); } finally { _lifecycle.Release(); } } /// /// /// Stops the /sample stream, disposes the Agent client, and drops the probe cache and /// the observation index — values sourced from a torn-down connection must never keep /// reading Good. LastSuccessfulRead is retained because it is diagnostic ("when did /// this driver last hear from the Agent"), not served data. Idempotent. /// public async Task ShutdownAsync(CancellationToken cancellationToken) { await _lifecycle.WaitAsync(cancellationToken).ConfigureAwait(false); try { var lastRead = ReadHealth().LastSuccessfulRead; await StopSampleStreamAsync(cancellationToken).ConfigureAwait(false); await StopProbeLoopAsync(cancellationToken).ConfigureAwait(false); await TeardownCoreAsync().ConfigureAwait(false); WriteHealth(new DriverHealth(DriverState.Unknown, lastRead, null)); } finally { _lifecycle.Release(); } } /// public DriverHealth GetHealth() => ReadHealth(); /// /// /// Sums the three caches that scale with plant size: the /probe device model (by /// declared data item), the live observation index (by indexed data item), and the authored /// tag table. See for why the weights are coarse /// constants. /// public long GetMemoryFootprint() => ((long)(Volatile.Read(ref _session)?.ProbeModelDataItemCount ?? 0) * ProbeModelBytesPerDataItem) + ((long)ObservationIndex.Count * IndexBytesPerEntry) + ((long)Volatile.Read(ref _binding).IndexTags.Count * TagBytesPerEntry); /// /// /// Drops the /probe device model — the browse cache, which is re-derivable from the /// Agent at any time — and keeps the observation index, which is required for correctness /// (it is the read surface, and re-deriving it would mean a full re-prime the caller did not /// ask for). Flushing cannot wedge the driver: every consumer of the model goes through /// , which re-fetches and re-caches on a miss. /// /// The drop is a single of a whole new /// , so the model and its data-item count can never be seen /// out of step (which would skew ), and a flush that /// races a concurrent teardown or re-initialize is discarded rather than resurrecting /// the session it was reading. /// /// public Task FlushOptionalCachesAsync(CancellationToken cancellationToken) { var session = Volatile.Read(ref _session); if (session?.ProbeModel is null) { return Task.CompletedTask; } var flushed = session with { ProbeModel = null, ProbeModelDataItemCount = 0 }; if (!ReferenceEquals(Interlocked.CompareExchange(ref _session, flushed, session), session)) { // Someone swapped the session while we were building the flushed copy — their state is // newer than ours, and publishing would undo it. Whatever they installed is either a // teardown (no cache to flush) or a fresh initialize (a cache the operator's budget // breach never saw), so dropping this flush is correct, not merely safe. return Task.CompletedTask; } _logger.LogInformation( "MTConnect driver {DriverInstanceId} flushed its cached /probe device model ({DataItemCount} data items); it is re-fetched on the next browse.", _driverInstanceId, session.ProbeModelDataItemCount); return Task.CompletedTask; } /// /// /// /// One /current per batch, whatever the batch size. MTConnect's /// /current answers for the whole device (or the whole Agent), so N references /// cost one deadline-bounded round-trip and are answered out of that one document. There /// is no per-reference request to make and none is made. /// /// /// The response is indexed into a per-call snapshot, NOT into the driver's shared /// observation index. That index is written by the /sample pump (Task 11) and /// is deliberately last-write-wins with no timestamp arbitration. A read that folded its /// /current into it would be a second writer racing the first, and the /// /current document is frequently the older of the two — the Agent /// composes it while the pump's newer delta is already in flight. Losing that race rolls /// a subscribed value backwards and leaves it there until the data item next /// changes, which for an EVENT such as Execution can be hours. Keeping the /// read path a pure function of (document, authored tags) costs one index construction /// per batch — negligible beside the HTTP round-trip that precedes it — and makes it /// impossible for a read to corrupt subscription state. The reference source-of-truth is /// unchanged either way: both paths coerce through the same /// logic, so a read and a subscription report a /// given observation identically. /// /// /// Nothing but genuine caller cancellation crosses this boundary. An unreachable /// Agent, a blown deadline, an unusable answer — each fails the batch's values /// (every reference codes ) rather than the batch: an /// exception here would also fail the references that were perfectly readable, and the /// node-manager cannot tell a thrown read from a broken driver. This is the opposite /// posture to , which is specified to throw so the /// runtime's retry/backoff loop can own recovery. /// /// /// The status-code distinctions the index draws — BadWaitingForInitialData for an /// authored-but-unreported id, BadNodeIdUnknown for an unknown one, /// BadNoCommunication for the Agent's UNAVAILABLE — are passed through /// unmodified. They tell an operator three different things. /// /// /// /// is null — a caller bug, not device data, and the /// one thing this method is loud about (matching 's /// documented posture). /// public async Task> ReadAsync( IReadOnlyList fullReferences, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(fullReferences); // Checked before anything else: an empty batch must cost the Agent nothing, and there is no // failure it could report. if (fullReferences.Count == 0) { return []; } // One read of the session reference: the client used below is the one that was live at this // instant, never a half-swapped pairing (see the _session remarks). var client = Volatile.Read(ref _session)?.Client; var options = _options; // Read once, for the same reason the session is: every reference in this batch must be // resolved against ONE tag table, not half against the table a concurrent re-initialize // just replaced. var binding = Volatile.Read(ref _binding); if (client is null) { // Before InitializeAsync, after a faulted one, or after ShutdownAsync. Health is left // exactly as it is: a read attempted against a driver that was never started (or was // deliberately stopped) is not evidence of degradation. return BadBatch(fullReferences.Count, BadNotConnected); } try { var current = await BoundedAsync(client.CurrentAsync, "/current", options.RequestTimeoutMs, cancellationToken) .ConfigureAwait(false); // Built from the SAME authored tags as the shared index, so the authored-vs-unknown // distinction and every coercion behave identically — it is a snapshot, not a variant. var snapshot = new MTConnectObservationIndex(binding.IndexTags); snapshot.Apply(current); var results = new DataValueSnapshot[fullReferences.Count]; for (var i = 0; i < results.Length; i++) { // Each reference is a RawPath, translated to the DataItem id the index is keyed by. // An unmapped reference falls through unchanged and misses — Bad, never a throw. // Get never throws and never returns null, including for a null/blank reference. results[i] = snapshot.Get(binding.ToDataItemId(fullReferences[i])); } _readDegraded = false; PublishHealthy(DateTime.UtcNow); return results; } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { // The caller asked us to stop. This is the ONLY exception that may propagate: a blown // per-call deadline arrives as the TimeoutException BoundedAsync raises instead, and is // handled below as the Agent failure it is. throw; } catch (Exception ex) { DegradeAfterReadFailure(ex); return BadBatch(fullReferences.Count, BadCommunicationError); } } // ---- ISubscribable: ONE shared /sample long poll behind every handle ---- /// public event EventHandler? OnDataChange; /// /// /// /// Initial data comes from the primed /current, not from the wire. Every /// subscribed reference gets one callback immediately (the OPC UA Part 4 convention), /// answered out of the observation index — including the references with no value yet, /// which report BadWaitingForInitialData. Waiting for the Agent's next chunk /// instead would leave a fresh subscription silent for up to a whole heartbeat, /// and permanently silent for any data item that simply is not changing. /// /// /// Subscribe does not wait for the stream. The first subscription starts the /// shared pump; the pump then opens /sample on its own task under its own /// cancellation token. This is deliberate on both counts. The opening handshake cannot /// be awaited inside the caller's Subscribe timeout in any useful sense — the first /// chunk may legitimately be a heartbeat away — and the long-lived poll must NOT run /// under the caller's token, which belongs to one Subscribe call and is disposed the /// moment it returns. What bounds the running stream is the client's own heartbeat /// watchdog, and what stops it is or the lifecycle. /// /// /// is not honoured per subscription, and /// cannot be: MTConnect's publish cadence is the Agent-side interval query /// parameter of the single shared /sample request /// (), fixed when the client is /// built. Honouring a per-handle interval would mean one stream per subscription — the /// Agent load this driver exists to avoid — or client-side throttling that would drop /// transient EVENT values the Agent bothered to send. The parameter is accepted and /// ignored, which is also what the natively-subscribing sibling drivers do. /// /// /// Subscribing before is legal: the interest is recorded, /// every reference reports "no value yet", and the initialize that follows starts the /// pump for it. It does not throw and it does not dial an Agent that is not there yet. /// /// /// /// is null — a caller bug, and the only thing this /// method is loud about, matching 's posture. /// public Task SubscribeAsync( IReadOnlyList fullReferences, TimeSpan publishingInterval, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(fullReferences); // De-duplicated: a reference asked for twice in one subscription is one subscribed value, // and would otherwise raise two identical callbacks per chunk forever. var references = new HashSet(StringComparer.Ordinal); foreach (var reference in fullReferences) { if (reference is not null) { references.Add(reference); } } var handle = new MTConnectSampleHandle(Interlocked.Increment(ref _nextSubscriptionId)); lock (_subscriptionLock) { _subscriptions[handle.SubscriptionId] = new SubscriptionState(handle, references); // No-op when a pump is already running: every handle shares the one stream. StartSampleStreamCore(republishOnStart: false); } // Raised OUTSIDE the lock — a handler is caller code (see the _subscriptionLock remarks). // The callback carries the caller's OWN reference (the RawPath); only the index lookup // behind it is translated to the DataItem id. var index = ObservationIndex; var binding = Volatile.Read(ref _binding); foreach (var reference in references) { RaiseDataChange(handle, reference, index.Get(binding.ToDataItemId(reference))); } _logger.LogDebug( "MTConnect driver {DriverInstanceId} opened subscription {SubscriptionId} over {ReferenceCount} reference(s).", _driverInstanceId, handle.DiagnosticId, references.Count); return Task.FromResult(handle); } /// /// /// Drops the handle's references and stops the shared stream when the last subscription /// goes. An unknown handle — already unsubscribed, or issued by a different driver — is a /// no-op rather than an error: this runs on teardown paths, including failure paths, where /// throwing would mask whatever prompted the teardown. /// /// is null. public async Task UnsubscribeAsync(ISubscriptionHandle handle, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(handle); if (handle is not MTConnectSampleHandle ours) { return; } bool stopStream; lock (_subscriptionLock) { if (!_subscriptions.Remove(ours.SubscriptionId)) { return; } stopStream = !HasSubscribedReferencesCore(); } if (stopStream) { // Outside the lock: StopSampleStreamAsync awaits the pump, and the pump takes this same // lock to read the subscription registry. Holding it here would deadlock the two. // The caller's token bounds the wait: DriverInstanceActor puts a 5 s budget around this // call, and until now that budget was inert. await StopSampleStreamAsync(cancellationToken).ConfigureAwait(false); } _logger.LogDebug( "MTConnect driver {DriverInstanceId} closed subscription {SubscriptionId}{StreamNote}.", _driverInstanceId, ours.DiagnosticId, stopStream ? " and stopped the shared /sample stream" : string.Empty); } // ---- ITagDiscovery: the /probe device model, streamed as folders + variables ---- /// /// /// This one member is the whole browse opt-in. The Wave-0 universal /// DiscoveryDriverBrowser turns any driver that answers true here into a live /// address picker by capturing — so this driver ships no browser /// code, no per-driver picker component, and no entry in the browser's browse-patch table /// (its discovery is unconditional; there is no "enable browse" option to turn on first). /// It is answered on an un-initialized instance without dialling anything, because /// CanBrowse constructs a throwaway driver purely to read it — which is also why the /// constructor opens nothing. /// public bool SupportsOnlineDiscovery => true; /// /// /// Once, not the default UntilStable: streams the /// complete device model within the single call — /probe is one synchronous request /// that either answers with the whole hierarchy or fails — so nothing fills in asynchronously /// after connect and a settle loop would only re-fetch an identical answer. (The one thing /// that does change the model is an Agent restart, and that is /// IRediscoverable's job in Task 13, not a retry cadence's.) /// public DiscoveryRediscoverPolicy RediscoverPolicy => DiscoveryRediscoverPolicy.Once; /// /// /// /// Streams the Agent's /probe hierarchy verbatim: one folder per Device, /// one per nested Component to arbitrary depth, and one variable per /// DataItem under whichever container declares it — including the data items a /// device declares directly, outside any component (an availability item almost /// always lives there, and a walker that only recursed Components would silently /// drop it). /// /// /// is the DataItem@id, never the /// browse name. The universal browser commits a browsed leaf's FullName /// straight into TagConfig.FullName, and that is the key /// and the /sample pump resolve an observation against /// ( is keyed by DataItem@id). Committing /// the name instead would produce a picker that looks right and authors tags that /// can never report a value. The browse name — cosmetic — prefers the DataItem's /// name and falls back to its id, which for this protocol is the common /// case rather than an edge case. /// /// /// The type shape comes from 's returned /// record, in full and /// included. Recomputing the array shape /// locally from representation/sampleCount is the tempting shortcut and it /// is wrong: the inference also demotes DATA_SET/TABLE to /// even on a SAMPLE, and honours /// TIME_SERIES only on a SAMPLE. The AdminUI typed editor calls the same /// function, so any divergence here makes the picker and the editor disagree about one /// data item. The type attribute is passed verbatim for the same reason: /// the probe document spells it UPPER_SNAKE (PART_COUNT) where the streams /// document uses PascalCase, and the inference is separator-insensitive precisely to /// bridge the two — pre-normalising it here would defeat that. /// /// /// is this seam's job, deliberately not /// the inference's: a CONDITION data item is an alarm. /// IVariableHandle.MarkAsAlarmCondition is not called — that annotation /// needs the driver-side refs of an alarm's sibling attributes (in-alarm flag, priority, /// ack target), and an MTConnect condition has none: it is one text observation whose /// value is the state word. Core's generic node manager enriches on the flag alone. /// /// /// The model comes from , never from /// . may have /// dropped the cache under memory pressure; reading the field directly would make a /// flushed driver browse as an Agent with no data items, permanently. /// /// /// This throws, and that is the opposite of 's posture on /// purpose. A read has a per-reference Bad status code to report a failure in-band; a /// discovery stream has no such channel, so the only way to "fail quietly" would be to /// emit an empty tree — indistinguishable from an Agent that genuinely declares nothing, /// and read by an operator in the picker as a working connection to an empty machine /// (#485: empty is not an answer). /// /// /// What the callers actually do with the throw — worth stating precisely, because /// the reassuring version ("it is logged and retried") is not true here. The live /// consumer is the /raw browse-commit path /// (DiscoveryDriverBrowser/BrowserSessionService), which surfaces it to the /// operator as a failed browse — that is the case this posture is chosen for. /// DriverInstanceActor.HandleRediscoverAsync catches it, logs a Warning, and /// publishes an EMPTY node set; it does not reschedule, because retrying is a /// behaviour and this driver is /// . And that publish currently reaches /// nothing: DriverHostActor.HandleDiscoveredNodes short-circuits unconditionally /// (discovered-node injection is dormant in v3 — raw tags are authored through /// browse-commit instead), so at that seam the throw-versus-empty choice has no runtime /// effect today. It is made for the browse path, and for whatever re-enables injection. /// /// /// Nothing here writes the observation index. Discovery is a pure read of the /// device model — the /sample pump remains its sole writer. /// /// /// is null. /// /// The driver is not initialized, or it was torn down while the /probe fetch was in /// flight (see ). /// public async Task DiscoverAsync(IAddressSpaceBuilder builder, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(builder); var model = await GetOrFetchProbeModelAsync(cancellationToken).ConfigureAwait(false); // Read once: a concurrent re-initialize must not scope half this walk to one device and half // to another. var deviceScope = _options.DeviceName?.Trim(); var devices = 0; var variables = 0; foreach (var device in model.Devices) { if (device is null || !InDeviceScope(device, deviceScope)) { continue; } var name = FolderName(device.Name, device.Id); if (string.IsNullOrWhiteSpace(name)) { // Neither a name nor an id: there is no browse path that could identify this device, // and folding it in would emit a blank tree segment (and a blank NodeId component // downstream). Its data items are unreachable either way. _logger.LogWarning( "MTConnect driver {DriverInstanceId} skipped a device in {AgentUri}'s /probe model that declares neither a name nor an id; it cannot be given a browse path.", _driverInstanceId, _options.AgentUri); continue; } devices++; variables += StreamContainer(device, builder.Folder(name, name)); } // A configured scope that matches NO device is an authoring error — almost always a typo or // a device renamed in the Agent — and its only symptom is a browse that succeeds and shows // nothing. Reported at Warning, because at Debug the operator's evidence is an empty picker // and no explanation. An agent-wide browse that finds nothing is a different statement (the // Agent declares no devices) and stays a Debug tally. if (devices == 0 && !string.IsNullOrEmpty(deviceScope)) { _logger.LogWarning( "MTConnect driver {DriverInstanceId} is scoped to device '{DeviceScope}', but {AgentUri}'s /probe model declares no device with that name or id ({DeclaredCount} declared). The browse is empty because the scope matched nothing, not because the Agent has nothing.", _driverInstanceId, deviceScope, _options.AgentUri, model.Devices.Count); return; } _logger.LogDebug( "MTConnect driver {DriverInstanceId} streamed {DeviceCount} device(s) and {VariableCount} data item(s) from {AgentUri}'s /probe model{DeviceScope}.", _driverInstanceId, devices, variables, _options.AgentUri, deviceScope is null or "" ? string.Empty : $" (scoped to '{deviceScope}')"); } /// /// Streams one device/component's own data items into , then recurses /// into each nested component under a folder of its own. Returns how many variables it /// streamed, for the diagnostic tally. /// /// /// Walks the tree rather than using AllDataItems(): that helper flattens, and the /// folder structure — which item belongs to which component — is precisely what a browse /// picker exists to show. Recursion depth is bounded by the parser, which refuses a probe /// document nested beyond its own limit. /// private static int StreamContainer(IMTConnectComponentContainer container, IAddressSpaceBuilder scope) { var streamed = 0; foreach (var dataItem in container.DataItems) { // An id-less data item cannot be read, subscribed, or committed as a tag — there would be // nothing to key it by. The parser already requires the attribute; this is the guard for a // model built any other way. if (dataItem is null || string.IsNullOrWhiteSpace(dataItem.Id)) { continue; } // Every field the DataItem declares, straight through — the array shape is READ OFF the // result, never recomputed here. See the DiscoverAsync remarks. var inferred = MTConnectDataTypeInference.Infer( dataItem.Category, dataItem.Type, dataItem.Units, dataItem.Representation, dataItem.SampleCount); var browseName = string.IsNullOrWhiteSpace(dataItem.Name) ? dataItem.Id : dataItem.Name; scope.Variable( browseName, browseName, new DriverAttributeInfo( FullName: dataItem.Id, DriverDataType: inferred.DataType, IsArray: inferred.IsArray, ArrayDim: inferred.ArrayDim, // MTConnect is read-only and this driver implements no IWritable, so the tier // that is read-only from OPC UA is the only honest one. SecurityClass: SecurityClassification.ViewOnly, // Historization is authored per tag, never inferred from a device model. IsHistorized: false, IsAlarm: IsCondition(dataItem.Category))); streamed++; } foreach (var component in container.Components) { if (component is null) { continue; } // Guarded on the RESOLVED folder name, not on the id alone: a component with a blank id // but a real name browses perfectly well (the id is only the fallback), while one with // neither would open a folder named "" — a blank segment in every browse path beneath // it. The data-item guard above is stricter for a different reason: THAT id is the // correlation key an observation is resolved by, so a blank one is unreadable whatever // it is called. var name = FolderName(component.Name, component.Id); if (string.IsNullOrWhiteSpace(name)) { continue; } streamed += StreamContainer(component, scope.Folder(name, name)); } return streamed; } /// /// Whether a device is in this instance's configured scope. An unset scope is agent-wide. /// /// /// /// Matched against the device's name or its id — a device model may /// omit the name, and its id is then what an operator would put in the config and what a /// scoped Agent URL would carry. /// /// /// Belt-and-braces, not dead code. The production client already narrows the /// request to {AgentUri}/{DeviceName}/probe, so a conforming Agent answers with the /// one device anyway. But an Agent (or an intermediary) that ignored the path segment and /// answered agent-wide would otherwise publish every other machine's data items into an /// instance scoped to one, and the operator's only evidence would be a picker full of ids /// they never configured. Matching is exact for the same reason the URL is: a /// differently-cased device name is not the device the Agent would have served. /// /// private static bool InDeviceScope(MTConnectDevice device, string? deviceScope) => string.IsNullOrEmpty(deviceScope) || string.Equals(device.Name, deviceScope, StringComparison.Ordinal) || string.Equals(device.Id, deviceScope, StringComparison.Ordinal); /// The folder name for a device/component: its name, or its id when it declares none. private static string FolderName(string? name, string id) => string.IsNullOrWhiteSpace(name) ? id : name; /// /// Whether a data item's category makes it an alarm. Lives here rather than in /// because that function answers a type /// question; this is an address-space one. /// private static bool IsCondition(string? category) => category.AsSpan().Trim().Equals("CONDITION", StringComparison.OrdinalIgnoreCase); // ---- IHostConnectivityProbe: one host per Agent, on its own cheap tick ---- /// /// /// Raised from the probe loop's task, never under a lock and never under /// . A handler is arbitrary caller code: one that blocks forever /// blocks the loop with it (and therefore the teardown that awaits it), exactly as an /// handler does to the pump. One that throws is absorbed. /// public event EventHandler? OnHostStatusChanged; /// /// The single host this driver instance talks to, named by its configured /// . Mirrors ModbusDriver.HostName's /// host:port convention: several MTConnect drivers in one server disambiguate on the /// Admin dashboard by endpoint rather than by driver-instance id. /// public string HostName => _options.AgentUri; /// /// /// One entry, always — an Agent is a single host, and the driver knows of it whether or not /// it has been reached. Before the first probe tick (and for the whole life of an instance /// with Probe.Enabled = false) that entry reads , which /// is the truth: nothing has measured it. See for why a successful /// initialize does not seed it . /// public IReadOnlyList GetHostStatuses() { lock (_probeLock) { return [new HostConnectivityStatus(HostName, _hostState, _hostStateChangedUtc)]; } } /// /// Starts the periodic connectivity probe, if the operator asked for one. Caller must hold /// . /// /// The client this session owns; the loop dials that one, never a later one. /// The options this session started with. private void StartProbeLoopCore(IMTConnectAgentClient client, MTConnectDriverOptions options) { if (!options.Probe.Enabled) { return; } var cts = new CancellationTokenSource(); _probeCts = cts; Volatile.Write(ref _probeTask, Task.Run(() => RunProbeLoopAsync(client, options, cts.Token), CancellationToken.None)); } /// /// The connectivity probe: wait, ask the Agent whether it is answering, publish the verdict, /// repeat. Never throws — it is a detached background task, and an escaping exception would /// silently end host reporting while kept serving the last /// thing it saw. /// /// /// /// It calls directly and publishes /// nothing. Both halves of that are deliberate. /// /// /// Routing it through is the obvious shortcut and /// it is inert: that accessor answers from the cache /// already warmed, so the loop would issue no request at all and report /// forever no matter what the Agent was doing — a feature /// that ticks, logs, and measures nothing. /// /// /// Publishing the answer into is the opposite /// mistake: it would put a second writer on a field the lifecycle owns under a /// compare-and-swap, and would silently re-warm a cache that /// — or an Agent restart, see /// — deliberately dropped. The parsed model is /// discarded; only the fact that a request completed is used. /// /// /// Why /probe rather than the /sample heartbeat. The heartbeat is /// free, but it only exists while something is subscribed — an instance with no /// subscriptions would report indefinitely, and host /// reachability would become a function of OPC UA client behaviour. A dedicated tick is /// the one signal that is true whether or not anyone is watching. Its cost is a parse of /// the device model per interval, which is why /// exists. /// /// /// It waits first, then probes. has just completed a /// successful /probe one moment earlier; a tick that fired immediately would spend /// a second identical request to learn nothing, and would do it on every re-initialize. /// /// /// It never calls a lifecycle method — see the remarks. /// awaits this task while holding that semaphore, so /// reaching back into one would be a permanent hang with no exception and no log line. /// /// private async Task RunProbeLoopAsync( IMTConnectAgentClient client, MTConnectDriverOptions options, CancellationToken ct) { try { while (!ct.IsCancellationRequested) { try { await _probeScheduler(options.Probe.Interval, ct).ConfigureAwait(false); } catch (OperationCanceledException) { return; } if (ct.IsCancellationRequested) { return; } bool reachable; try { using var deadline = CancellationTokenSource.CreateLinkedTokenSource(ct); deadline.CancelAfter(options.Probe.Timeout); // The answer is deliberately discarded — see the remarks. _ = await client.ProbeAsync(deadline.Token).ConfigureAwait(false); reachable = true; } catch (OperationCanceledException) when (ct.IsCancellationRequested) { return; } catch (Exception ex) { // Unreachable, unusable, disposed, or past its deadline — from the host's point // of view these are one fact, and the data paths report the detail. reachable = false; _logger.LogDebug( ex, "MTConnect driver {DriverInstanceId} could not reach {AgentUri} on its connectivity probe.", _driverInstanceId, options.AgentUri); } TransitionHostTo(reachable ? HostState.Running : HostState.Stopped); } } catch (Exception ex) { _logger.LogError( ex, "MTConnect driver {DriverInstanceId} stopped its connectivity probe on an unexpected fault; host status for {AgentUri} will not be updated until the driver is re-initialized.", _driverInstanceId, options.AgentUri); } } /// /// Publishes a host state and announces it — only when it actually changed. Raising on /// every tick would emit one event per interval per Agent forever: Core would re-fan Bad /// quality across the host's subtree each time, and an operator watching the feed could not /// tell a new outage from an ongoing one. /// private void TransitionHostTo(HostState newState) { HostState old; lock (_probeLock) { old = _hostState; if (old == newState) { return; } _hostState = newState; _hostStateChangedUtc = DateTime.UtcNow; } _logger.LogInformation( "MTConnect driver {DriverInstanceId} host {AgentUri} went {Previous} -> {Current}.", _driverInstanceId, HostName, old, newState); // Outside the lock: a handler is caller code. try { OnHostStatusChanged?.Invoke(this, new HostStatusChangedEventArgs(HostName, old, newState)); } catch (Exception ex) { _logger.LogWarning( ex, "MTConnect driver {DriverInstanceId} caught a fault from a host-status subscriber for {AgentUri}; the probe continues.", _driverInstanceId, HostName); } } /// /// Stops the connectivity probe loop and waits for it to be gone. Idempotent, and never /// throws — it runs on teardown paths where a second exception would mask the first. /// /// /// ⚠️ This runs while is held and it awaits the loop, on /// exactly the same terms as : safe only because the loop /// never touches the lifecycle and every await it makes observes the token cancelled below. /// The wait is bounded by and the caller's token for the same /// reason: an handler that never returns would otherwise /// wedge the Akka dispatcher thread DriverInstanceActor.PostStop blocks on. /// /// The caller's deadline for the teardown wait. private async Task StopProbeLoopAsync(CancellationToken cancellationToken = default) { var cts = _probeCts; var task = Volatile.Read(ref _probeTask); _probeCts = null; Volatile.Write(ref _probeTask, null); if (cts is null) { return; } try { await cts.CancelAsync().ConfigureAwait(false); } catch (ObjectDisposedException) { // Raced another stop; the loop is already going away. } var stopped = true; if (task is not null) { try { await task.WaitAsync(TeardownWait, cancellationToken).ConfigureAwait(false); } catch (Exception ex) when (ex is TimeoutException or OperationCanceledException) { stopped = false; _logger.LogError( ex, "MTConnect driver {DriverInstanceId} gave up waiting for its connectivity probe to stop after {TeardownWait}; the usual cause is an OnHostStatusChanged subscriber that blocks. Teardown continues.", _driverInstanceId, TeardownWait); } catch (Exception ex) { // RunProbeLoopAsync is written not to throw, so this is belt-and-braces. _logger.LogDebug( ex, "MTConnect driver {DriverInstanceId} swallowed a fault from its connectivity probe while stopping it.", _driverInstanceId); } } if (stopped) { // Only once the loop is provably gone — see StopSampleStreamAsync for why. cts.Dispose(); } } // ---- IRediscoverable: the Agent instanceId watch ---- /// /// /// /// This event is what makes = /// safe. The runtime runs exactly one /// discovery pass per connect, so an Agent that restarts mid-run and renumbers, adds, or /// drops data items would otherwise leave a stale address space sitting behind a /// perfectly Healthy driver. /// /// /// Raised from the /sample pump's task, outside every lock. A handler that /// throws is absorbed and logged — a consumer bug must not tear down the stream serving /// every subscriber. A handler that blocks blocks the pump, and one that /// synchronously awaits a lifecycle method on this driver deadlocks it: teardown waits /// for the pump, and the pump is waiting inside the handler. Core's consumer schedules /// the rediscovery rather than running it inline, which is the only safe shape. /// /// public event EventHandler? OnRediscoveryNeeded; /// /// Handles an Agent restart observed on the wire: drop the device model it invalidated, then /// tell Core to re-discover — once per distinct instanceId. /// /// /// /// Dropping the cached /probe model is not housekeeping — it is what stops this /// whole capability being inert. Core's response to the event is to re-run /// , which reads /// ; with the cache still warm that call answers /// from a model describing a process that no longer exists, and the "rediscovery" would /// diff the stale tree against itself and change nothing. Dropped through the same /// compare-and-swap uses, so a lifecycle change /// landing mid-drop wins rather than being undone. /// /// /// Once per change, not once per chunk. Every chunk after a restart carries the new /// id, and a re-baseline that fails leaves the pump still comparing against the old one — /// so the guard is a latch on the id that was last announced, not on the pump's cursor /// state. A second, different restart is a second announcement; an id the Agent reuses /// after some other id is one too, because that is genuinely a new process. /// /// /// The instanceId the Agent's chunk reported. /// The options in force, for the log line's endpoint. private void AnnounceAgentRestart(long instanceId, MTConnectDriverOptions options) { if (Interlocked.Exchange(ref _announcedInstanceId, instanceId) == instanceId) { return; } InvalidateProbeModel(); _logger.LogInformation( "MTConnect driver {DriverInstanceId} dropped its cached /probe device model for {AgentUri} and raised rediscovery: agent instanceId is now {InstanceId}.", _driverInstanceId, options.AgentUri, instanceId); try { OnRediscoveryNeeded?.Invoke(this, new RediscoveryEventArgs("MTConnect agent instanceId changed", null)); } catch (Exception ex) { _logger.LogWarning( ex, "MTConnect driver {DriverInstanceId} caught a fault from a rediscovery subscriber; the /sample stream continues.", _driverInstanceId); } } /// /// Drops the cached /probe device model so the next consumer re-fetches it. Same /// compare-and-swap discipline as — retried, because /// unlike a flush this drop is a correctness requirement rather than a courtesy: losing the /// race must not leave a model from a dead Agent installed. /// private void InvalidateProbeModel() { while (true) { var session = Volatile.Read(ref _session); if (session?.ProbeModel is null) { return; } var dropped = session with { ProbeModel = null, ProbeModelDataItemCount = 0 }; if (ReferenceEquals(Interlocked.CompareExchange(ref _session, dropped, session), session)) { return; } } } /// /// Starts the shared /sample pump, if there is anything to pump and anything to pump /// from. Caller must hold . /// /// /// Whether the pump should republish the observation index to every subscriber before it /// opens the stream — true when a lifecycle start replaced the baseline the subscribers' /// current values came from. /// private void StartSampleStreamCore(bool republishOnStart) { // `IsCompleted: false`, not `is not null`: a pump that has already exited is not a running // stream, and conflating the two means a pump killed by an unexpected fault could never be // revived by a resubscribe. The genuinely-unrepeatable case is the latch below, not this. if (_pumpTask is { IsCompleted: false }) { return; } if (_streamUnsupported) { // Refusing IS correct — reconnecting reproduces the identical answer forever — but it // must never be silent. DriverInstanceActor handles a tag-set change as // Unsubscribe-then-Subscribe with NO re-initialize, and the unsubscribe stops a stream // that is not running; without this, the resubscribe returns a handle, the caller is // told SubscriptionEstablished, and the next successful read reports a perfectly green // driver over a subscription that can never deliver a value (#485). _streamDegraded = true; Degrade($"The MTConnect Agent at '{_options.AgentUri}' does not serve a framed /sample stream; subscriptions cannot deliver until the driver is re-initialized."); _logger.LogWarning( "MTConnect driver {DriverInstanceId} refused to start a /sample stream for {AgentUri}: the endpoint was already found not to stream. The subscription is registered but will receive NO updates until the driver is re-initialized.", _driverInstanceId, _options.AgentUri); return; } if (!HasSubscribedReferencesCore()) { return; } // One read of the session: the client the pump enumerates and the cursor it opens at are // guaranteed to belong to the same initialize. var session = Volatile.Read(ref _session); if (session is null) { return; } // A previous pump ran to completion; its cancellation source is ours to release before the // field is overwritten, or every revived stream leaks one. _pumpCts?.Dispose(); _pumpTask = null; // A new stream generation gets a fresh subscriber-fault Warning budget (see // _subscriberFaultLogged). Interlocked.Exchange(ref _subscriberFaultLogged, 0); var options = _options; var cts = new CancellationTokenSource(); _pumpCts = cts; _pumpTask = Task.Run( () => RunSampleStreamAsync( session.Client, options, session.NextSequence, session.InstanceId, republishOnStart, cts.Token), CancellationToken.None); } /// /// The shared /sample pump: reads chunks, keeps the observation index and the /// subscribers up to date, and owns every way the stream can stop. /// /// /// /// Never throws. It is a detached background task; an escaping exception would be /// an unobserved fault that silently ends the driver's subscription surface while every /// handle still looks alive. /// /// /// Never calls a lifecycle method. Its re-baseline goes straight to /// on the client it already holds — see /// and the remarks for why /// "just re-initialize" would be a silent permanent deadlock. /// /// /// Backoff applies to failures only. A re-baseline is a protocol event, not a /// fault, so it reopens the stream immediately and resets the attempt count; only a /// genuine failure (dropped connection, blown heartbeat, a re-baseline that itself /// failed) advances . /// /// private async Task RunSampleStreamAsync( IMTConnectAgentClient client, MTConnectDriverOptions options, long from, long instanceId, bool republishOnStart, CancellationToken ct) { var cursor = from; var attempt = 0; _reconnectAttempts = 0; try { if (republishOnStart) { FanOut(SubscribedReferences(), ObservationIndex); } while (!ct.IsCancellationRequested) { if (attempt > 0) { var delay = BackoffFor(attempt, options.Reconnect); if (delay > TimeSpan.Zero) { try { await Task.Delay(delay, ct).ConfigureAwait(false); } catch (OperationCanceledException) { return; } } } var outcome = await ConsumeStreamAsync(client, options, cursor, instanceId, ct).ConfigureAwait(false); cursor = outcome.Cursor; instanceId = outcome.InstanceId; switch (outcome.Verdict) { case StreamVerdict.Cancelled: return; case StreamVerdict.Unsupported: // Configuration, not weather: reconnecting reproduces the identical answer // forever, so retrying would be an unbounded hot loop against an endpoint // that will never stream until an operator changes something. Latch it, say // so loudly, and stop. A re-initialize clears the latch. lock (_subscriptionLock) { _streamUnsupported = true; } _streamDegraded = true; Degrade(outcome.Error!); _logger.LogError( outcome.Error, "MTConnect driver {DriverInstanceId} cannot stream /sample from {AgentUri}: the endpoint did not answer with a framed multipart stream. Subscriptions stay registered but will receive no updates until the driver is re-initialized; check for a proxy or load balancer in front of the Agent.", _driverInstanceId, options.AgentUri); return; case StreamVerdict.Reopen: attempt = 0; _reconnectAttempts = 0; break; default: _streamDegraded = true; Degrade(outcome.Error!); // A stream that DELIVERED before it dropped proves the endpoint works, so // its failure starts a fresh backoff ladder. Without this the counter is // monotonic for the life of the process: a healthy agent that drops a // connection once an hour would be pinned at MaxBackoffMs within a day, and // "the first retry is immediate" would be true exactly once ever. attempt = outcome.DeliveredAny ? 1 : attempt + 1; _reconnectAttempts = attempt; _logger.LogWarning( outcome.Error, "MTConnect driver {DriverInstanceId} lost the /sample stream from {AgentUri}; reconnecting from sequence {Cursor} (attempt {Attempt}, delivered={Delivered}).", _driverInstanceId, options.AgentUri, cursor, attempt, outcome.DeliveredAny); break; } } } catch (Exception ex) { _logger.LogError( ex, "MTConnect driver {DriverInstanceId} stopped its /sample pump on an unexpected fault; subscriptions will receive no further updates until the driver is re-initialized.", _driverInstanceId); _streamDegraded = true; Degrade(ex); } finally { // The session's cursor is the opening `from` of the NEXT pump, and until now it only // ever held the priming /current's. A last-unsubscribe + resubscribe (no re-initialize) // would therefore reopen at a sequence the Agent has long since evicted: it self-heals // via OUT_OF_RANGE, but only after a wasted round trip, and a still-buffered stale // cursor replays historical observations to subscribers first. One publish per pump // lifetime, dropped if a lifecycle change has already installed a newer session. PublishAgentCursor(instanceId, cursor); } } /// /// Enumerates one /sample connection until it ends, and classifies how it ended. /// Re-baselines are performed here, inside the enumeration, so the index and the /// subscribers are up to date before the stream is reopened. /// private async Task ConsumeStreamAsync( IMTConnectAgentClient client, MTConnectDriverOptions options, long cursor, long instanceId, CancellationToken ct) { // The RUNNING cursor the gap check is made against: the previous chunk's nextSequence, equal // to the opening `from` only for the first chunk. Comparing every chunk against the opening // `from` instead reports a gap on every chunk once the Agent's buffer rolls past it — an // endless /current re-baseline storm against a perfectly healthy stream. var expected = cursor; // Whether THIS connection ever handed us a chunk. A stream that delivered and then dropped // is a working endpoint having a bad moment; one that never delivered may be an endpoint // that cannot work at all. They must not share a backoff ladder — see RunSampleStreamAsync. var delivered = false; try { await foreach (var chunk in client.SampleAsync(cursor, ct).ConfigureAwait(false)) { if (chunk is null) { continue; } // CHECKED BEFORE THE GAP: an Agent restart changes instanceId *and* resets // sequences, so a restart usually trips IsSequenceGap too — and the two need // different handling. A gap keeps the held values (they are still this device's); // a restart invalidates every one of them, because the device model they describe // no longer exists. Testing the gap first would misdiagnose the restart and keep // serving values the new Agent may never report again. if (chunk.InstanceId != instanceId) { _logger.LogWarning( "MTConnect driver {DriverInstanceId} saw {AgentUri} restart (instanceId {Previous} -> {Current}); dropping every held observation and re-baselining.", _driverInstanceId, options.AgentUri, instanceId, chunk.InstanceId); // Dropped BEFORE the re-baseline, and the subscribers are told: if the /current // then fails, holding stale values from a dead device model would be worse than // holding none. ObservationIndex.Clear(); FanOut(SubscribedReferences(), ObservationIndex); // The driver's own state is consistent by this point, so it is safe to hand the // restart to arbitrary caller code — which may call straight back into // DiscoverAsync. AnnounceAgentRestart(chunk.InstanceId, options); return await RebaselineAsync(client, options, expected, instanceId, delivered, ct) .ConfigureAwait(false); } if (IMTConnectAgentClient.IsSequenceGap(expected, chunk)) { _logger.LogWarning( "MTConnect driver {DriverInstanceId} fell out of {AgentUri}'s buffer (expected sequence {Expected}, oldest retained {FirstSequence}); re-baselining from /current.", _driverInstanceId, options.AgentUri, expected, chunk.FirstSequence); return await RebaselineAsync(client, options, expected, instanceId, delivered, ct) .ConfigureAwait(false); } ApplyAndFanOut(chunk); delivered = true; // EVERY chunk advances the cursor, including an observation-free heartbeat: the // Agent sends those precisely so a quiet connection can be told from a dead one, and // they carry the sequence forward. Not advancing on one makes the NEXT chunk look // like a gap. expected = chunk.NextSequence; _streamDegraded = false; PublishHealthy(DateTime.UtcNow); } // The seam contracts that cancellation is the ONLY way this enumeration ends quietly, so // a quiet end under a live token is a non-conforming client — reported as the lost // stream it is rather than treated as "the subscription is simply idle" (#485). return ct.IsCancellationRequested ? new StreamOutcome(StreamVerdict.Cancelled, expected, instanceId, delivered, null) : new StreamOutcome( StreamVerdict.Transient, expected, instanceId, delivered, new MTConnectStreamEndedException( MTConnectStreamEndReason.ConnectionClosed, options.AgentUri, 0)); } catch (OperationCanceledException) when (ct.IsCancellationRequested) { return new StreamOutcome(StreamVerdict.Cancelled, expected, instanceId, delivered, null); } catch (MTConnectStreamNotSupportedException ex) { return new StreamOutcome(StreamVerdict.Unsupported, expected, instanceId, delivered, ex); } catch (InvalidDataException ex) { // THE OTHER HALF OF RING-BUFFER OVERFLOW, and the one IsSequenceGap cannot see: a real // cppagent answers a `from` that has already fallen out of its buffer with an // MTConnectError / OUT_OF_RANGE document served under HTTP 200, which the client // surfaces here rather than as a gap-bearing chunk. Re-baselining on the gap alone would // mean the driver recovers from falling a little behind but not from falling a lot. _logger.LogWarning( ex, "MTConnect driver {DriverInstanceId} could not resume {AgentUri}'s /sample stream at sequence {Cursor} (the Agent rejected it as out of range); re-baselining from /current.", _driverInstanceId, options.AgentUri, cursor); return await RebaselineAsync(client, options, expected, instanceId, delivered, ct) .ConfigureAwait(false); } catch (Exception ex) { // Transient by default: a dropped connection, a closing boundary, the heartbeat // watchdog's TimeoutException. Reconnect under backoff. return new StreamOutcome(StreamVerdict.Transient, expected, instanceId, delivered, ex); } } /// /// Re-primes the observation index from a fresh /current and reports where the stream /// should resume. The recovery path for every way the incremental sequence can be broken: /// a ring-buffer overflow (detected as a gap, or refused outright as OUT_OF_RANGE) /// and an Agent restart. /// /// /// This calls the client directly and takes no lock. Re-baselining by calling /// — the obvious way to write it, since that method already /// does exactly this work — would take the non-reentrant semaphore /// from inside a pump that a lifecycle method may already be waiting on, and hang the driver /// permanently with no exception and no log line. /// private async Task RebaselineAsync( IMTConnectAgentClient client, MTConnectDriverOptions options, long fallbackCursor, long fallbackInstanceId, bool deliveredAny, CancellationToken ct) { try { var current = await BoundedAsync(client.CurrentAsync, "/current", options.RequestTimeoutMs, ct) .ConfigureAwait(false); PublishAgentCursor(current.InstanceId, current.NextSequence); ApplyAndFanOut(current); _streamDegraded = false; PublishHealthy(DateTime.UtcNow); _logger.LogInformation( "MTConnect driver {DriverInstanceId} re-baselined from {AgentUri}'s /current ({ObservationCount} observation(s), agent instanceId {InstanceId}) and resumes /sample from sequence {NextSequence}.", _driverInstanceId, options.AgentUri, current.Observations.Count, current.InstanceId, current.NextSequence); return new StreamOutcome( StreamVerdict.Reopen, current.NextSequence, current.InstanceId, deliveredAny, null); } catch (OperationCanceledException) when (ct.IsCancellationRequested) { return new StreamOutcome( StreamVerdict.Cancelled, fallbackCursor, fallbackInstanceId, deliveredAny, null); } catch (Exception ex) { // The Agent is unreachable or unusable right now. Treated as an ordinary transient // failure so the reconnect backoff applies — without it, an Agent that is down would be // re-baselined against as fast as the loop can spin. The cursor is deliberately left // where it was: reopening there either works or lands right back here. return new StreamOutcome( StreamVerdict.Transient, fallbackCursor, fallbackInstanceId, deliveredAny, ex); } } /// /// The reconnect delay before attempt , as a pure function — the /// backoff policy is asserted directly rather than inferred from how long a test slept. /// /// /// The first retry honours (0 by /// default = immediate, which is what an operator wants after a one-off blip); every later /// one multiplies, from a floor, up to /// . Operator-authored nonsense cannot /// un-bound the loop: a multiplier that cannot grow the delay falls back to /// , a non-positive cap falls back to /// (clamping to it would make EVERY delay zero — a spin, /// not a backoff), and a negative minimum clamps to zero. /// /// The 1-based reconnect attempt about to be made. /// The authored backoff options. /// How long to wait before that attempt. internal static TimeSpan BackoffFor(int attempt, MTConnectReconnectOptions reconnect) { ArgumentNullException.ThrowIfNull(reconnect); // A non-positive cap is "unset", NOT "no cap". Math.Max(0, ...) would clamp every delay to // zero and turn the reconnect loop into a spin against a refused connection. var capMs = reconnect.MaxBackoffMs > 0 ? reconnect.MaxBackoffMs : DefaultMaxBackoffMs; var minMs = Math.Max(0, reconnect.MinBackoffMs); if (attempt <= 1) { return TimeSpan.FromMilliseconds(Math.Min(minMs, capMs)); } var multiplier = reconnect.BackoffMultiplier > 1.0 ? reconnect.BackoffMultiplier : DefaultBackoffMultiplier; var delayMs = (double)Math.Max(minMs, BackoffGrowthFloorMs); for (var step = 2; step <= attempt; step++) { delayMs *= multiplier; if (delayMs >= capMs) { return TimeSpan.FromMilliseconds(capMs); } } return TimeSpan.FromMilliseconds(Math.Min(delayMs, capMs)); } /// /// Folds a /current snapshot or one /sample chunk into the shared index and /// publishes every reference it touched to the handles that subscribe it. /// private void ApplyAndFanOut(MTConnectStreamsResult document) { var index = ObservationIndex; index.Apply(document); if (document.Observations is not { Count: > 0 }) { return; } // De-duplicated in document order: one Agent document may report a data item more than once // (a Condition container carries several simultaneously-active states), and the index has // already reconciled those into a single snapshot. // // The Agent reports DataItem ids; subscribers hold RawPaths. Each touched id is therefore // EXPANDED into the references it is published under — every RawPath bound to it (a DataItem // may legitimately back more than one authored raw tag), plus the id itself for the // dataItemId-keyed authoring surface (the driver CLI). Publishing the id instead of the // RawPath would miss DriverHostActor's NodeId table and drop the value silently. var binding = Volatile.Read(ref _binding); var touched = new List(document.Observations.Count); var seen = new HashSet(StringComparer.Ordinal); foreach (var observation in document.Observations) { if (observation is null || string.IsNullOrWhiteSpace(observation.DataItemId)) { continue; } if (seen.Add(observation.DataItemId)) { binding.AppendReferences(observation.DataItemId, touched); } } FanOut(touched, index); } /// /// Publishes the index's current snapshot of to every handle /// that subscribes them. References nobody subscribed raise nothing: the Agent streams the /// whole device, and republishing all of it would flood the server with values no client /// asked for. /// /// /// are subscriber references — RawPaths under v3 — and /// each is carried through to the callback unchanged; only the index lookup behind it is /// translated to the DataItem id the Agent reports under. That asymmetry is the point: the /// value must arrive under the reference the caller subscribed, or it is routed nowhere. /// private void FanOut(IReadOnlyList references, MTConnectObservationIndex index) { if (references.Count == 0) { return; } SubscriptionState[] states; lock (_subscriptionLock) { if (_subscriptions.Count == 0) { return; } states = [.. _subscriptions.Values]; } var binding = Volatile.Read(ref _binding); // The lock is released before any callback: handlers are caller code. foreach (var reference in references) { DataValueSnapshot? snapshot = null; foreach (var state in states) { if (!state.References.Contains(reference)) { continue; } // Read once per reference, so every handle sees the identical snapshot. snapshot ??= index.Get(binding.ToDataItemId(reference)); RaiseDataChange(state.Handle, reference, snapshot); } } } /// /// Raises one data-change callback to every subscriber, absorbing anything each one /// throws. A consumer bug must not tear down the pump that serves the others — nor rob them /// of the value. /// /// /// /// The invocation list is walked by hand, with the try/catch INSIDE the loop. A /// plain OnDataChange?.Invoke(...) is a multicast call: the first handler that /// throws aborts the rest of the list, so one broken consumer silently starves every /// other subscriber of that observation — and a single try/catch around the whole /// invoke absorbs the exception while leaving the starvation in place, which reads as /// "isolated" in a log and is not. /// /// /// One is shared across the list deliberately: it is /// an immutable record, and allocating per handler on the hot path would buy nothing. /// /// private void RaiseDataChange(MTConnectSampleHandle handle, string reference, DataValueSnapshot snapshot) { var handlers = OnDataChange; if (handlers is null) { return; } var args = new DataChangeEventArgs(handle, reference, snapshot); foreach (var registration in handlers.GetInvocationList()) { try { ((EventHandler)registration).Invoke(this, args); } catch (Exception ex) { LogSubscriberFault(ex, reference, handle); } } } /// /// Reports a faulting data-change subscriber once per stream generation at Warning, /// and at Debug thereafter. /// /// /// A consumer that throws on every value would otherwise emit one Warning per observation /// per handle — on a busy Agent, thousands a second, which buries the incident it is /// reporting along with everything else in the log. Nothing is lost: the later faults are /// still emitted, at Debug, and the latch is cleared whenever the pump (re)starts. /// private void LogSubscriberFault(Exception ex, string reference, MTConnectSampleHandle handle) { if (Interlocked.Exchange(ref _subscriberFaultLogged, 1) == 0) { _logger.LogWarning( ex, "MTConnect driver {DriverInstanceId} caught a fault from a data-change subscriber for {Reference} on {SubscriptionId}; the stream continues and every other subscriber still received the value. Further subscriber faults are logged at Debug until the stream restarts.", _driverInstanceId, reference, handle.DiagnosticId); return; } _logger.LogDebug( ex, "MTConnect driver {DriverInstanceId} caught a further fault from a data-change subscriber for {Reference} on {SubscriptionId}.", _driverInstanceId, reference, handle.DiagnosticId); } /// Every reference any live subscription is interested in. private List SubscribedReferences() { lock (_subscriptionLock) { var all = new HashSet(StringComparer.Ordinal); foreach (var state in _subscriptions.Values) { all.UnionWith(state.References); } return [.. all]; } } /// /// Whether any live subscription actually names a reference. Caller must hold /// . A subscription over an empty reference list is legal /// (a caller that adds tags later) and must not open a stream by itself. /// private bool HasSubscribedReferencesCore() { foreach (var state in _subscriptions.Values) { if (state.References.Count > 0) { return true; } } return false; } /// /// Records where the Agent now is — its instanceId and the sequence a stream /// should resume from — onto whichever session is installed. /// /// /// /// Both fields move together, in one compare-and-swap. They are one fact about /// one Agent process, and the session's own contract is that a reader gets the client /// and its cursor as a consistent pair. Publishing the id alone (as this did) left the /// session holding a NEW instanceId beside the cursor of the process that had just /// died — the exact tear the packing exists to prevent. /// /// /// There is no "the id did not change, so skip" early return, because the two /// re-baselines that do NOT change the id — a sequence gap and an OUT_OF_RANGE /// rejection — are precisely the ones whose whole purpose is to move the cursor. /// Suppressing the write there would leave the next pump reopening at the stale /// sequence that had just been rejected. /// /// /// Compare-and-swap for the same reason uses one: /// a lifecycle change landing mid-update has newer state, so this write is dropped /// rather than replayed over it. /// /// /// The Agent instanceId last observed. /// The sequence a /sample stream should resume from. private void PublishAgentCursor(long instanceId, long nextSequence) { while (true) { var session = Volatile.Read(ref _session); if (session is null || (session.InstanceId == instanceId && session.NextSequence == nextSequence)) { return; } var updated = session with { InstanceId = instanceId, NextSequence = nextSequence }; if (ReferenceEquals(Interlocked.CompareExchange(ref _session, updated, session), session)) { return; } } } /// /// The /probe device model: the cached copy when warm, otherwise a fresh /// deadline-bounded /probe that re-populates the cache. This is the accessor every /// model consumer (Task 12's DiscoverAsync) must use, so that /// can never leave browse permanently disabled. /// /// /// Lock-free by design — see the remarks for why this does not take /// . It reads the session once, fetches against that /// session's client, and re-caches only if that same session is still installed, so a /// concurrent / can neither be /// undone by a late re-cache nor leak a model belonging to a client that is gone. A teardown /// that lands mid-fetch surfaces as the below rather /// than as a raw from the disposed client. /// /// Cancellation token. /// The Agent's device model. /// /// The driver is not initialized, or it was torn down while this fetch was in flight. /// internal async Task GetOrFetchProbeModelAsync(CancellationToken cancellationToken) { var session = Volatile.Read(ref _session) ?? throw NotConnected(); if (session.ProbeModel is not null) { return session.ProbeModel; } MTConnectProbeModel model; try { model = await BoundedAsync(session.Client.ProbeAsync, "/probe", _options.RequestTimeoutMs, cancellationToken) .ConfigureAwait(false); } catch (ObjectDisposedException ex) { // The lock-free window the _session remarks name: a teardown disposed this client after // we captured it. Reported as the not-connected error the caller already handles, so // browse has exactly one failure mode to reason about. throw NotConnected(ex); } // Publish only onto the session we actually fetched against. var recached = session with { ProbeModel = model, ProbeModelDataItemCount = CountDataItems(model) }; _ = Interlocked.CompareExchange(ref _session, recached, session); return model; } private InvalidOperationException NotConnected(Exception? inner = null) => new($"MTConnect driver '{_driverInstanceId}' has no live agent client; the /probe device model cannot be fetched unless InitializeAsync has succeeded and the driver has not since been shut down or re-initialized.", inner); /// /// Builds driver options from a DriverConfig JSON document. /// /// /// Lives here rather than on the factory because the driver itself needs it: the runtime /// delivers config changes as JSON to a live instance (see the class remarks). Task 15's /// MTConnectDriverFactoryExtensions.CreateInstance should delegate to this method /// rather than deserializing a second DTO — two parsers for one document drift. /// /// Driver instance id, used only in error text. /// The driver's DriverConfig JSON. /// The parsed options. /// The document is unusable or omits agentUri. internal static MTConnectDriverOptions ParseOptions(string driverInstanceId, string driverConfigJson) { ArgumentException.ThrowIfNullOrWhiteSpace(driverInstanceId); ArgumentException.ThrowIfNullOrWhiteSpace(driverConfigJson); MTConnectDriverConfigDto? dto; try { dto = JsonSerializer.Deserialize(driverConfigJson, JsonOptions); } catch (JsonException ex) { throw new InvalidOperationException( $"MTConnect driver config for '{driverInstanceId}' is not valid JSON: {ex.Message}", ex); } if (dto is null) { throw new InvalidOperationException( $"MTConnect driver config for '{driverInstanceId}' deserialised to null"); } if (string.IsNullOrWhiteSpace(dto.AgentUri)) { throw new InvalidOperationException( $"MTConnect driver config for '{driverInstanceId}' missing required AgentUri"); } return new MTConnectDriverOptions { AgentUri = dto.AgentUri.Trim(), DeviceName = string.IsNullOrWhiteSpace(dto.DeviceName) ? null : dto.DeviceName.Trim(), RequestTimeoutMs = dto.RequestTimeoutMs ?? 5_000, SampleIntervalMs = dto.SampleIntervalMs ?? 1_000, SampleCount = dto.SampleCount ?? 1_000, HeartbeatMs = dto.HeartbeatMs ?? 10_000, Tags = dto.Tags is { Count: > 0 } ? [.. dto.Tags.Select(t => BuildTag(t, driverInstanceId))] : [], // v3: the deploy artifact injects the driver's authored raw tags here (see // DriverDeviceConfigMerger). Bound verbatim — the per-entry TagConfig blob is mapped // later, at binding-build time, where a bad blob can be skipped rather than failing the // whole deployment over one tag. RawTags = dto.RawTags is { Count: > 0 } ? [.. dto.RawTags] : [], Probe = new MTConnectProbeOptions { Enabled = dto.Probe?.Enabled ?? true, Interval = TimeSpan.FromMilliseconds(dto.Probe?.IntervalMs ?? 5_000), Timeout = TimeSpan.FromMilliseconds(dto.Probe?.TimeoutMs ?? 2_000), }, Reconnect = new MTConnectReconnectOptions { MinBackoffMs = dto.Reconnect?.MinBackoffMs ?? 0, MaxBackoffMs = dto.Reconnect?.MaxBackoffMs ?? 30_000, BackoffMultiplier = dto.Reconnect?.BackoffMultiplier ?? 2.0, }, }; } // ---- lifecycle internals (all called with _lifecycle held) ---- /// /// Brings the driver up against , reusing /// when the caller determined the transport is unchanged. /// Nothing is published until BOTH agent calls have succeeded, so a failure can never leave a /// half-initialized instance behind. /// private async Task StartCoreAsync( MTConnectDriverOptions options, IMTConnectAgentClient? existingClient, CancellationToken ct) { var lastRead = ReadHealth().LastSuccessfulRead; WriteHealth(new DriverHealth(DriverState.Initializing, lastRead, null)); IMTConnectAgentClient? built = null; try { // arch-review 01/S-6: a 0 authored here must never come to mean "wait forever". Checked // in the driver, not only in the production client's constructor, so the guard holds for // any client implementation behind the seam. RequirePositive(options.RequestTimeoutMs, nameof(MTConnectDriverOptions.RequestTimeoutMs)); RequirePositive(options.HeartbeatMs, nameof(MTConnectDriverOptions.HeartbeatMs)); RequirePositive(options.SampleIntervalMs, nameof(MTConnectDriverOptions.SampleIntervalMs)); RequirePositive(options.SampleCount, nameof(MTConnectDriverOptions.SampleCount)); // Same 01/S-6 guard for Task 13's two knobs, gated on the feature that reads them. A 0 // interval turns the "periodic" probe into an unthrottled hot loop against the Agent, and // a 0 (or negative) timeout cancels every probe before it can be answered — pinning the // host to Stopped forever and reporting an outage that does not exist. Checked only when // probing is ON because a disabled probe reads neither value, and faulting a deployment // over a field the driver never touches would buy nothing; flipping Enabled back on is a // config edit, and ReinitializeAsync runs this same validation. if (options.Probe.Enabled) { RequirePositive(options.Probe.Interval, $"{nameof(MTConnectDriverOptions.Probe)}.{nameof(MTConnectProbeOptions.Interval)}"); RequirePositive(options.Probe.Timeout, $"{nameof(MTConnectDriverOptions.Probe)}.{nameof(MTConnectProbeOptions.Timeout)}"); } var client = existingClient; if (client is null) { built = _agentClientFactory(options) ?? throw new InvalidOperationException( $"The MTConnect agent-client factory for driver '{_driverInstanceId}' returned null."); client = built; } var probe = await BoundedAsync(client.ProbeAsync, "/probe", options.RequestTimeoutMs, ct) .ConfigureAwait(false); var current = await BoundedAsync(client.CurrentAsync, "/current", options.RequestTimeoutMs, ct) .ConfigureAwait(false); // ---- commit point: everything below is non-throwing bookkeeping ---- _options = options; // Derived from the options AND the /probe model just fetched: a raw tag whose TagConfig // blob declares no type takes the type the Agent itself declares for that DataItem, so a // browse-committed tag (whose blob carries only the address) publishes as the number it // is rather than as text. Published before the index, which is built from it. var binding = BuildBinding(options, probe); Volatile.Write(ref _binding, binding); built = null; // The Agent this session talks to, as of now. A restart is measured against THIS, so a // re-initialize against a different Agent process never re-announces the change the // operator's own action caused. Volatile.Write(ref _announcedInstanceId, current.InstanceId); // One publish, so a concurrent lock-free reader sees the new client, the new probe // cache, and the pump's opening cursor together or none of them. Volatile.Write( ref _session, new AgentSession(client, probe, CountDataItems(probe), current.InstanceId, current.NextSequence)); var index = new MTConnectObservationIndex(binding.IndexTags); index.Apply(current); Volatile.Write(ref _index, index); // A fresh session has no failures behind it — including the /sample verdict, because a // re-initialize is precisely the event that means an operator changed something. _streamDegraded = false; _readDegraded = false; WriteHealth(new DriverHealth(DriverState.Healthy, DateTime.UtcNow, null)); // Standing subscriptions survive a lifecycle restart; everything they hold came from the // baseline just thrown away, so the restarted pump republishes as its first act. It is // the PUMP that republishes rather than this method, deliberately: raising subscriber // callbacks from here would run caller code on the thread holding the non-reentrant // lifecycle semaphore, and a handler that touched the driver's lifecycle would deadlock. lock (_subscriptionLock) { _streamUnsupported = false; StartSampleStreamCore(republishOnStart: true); } StartProbeLoopCore(client, options); _logger.LogInformation( "MTConnect driver {DriverInstanceId} initialized against {AgentUri}{DeviceScope}: {DeviceCount} device(s), {ObservationCount} primed observation(s), agent instanceId {InstanceId}.", _driverInstanceId, options.AgentUri, options.DeviceName is null ? string.Empty : $"/{options.DeviceName}", probe.Devices.Count, index.Count, current.InstanceId); } catch (Exception ex) { // A client this attempt built is unreachable from anywhere else; one the caller handed in // is disposed by the teardown below. Either way nothing survives a failed start. SafeDispose(built); await TeardownCoreAsync().ConfigureAwait(false); WriteHealth(new DriverHealth(DriverState.Faulted, lastRead, ex.Message)); _logger.LogError( ex, "MTConnect driver {DriverInstanceId} failed to initialize against {AgentUri}; the driver is Faulted and serves no values.", _driverInstanceId, options.AgentUri); throw; } } /// /// Releases the agent client and drops every piece of served state. Never throws — it runs /// on the failure path of , where a second exception would mask /// the real one. /// private Task TeardownCoreAsync() { // Retire the whole session in one write, BEFORE disposing: a lock-free reader that still // holds the old reference gets a disposed client (caught + named at each seam), never a // session that is half-torn-down. var session = Volatile.Read(ref _session); Volatile.Write(ref _session, null); SafeDispose(session?.Client); // A fresh index rather than Clear(): the tag table it coerces against belongs to the options // that are being torn down. The binding is left in place — it is a pure function of the // options still installed, and a reference must keep resolving to the tag it names so a read // against a torn-down driver answers BadWaitingForInitialData ("authored, no value") rather // than BadNodeIdUnknown ("no such tag"). Volatile.Write(ref _index, new MTConnectObservationIndex(Volatile.Read(ref _binding).IndexTags)); return Task.CompletedTask; } /// /// Stops the shared /sample long-poll stream and waits for the pump to actually be /// gone. Idempotent, and never throws — it runs on teardown paths where a second exception /// would mask the first. /// /// /// /// It waits, rather than merely signalling. Every caller is about to dispose the /// client the pump is enumerating (or install a new observation index the old pump would /// write into). Returning while the pump is still unwinding would let it publish one more /// value into state its owner has already retired — the exact "torn down but still /// reporting" shape the lifecycle is written to prevent. /// /// /// ⚠️ This runs while is held (from /// , and /// ) — and it awaits the pump. That is only safe because the /// pump never touches : its re-baseline calls /// directly (see /// ) rather than re-entering a lifecycle method, and every /// await it makes observes the pump token cancelled below. Break either property and /// this await is a permanent hang with no exception and no log line. /// /// /// The wait is bounded (, and the caller's token) — the /// one thing outside the guarantee above is a subscriber callback, and an /// handler that never returns would otherwise block teardown /// forever. That is not merely slow: DriverInstanceActor.PostStop blocks an Akka /// dispatcher thread on while this driver holds /// , so an unbounded wait wedges an actor-system thread and /// every later lifecycle call with it. Abandoning the wait is safe because the pump's /// registry slots are cleared BEFORE it: a pump left running cannot be resurrected into /// , and it is already cancelled. /// /// /// The caller's deadline for the teardown wait. private async Task StopSampleStreamAsync(CancellationToken cancellationToken = default) { CancellationTokenSource? cts; Task? task; lock (_subscriptionLock) { cts = _pumpCts; task = _pumpTask; _pumpCts = null; _pumpTask = null; } if (cts is null) { return; } try { await cts.CancelAsync().ConfigureAwait(false); } catch (ObjectDisposedException) { // Raced another stop; the pump is already going away. } var stopped = true; if (task is not null) { try { await task.WaitAsync(TeardownWait, cancellationToken).ConfigureAwait(false); } catch (Exception ex) when (ex is TimeoutException or OperationCanceledException) { stopped = false; _logger.LogError( ex, "MTConnect driver {DriverInstanceId} gave up waiting for its /sample pump to stop after {TeardownWait}. The pump is cancelled and can publish nothing further, but something inside it is not returning — the usual cause is an OnDataChange subscriber that blocks. Teardown continues.", _driverInstanceId, TeardownWait); } catch (Exception ex) { // RunSampleStreamAsync is written not to throw, so this is belt-and-braces — but a // teardown that rethrew would replace whatever failure prompted it. _logger.LogDebug( ex, "MTConnect driver {DriverInstanceId} swallowed a fault from its /sample pump while stopping the stream.", _driverInstanceId); } } if (stopped) { // Disposed ONLY once the pump is provably gone: an abandoned pump still holds this // source's token, and disposing underneath it turns its next linked-token creation into // an ObjectDisposedException. One leaked CTS beats a mystery exception. cts.Dispose(); } // Nothing to be degraded about once the stream is gone — UNLESS the endpoint is latched as // unable to stream at all, in which case there is no prospect of one coming back, and // clearing this would let the next successful read report a green driver over a subscription // that can never deliver (#485). See StartSampleStreamCore. lock (_subscriptionLock) { if (!_streamUnsupported) { _streamDegraded = false; } } } /// /// Resolves the options a lifecycle call should apply, reporting an unreadable document /// consistently for both entry points before any state is destroyed. /// /// /// /// One rule, deliberately shared by and /// : a config document that cannot be read never /// destroys working state; it faults only a driver that had none. /// /// /// With a live session the driver keeps serving its previous, valid configuration and /// health is left alone — rejecting the edit already fails the deployment /// (DriverInstanceActor turns the throw into ApplyResult(false)), and /// downing a working driver over an unreadable document is a far larger blast radius /// than the change that was refused. Publishing while /// still serving values would also be a lie in the other direction. /// /// /// With no live session there is nothing to protect and every tag is already dark, so /// the failure must reach as — /// otherwise the dashboard keeps showing whatever state the driver was in before, which /// is the failure-that-looks-like-success shape this codebase exists to avoid. /// /// private MTConnectDriverOptions ResolveIncomingOptions(string driverConfigJson) { try { return HasConfigBody(driverConfigJson) ? ParseOptions(_driverInstanceId, driverConfigJson) : _options; } catch (Exception ex) { if (Volatile.Read(ref _session) is not null) { _logger.LogError( ex, "MTConnect driver {DriverInstanceId} rejected a new driver config document and kept running its previous configuration.", _driverInstanceId); } else { WriteHealth(new DriverHealth(DriverState.Faulted, ReadHealth().LastSuccessfulRead, ex.Message)); _logger.LogError( ex, "MTConnect driver {DriverInstanceId} could not read its driver config document and has no running configuration to fall back on; the driver is Faulted.", _driverInstanceId); } throw; } } /// /// True when carries a real config body — i.e. a JSON /// object or array with at least one element. The bootstrapper always passes a /// populated document; tests and probe-only callers pass a semantically empty one, which /// keeps the constructor-supplied options. /// /// /// /// Emptiness is decided by parsing, not by string comparison. Matching the /// literals "{}" / "[]" misses every other spelling of the same document — /// "{ }", a pretty-printed "{\n}", "{ /* nothing yet */ }" — and /// each of those would then reach , fail the required-AgentUri /// check, and turn a semantically empty document into a cold-start fault or a rejected /// deployment. That is precisely the asymmetry the keep-the-constructor-options rule /// exists to prevent. /// /// /// A malformed document is emphatically NOT empty and returns true, so /// runs and produces the real, quoted parse error. Treating /// unreadable text as "no config supplied" would silently swallow a corrupt document and /// start the driver on stale options — a failure wearing success's clothes. /// /// private static bool HasConfigBody(string? driverConfigJson) { if (string.IsNullOrWhiteSpace(driverConfigJson)) { return false; } try { using var document = JsonDocument.Parse(driverConfigJson, DocumentOptions); return document.RootElement.ValueKind switch { JsonValueKind.Object => document.RootElement.EnumerateObject().Any(), JsonValueKind.Array => document.RootElement.EnumerateArray().Any(), // A bare `null` document is the JSON spelling of "nothing supplied". JsonValueKind.Null => false, // A scalar root is not a config object at all; let ParseOptions say so. _ => true, }; } catch (JsonException) { return true; } } /// /// Whether the new options change anything reads at /// construction, and therefore cannot be applied to a live client. See /// for why this is wider than the plan's AgentUri/DeviceName. /// private static bool RequiresClientRebuild(MTConnectDriverOptions current, MTConnectDriverOptions incoming) => !string.Equals(current.AgentUri, incoming.AgentUri, StringComparison.Ordinal) || !string.Equals(current.DeviceName, incoming.DeviceName, StringComparison.Ordinal) || current.RequestTimeoutMs != incoming.RequestTimeoutMs || current.HeartbeatMs != incoming.HeartbeatMs || current.SampleIntervalMs != incoming.SampleIntervalMs || current.SampleCount != incoming.SampleCount; /// /// Counts every data item the device model declares. Done once, when the model is cached, /// rather than re-walking the tree on every 30 s poll. /// private static int CountDataItems(MTConnectProbeModel model) => model.Devices.Sum(d => d.AllDataItems().Count()); /// /// Runs one unary agent call under a deadline linked to the caller's token. Belt-and-braces /// over the production client's own per-call deadline: the seam does not oblige an /// implementation to bound itself, and an unbounded initialize would wedge the driver-host /// actor (arch-review R2-01). /// private async Task BoundedAsync( Func> call, string request, int timeoutMs, CancellationToken ct) { using var deadline = CancellationTokenSource.CreateLinkedTokenSource(ct); deadline.CancelAfter(TimeSpan.FromMilliseconds(timeoutMs)); try { return await call(deadline.Token).ConfigureAwait(false); } catch (OperationCanceledException) when (!ct.IsCancellationRequested) { // A bare "A task was canceled" names neither the Agent nor the request that stalled. throw new TimeoutException( $"MTConnect {request} request for driver '{_driverInstanceId}' against '{_options.AgentUri}' did not complete within {timeoutMs} ms."); } } /// /// One Bad-coded snapshot per requested reference — the shape every read failure degrades /// to. Arity is preserved because the caller indexes the result positionally against the /// references it asked for. /// private static DataValueSnapshot[] BadBatch(int count, uint status) { // No source timestamp: nothing was observed. The server timestamp is when we gave up. var serverTimestamp = DateTime.UtcNow; var results = new DataValueSnapshot[count]; for (var i = 0; i < count; i++) { results[i] = new DataValueSnapshot(null, status, null, serverTimestamp); } return results; } /// /// Routes a read failure to the health surface (05/STAB-9's fleet-wide shape): log it and /// degrade, preserving LastSuccessfulRead — the "when did this driver last hear from /// the Agent" diagnostic. Never downgrades , which is a /// stronger, config-level signal that a transient read failure must not paper over. /// private void DegradeAfterReadFailure(Exception ex) { _readDegraded = true; Degrade(ex); _logger.LogWarning( ex, "MTConnect driver {DriverInstanceId} could not read /current from {AgentUri}; the batch is reported Bad.", _driverInstanceId, _options.AgentUri); } /// Publishes without downgrading a Faulted driver. private void Degrade(Exception ex) => Degrade(ex.Message); /// /// Publishes with an explanatory message, without /// downgrading a Faulted driver. /// /// /// Deliberately not , even for the "an operator must fix /// this" cases such as an endpoint that cannot stream. In this codebase Faulted is not /// merely a stronger word: DriverHealthReport maps any Faulted driver to a /// /readyz 503, taking the whole server node out of readiness. A driver whose /// reads are perfectly healthy and whose subscriptions are dead is degraded, not unready — /// de-readying the node over one misconfigured Agent would be a far larger blast radius than /// the fault it reports. Faulted stays reserved for the case that earns it: an initialize /// that failed, after which the driver serves nothing at all. /// private void Degrade(string message) { var current = ReadHealth(); if (current.State != DriverState.Faulted) { WriteHealth(new DriverHealth(DriverState.Degraded, current.LastSuccessfulRead, message)); } } /// /// Publishes unless the other data path is currently /// failing. /// /// /// /// Health precedence, stated deliberately. This driver has two independent Agent /// data paths: /current (the batch read) and /// /sample (the subscription pump). They are different HTTP requests and fail /// separately — a proxy that collapses multipart/x-mixed-replace breaks streaming /// while reads stay perfect; an Agent that stalls composing a whole-device snapshot /// breaks reads while the stream keeps flowing. /// /// /// The naive last-writer-wins publish would let either path paper over the other's /// failure, and "Healthy" while half the surface returns Bad is precisely the /// failure-wearing-success shape this codebase exists to avoid. So each path owns a /// flag, clears only its own, and Healthy requires both clear: the driver reports /// Degraded while anything is failing and recovers as soon as the failing path /// itself recovers. — a config-level verdict — is /// never downgraded by either. /// /// /// When the Agent was last successfully heard from. private void PublishHealthy(DateTime at) { var current = ReadHealth(); if (_streamDegraded || _readDegraded) { if (current.State != DriverState.Faulted) { WriteHealth(new DriverHealth(DriverState.Degraded, at, current.LastError)); } return; } WriteHealth(new DriverHealth(DriverState.Healthy, at, null)); } /// /// Disposes an agent client without letting a disposal fault escape. Teardown runs on the /// failure path of , where a second exception would replace the /// real one — but the swallowed exception is still traced rather than vanishing: a /// client that cannot be released is how a socket-handle leak starts, and an empty catch is /// the only place in this driver where a failure would leave no evidence at all. Debug /// level, because it is diagnostic detail about an already-reported failure. /// private void SafeDispose(IMTConnectAgentClient? client) { if (client is null) { return; } try { client.Dispose(); } catch (Exception ex) { _logger.LogDebug( ex, "MTConnect driver {DriverInstanceId} swallowed a fault while disposing its agent client during teardown; the client may not have released its connections.", _driverInstanceId); } } private static void RequirePositive(int value, string optionName) { if (value <= 0) { throw new ArgumentOutOfRangeException( nameof(MTConnectDriverOptions), value, $"{optionName} must be positive; a non-positive value would either un-bound a deadline or ask the Agent for nothing."); } } /// /// The half of the 01/S-6 guard, for the connectivity probe's two /// knobs. A non-positive interval is an unthrottled loop; a non-positive timeout cancels /// every request before it can be answered. /// private static void RequirePositive(TimeSpan value, string optionName) { if (value <= TimeSpan.Zero) { throw new ArgumentOutOfRangeException( nameof(MTConnectDriverOptions), value, $"{optionName} must be positive; a non-positive value would either spin the connectivity probe against the Agent or cancel every probe before it could be answered."); } } /// Barrier-protected read of the cross-thread health snapshot. private DriverHealth ReadHealth() => Volatile.Read(ref _health); /// Barrier-protected publish of a new health snapshot. private void WriteHealth(DriverHealth value) => Volatile.Write(ref _health, value); private static MTConnectTagDefinition BuildTag(MTConnectTagDto dto, string driverInstanceId) { if (string.IsNullOrWhiteSpace(dto.FullName)) { throw new InvalidOperationException( $"MTConnect driver config for '{driverInstanceId}' has a tag with no FullName (the Agent DataItem id it binds)."); } return new MTConnectTagDefinition( dto.FullName.Trim(), // Absent means "no declared type": String is the one coercion that cannot fail, and it // carries the Agent's text through untouched. dto.DriverDataType is null ? DriverDataType.String : ParseEnum(dto.DriverDataType, dto.FullName, driverInstanceId, nameof(MTConnectTagDefinition.DriverDataType)), dto.IsArray ?? false, dto.ArrayDim ?? 0, dto.MtCategory, dto.MtType, dto.MtSubType, dto.Units); } // ---- v3 tag binding: RawPath ↔ dataItemId, and the definitions the index coerces against ---- /// /// Derives the driver's tag tables from an options document — the RawPath → dataItemId /// translation the data plane runs through, its reverse, and the definition set the /// observation index coerces against. /// /// /// /// A bad raw tag is skipped, never thrown. One unreadable TagConfig blob /// must not fail a whole driver — and therefore a whole deployment — over a tag the rest /// of the plant does not depend on. The skipped tag's RawPath then resolves to nothing /// and reports BadNodeIdUnknown, which is the documented /// posture ("a miss is a miss") and is /// visible to an operator, unlike a silent default. /// /// /// Coercion type precedence (see ): /// the blob's own driverDataType, else a /// entry naming the same DataItem id, else the Agent's /probe declaration for that /// data item, else . The third step is what makes a /// browse-committed tag work: RawBrowseCommitMapper writes only the address into /// the blob, so without it every value from the address picker would publish as text /// under a numerically-typed OPC UA node. /// /// /// The options document to derive from. /// The Agent's device model when one has been fetched; else null. private TagBinding BuildBinding(MTConnectDriverOptions options, MTConnectProbeModel? probeModel) { // The dataItemId-keyed authoring surface first (the driver CLI, and every pre-v3 caller). // Last-wins on a duplicate id, matching MTConnectObservationIndex's own documented rule. var definitions = new Dictionary(StringComparer.Ordinal); foreach (var tag in options.Tags) { if (tag is null || string.IsNullOrWhiteSpace(tag.FullName)) continue; definitions[tag.FullName] = tag; } if (options.RawTags.Count == 0) { return new TagBinding( FrozenDictionary.Empty, FrozenDictionary>.Empty, [.. definitions.Values]); } var declared = IndexProbeDataItems(probeModel); var idByRawPath = new Dictionary(options.RawTags.Count, StringComparer.Ordinal); var rawPathsById = new Dictionary>(StringComparer.Ordinal); foreach (var entry in options.RawTags) { if (entry is null || string.IsNullOrWhiteSpace(entry.RawPath)) continue; if (!TryReadRawTag(entry, out var dataItemId, out var authored)) { _logger.LogWarning( "MTConnect driver {DriverInstanceId} could not map the raw tag {RawPath} to an Agent DataItem; it names no readable id (expected 'fullName', 'dataItemId' or 'address') or declares an unrecognised driverDataType. The tag is skipped and reports BadNodeIdUnknown; every other tag is unaffected.", _driverInstanceId, entry.RawPath); continue; } idByRawPath[entry.RawPath] = dataItemId; if (!rawPathsById.TryGetValue(dataItemId, out var paths)) { rawPathsById[dataItemId] = paths = []; } if (!paths.Contains(entry.RawPath, StringComparer.Ordinal)) { // One DataItem may legitimately back several authored raw tags; a duplicate RawPath // (the same tag listed twice) must still publish once. paths.Add(entry.RawPath); } if (authored is not null) { // The blob declared a type: it is the deploy-time authority and supersedes an older // Tags entry for the same DataItem. definitions[dataItemId] = authored; } else if (!definitions.ContainsKey(dataItemId)) { definitions[dataItemId] = InferDefinition(dataItemId, declared); } } return new TagBinding( idByRawPath.ToFrozenDictionary(StringComparer.Ordinal), rawPathsById.ToFrozenDictionary( pair => pair.Key, pair => (IReadOnlyList)pair.Value, StringComparer.Ordinal), [.. definitions.Values]); } /// /// Reads one authored raw tag: the DataItem id it binds, and the definition its blob declares /// (null when the blob declares no type and the caller must fall back). /// /// /// when the blob is unreadable, names no DataItem id, or carries a /// driverDataType that is not a — a present-but-invalid /// type rejects the tag rather than silently defaulting, which would publish the Agent's text /// as a Good value of the wrong type (the sibling drivers' strict-enum posture). /// private static bool TryReadRawTag( RawTagEntry entry, out string dataItemId, out MTConnectTagDefinition? authored) { dataItemId = string.Empty; authored = null; if (string.IsNullOrWhiteSpace(entry.TagConfig)) return false; MTConnectRawTagConfigDto? blob; try { blob = JsonSerializer.Deserialize(entry.TagConfig, JsonOptions); } catch (JsonException) { return false; } if (blob is null) return false; var id = FirstNonBlank(blob.FullName, blob.DataItemId, blob.Address); if (id is null) return false; dataItemId = id; // "driverDataType" (the driver's own tags[] field) or "dataType" (what the AdminUI's typed // MTConnectTagConfigModel writes) — see the DTO remarks. var authoredType = FirstNonBlank(blob.DriverDataType, blob.DataType); if (authoredType is null) return true; // A digits-only type is rejected outright, BEFORE the parse: Enum.TryParse accepts "8" as // readily as "Float64" and returns the member with that ordinal, so a numerically-serialized // enum (the repo-wide authoring trap) would otherwise coerce every value of that tag to a // type nobody chose, under Good quality. The AdminUI's MTConnectTagConfigModel.Validate // refuses the same shape; this is the runtime half of that guard. if (authoredType.All(char.IsAsciiDigit) || !Enum.TryParse(authoredType, ignoreCase: true, out var parsed) || !Enum.IsDefined(parsed)) { return false; } authored = new MTConnectTagDefinition( dataItemId, parsed, blob.IsArray ?? false, blob.ArrayDim ?? 0, blob.MtCategory, blob.MtType, blob.MtSubType, blob.Units); return true; } /// /// The definition for a raw tag whose blob declared no type: the Agent's own /probe /// declaration when the model is in hand, else — the one /// coercion that cannot fail, carrying the Agent's text through untouched. /// private static MTConnectTagDefinition InferDefinition( string dataItemId, FrozenDictionary declared) { if (!declared.TryGetValue(dataItemId, out var dataItem)) { return new MTConnectTagDefinition(dataItemId, DriverDataType.String); } var inferred = MTConnectDataTypeInference.Infer( dataItem.Category, dataItem.Type, dataItem.Units, dataItem.Representation, dataItem.SampleCount); return new MTConnectTagDefinition( dataItemId, inferred.DataType, inferred.IsArray, (int)(inferred.ArrayDim ?? 0), dataItem.Category, dataItem.Type, dataItem.SubType, dataItem.Units); } /// Flattens a device model into a dataItemId → DataItem lookup; null ⇒ empty. private static FrozenDictionary IndexProbeDataItems(MTConnectProbeModel? probeModel) { if (probeModel is null) return FrozenDictionary.Empty; var map = new Dictionary(StringComparer.Ordinal); foreach (var device in probeModel.Devices) { if (device is null) continue; foreach (var dataItem in device.AllDataItems()) { if (dataItem is null || string.IsNullOrWhiteSpace(dataItem.Id)) continue; map[dataItem.Id] = dataItem; } } return map.ToFrozenDictionary(StringComparer.Ordinal); } /// The first of that carries text, trimmed; else null. private static string? FirstNonBlank(params string?[] candidates) { foreach (var candidate in candidates) { if (!string.IsNullOrWhiteSpace(candidate)) return candidate.Trim(); } return null; } /// /// Parses an enum authored as a name. Config surfaces serialize enums as names /// project-wide; a numerically-serialized enum is the known trap that faults a driver at /// deploy, so this reports the offending field rather than silently taking member 0. /// private static T ParseEnum(string raw, string tagName, string driverInstanceId, string field) where T : struct, Enum => Enum.TryParse(raw, ignoreCase: true, out var parsed) && Enum.IsDefined(parsed) ? parsed : throw new InvalidOperationException(string.Create( CultureInfo.InvariantCulture, $"MTConnect driver config for '{driverInstanceId}' tag '{tagName}' has an unrecognised {field} '{raw}'. Expected one of: {string.Join(", ", Enum.GetNames())}.")); /// /// The live agent client and the /probe cache derived from it, as one immutable /// unit. Immutable so that publishing a change is a single reference swap — see the /// remarks for why the three values must move together. /// /// The agent client this session owns and will dispose on teardown. /// /// The cached device model, or null once has /// dropped it. Never null at the moment a session is first published. /// /// /// Data items declared by , counted once at cache time and /// always 0 when it is null. /// /// /// The Agent's instanceId as last observed. Updated in place (by /// ) when the pump sees the Agent restart, so the /// session always names the Agent process the client is actually talking to. /// /// /// The sequence the /sample pump should open at — the priming /current's /// nextSequence. Lives here so that starting the pump reads the client and its cursor /// as one consistent pair; a torn read would open the stream at another session's sequence. /// private sealed record AgentSession( IMTConnectAgentClient Client, MTConnectProbeModel? ProbeModel, int ProbeModelDataItemCount, long InstanceId, long NextSequence); /// /// The tag tables one options document derives, as one immutable unit — so installing new /// options is a single reference swap and no reader can pair one document's RawPaths with /// another's definitions. /// /// /// The v3 wire translation: the RawPath the server subscribes / reads / routes by, mapped to /// the Agent DataItem id the observation index is keyed by. /// /// /// The reverse, one-to-many: every RawPath a DataItem must be published under. A DataItem may /// legitimately back more than one authored raw tag, and each of those tags is a separate /// OPC UA node that must receive the value. /// /// /// The definition set coerces against — the union of /// the authored and the definitions derived from /// . This is how a raw tag's declared (or /// probe-inferred) DriverDataType reaches the coercion. /// private sealed record TagBinding( FrozenDictionary DataItemIdByRawPath, FrozenDictionary> RawPathsByDataItemId, IReadOnlyList IndexTags) { /// /// Translates a subscriber/read reference to the DataItem id the index is keyed by. /// /// /// An unmapped reference is returned unchanged rather than rejected, which serves /// two callers at once: the dataItemId-keyed authoring surface (the driver CLI, and a /// driver configured with Tags and no RawTags) resolves normally, and a /// RawPath that no authored raw tag claims falls through to an index miss — /// BadNodeIdUnknown, never an exception. There is no third behaviour to add here: /// the index already distinguishes "authored but not yet reported" from "unknown". /// /// The caller's reference (a RawPath under v3). /// The DataItem id to look up. public string ToDataItemId(string? reference) => reference is not null && DataItemIdByRawPath.TryGetValue(reference, out var dataItemId) ? dataItemId : reference ?? string.Empty; /// /// Appends every reference must be published under: each /// RawPath bound to it, plus the id itself for the dataItemId-keyed authoring surface. /// /// The DataItem the Agent just reported. /// The reference list being accumulated for the fan-out. public void AppendReferences(string dataItemId, List into) { if (RawPathsByDataItemId.TryGetValue(dataItemId, out var rawPaths)) { into.AddRange(rawPaths); // Only when a RawPath is spelled identically to the id — otherwise a subscriber // holding that one reference would be raised twice for one observation. if (rawPaths.Contains(dataItemId, StringComparer.Ordinal)) return; } into.Add(dataItemId); } } /// /// One live subscription: its handle and the references it wants. Both are effectively /// immutable once registered, so the pump reads without a lock — /// it takes only to snapshot the registry itself. /// /// The identity handed to the caller and stamped on every event. /// The de-duplicated DataItem ids this subscription publishes. private sealed record SubscriptionState(MTConnectSampleHandle Handle, HashSet References); /// How one /sample connection ended, and therefore what the pump does next. private enum StreamVerdict { /// Cancelled by teardown — stop, publish nothing, report nothing. Cancelled = 0, /// A transient failure: degrade and reconnect under the backoff. Transient = 1, /// Re-baselined successfully: reopen immediately at the new cursor, no backoff. Reopen = 2, /// The endpoint cannot stream at all: latch, report loudly, and stop retrying. Unsupported = 3 } /// How the connection ended. /// The sequence the next connection must open at. /// The Agent instance the cursor belongs to. /// /// Whether this connection yielded at least one chunk before it ended. Distinguishes a /// working endpoint having a bad moment from one that has never worked, which is what keeps /// the reconnect backoff from ratcheting monotonically over a process lifetime. /// /// The failure, when there was one. private sealed record StreamOutcome( StreamVerdict Verdict, long Cursor, long InstanceId, bool DeliveredAny, Exception? Error); // ---- config DTOs ---- /// /// Nullable-init mirror of the driver's DriverConfig JSON. Every property is nullable /// so an omitted key means "use the default" rather than "reset to zero". /// private sealed class MTConnectDriverConfigDto { public string? AgentUri { get; init; } public string? DeviceName { get; init; } public int? RequestTimeoutMs { get; init; } public int? SampleIntervalMs { get; init; } public int? SampleCount { get; init; } public int? HeartbeatMs { get; init; } public List? Tags { get; init; } /// /// The v3 authored raw tags the deploy artifact injects (RawPath + driver TagConfig /// blob). Bound verbatim into ; this is the /// collection a deployed instance's data plane is keyed by. /// public List? RawTags { get; init; } public MTConnectProbeDto? Probe { get; init; } public MTConnectReconnectDto? Reconnect { get; init; } } /// /// A raw tag's TagConfig blob: the driver-specific address of one authored raw tag. /// Every field is optional — the blob is operator-authorable and reaches the driver through /// several authoring surfaces that populate different subsets of it. /// /// /// /// Three accepted spellings of the DataItem id, tried in order: /// fullName (the driver's own tag shape, and what the typed editor writes), /// dataItemId (the protocol's own word for it, and the spelling an operator /// hand-authoring the blob reaches for first), and address (what /// RawBrowseCommitMapper writes for a driver with no typed editor of its own — the /// shape every browse-committed MTConnect tag currently arrives in). Accepting all three /// costs two dictionary probes; rejecting two of them would make a tag committed from the /// address picker silently unreadable. /// /// /// Two accepted spellings of the coercion type, likewise: dataType — what /// the AdminUI's MTConnectTagConfigModel writes, and therefore what every tag /// authored through the typed editor carries — and driverDataType, matching the /// driver's own tags[] field name. Reading only one of them would make the editor /// and the driver disagree about the type of every tag. /// /// /// Kept separate from rather than reusing it: that DTO is /// the strict tags[] surface, where a missing fullName is a config error /// worth failing the deployment over. A raw tag's blob is authored elsewhere and its /// failure mode is per-tag, not per-driver. /// /// private sealed class MTConnectRawTagConfigDto { public string? FullName { get; init; } public string? DataItemId { get; init; } public string? Address { get; init; } public string? DriverDataType { get; init; } public string? DataType { get; init; } public bool? IsArray { get; init; } public int? ArrayDim { get; init; } public string? MtCategory { get; init; } public string? MtType { get; init; } public string? MtSubType { get; init; } public string? Units { get; init; } } /// One authored tag. DriverDataType stays a string — see . private sealed class MTConnectTagDto { public string? FullName { get; init; } public string? DriverDataType { get; init; } public bool? IsArray { get; init; } public int? ArrayDim { get; init; } public string? MtCategory { get; init; } public string? MtType { get; init; } public string? MtSubType { get; init; } public string? Units { get; init; } } /// Background connectivity-probe knobs (Task 13 consumes them). private sealed class MTConnectProbeDto { public bool? Enabled { get; init; } public int? IntervalMs { get; init; } public int? TimeoutMs { get; init; } } /// Reconnect-backoff knobs (Task 11's pump consumes them). private sealed class MTConnectReconnectDto { public int? MinBackoffMs { get; init; } public int? MaxBackoffMs { get; init; } public double? BackoffMultiplier { get; init; } } }