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