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"; } }