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;
///
/// The DataItem ids the authored raw tags bind. DiscoverAsync emits
/// FullName: dataItem.Id from the Agent's probe model, so the two sets are the same
/// vocabulary.
public IReadOnlyCollection? AuthoredDiscoveryRefs =>
Volatile.Read(ref _binding).RawPathsByDataItemId.Keys;
///
/// False — DiscoverAsync is built purely from the Agent's probe model and never
/// re-emits authored tags, so an authored DataItem that the Agent stopped publishing IS visible as
/// missing. MTConnect is currently the only driver where both drift directions are detectable.
public bool DiscoveryStreamIncludesAuthoredTags => false;
///
///
/// 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; }
}
}