From 9a67ed918a00cc3993f918f20027deda820c742f Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Fri, 24 Jul 2026 16:29:17 -0400 Subject: [PATCH] feat(mtconnect): ISubscribable /sample pump + ring-buffer re-baseline (Task 11) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One shared /sample long poll behind every subscription handle. Subscribe fires initial data per reference from the primed /current (Part 4 convention) and returns without waiting on the stream; the pump runs on its own task under its own token — never the caller's Subscribe token, which is disposed the moment that call returns. The pump owns every way the sequence can break: - sequence gap, measured against the RUNNING cursor (the previous chunk's nextSequence). Comparing against the opening `from` reports a gap on every chunk once the buffer rolls — a /current re-baseline storm. - OUT_OF_RANGE under HTTP 200 (InvalidDataException out of SampleAsync), which IsSequenceGap cannot see at all — without it a driver recovers from falling a little behind but not from falling a lot. - agent restart, checked BEFORE the gap because a restart resets sequences and so usually trips the gap check too; it clears the index (the held values describe a device model that no longer exists) and tells the subscribers. - every chunk advances the cursor, heartbeats included. Re-baseline calls CurrentAsync directly on the client the pump already holds: routing it through ReinitializeAsync would take the non-reentrant lifecycle semaphore from inside a pump a lifecycle method may already be awaiting, and hang the driver with no exception and no log. StopSampleStreamAsync is filled in (same call sites) and InitializeAsync now stops the stream before teardown too, so a re-Initialize on a live instance cannot leave a pump enumerating a disposed client. MTConnectStreamEnded/TimeoutException reconnect under a geometric backoff with a 100 ms growth floor (MinBackoffMs defaults to 0, and 0 x multiplier is still 0). MTConnectStreamNotSupportedException does NOT retry — it latches, reports loudly, and is cleared only by a re-initialize. Health precedence is explicit: /current and /sample fail independently, so each path owns a flag, clears only its own, and Healthy requires both clear. Neither can paper over the other's failure. 29 new tests, all deterministic — no sleeps, no wall-clock assertions. The fake grew per-generation sample streams (so a reconnect has somewhere to land), scripted /current answers, scripted sample failures, and a chunk-consumed signal that also fires when the consumer abandons the stream. 375/375 green. --- .../MTConnectDriver.cs | 967 +++++++++++++++++- .../MTConnectSampleHandle.cs | 29 + .../CannedAgentClient.cs | 165 ++- .../MTConnectSubscribeTests.cs | 827 +++++++++++++++ 4 files changed, 1943 insertions(+), 45 deletions(-) create mode 100644 src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectSampleHandle.cs create mode 100644 tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/MTConnectSubscribeTests.cs diff --git a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectDriver.cs b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectDriver.cs index 8dcc9f68..36214f5a 100644 --- a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectDriver.cs +++ b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectDriver.cs @@ -42,19 +42,19 @@ namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect; /// nextSequence. /// /// -/// Scope. This type currently implements and -/// . ISubscribable (Task 11), ITagDiscovery (Task 12) +/// Scope. This type currently implements , +/// and . ITagDiscovery (Task 12) /// and IHostConnectivityProbe/IRediscoverable (Task 13) are added on top of /// the state this class already caches — the probe model, the Agent instanceId, and /// the observation index. /// /// /// The read path never writes the shared observation index — see -/// . That index has exactly one writer, and Task 11's /sample -/// pump is it. +/// . That index has exactly one writer: the /sample pump +/// (). /// /// -public sealed class MTConnectDriver : IDriver, IReadable +public sealed class MTConnectDriver : IDriver, IReadable, ISubscribable { /// /// Coarse per-entry byte weights behind . The contract asks @@ -85,6 +85,19 @@ public sealed class MTConnectDriver : IDriver, IReadable /// 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; + /// /// Config-JSON reader options, mirroring the sibling driver factories. Note there is /// deliberately no JsonStringEnumConverter: enum-carrying DTO fields stay @@ -145,10 +158,58 @@ public sealed class MTConnectDriver : IDriver, IReadable /// 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 = []; + private MTConnectDriverOptions _options; private MTConnectObservationIndex _index; - private long? _agentInstanceId; 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; /// /// The live agent client plus the /probe cache derived from it, held as one @@ -247,10 +308,28 @@ public sealed class MTConnectDriver : IDriver, IReadable internal MTConnectProbeModel? CachedProbeModel => Volatile.Read(ref _session)?.ProbeModel; /// - /// The Agent's instanceId as of the last successful /current. Task 13 watches - /// this for change to raise rediscovery (an Agent restart invalidates every cached id). + /// 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 => _agentInstanceId; + internal long? AgentInstanceId => Volatile.Read(ref _session)?.InstanceId; + + /// + /// 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 options currently in force — the constructor's, or the last document applied. internal MTConnectDriverOptions EffectiveOptions => _options; @@ -275,7 +354,12 @@ public sealed class MTConnectDriver : IDriver, IReadable // 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. + // 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().ConfigureAwait(false); await TeardownCoreAsync().ConfigureAwait(false); await StartCoreAsync(options, existingClient: null, cancellationToken).ConfigureAwait(false); @@ -504,7 +588,8 @@ public sealed class MTConnectDriver : IDriver, IReadable results[i] = snapshot.Get(fullReferences[i]); } - WriteHealth(new DriverHealth(DriverState.Healthy, DateTime.UtcNow, null)); + _readDegraded = false; + PublishHealthy(DateTime.UtcNow); return results; } @@ -523,6 +608,648 @@ public sealed class MTConnectDriver : IDriver, IReadable } } + // ---- 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). + var index = ObservationIndex; + foreach (var reference in references) + { + RaiseDataChange(handle, reference, index.Get(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. + await StopSampleStreamAsync().ConfigureAwait(false); + } + + _logger.LogDebug( + "MTConnect driver {DriverInstanceId} closed subscription {SubscriptionId}{StreamNote}.", + _driverInstanceId, ours.DiagnosticId, stopStream ? " and stopped the shared /sample stream" : string.Empty); + } + + /// + /// 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) + { + if (_pumpTask is not null || _streamUnsupported || !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; + } + + 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; + + 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; + + break; + + default: + _streamDegraded = true; + Degrade(outcome.Error!); + attempt++; + + _logger.LogWarning( + outcome.Error, + "MTConnect driver {DriverInstanceId} lost the /sample stream from {AgentUri}; reconnecting from sequence {Cursor} (attempt {Attempt}).", + _driverInstanceId, options.AgentUri, cursor, attempt); + + 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); + } + } + + /// + /// 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; + + 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); + + return await RebaselineAsync(client, options, expected, instanceId, 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, ct).ConfigureAwait(false); + } + + ApplyAndFanOut(chunk); + + // 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, null) + : new StreamOutcome( + StreamVerdict.Transient, + expected, + instanceId, + new MTConnectStreamEndedException( + MTConnectStreamEndReason.ConnectionClosed, options.AgentUri, 0)); + } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + return new StreamOutcome(StreamVerdict.Cancelled, expected, instanceId, null); + } + catch (MTConnectStreamNotSupportedException ex) + { + return new StreamOutcome(StreamVerdict.Unsupported, expected, instanceId, 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, 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, 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, + CancellationToken ct) + { + try + { + var current = await BoundedAsync(client.CurrentAsync, "/current", options.RequestTimeoutMs, ct) + .ConfigureAwait(false); + + PublishAgentInstanceId(current.InstanceId); + 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, null); + } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + return new StreamOutcome(StreamVerdict.Cancelled, fallbackCursor, fallbackInstanceId, 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, 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 + /// , and negative values clamp 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); + + var capMs = Math.Max(0, reconnect.MaxBackoffMs); + 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. + 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)) + { + touched.Add(observation.DataItemId); + } + } + + 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. + /// + private void FanOut(IReadOnlyList references, MTConnectObservationIndex index) + { + if (references.Count == 0) + { + return; + } + + SubscriptionState[] states; + lock (_subscriptionLock) + { + if (_subscriptions.Count == 0) + { + return; + } + + states = [.. _subscriptions.Values]; + } + + // 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(reference); + RaiseDataChange(state.Handle, reference, snapshot); + } + } + } + + /// + /// Raises one data-change callback, absorbing anything the subscriber throws. A consumer bug + /// must not tear down the pump that serves every OTHER subscriber — and the exception is + /// logged rather than swallowed silently, because a handler that throws on every value is + /// otherwise invisible. + /// + private void RaiseDataChange(MTConnectSampleHandle handle, string reference, DataValueSnapshot snapshot) + { + try + { + OnDataChange?.Invoke(this, new DataChangeEventArgs(handle, reference, snapshot)); + } + catch (Exception ex) + { + _logger.LogWarning( + ex, + "MTConnect driver {DriverInstanceId} caught a fault from a data-change subscriber for {Reference} on {SubscriptionId}; the stream continues.", + _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 the Agent instanceId a re-baseline observed, onto whichever session is + /// installed. A compare-and-swap for the same reason + /// uses one: if a lifecycle change lands mid-update, + /// its state is newer and this write must be dropped, not replayed over it. + /// + private void PublishAgentInstanceId(long instanceId) + { + while (true) + { + var session = Volatile.Read(ref _session); + if (session is null || session.InstanceId == instanceId) + { + return; + } + + var updated = session with { InstanceId = instanceId }; + 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 @@ -687,17 +1414,34 @@ public sealed class MTConnectDriver : IDriver, IReadable _options = options; built = null; - // One publish, so a concurrent lock-free reader sees the new client and the new probe - // cache together or neither of them. - Volatile.Write(ref _session, new AgentSession(client, probe, CountDataItems(probe))); - _agentInstanceId = 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(options.Tags); 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); + } + _logger.LogInformation( "MTConnect driver {DriverInstanceId} initialized against {AgentUri}{DeviceScope}: {DeviceCount} device(s), {ObservationCount} primed observation(s), agent instanceId {InstanceId}.", _driverInstanceId, @@ -739,8 +1483,6 @@ public sealed class MTConnectDriver : IDriver, IReadable Volatile.Write(ref _session, null); SafeDispose(session?.Client); - _agentInstanceId = null; - // A fresh index rather than Clear(): the tag table it coerces against belongs to the options // that are being torn down. Volatile.Write(ref _index, new MTConnectObservationIndex(_options.Tags)); @@ -749,21 +1491,85 @@ public sealed class MTConnectDriver : IDriver, IReadable } /// - /// Stops the shared /sample long-poll stream. A no-op until Task 11 introduces the - /// pump; it exists now because both and - /// must stop the stream before the client they are about - /// to dispose goes away, and that ordering is easy to lose when the pump is bolted on later. + /// 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. /// /// - /// ⚠️ This runs while is held. Whatever Task 11 fills it with - /// — cancelling the pump's token, awaiting its task — must not call back into - /// , or - /// , directly or by awaiting a pump that is itself blocked on one - /// of them: the semaphore is non-reentrant and the result is a permanent hang with no - /// exception and no log. See the remarks for the re-baseline - /// pattern that avoids this. + /// + /// 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 one thing outside that guarantee is a subscriber callback: an + /// handler that blocks forever blocks teardown with it. The + /// pump raises callbacks outside every lock precisely so that a slow handler only + /// delays; an infinite one is a consumer bug this driver cannot paper over. + /// /// - private static Task StopSampleStreamAsync() => Task.CompletedTask; + private async Task StopSampleStreamAsync() + { + 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. + } + + if (task is not null) + { + try + { + await task.ConfigureAwait(false); + } + 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); + } + } + + cts.Dispose(); + + // There is no stream any more, so there is nothing for it to be degraded about. Leaving the + // flag set would pin a later successful read to Degraded forever. + _streamDegraded = false; + } /// /// Resolves the options a lifecycle call should apply, reporting an unreadable document @@ -943,11 +1749,8 @@ public sealed class MTConnectDriver : IDriver, IReadable /// private void DegradeAfterReadFailure(Exception ex) { - var current = ReadHealth(); - if (current.State != DriverState.Faulted) - { - WriteHealth(new DriverHealth(DriverState.Degraded, current.LastSuccessfulRead, ex.Message)); - } + _readDegraded = true; + Degrade(ex); _logger.LogWarning( ex, @@ -955,6 +1758,57 @@ public sealed class MTConnectDriver : IDriver, IReadable _driverInstanceId, _options.AgentUri); } + /// Publishes without downgrading a Faulted driver. + private void Degrade(Exception ex) + { + var current = ReadHealth(); + if (current.State != DriverState.Faulted) + { + WriteHealth(new DriverHealth(DriverState.Degraded, current.LastSuccessfulRead, ex.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 @@ -1050,10 +1904,53 @@ public sealed class MTConnectDriver : IDriver, IReadable /// 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); + int ProbeModelDataItemCount, + long InstanceId, + long NextSequence); + + /// + /// 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. + /// The failure, when there was one. + private sealed record StreamOutcome(StreamVerdict Verdict, long Cursor, long InstanceId, Exception? Error); // ---- config DTOs ---- diff --git a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectSampleHandle.cs b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectSampleHandle.cs new file mode 100644 index 00000000..ee0119ed --- /dev/null +++ b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect/MTConnectSampleHandle.cs @@ -0,0 +1,29 @@ +using ZB.MOM.WW.OtOpcUa.Core.Abstractions; + +namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect; + +/// +/// Driver-internal identity of one /sample subscription, matching the sibling drivers' +/// shape (Galaxy's GalaxySubscriptionHandle, the shared PollGroupEngine's +/// PollSubscriptionHandle): a monotonic per-driver id plus a diagnostic string that +/// carries it into logs. +/// +/// +/// +/// Every handle shares ONE Agent stream. Unlike a polled driver, where each +/// subscription owns a loop, MTConnect's /sample long poll is per-Agent: the driver +/// opens exactly one and fans each chunk out to whichever handles subscribe the reporting +/// DataItem. The handle is therefore purely an identity — it owns no connection, no task, +/// and no cursor — and dropping one only stops the stream when it was the last. +/// +/// +/// A for value equality on the id, so a handle that has round-tripped +/// through the caller still resolves. The id is never reused within a driver instance. +/// +/// +/// The monotonic per-driver subscription id. +internal sealed record MTConnectSampleHandle(long SubscriptionId) : ISubscriptionHandle +{ + /// + public string DiagnosticId => $"mtconnect-sub-{SubscriptionId}"; +} diff --git a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/CannedAgentClient.cs b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/CannedAgentClient.cs index 71fe5e52..d3c2d20d 100644 --- a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/CannedAgentClient.cs +++ b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/CannedAgentClient.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using System.Runtime.CompilerServices; using System.Threading.Channels; @@ -34,9 +35,19 @@ namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests; /// internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable { - private readonly Channel _chunks = Channel.CreateUnbounded(); private readonly CancellationTokenSource _disposeCts = new(); + /// Chunks a test has scripted but not yet pumped — see . + private readonly ConcurrentQueue _script = new(); + + /// + /// The current /sample stream generation. Replaced (not merely completed) by + /// so that a pump which reconnects after a dropped stream gets a + /// live channel to read from instead of one that is permanently closed — without that, a + /// reconnect test could only ever observe the reconnect failing. + /// + private Channel _chunks = Channel.CreateUnbounded(); + private int _probeCallCount; private int _currentCallCount; private int _sampleCallCount; @@ -66,6 +77,44 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable public static MTConnectStreamsResult Chunk(string fixture) => MTConnectStreamsParser.Parse(File.ReadAllText(fixture)); + /// + /// An Agent scripted for the ring-buffer-overflow story: the driver primes from + /// Fixtures/current.xml (nextSequence 108), then the first scripted chunk is + /// Fixtures/sample-gap.xml — whose firstSequence (5000) is strictly newer than + /// the cursor, i.e. the buffer rolled past the driver — and the second is an ordinary chunk + /// that continues contiguously from the re-baseline. + /// + /// + /// The second /current answer is what makes this a story rather than a single + /// event: it advertises the buffer the gap chunk revealed (firstSequence 5000 / + /// nextSequence 5005), so a driver that re-baselines and resumes from that answer sees + /// the follow-up chunk as contiguous. A driver that re-baselined but resumed from the stale + /// cursor would report a gap on every subsequent chunk instead. + /// + public static CannedAgentClient WithGapThenResume() + { + var client = FromFixtures(); + var gap = Chunk("Fixtures/sample-gap.xml"); + + // 1st /current = the InitializeAsync prime (the fixture). 2nd = the post-gap re-baseline. + client.CurrentAnswers.Enqueue(client.Current); + client.CurrentAnswers.Enqueue( + new MTConnectStreamsResult(gap.InstanceId, gap.NextSequence, gap.FirstSequence, gap.Observations)); + + client.ScriptChunks( + gap, + new MTConnectStreamsResult( + gap.InstanceId, + gap.NextSequence + 5, + gap.NextSequence, + [ + new MTConnectObservation( + "dev1_pos", "201.5000", new DateTime(2026, 7, 24, 12, 6, 0, DateTimeKind.Utc)), + ])); + + return client; + } + /// The document /probe answers with. Settable so a test can re-shape the model. public MTConnectProbeModel Probe { get; set; } @@ -81,6 +130,31 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable /// When set, /current throws this instead of answering. public Exception? CurrentFailure { get; set; } + /// + /// Scripted /current answers, consumed in order: each call dequeues the next one and + /// promotes it to , so once the script runs out the Agent keeps + /// answering with the last thing it said rather than travelling back in time. This is how a + /// test scripts "the prime, then the post-gap re-baseline" without racing the pump for a + /// property setter. + /// + public ConcurrentQueue CurrentAnswers { get; } = new(); + + /// + /// Scripted one-shot /sample failures: each enumeration dequeues one and throws it. + /// Models a transient stream fault the Agent recovers from — e.g. the + /// a real cppagent's OUT_OF_RANGE error document + /// (served under HTTP 200) surfaces as. + /// + public ConcurrentQueue SampleFailures { get; } = new(); + + /// + /// A sticky /sample failure: every enumeration throws it, forever. Models a + /// configuration-level fault such as + /// , where reconnecting reproduces the + /// identical answer — the shape a pump must NOT retry-loop against. + /// + public Exception? SampleFailure { get; set; } + /// /// When set, the next /probe parks here until the test completes it, then the /// gate clears itself so later calls answer immediately. @@ -101,6 +175,14 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable /// public TaskCompletionSource? CurrentGate { get; set; } + /// + /// Signalled the moment a /current request lands — before + /// parks it. This is the deterministic "the caller is now inside + /// the request" barrier a test needs before landing a concurrent lifecycle change on it; + /// without it the only alternative is polling on a timer. + /// + public TaskCompletionSource? CurrentEntered { get; set; } + /// Number of /probe requests issued against this client. public int ProbeCallCount => Volatile.Read(ref _probeCallCount); @@ -152,6 +234,7 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable ObjectDisposedException.ThrowIf(IsDisposed, this); ct.ThrowIfCancellationRequested(); Interlocked.Increment(ref _currentCallCount); + CurrentEntered?.TrySetResult(); // The request was accepted; the answer is withheld until the test releases the gate or the // caller's deadline cancels the token. @@ -161,6 +244,12 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable await gate.Task.WaitAsync(ct).ConfigureAwait(false); } + // A scripted answer supersedes — and then becomes — the standing one. + if (CurrentAnswers.TryDequeue(out var scripted)) + { + Current = scripted; + } + return CurrentFailure is null ? Current : throw CurrentFailure; } @@ -172,6 +261,21 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable Interlocked.Increment(ref _sampleCallCount); LastSampleFrom = from; + // Thrown before the first chunk, exactly like a real client that fails its /sample request: + // one-shot scripts first, then the sticky "this endpoint will never stream" failure. + if (SampleFailures.TryDequeue(out var oneShot)) + { + throw oneShot; + } + + if (SampleFailure is { } sticky) + { + throw sticky; + } + + // Bound to THIS generation: EndStream swaps in a fresh channel, so a consumer that + // reconnects reads the new one rather than the closed one it just lost. + var chunks = Volatile.Read(ref _chunks); using var lifetime = CancellationTokenSource.CreateLinkedTokenSource(ct, _disposeCts.Token); var delivered = 0L; @@ -180,7 +284,7 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable ScriptedChunk scripted; try { - scripted = await _chunks.Reader.ReadAsync(lifetime.Token).ConfigureAwait(false); + scripted = await chunks.Reader.ReadAsync(lifetime.Token).ConfigureAwait(false); } catch (ChannelClosedException) { @@ -192,9 +296,18 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable delivered++; - yield return scripted.Result; - - scripted.Consumed.TrySetResult(); + try + { + yield return scripted.Result; + } + finally + { + // Signalled when the consumer is DONE with this chunk — whether it came back for the + // next one or abandoned the enumeration. The finally is load-bearing: a chunk that + // makes the pump break out (a sequence gap, an agent restart) is never resumed past, + // so signalling only after the resume would hang every re-baseline test. + scripted.Consumed.TrySetResult(); + } } } @@ -212,7 +325,7 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable var scripted = new ScriptedChunk( chunk, new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously)); - if (!_chunks.Writer.TryWrite(scripted)) + if (!Volatile.Read(ref _chunks).Writer.TryWrite(scripted)) { throw new InvalidOperationException("The scripted /sample stream is already closed."); } @@ -221,10 +334,42 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable } /// - /// Closes the scripted /sample stream, which surfaces to the consumer as an - /// — the transient, reconnectable end. + /// Queues chunks for to deliver one at a time. Nothing is sent until + /// a test asks for it — the script is a plan, not a schedule. /// - public void EndStream() => _chunks.Writer.TryComplete(); + /// The chunks the Agent should send, in order. + public void ScriptChunks(params MTConnectStreamsResult[] chunks) + { + ArgumentNullException.ThrowIfNull(chunks); + + foreach (var chunk in chunks) + { + _script.Enqueue(chunk); + } + } + + /// + /// Pumps the next scripted chunk (see ) and waits until the + /// consumer has finished with it. + /// + /// A task completing once the consumer has processed — or abandoned the stream on — that chunk. + /// The script is exhausted. + public Task PumpOnce() => + _script.TryDequeue(out var next) + ? PumpAsync(next) + : throw new InvalidOperationException( + "No scripted /sample chunk left to pump; call ScriptChunks first."); + + /// + /// Closes the current scripted /sample stream, which surfaces to its consumer as an + /// — the transient, reconnectable end — and opens + /// a fresh one so a reconnecting consumer has somewhere to land. + /// + public void EndStream() + { + var closed = Interlocked.Exchange(ref _chunks, Channel.CreateUnbounded()); + closed.Writer.TryComplete(); + } /// public void Dispose() @@ -237,7 +382,7 @@ internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable } _disposeCts.Cancel(); - _chunks.Writer.TryComplete(); + Volatile.Read(ref _chunks).Writer.TryComplete(); _disposeCts.Dispose(); } diff --git a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/MTConnectSubscribeTests.cs b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/MTConnectSubscribeTests.cs new file mode 100644 index 00000000..d7b950b7 --- /dev/null +++ b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests/MTConnectSubscribeTests.cs @@ -0,0 +1,827 @@ +using System.Collections.Concurrent; +using Shouldly; +using Xunit; +using ZB.MOM.WW.OtOpcUa.Core.Abstractions; + +namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests; + +/// +/// Task 11 — 's half: one shared +/// /sample long-poll pump behind every handle, the ring-buffer re-baseline, and the +/// teardown that must leave nothing running. Every test runs against +/// ; no socket is opened anywhere in this file. +/// +/// +/// +/// Nothing here waits on wall-clock time. The fake completes a +/// only once the pump has finished with that +/// chunk, and exposes the pump's own task +/// () as the teardown barrier — so every +/// "the pump has now done X" statement is a fact, not a sleep. The one +/// below is a failure guard that turns a genuine hang (an +/// un-terminated retry loop, a deadlocked teardown) into a clean red test; no assertion +/// is ever made about elapsed time. +/// +/// +/// The load-bearing tests are the ones about a stream that has gone wrong. A pump +/// that only ever meets contiguous chunks is a dozen lines; the defects live in the four +/// ways it can be knocked off the sequence — the buffer rolling past the cursor +/// (IsSequenceGap), the Agent refusing an evicted from outright +/// (OUT_OF_RANGE under HTTP 200 ⇒ ), the Agent +/// restarting (a new instanceId, which ALSO looks like a gap), and the connection +/// dropping — plus the one failure it must NOT retry +/// (, where reconnecting reproduces the +/// identical answer forever). +/// +/// +public sealed class MTConnectSubscribeTests +{ + private const string AgentUri = "http://fixture-agent:5000"; + + /// The instanceId every canned fixture shares. + private const long FixtureInstanceId = 1655000000L; + + /// A DIFFERENT instanceId — the Agent came back as a new process. + private const long RestartedInstanceId = 1655999999L; + + /// The cursor Fixtures/current.xml primes the pump with. + private const long PrimedNextSequence = 108L; + + // Canonical Opc.Ua.StatusCodes numerics, restated so the test asserts the wire value a client + // actually sees rather than a driver-private constant it could drift with. + private const uint Good = 0x00000000u; + private const uint BadNoCommunication = 0x80310000u; + private const uint BadWaitingForInitialData = 0x80320000u; + private const uint BadNodeIdUnknown = 0x80340000u; + + /// + /// Failure guard for the handful of awaits that a defect could turn into a permanent hang + /// (a NotSupported retry loop, a teardown that deadlocks on the lifecycle semaphore). + /// Deliberately far longer than any of these operations could legitimately take — it exists + /// to make a hang red, not to assert a duration. + /// + private static readonly TimeSpan Watchdog = TimeSpan.FromSeconds(10); + + private static readonly DateTime ObservedAt = new(2026, 7, 24, 12, 30, 0, DateTimeKind.Utc); + + private static CancellationToken Ct => TestContext.Current.CancellationToken; + + private static MTConnectDriverOptions Opts(MTConnectReconnectOptions? reconnect = null) => + new() + { + AgentUri = AgentUri, + RequestTimeoutMs = 5000, + Reconnect = reconnect ?? new MTConnectReconnectOptions(), + Tags = + [ + new MTConnectTagDefinition("dev1_pos", DriverDataType.Float64), + new MTConnectTagDefinition("dev1_execution", DriverDataType.String), + new MTConnectTagDefinition("dev1_partcount", DriverDataType.Int32), + new MTConnectTagDefinition("dev1_never_reported", DriverDataType.Int32), + ], + }; + + // ---- fixtures ---- + + private static (MTConnectDriver Driver, CannedAgentClient Client) NewDriver( + CannedAgentClient? client = null, MTConnectDriverOptions? options = null) + { + var agent = client ?? CannedAgentClient.FromFixtures(); + + return (new MTConnectDriver(options ?? Opts(), "mt1", _ => agent), agent); + } + + private static async Task<(MTConnectDriver Driver, CannedAgentClient Client)> InitializedDriverAsync( + CannedAgentClient? client = null, MTConnectDriverOptions? options = null) + { + var (driver, agent) = NewDriver(client, options); + await driver.InitializeAsync("{}", Ct); + + return (driver, agent); + } + + /// Attaches a recorder to the driver's data-change event and returns its log. + private static ConcurrentQueue Record(MTConnectDriver driver) + { + var seen = new ConcurrentQueue(); + driver.OnDataChange += (_, e) => seen.Enqueue(e); + + return seen; + } + + /// Builds a /sample chunk (or /current answer) for the fixture Agent. + private static MTConnectStreamsResult Chunk( + long firstSequence, long nextSequence, params (string Id, string Value)[] observations) => + ChunkFrom(FixtureInstanceId, firstSequence, nextSequence, observations); + + private static MTConnectStreamsResult ChunkFrom( + long instanceId, long firstSequence, long nextSequence, params (string Id, string Value)[] observations) => + new( + instanceId, + nextSequence, + firstSequence, + [.. observations.Select(o => new MTConnectObservation(o.Id, o.Value, ObservedAt))]); + + private static IReadOnlyList For( + ConcurrentQueue seen, string reference) => + [.. seen.Where(e => e.FullReference == reference)]; + + // ---- initial data ---- + + /// + /// The plan's first TDD case, and the OPC UA Part 4 convention: a subscription reports the + /// current value immediately, out of the index the priming /current already filled — + /// it does not make the caller wait for the Agent's next chunk, which may be a whole + /// heartbeat away. + /// + [Fact] + public async Task Subscribe_fires_initial_data_from_current() + { + var (driver, _) = await InitializedDriverAsync(); + var seen = Record(driver); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + seen.Select(e => e.FullReference).ShouldContain("dev1_pos"); + var initial = seen.Single(); + initial.Snapshot.StatusCode.ShouldBe(Good); + initial.Snapshot.Value.ShouldBe(123.4567d); + } + + /// + /// Initial data fires for EVERY subscribed reference, including the ones with no value yet. + /// A driver that fired only for the references it happens to hold a value for would leave + /// the others with no snapshot at all — indistinguishable, from the server's side, from a + /// subscription that never established. + /// + [Fact] + public async Task Subscribe_fires_initial_data_for_every_ref_including_the_valueless_ones() + { + var (driver, _) = await InitializedDriverAsync(); + var seen = Record(driver); + + await driver.SubscribeAsync( + ["dev1_pos", "dev1_partcount", "dev1_never_reported", "nope"], TimeSpan.FromMilliseconds(50), Ct); + + seen.Count.ShouldBe(4); + For(seen, "dev1_pos")[0].Snapshot.StatusCode.ShouldBe(Good); + For(seen, "dev1_partcount")[0].Snapshot.StatusCode.ShouldBe(BadNoCommunication); + For(seen, "dev1_never_reported")[0].Snapshot.StatusCode.ShouldBe(BadWaitingForInitialData); + For(seen, "nope")[0].Snapshot.StatusCode.ShouldBe(BadNodeIdUnknown); + } + + /// Every event carries the handle the caller was given — that is how it demultiplexes. + [Fact] + public async Task Subscribe_returns_a_distinct_handle_per_call_and_stamps_it_on_every_event() + { + var (driver, _) = await InitializedDriverAsync(); + var seen = Record(driver); + + var a = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + var b = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + a.ShouldNotBe(b); + a.DiagnosticId.ShouldNotBe(b.DiagnosticId); + a.DiagnosticId.ShouldNotBeNullOrWhiteSpace(); + For(seen, "dev1_pos").Select(e => e.SubscriptionHandle).ShouldBe([a, b]); + } + + /// Subscribing to nothing costs the Agent nothing — no stream is opened for it. + [Fact] + public async Task Subscribe_of_an_empty_ref_list_opens_no_stream() + { + var (driver, client) = await InitializedDriverAsync(); + + var handle = await driver.SubscribeAsync([], TimeSpan.FromMilliseconds(50), Ct); + + handle.ShouldNotBeNull(); + client.SampleCallCount.ShouldBe(0); + + // Asserted on the pump itself, not only on the call count: a pump that HAD been started + // might simply not have reached its first request yet, and would pass the count alone. + driver.SampleStreamTask.ShouldBeNull(); + + // …and it still unsubscribes cleanly, so a caller that conditionally adds refs later is symmetric. + await driver.UnsubscribeAsync(handle, Ct); + } + + // ---- the shared stream ---- + + /// + /// One Agent, one /sample long poll — no matter how many handles. MTConnect streams + /// the whole device, so a stream per subscription would multiply the Agent's connection load + /// by the number of OPC UA subscriptions for identical data. + /// + [Fact] + public async Task Two_subscriptions_share_exactly_one_sample_stream() + { + var (driver, client) = await InitializedDriverAsync(); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + await driver.SubscribeAsync(["dev1_execution"], TimeSpan.FromMilliseconds(50), Ct); + + // Barrier: a chunk can only be consumed by a running enumeration, so both subscriptions have + // demonstrably landed on the same one by the time this returns. + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(1); + client.LastSampleFrom.ShouldBe(PrimedNextSequence); + } + + /// + /// A chunk updates the index for everything it carries, but only the subscribed references + /// raise a callback. The Agent streams every DataItem on the device; publishing all of them + /// would flood the server with values nothing asked for. + /// + [Fact] + public async Task Pump_raises_OnDataChange_only_for_subscribed_refs() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client.PumpAsync(CannedAgentClient.Chunk("Fixtures/sample.xml")).WaitAsync(Watchdog, Ct); + + seen.Select(e => e.FullReference).Distinct().ShouldBe(["dev1_pos"]); + For(seen, "dev1_pos")[0].Snapshot.Value.ShouldBe(124.01d); + + // The index took the whole chunk even though only one ref was published. + driver.ObservationIndex.Get("dev1_partcount").Value.ShouldBe(42); + } + + /// Each handle hears about the references IT subscribed, and only those. + [Fact] + public async Task Pump_fans_each_observation_to_every_handle_that_subscribes_it() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + var a = await driver.SubscribeAsync(["dev1_pos", "dev1_execution"], TimeSpan.FromMilliseconds(50), Ct); + var b = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client + .PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"), ("dev1_execution", "READY"))) + .WaitAsync(Watchdog, Ct); + + For(seen, "dev1_pos").Select(e => e.SubscriptionHandle).ShouldBe([a, b], ignoreOrder: true); + For(seen, "dev1_execution").Select(e => e.SubscriptionHandle).ShouldBe([a]); + } + + /// + /// The running-cursor pin. An observation-free heartbeat still advances the sequence — + /// the Agent sends one precisely so a quiet connection can be told from a dead one — and the + /// gap check must be made against the PREVIOUS chunk's nextSequence, never against the + /// from the stream was opened with. A driver that compares against the opening + /// from, or that skips an empty chunk when advancing, reports a gap on the third chunk + /// here and re-baselines against a perfectly healthy stream. + /// + [Fact] + public async Task Pump_advances_the_cursor_on_every_chunk_including_an_observation_free_one() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + // A heartbeat: no observations, and firstSequence is the buffer floor, not the cursor. + await client.PumpAsync(Chunk(1, 120)).WaitAsync(Watchdog, Ct); + + // Contiguous with the heartbeat's nextSequence — a gap ONLY if the heartbeat was ignored. + await client.PumpAsync(Chunk(120, 125, ("dev1_pos", "2.5"))).WaitAsync(Watchdog, Ct); + + client.CurrentCallCount.ShouldBe(1); // the initialize prime, and nothing else + client.SampleCallCount.ShouldBe(1); // one uninterrupted stream + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(2.5d); + } + + // ---- ring-buffer overflow ---- + + /// + /// The plan's second TDD case. The Agent's buffer rolled past the cursor + /// (firstSequence 5000 > the driver's 108), so the observations in between are + /// gone: the driver must re-baseline from /current rather than trust an incremental + /// update — and must do it once, then resume, not once per chunk forever. + /// + [Fact] + public async Task Sequence_gap_triggers_exactly_one_current_rebaseline_then_resumes() + { + var client = CannedAgentClient.WithGapThenResume(); + var (driver, _) = await InitializedDriverAsync(client); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client.PumpOnce().WaitAsync(Watchdog, Ct); // the gap chunk + + client.CurrentCallCount.ShouldBeGreaterThan(1); // the plan's assertion + client.CurrentCallCount.ShouldBe(2); // prime + exactly one re-baseline + + await client.PumpOnce().WaitAsync(Watchdog, Ct); // the contiguous follow-up + + client.CurrentCallCount.ShouldBe(2); // no re-baseline storm + client.SampleCallCount.ShouldBe(2); // the stream was reopened once + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(201.5d); + } + + /// + /// …and it resumes from the re-baselined nextSequence, not from the stale + /// cursor the Agent has already evicted. Reopening at the old cursor would earn the identical + /// rejection immediately and forever. + /// + [Fact] + public async Task Sequence_gap_resumes_the_stream_from_the_rebaselined_next_sequence() + { + var client = CannedAgentClient.WithGapThenResume(); + var (driver, _) = await InitializedDriverAsync(client); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + await client.PumpOnce().WaitAsync(Watchdog, Ct); // gap + await client.PumpOnce().WaitAsync(Watchdog, Ct); // resume — only a reopened stream delivers this + + client.LastSampleFrom.ShouldBe(5005L); + } + + /// + /// The other half of the overflow story. 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 as an + /// , never as a gap-bearing chunk. A pump that re-baselines + /// only on IsSequenceGap therefore fails precisely when it has fallen furthest behind. + /// + [Fact] + public async Task Out_of_range_stream_error_also_triggers_a_rebaseline() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + + client.SampleFailures.Enqueue(new InvalidDataException( + "MTConnect /sample answered an MTConnectError document (OUT_OF_RANGE) under HTTP 200")); + client.CurrentAnswers.Enqueue(Chunk(490, 500, ("dev1_pos", "50.5"))); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + // Only a pump that re-baselined and reopened the stream can consume this. + await client.PumpAsync(Chunk(500, 505, ("dev1_pos", "77.5"))).WaitAsync(Watchdog, Ct); + + client.CurrentCallCount.ShouldBe(2); + client.SampleCallCount.ShouldBe(2); + client.LastSampleFrom.ShouldBe(500L); + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(77.5d); + } + + // ---- agent restart ---- + + /// + /// The instanceId check must come BEFORE the gap check. An Agent restart changes + /// instanceId AND resets sequences, so a restart usually trips + /// IsSequenceGap too — and the two demand 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. This chunk is deliberately BOTH: new + /// instanceId and a firstSequence far past the cursor. A driver that tests the gap + /// first keeps serving dev1_execution's pre-restart value indefinitely — the new + /// Agent never reports it again, so nothing would ever overwrite it. + /// + [Fact] + public async Task Agent_restart_clears_the_index_and_is_detected_before_the_gap_path() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos", "dev1_execution"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + // The restarted Agent's /current knows nothing about dev1_execution. + client.CurrentAnswers.Enqueue(ChunkFrom(RestartedInstanceId, 40, 42, ("dev1_pos", "9.5"))); + + await client + .PumpAsync(ChunkFrom(RestartedInstanceId, 5000, 5005, ("dev1_pos", "5.5"))) + .WaitAsync(Watchdog, Ct); + + driver.AgentInstanceId.ShouldBe(RestartedInstanceId); + client.CurrentCallCount.ShouldBe(2); + + // Cleared: the pre-restart value is gone, not merely shadowed. + driver.ObservationIndex.Get("dev1_execution").StatusCode.ShouldBe(BadWaitingForInitialData); + driver.ObservationIndex.Get("dev1_pos").Value.ShouldBe(9.5d); + + // And the subscriber was TOLD its value went away, rather than being left holding a stale Good. + For(seen, "dev1_execution").Last().Snapshot.StatusCode.ShouldBe(BadWaitingForInitialData); + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(9.5d); + } + + // ---- stream ends ---- + + /// + /// A dropped connection is transient: reconnect from the cursor and keep going. It is + /// emphatically NOT a re-baseline — the sequence is still valid, and spending a + /// /current on every network blip would be a self-inflicted load. + /// + [Fact] + public async Task Stream_ended_reconnects_from_the_cursor_without_a_rebaseline() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + client.EndStream(); + + // Only a reconnected pump can consume a chunk written to the NEW stream generation. + await client.PumpAsync(Chunk(113, 118, ("dev1_pos", "2.5"))).WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(2); + client.LastSampleFrom.ShouldBe(113L); + client.CurrentCallCount.ShouldBe(1); + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(2.5d); + } + + /// + /// The hot-loop pin. means the + /// endpoint will never stream to this request — a reverse proxy, a health page, an Agent + /// fronted by something that collapses multipart/x-mixed-replace. Reconnecting + /// reproduces the identical answer forever, so the pump must give up loudly instead of + /// hammering the endpoint until an operator notices. A later subscription must not restart + /// the loop either: the answer has not changed just because someone asked again. + /// + [Fact] + public async Task Stream_not_supported_stops_the_pump_and_never_retries() + { + var (driver, client) = await InitializedDriverAsync(); + client.SampleFailure = new MTConnectStreamNotSupportedException( + "the Agent answered /sample with a single non-multipart document"); + + var handle = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + var pump = driver.SampleStreamTask.ShouldNotBeNull(); + + await pump.WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(1); + driver.GetHealth().State.ShouldBe(DriverState.Degraded); + driver.GetHealth().LastError.ShouldNotBeNullOrWhiteSpace(); + + // A full unsubscribe/resubscribe cycle must not re-open the wound either. + await driver.UnsubscribeAsync(handle, Ct).WaitAsync(Watchdog, Ct); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + client.SampleCallCount.ShouldBe(1); + } + + /// + /// …but a re-initialize IS an operator changing something, so it clears the verdict and + /// tries again. Otherwise fixing the proxy would need a process restart. + /// + [Fact] + public async Task Reinitialize_clears_the_stream_unsupported_verdict() + { + var (driver, client) = await InitializedDriverAsync(); + client.SampleFailure = new MTConnectStreamNotSupportedException("not multipart"); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + await driver.SampleStreamTask!.WaitAsync(Watchdog, Ct); + + client.SampleFailure = null; + await driver.ReinitializeAsync("{}", Ct).WaitAsync(Watchdog, Ct); + + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(2); + } + + // ---- backoff ---- + + /// + /// The reconnect backoff as a pure function of the attempt number: the first retry honours + /// MinBackoffMs (zero = immediate), and every later one grows geometrically to the + /// cap. The 100 ms growth floor matters more than it looks — the default + /// MinBackoffMs is 0, and 0 × multiplier is still 0, so without a floor + /// the "geometric" backoff would be an unbounded hot loop against a dead Agent. + /// + [Fact] + public void Reconnect_backoff_starts_at_min_grows_geometrically_and_caps() + { + var options = new MTConnectReconnectOptions + { + MinBackoffMs = 0, MaxBackoffMs = 1000, BackoffMultiplier = 2.0, + }; + + MTConnectDriver.BackoffFor(1, options).ShouldBe(TimeSpan.Zero); + MTConnectDriver.BackoffFor(2, options).ShouldBe(TimeSpan.FromMilliseconds(200)); + MTConnectDriver.BackoffFor(3, options).ShouldBe(TimeSpan.FromMilliseconds(400)); + MTConnectDriver.BackoffFor(4, options).ShouldBe(TimeSpan.FromMilliseconds(800)); + MTConnectDriver.BackoffFor(5, options).ShouldBe(TimeSpan.FromMilliseconds(1000)); + MTConnectDriver.BackoffFor(500, options).ShouldBe(TimeSpan.FromMilliseconds(1000)); + } + + /// + /// Operator-authored nonsense must not un-bound the loop: a multiplier of 1 (or less) would + /// never grow, and a negative delay is not a delay. Both fall back to the shipped default. + /// + [Fact] + public void Reconnect_backoff_survives_a_degenerate_multiplier() + { + var options = new MTConnectReconnectOptions + { + MinBackoffMs = -5, MaxBackoffMs = 10_000, BackoffMultiplier = 0.5, + }; + + MTConnectDriver.BackoffFor(1, options).ShouldBe(TimeSpan.Zero); + MTConnectDriver.BackoffFor(2, options).ShouldBeGreaterThan(TimeSpan.Zero); + MTConnectDriver.BackoffFor(3, options) + .ShouldBeGreaterThan(MTConnectDriver.BackoffFor(2, options)); + } + + // ---- unsubscribe ---- + + /// + /// Dropping one handle drops only its references; the shared stream keeps running for the + /// others. Tearing the stream down on the first unsubscribe would silently blind every + /// remaining subscription. + /// + [Fact] + public async Task Unsubscribe_drops_only_that_handles_refs_and_keeps_the_shared_stream() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + var a = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + var b = await driver.SubscribeAsync(["dev1_execution"], TimeSpan.FromMilliseconds(50), Ct); + var pump = driver.SampleStreamTask.ShouldNotBeNull(); + + await driver.UnsubscribeAsync(a, Ct).WaitAsync(Watchdog, Ct); + seen.Clear(); + + await client + .PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"), ("dev1_execution", "READY"))) + .WaitAsync(Watchdog, Ct); + + pump.IsCompleted.ShouldBeFalse(); + seen.Select(e => e.FullReference).Distinct().ShouldBe(["dev1_execution"]); + + await driver.UnsubscribeAsync(b, Ct).WaitAsync(Watchdog, Ct); + + pump.IsCompleted.ShouldBeTrue(); + driver.SampleStreamTask.ShouldBeNull(); + } + + /// The last unsubscribe stops the stream — and nothing is published after it. + [Fact] + public async Task Unsubscribe_of_the_last_handle_stops_the_stream_and_silences_the_driver() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + var handle = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + await driver.UnsubscribeAsync(handle, Ct).WaitAsync(Watchdog, Ct); + seen.Clear(); + + // The stream is gone, so this chunk has no consumer and nothing may reach the recorder. + client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).IsCompleted.ShouldBeFalse(); + + seen.ShouldBeEmpty(); + driver.SampleStreamTask.ShouldBeNull(); + } + + /// + /// Unsubscribing twice, or with a handle this driver never issued, is a no-op — not a throw. + /// Teardown paths run on failure paths, and an exception there would mask the real fault. + /// + [Fact] + public async Task Unsubscribe_is_idempotent_and_ignores_a_foreign_handle() + { + var (driver, _) = await InitializedDriverAsync(); + var handle = await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + await driver.UnsubscribeAsync(handle, Ct).WaitAsync(Watchdog, Ct); + await driver.UnsubscribeAsync(handle, Ct).WaitAsync(Watchdog, Ct); + await driver.UnsubscribeAsync(new ForeignHandle(), Ct).WaitAsync(Watchdog, Ct); + } + + /// A null argument is a caller bug, and the one thing these methods are loud about. + [Fact] + public async Task Null_arguments_throw_ArgumentNullException() + { + var (driver, _) = await InitializedDriverAsync(); + + await Should.ThrowAsync( + () => driver.SubscribeAsync(null!, TimeSpan.FromMilliseconds(50), Ct)); + await Should.ThrowAsync(() => driver.UnsubscribeAsync(null!, Ct)); + } + + // ---- lifecycle ---- + + /// + /// Subscribing before the driver is initialized registers the interest and reports the + /// honest "no value yet" — it does not throw, and it does not dial an Agent that does not + /// exist yet. Initialize then picks the standing subscription up. + /// + [Fact] + public async Task Subscribe_before_initialize_registers_and_initialize_starts_the_stream() + { + var (driver, client) = NewDriver(); + var seen = Record(driver); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + client.SampleCallCount.ShouldBe(0); + driver.SampleStreamTask.ShouldBeNull(); + seen.Single().Snapshot.StatusCode.ShouldBe(BadWaitingForInitialData); + + await driver.InitializeAsync("{}", Ct); + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(1); + For(seen, "dev1_pos").Last().Snapshot.Value.ShouldBe(1.5d); + } + + /// + /// A re-initialize replaces the served state, so it must replace the pump with it: the old + /// one is stopped (it holds the old cursor, and possibly the old client) and a new one starts + /// from the fresh /current. The standing subscription survives — the caller did not + /// ask for it to be dropped — and is re-primed, because every value it holds came from a + /// baseline that has just been thrown away. + /// + [Fact] + public async Task Reinitialize_stops_the_old_pump_starts_a_new_one_and_republishes() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + var first = driver.SampleStreamTask.ShouldNotBeNull(); + + // Barrier: the first pump is demonstrably enumerating before the re-initialize lands on it. + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + client.CurrentAnswers.Enqueue(Chunk(300, 310, ("dev1_pos", "310.5"))); + seen.Clear(); + + await driver.ReinitializeAsync("{}", Ct).WaitAsync(Watchdog, Ct); + + first.IsCompleted.ShouldBeTrue(); + var second = driver.SampleStreamTask.ShouldNotBeNull(); + second.ShouldNotBeSameAs(first); + + await client.PumpAsync(Chunk(310, 315, ("dev1_pos", "315.5"))).WaitAsync(Watchdog, Ct); + + client.SampleCallCount.ShouldBe(2); + client.LastSampleFrom.ShouldBe(310L); + For(seen, "dev1_pos").Select(e => e.Snapshot.Value).ShouldBe([310.5d, 315.5d]); + } + + /// + /// Shutdown stops the pump BEFORE disposing the client. Getting that order wrong leaves the + /// pump enumerating a disposed client — which it would read as a dropped connection and try + /// to reconnect, forever, against an object that no longer exists. + /// + [Fact] + public async Task Shutdown_stops_the_pump_before_disposing_the_client_and_leaks_nothing() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + var pump = driver.SampleStreamTask.ShouldNotBeNull(); + + // Barrier: the stream is demonstrably open before the shutdown lands, so the call count + // below is a statement about reconnect churn rather than about who won a race. + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + seen.Clear(); + + await driver.ShutdownAsync(Ct).WaitAsync(Watchdog, Ct); + + pump.IsCompleted.ShouldBeTrue(); + pump.IsFaulted.ShouldBeFalse(); + driver.SampleStreamTask.ShouldBeNull(); + client.IsDisposed.ShouldBeTrue(); + + // No reconnect churn against the disposed client, and nothing published after teardown. + client.SampleCallCount.ShouldBe(1); + seen.ShouldBeEmpty(); + } + + /// + /// The deadlock pin. The lifecycle semaphore is non-reentrant, so a pump that reached + /// back into InitializeAsync/ReinitializeAsync/ShutdownAsync — the + /// natural way to write "on a gap, just re-prime" — would hang the driver permanently, with + /// no exception and no log. Here the pump is parked INSIDE its re-baseline /current + /// when a shutdown lands and takes that semaphore: the shutdown must cancel the pump and + /// complete. If the pump were waiting on the lifecycle lock instead, neither side would ever + /// move. + /// + [Fact] + public async Task Shutdown_while_the_pump_is_mid_rebaseline_does_not_deadlock() + { + var (driver, client) = await InitializedDriverAsync(); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + client.CurrentEntered = entered; + client.CurrentGate = release; + + // A gap forces the re-baseline; the gate holds the pump inside the /current request. + _ = client.PumpAsync(Chunk(5000, 5005, ("dev1_pos", "5.5"))); + await entered.Task.WaitAsync(Watchdog, Ct); + + try + { + await driver.ShutdownAsync(Ct).WaitAsync(Watchdog, Ct); + } + finally + { + release.TrySetResult(); + } + + driver.SampleStreamTask.ShouldBeNull(); + client.IsDisposed.ShouldBeTrue(); + } + + // ---- robustness + health ---- + + /// + /// A subscriber that throws is a bug in the consumer, not a reason to stop delivering data to + /// every other subscriber. The pump absorbs it and keeps pumping. + /// + [Fact] + public async Task A_throwing_subscriber_does_not_kill_the_pump() + { + var (driver, client) = await InitializedDriverAsync(); + var seen = Record(driver); + driver.OnDataChange += (_, _) => throw new InvalidOperationException("subscriber blew up"); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + seen.Clear(); + + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + await client.PumpAsync(Chunk(113, 118, ("dev1_pos", "2.5"))).WaitAsync(Watchdog, Ct); + + driver.SampleStreamTask!.IsCompleted.ShouldBeFalse(); + For(seen, "dev1_pos").Select(e => e.Snapshot.Value).ShouldBe([1.5d, 2.5d]); + } + + /// + /// Health precedence, read side. A healthy stream must not paper over a broken read + /// path: /sample and /current are different Agent requests and can fail + /// independently, and "Healthy" while every OPC UA Read returns Bad is exactly the + /// failure-wearing-success shape this driver is written against. + /// + [Fact] + public async Task A_pumped_chunk_does_not_paper_over_a_failed_read() + { + var (driver, client) = await InitializedDriverAsync(); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + + client.CurrentFailure = new HttpRequestException("connection refused"); + await driver.ReadAsync(["dev1_pos"], Ct); + driver.GetHealth().State.ShouldBe(DriverState.Degraded); + + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + driver.GetHealth().State.ShouldBe(DriverState.Degraded); + + // …and clearing the read failure does restore it, because the stream never failed. + client.CurrentFailure = null; + await driver.ReadAsync(["dev1_pos"], Ct); + driver.GetHealth().State.ShouldBe(DriverState.Healthy); + } + + /// + /// Health precedence, stream side. The mirror image: a successful read must not paper + /// over a stream that has stopped delivering. Only the path that failed may clear its own + /// degradation. + /// + [Fact] + public async Task A_successful_read_does_not_paper_over_a_broken_stream() + { + var (driver, client) = await InitializedDriverAsync(); + client.SampleFailure = new MTConnectStreamNotSupportedException("not multipart"); + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + await driver.SampleStreamTask!.WaitAsync(Watchdog, Ct); + + var res = await driver.ReadAsync(["dev1_pos"], Ct); + + res[0].StatusCode.ShouldBe(Good); + driver.GetHealth().State.ShouldBe(DriverState.Degraded); + } + + /// A recovered stream clears its own degradation once chunks flow again. + [Fact] + public async Task A_recovered_stream_restores_health() + { + var (driver, client) = await InitializedDriverAsync(); + client.SampleFailures.Enqueue(new TimeoutException("no chunk within the heartbeat window")); + + await driver.SubscribeAsync(["dev1_pos"], TimeSpan.FromMilliseconds(50), Ct); + await client.PumpAsync(Chunk(PrimedNextSequence, 113, ("dev1_pos", "1.5"))).WaitAsync(Watchdog, Ct); + + driver.GetHealth().State.ShouldBe(DriverState.Healthy); + client.SampleCallCount.ShouldBe(2); + } + + /// A handle from some other driver — the type check must not be a cast. + private sealed class ForeignHandle : ISubscriptionHandle + { + public string DiagnosticId => "not-ours"; + } +}