26eb511a7e
Important 1 - the probe-model triple was read outside _lifecycle while every
write happened under it. Packed (Client, ProbeModel, ProbeModelDataItemCount)
into one immutable AgentSession swapped by a single Volatile.Write, matching
what _index already does, rather than locking the readers: both await the
network, so taking _lifecycle would stall a shutdown for a whole
RequestTimeoutMs and would widen the non-reentrancy hazard. GetOrFetch and
Flush re-publish via CompareExchange so a late re-cache cannot undo a
concurrent teardown/re-init, and a client disposed mid-fetch now surfaces as
the same "not connected" InvalidOperationException browse already handles
instead of a raw ObjectDisposedException.
Important 2 - HasConfigBody decided emptiness by matching the literals "{}" /
"[]", so "{ }", "{\n}" and every pretty-printed empty document fell through to
ParseOptions, failed the required-AgentUri check, and turned a semantically
empty config into a cold-start fault. Now parsed: an object/array with no
elements (or a bare null) is empty; malformed text is deliberately NOT empty so
ParseOptions produces the real quoted error rather than silently starting on
stale options.
Important 3 - documented the _lifecycle non-reentrancy hazard in the code, on
both the semaphore and StopSampleStreamAsync, naming Task 11's re-baseline as
the specific path that would deadlock and giving the CurrentAsync-directly
pattern that avoids it.
Important 4 - removed the Initialize/Reinitialize asymmetry rather than
documenting it. Both now share one rule via ResolveIncomingOptions: an
unreadable config document never destroys working state, and faults only a
driver that had none. InitializeAsync previously tore down before parsing, so a
bad document on a live instance destroyed a healthy client - the exact outcome
ReinitializeAsync was written to avoid.
Minors: SafeDispose traces the swallowed disposal fault at Debug (a client that
cannot be released is how a handle leak starts); the RequirePositive test is a
Theory over all four timing knobs, not just RequestTimeoutMs (arch-review
01/S-6); the footprint test no longer overclaims "is zero before initialize".
CannedAgentClient gains a one-shot ProbeGate so a lifecycle change landing
mid-request is deterministic - no timers. 346/346 (329 + 17). Six mutations
verified: fetch-CAS, disposed-translation, flush retired-session guard,
literal HasConfigBody, teardown-before-parse, and two dropped RequirePositive
calls each fail only their own tests.
246 lines
10 KiB
C#
246 lines
10 KiB
C#
using System.Runtime.CompilerServices;
|
|
using System.Threading.Channels;
|
|
|
|
namespace ZB.MOM.WW.OtOpcUa.Driver.MTConnect.Tests;
|
|
|
|
/// <summary>
|
|
/// The shared <see cref="IMTConnectAgentClient"/> test double: an Agent that serves the canned
|
|
/// <c>Fixtures/probe.xml</c> + <c>Fixtures/current.xml</c> documents, counts every call, records
|
|
/// its own disposal, and lets a test drive the <c>/sample</c> chunk sequence by hand.
|
|
/// </summary>
|
|
/// <remarks>
|
|
/// <para>
|
|
/// <b>Deterministic by construction — no timers, no sleeps, no polling.</b> The
|
|
/// <c>/sample</c> leg reads from an unbounded <see cref="Channel{T}"/> that only a test
|
|
/// fills, and <see cref="PumpAsync"/> does not return until the driver has actually consumed
|
|
/// the chunk it wrote (each chunk carries its own completion source, signalled after the
|
|
/// enumerator's <c>yield return</c> resumes). A subscription test can therefore say "one
|
|
/// chunk has now been fully processed" as a fact rather than as a timing hope.
|
|
/// </para>
|
|
/// <para>
|
|
/// <b>It honours the seam's stream-end contract</b> (see
|
|
/// <see cref="IMTConnectAgentClient.SampleAsync"/>): cancelling the token is the only way the
|
|
/// enumeration ends without throwing. Closing the scripted stream via
|
|
/// <see cref="EndStream"/> raises <see cref="MTConnectStreamEndedException"/>, exactly as the
|
|
/// production client does when an Agent drops the connection — so a pump written against the
|
|
/// fake cannot silently pass while mishandling the real one.
|
|
/// </para>
|
|
/// <para>
|
|
/// <b>Disposal is observable</b> (<see cref="DisposeCount"/>) because "did the driver
|
|
/// actually release the client?" is a behaviour, not an implementation detail: a
|
|
/// <c>ReinitializeAsync</c> that re-points the Agent but leaks the old client keeps a live
|
|
/// connection pool per re-deploy, and nothing else in the test surface would notice.
|
|
/// </para>
|
|
/// </remarks>
|
|
internal sealed class CannedAgentClient : IMTConnectAgentClient, IDisposable
|
|
{
|
|
private readonly Channel<ScriptedChunk> _chunks = Channel.CreateUnbounded<ScriptedChunk>();
|
|
private readonly CancellationTokenSource _disposeCts = new();
|
|
|
|
private int _probeCallCount;
|
|
private int _currentCallCount;
|
|
private int _sampleCallCount;
|
|
private int _disposeCount;
|
|
|
|
private CannedAgentClient(MTConnectProbeModel probe, MTConnectStreamsResult current)
|
|
{
|
|
Probe = probe;
|
|
Current = current;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Builds a client serving the repo's canned fixtures — <c>Fixtures/probe.xml</c> for
|
|
/// <c>/probe</c> and <c>Fixtures/current.xml</c> for <c>/current</c>.
|
|
/// </summary>
|
|
/// <param name="probeFixture">Probe-document fixture path, relative to the test binaries.</param>
|
|
/// <param name="currentFixture">Streams-document fixture path, relative to the test binaries.</param>
|
|
public static CannedAgentClient FromFixtures(
|
|
string probeFixture = "Fixtures/probe.xml",
|
|
string currentFixture = "Fixtures/current.xml") =>
|
|
new(
|
|
MTConnectProbeParser.Parse(File.ReadAllText(probeFixture)),
|
|
MTConnectStreamsParser.Parse(File.ReadAllText(currentFixture)));
|
|
|
|
/// <summary>Parses a streams fixture into a chunk a test can script onto the sample stream.</summary>
|
|
/// <param name="fixture">Streams-document fixture path, relative to the test binaries.</param>
|
|
public static MTConnectStreamsResult Chunk(string fixture) =>
|
|
MTConnectStreamsParser.Parse(File.ReadAllText(fixture));
|
|
|
|
/// <summary>The document <c>/probe</c> answers with. Settable so a test can re-shape the model.</summary>
|
|
public MTConnectProbeModel Probe { get; set; }
|
|
|
|
/// <summary>
|
|
/// The document <c>/current</c> answers with. Settable so a re-baseline (or an agent-restart
|
|
/// <c>instanceId</c> change) can be scripted between calls.
|
|
/// </summary>
|
|
public MTConnectStreamsResult Current { get; set; }
|
|
|
|
/// <summary>When set, <c>/probe</c> throws this instead of answering.</summary>
|
|
public Exception? ProbeFailure { get; set; }
|
|
|
|
/// <summary>When set, <c>/current</c> throws this instead of answering.</summary>
|
|
public Exception? CurrentFailure { get; set; }
|
|
|
|
/// <summary>
|
|
/// When set, the <b>next</b> <c>/probe</c> parks here until the test completes it, then the
|
|
/// gate clears itself so later calls answer immediately.
|
|
/// </summary>
|
|
/// <remarks>
|
|
/// This is how a test holds a request "in flight" across a concurrent lifecycle change —
|
|
/// a shutdown or a re-initialize landing mid-fetch — with no timers and no races of its own.
|
|
/// The answer is captured when the request lands, not when it completes, so a test can
|
|
/// change <see cref="Probe"/> meanwhile and still tell the two documents apart.
|
|
/// </remarks>
|
|
public TaskCompletionSource? ProbeGate { get; set; }
|
|
|
|
/// <summary>
|
|
/// When set, <c>/current</c> does not answer until this source completes — an Agent that
|
|
/// accepted the request and then went quiet. The wait is cancellation-observing, so the
|
|
/// caller's own deadline is what ends it; the test never sleeps and never races a timer it
|
|
/// did not set. Leave <c>null</c> for the ordinary immediate answer.
|
|
/// </summary>
|
|
public TaskCompletionSource? CurrentGate { get; set; }
|
|
|
|
/// <summary>Number of <c>/probe</c> requests issued against this client.</summary>
|
|
public int ProbeCallCount => Volatile.Read(ref _probeCallCount);
|
|
|
|
/// <summary>Number of <c>/current</c> requests issued against this client.</summary>
|
|
public int CurrentCallCount => Volatile.Read(ref _currentCallCount);
|
|
|
|
/// <summary>Number of <c>/sample</c> enumerations started against this client.</summary>
|
|
public int SampleCallCount => Volatile.Read(ref _sampleCallCount);
|
|
|
|
/// <summary>Number of <see cref="Dispose"/> calls (not clamped — a double-dispose is visible).</summary>
|
|
public int DisposeCount => Volatile.Read(ref _disposeCount);
|
|
|
|
/// <summary>Whether the client has been disposed at least once.</summary>
|
|
public bool IsDisposed => DisposeCount > 0;
|
|
|
|
/// <summary>The <c>from</c> sequence the most recent <c>/sample</c> enumeration was opened at.</summary>
|
|
public long? LastSampleFrom { get; private set; }
|
|
|
|
/// <inheritdoc/>
|
|
public async Task<MTConnectProbeModel> ProbeAsync(CancellationToken ct)
|
|
{
|
|
ObjectDisposedException.ThrowIf(IsDisposed, this);
|
|
ct.ThrowIfCancellationRequested();
|
|
Interlocked.Increment(ref _probeCallCount);
|
|
|
|
if (ProbeFailure is not null)
|
|
{
|
|
throw ProbeFailure;
|
|
}
|
|
|
|
// Captured now: the Agent answers with the document it held when the request landed.
|
|
var answer = Probe;
|
|
|
|
if (ProbeGate is { } gate)
|
|
{
|
|
ProbeGate = null;
|
|
await gate.Task.WaitAsync(ct).ConfigureAwait(false);
|
|
|
|
// A real client whose handler was disposed mid-request fails the request; so does this.
|
|
ObjectDisposedException.ThrowIf(IsDisposed, this);
|
|
}
|
|
|
|
return answer;
|
|
}
|
|
|
|
/// <inheritdoc/>
|
|
public async Task<MTConnectStreamsResult> CurrentAsync(CancellationToken ct)
|
|
{
|
|
ObjectDisposedException.ThrowIf(IsDisposed, this);
|
|
ct.ThrowIfCancellationRequested();
|
|
Interlocked.Increment(ref _currentCallCount);
|
|
|
|
// The request was accepted; the answer is withheld until the test releases the gate or the
|
|
// caller's deadline cancels the token.
|
|
var gate = CurrentGate;
|
|
if (gate is not null)
|
|
{
|
|
await gate.Task.WaitAsync(ct).ConfigureAwait(false);
|
|
}
|
|
|
|
return CurrentFailure is null ? Current : throw CurrentFailure;
|
|
}
|
|
|
|
/// <inheritdoc/>
|
|
public async IAsyncEnumerable<MTConnectStreamsResult> SampleAsync(
|
|
long from, [EnumeratorCancellation] CancellationToken ct)
|
|
{
|
|
ObjectDisposedException.ThrowIf(IsDisposed, this);
|
|
Interlocked.Increment(ref _sampleCallCount);
|
|
LastSampleFrom = from;
|
|
|
|
using var lifetime = CancellationTokenSource.CreateLinkedTokenSource(ct, _disposeCts.Token);
|
|
var delivered = 0L;
|
|
|
|
while (true)
|
|
{
|
|
ScriptedChunk scripted;
|
|
try
|
|
{
|
|
scripted = await _chunks.Reader.ReadAsync(lifetime.Token).ConfigureAwait(false);
|
|
}
|
|
catch (ChannelClosedException)
|
|
{
|
|
// Mirrors the production client: a stream that stops for any reason other than the
|
|
// caller cancelling is an exception, never a quiet end of enumeration.
|
|
throw new MTConnectStreamEndedException(
|
|
MTConnectStreamEndReason.ConnectionClosed, "canned://agent", delivered);
|
|
}
|
|
|
|
delivered++;
|
|
|
|
yield return scripted.Result;
|
|
|
|
scripted.Consumed.TrySetResult();
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Queues a chunk onto the scripted <c>/sample</c> stream and waits until the consumer has
|
|
/// processed it. This is the deterministic replacement for "wait a bit and hope the pump
|
|
/// ran".
|
|
/// </summary>
|
|
/// <param name="chunk">The chunk the Agent should send next.</param>
|
|
/// <returns>A task completing once the enumerating consumer has moved past <paramref name="chunk"/>.</returns>
|
|
public Task PumpAsync(MTConnectStreamsResult chunk)
|
|
{
|
|
ArgumentNullException.ThrowIfNull(chunk);
|
|
|
|
var scripted = new ScriptedChunk(
|
|
chunk, new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously));
|
|
|
|
if (!_chunks.Writer.TryWrite(scripted))
|
|
{
|
|
throw new InvalidOperationException("The scripted /sample stream is already closed.");
|
|
}
|
|
|
|
return scripted.Consumed.Task;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Closes the scripted <c>/sample</c> stream, which surfaces to the consumer as an
|
|
/// <see cref="MTConnectStreamEndedException"/> — the transient, reconnectable end.
|
|
/// </summary>
|
|
public void EndStream() => _chunks.Writer.TryComplete();
|
|
|
|
/// <inheritdoc/>
|
|
public void Dispose()
|
|
{
|
|
Interlocked.Increment(ref _disposeCount);
|
|
|
|
if (DisposeCount > 1)
|
|
{
|
|
return;
|
|
}
|
|
|
|
_disposeCts.Cancel();
|
|
_chunks.Writer.TryComplete();
|
|
_disposeCts.Dispose();
|
|
}
|
|
|
|
private sealed record ScriptedChunk(MTConnectStreamsResult Result, TaskCompletionSource Consumed);
|
|
}
|