Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
230 changes: 219 additions & 11 deletions src/Daqifi.Core.Tests/Device/TextExchangeLineFramingTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,59 @@ public class TextExchangeLineFramingTests

private const int CompletionTimeoutMs = 150;

/// <summary>
/// Completion timeout for the fragmented-reply test, and the two constants below it are sized
/// against it. Two things have to be true at once there, and they pull in opposite directions:
/// </summary>
/// <remarks>
/// <para>
/// One fragment gap must stay well inside this window, or a wait loop that DOES restart its
/// window still looks like one that cut the reply short. That is the side this test kept
/// failing on. The gap used to be counted in the reader thread's own idle reads —
/// five <see cref="IdleReadSleepMs"/> sleeps — which made the gap the SUM of five scheduling
/// delays rather than one. On the macOS runner each 10 ms sleep came back at roughly 100 ms
/// under the load of the rest of the suite, so the nominal 50 ms gap arrived at ~500 ms, ate
/// the whole window, and the exchange returned six of the twelve lines. Widening the window
/// (300 ms -> 500 ms) did not fix it, because the stretch scales with the window.
/// </para>
/// <para>
/// The gap is now measured on the stream's own clock — the next fragment is handed over on the
/// first read that finds <see cref="StageGapMs"/> elapsed since the previous one drained — so a
/// stretched sleep no longer accumulates. The gap is <see cref="StageGapMs"/> plus at most ONE
/// late read, which means a single scheduling stall of 450 ms has to land inside one gap before
/// the test lies about the exchange, instead of five stalls of 90 ms.
/// </para>
/// <para>
/// The fragments must ALSO span longer than this window in total, or the test passes for a loop
/// that never restarts anything. That needs
/// (<see cref="StagedFragmentCount"/> - 1) x <see cref="StageGapMs"/> &gt; the window:
/// 15 x 50 ms = 750 ms nominal, and more on a slow machine, since delay can only lengthen a
/// gap. Both sides therefore hold in the direction a loaded runner pushes them — but the span
/// is asserted from the stream's record of it rather than left to this arithmetic, so pacing
/// that collapses fails instead of going vacuous.
/// </para>
/// </remarks>
private const int StagedCompletionTimeoutMs = 500;

/// <summary>
/// Nominal gap between fragments, measured on the stream's clock rather than counted in idle
/// reads; see <see cref="StagedCompletionTimeoutMs"/>.
/// </summary>
private const int StageGapMs = 50;

/// <summary>
/// How many fragments the staged reply is dribbled out in. Enough that the fragments outlast
/// one completion window even though each gap is a small fraction of it.
/// </summary>
private const int StagedFragmentCount = 16;

/// <summary>
/// How long an idle read parks the consumer thread for. Only the polling interval now: it sets
/// how soon after <see cref="StageGapMs"/> elapses the next fragment is noticed, and no longer
/// multiplies into the gap itself.
/// </summary>
private const int IdleReadSleepMs = 10;

[Fact]
public async Task ExecuteTextCommand_WhenTheDeviceAnswersWithABareLineFeed_ReturnsTheLine()
{
Expand Down Expand Up @@ -158,6 +211,80 @@ public async Task ExecuteTextCommand_WhenTheReplyMixesBothLineEndings_ReturnsEve
device.Disconnect();
}

[Fact]
public async Task ExecuteTextCommand_WhenTheDeviceAnswersInFragmentsSpacedApart_ReturnsEveryLine()
{
// The property the inactivity tail exists for, and the one an event-driven wait loop is
// most likely to get wrong: every line must restart the completion window, so a device
// that dribbles its answer out over more than one completion timeout in total is never cut
// off part-way. A loop that instead measured its deadline from when it started waiting
// would return the first few lines here and drop the rest.
//
// The fragments are paced on the stream's own clock — the next one is handed over on the
// first read that finds StageGapMs elapsed since the previous one drained — rather than by
// a delay on the test thread or by counting the reader thread's sleeps. A sleep is a floor,
// not a promise, and a loaded runner stretches one; timing the gap instead of counting
// sleeps keeps one stretched sleep from being multiplied into the whole gap. The margin
// between a gap and the completion window is what keeps this honest, and it is sized in the
// note on StagedCompletionTimeoutMs.
var expected = Enumerable.Range(1, StagedFragmentCount).Select(i => $"line{i}").ToArray();

using var transport = new ScriptedReplyTransport(
stageGapMs: StageGapMs,
expected.Select(line => $"{line}\r\n").ToArray());
using var device = new LineFramingTestableDevice("Fragmented Device", transport);

device.Connect();

var lines = await device.CallAsync(
() => transport.Release(),
completionTimeoutMs: StagedCompletionTimeoutMs);

Assert.Equal(expected, lines);
AssertRecognisedTheReply(device);

// Without this the test would still pass if the pacing ever collapsed to nothing — every
// line would arrive inside a single completion window and a wait loop that never restarted
// its window would look correct. Asserting the reply really did outlast one window is what
// makes the collection assertion above evidence for the property in the comment.
//
// Asked of the stream rather than of a stopwatch around the call: the call cannot return
// until a whole completion window has passed with no new line, so its own elapsed time is
// greater than the window whatever the pacing did, and asserting on it proves nothing. The
// span between the first and last fragment reaching the reader is the quantity that has to
// outlast the window, and only the stream can measure it.
var stagedSpanMs = transport.StagedSpan.TotalMilliseconds;
Assert.True(stagedSpanMs > StagedCompletionTimeoutMs,
$"The staged reply was handed over across {stagedSpanMs:F0} ms, which is inside one " +
$"{StagedCompletionTimeoutMs} ms completion window: the fragments were not spaced apart, " +
"so this test no longer exercises the window restarting on each line.");

device.Disconnect();
}

[Fact]
public async Task ExecuteTextCommand_WhenTheCallerCancelsWhileWaitingForAReply_ThrowsOperationCanceled()
{
// The wait for the device's reply is cancellable, and must stay cancellable however it is
// implemented. A silent device would otherwise hold the operation lock for the whole
// first-response window with nothing the caller could do about it.
using var transport = new ScriptedReplyTransport("\r\n");
using var device = new LineFramingTestableDevice("Cancelled Device", transport);

device.Connect();

using var cancellation = new CancellationTokenSource();

// Armed from inside the setup action so it fires while the exchange is waiting for a reply
// that is never released, rather than racing the validation that runs before it.
await Assert.ThrowsAnyAsync<OperationCanceledException>(
() => device.CallAsync(
() => cancellation.CancelAfter(50),
cancellationToken: cancellation.Token));

device.Disconnect();
}

private static void AssertRecognisedTheReply(LineFramingTestableDevice device)
{
// The exchange reports which branch its reply wait loop left by, so this asks it rather
Expand Down Expand Up @@ -207,12 +334,16 @@ public LineFramingTestableDevice(string name, IStreamTransport transport)
/// </summary>
internal override void OnReplyWaitCompleted(bool sawResponse) => RecognisedTheReply = sawResponse;

public Task<IReadOnlyList<string>> CallAsync(Action setupAction)
public Task<IReadOnlyList<string>> CallAsync(
Action setupAction,
int? completionTimeoutMs = null,
CancellationToken cancellationToken = default)
{
return ExecuteTextCommandAsync(
setupAction,
responseTimeoutMs: ResponseTimeoutMs,
completionTimeoutMs: CompletionTimeoutMs);
completionTimeoutMs: completionTimeoutMs ?? CompletionTimeoutMs,
cancellationToken: cancellationToken);
}
}

Expand All @@ -227,10 +358,29 @@ private sealed class ScriptedReplyTransport : IStreamTransport
private bool _disposed;

public ScriptedReplyTransport(string reply)
: this(stageGapMs: int.MaxValue, reply)
{
_stream = new ScriptedReplyStream(reply);
}

/// <summary>
/// A reply the device sends in fragments: each stage is handed over on the first read that
/// finds <paramref name="stageGapMs"/> elapsed since the previous stage drained, so the gap
/// between fragments is measured on the stream's clock rather than counted in the reader
/// thread's sleeps — one late read can lengthen a gap, but it cannot be multiplied into it.
/// </summary>
public ScriptedReplyTransport(int stageGapMs, params string[] stages)
{
_stream = new ScriptedReplyStream(stages, stageGapMs);
}

/// <summary>
/// How long the stream took to hand over every stage: the elapsed time from the first
/// fragment draining to the last one draining, or <see cref="TimeSpan.Zero"/> if fewer than
/// two fragments were ever handed over. This is what the fragmented-reply test asserts is
/// longer than one completion window.
/// </summary>
public TimeSpan StagedSpan => _stream.StagedSpan;

public Stream Stream => _disposed
? throw new ObjectDisposedException(nameof(ScriptedReplyTransport))
: _stream;
Expand Down Expand Up @@ -274,12 +424,40 @@ public void Dispose()

private sealed class ScriptedReplyStream : Stream
{
private readonly byte[] _reply;
private readonly byte[][] _stages;
private readonly TimeSpan _stageGap;
private readonly Stopwatch _clock = Stopwatch.StartNew();
private readonly object _gate = new();
private bool _released;
private int _stage;
private int _position;
private TimeSpan _stageDrainedAt;
private TimeSpan _firstStageDrainedAt;
private TimeSpan _lastStageDrainedAt;
private bool _anyStageDrained;

public ScriptedReplyStream(string[] stages, int stageGapMs)
{
_stages = new byte[stages.Length][];
for (var i = 0; i < stages.Length; i++)
{
_stages[i] = Encoding.ASCII.GetBytes(stages[i]);
}

_stageGap = TimeSpan.FromMilliseconds(stageGapMs);
}

public ScriptedReplyStream(string reply) => _reply = Encoding.ASCII.GetBytes(reply);
/// <summary>See <see cref="ScriptedReplyTransport.StagedSpan"/>.</summary>
public TimeSpan StagedSpan
{
get
{
lock (_gate)
{
return _lastStageDrainedAt - _firstStageDrainedAt;
}
}
}

public void Release()
{
Expand All @@ -301,17 +479,47 @@ public override int Read(byte[] buffer, int offset, int count)
{
lock (_gate)
{
if (_released && _position < _reply.Length)
if (_released && _stage < _stages.Length)
{
var toCopy = Math.Min(count, _reply.Length - _position);
Array.Copy(_reply, _position, buffer, offset, toCopy);
_position += toCopy;
return toCopy;
var current = _stages[_stage];
if (_position < current.Length)
{
var toCopy = Math.Min(count, current.Length - _position);
Array.Copy(current, _position, buffer, offset, toCopy);
_position += toCopy;

if (_position == current.Length)
{
// Handed over in full, so the gap to the next one starts here. The
// same instant is the one the reader turns into a line, which is
// what makes this the right clock to pace against.
_stageDrainedAt = _clock.Elapsed;
if (!_anyStageDrained)
{
_anyStageDrained = true;
_firstStageDrainedAt = _stageDrainedAt;
}

_lastStageDrainedAt = _stageDrainedAt;
}

return toCopy;
}

// This stage is drained. The next one is due once the gap has really
// elapsed on this clock — not once some number of sleeps have returned, a
// count that turns each stretched sleep on a loaded runner into a
// proportionally stretched gap.
if (_clock.Elapsed - _stageDrainedAt >= _stageGap)
{
_stage++;
_position = 0;
}
}
}

// Idle link: nothing to hand over, and no busy-spin in the reader thread.
Thread.Sleep(10);
Thread.Sleep(IdleReadSleepMs);
return 0;
}

Expand Down
Loading
Loading