diff --git a/src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj b/src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj index fa402d5..3511fb6 100644 --- a/src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj +++ b/src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj @@ -25,6 +25,14 @@ <_Parameter1>CanKit.Pro.CANopen + + + <_Parameter1>CanKit.Pro.IsoTp + + + <_Parameter1>CanKit.Pro.Uds + diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs index d23122a..cc39ddd 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs @@ -6,6 +6,7 @@ using System.Threading.Tasks; using CanKit.Abstractions.API.Can.Definitions; using CanKit.Abstractions.API.Common.Definitions; +using CanKit.Pro.Actor; using CanKit.Pro.RawCan; namespace CanKit.Pro.IsoTp; @@ -48,6 +49,8 @@ public sealed class IsoTpFunctionalClient : IDisposable private readonly ICanBusService _service; private readonly bool _ownsService; private readonly IsoTpFunctionalOptions _options; + private readonly ProtocolActor? _clock; + private readonly ITimeSource _time; // Endpoint used only for building the outbound SF payload (Normal addressing, no AE byte). private readonly IsoTpEndpoint _txEndpoint; @@ -63,13 +66,21 @@ public sealed class IsoTpFunctionalClient : IDisposable private int _disposed; + /// + /// A null is production: windows end with + /// and deadlines are + /// readings. An injected actor is the test seam (#171). Its timers + /// and this client's deadlines share that actor's clock, so a test ends a window by + /// advancing the clock instead of sleeping. The actor stays the caller's to dispose. + /// internal IsoTpFunctionalClient( ICanBusService service, uint functionalTxCanId, uint responseRxCanIdRangeStart, uint responseRxCanIdRangeEnd, IsoTpFunctionalOptions options, - bool ownsService) + bool ownsService, + ProtocolActor? clock = null) { if (responseRxCanIdRangeEnd < responseRxCanIdRangeStart) throw new ArgumentOutOfRangeException(nameof(responseRxCanIdRangeEnd), @@ -78,6 +89,8 @@ internal IsoTpFunctionalClient( _service = service ?? throw new ArgumentNullException(nameof(service)); _options = options ?? throw new ArgumentNullException(nameof(options)); _ownsService = ownsService; + _clock = clock; + _time = clock?.TimeSource ?? MonotonicTimeSource.Instance; _txEndpoint = IsoTpEndpoint.Normal(functionalTxCanId, 0, options.IsExtendedCanId); @@ -240,7 +253,7 @@ public async Task> CollectResponsesAsync( public IsoTpFunctionalListener Listen() { ThrowIfDisposed(); - return new IsoTpFunctionalListener(_service.Subscribe(_responseFilter, includeEcho: true)); + return new IsoTpFunctionalListener(_service.Subscribe(_responseFilter, includeEcho: true), _clock); } /// @@ -302,7 +315,7 @@ private async Task SendSingleFrameAsync(ReadOnlyMemory> CollectFromSubscriptionAsync( + private async Task> CollectFromSubscriptionAsync( ISubscription sub, TimeSpan window, CancellationToken cancellationToken) { var responses = new List(); @@ -311,16 +324,16 @@ private static async Task> CollectFromSub // The window's end is also held as an arrival stamp: the timer's callback and this // method's continuations are scheduling, and a frame that arrived after the deadline // but before they ran is not the window's (Codex on #150). - long deadline = Stopwatch.GetTimestamp() + (long)(window.TotalSeconds * Stopwatch.Frequency); - using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - windowCts.CancelAfter(window); - var windowToken = windowCts.Token; + long deadline = _time.GetTimestamp() + (long)(window.TotalSeconds * _time.Frequency); + using var windowEnd = new FunctionalWindow(_clock, window, cancellationToken); + var windowToken = windowEnd.Token; try { await foreach (var frameEvent in sub.Frames.WithCancellation(windowToken).ConfigureAwait(false)) { - if (TryParseFunctionalResponse(frameEvent, out var response) && response!.HostArrivalTimestamp <= deadline) + if (TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp) + && response!.HostArrivalTimestamp <= deadline) responses.Add(response); } } @@ -340,7 +353,8 @@ private static async Task> CollectFromSub sub.Dispose(); while (sub.TryRead(out var frameEvent)) { - if (TryParseFunctionalResponse(frameEvent, out var response) && response!.HostArrivalTimestamp <= deadline) + if (TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp) + && response!.HostArrivalTimestamp <= deadline) responses.Add(response); } } @@ -359,7 +373,7 @@ private static void DrainBuffered(ISubscription sub) } internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent, - out IsoTpFunctionalResponse? response) + out IsoTpFunctionalResponse? response, Func now) { var frame = frameEvent.Frame; var payload = frame.Data.ToArray(); @@ -389,8 +403,9 @@ internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent, var pdu = new byte[pci.Length]; Array.Copy(payload, pci.DataOffset, pdu, 0, pci.Length); - // Stamped by the demux at arrival; "now" only for an event built without a stamp. - var arrival = frameEvent.HostArrivalTimestamp > 0 ? frameEvent.HostArrivalTimestamp : Stopwatch.GetTimestamp(); + // Stamped by the demux at arrival; "now", on the collector's clock, only for an event + // built without a stamp. + var arrival = frameEvent.HostArrivalTimestamp > 0 ? frameEvent.HostArrivalTimestamp : now(); response = new IsoTpFunctionalResponse((uint)frame.ID, pdu, arrival); return true; } diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs index fa883e4..4a8305c 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs @@ -3,6 +3,7 @@ using System.Diagnostics; using System.Threading; using System.Threading.Tasks; +using CanKit.Pro.Actor; using CanKit.Pro.RawCan; namespace CanKit.Pro.IsoTp; @@ -27,15 +28,19 @@ public sealed class IsoTpFunctionalListener : IDisposable private static readonly TimeSpan MaxWindow = TimeSpan.FromMilliseconds(uint.MaxValue - 1); private readonly ISubscription _subscription; + private readonly ProtocolActor? _clock; + private readonly ITimeSource _time; // Arrived after a collection's deadline but before its drain read the buffer: kept for the // next collection, whose window they are in (Codex on #150). private readonly Queue _carried = new(); private bool _ended; private int _disposed; - internal IsoTpFunctionalListener(ISubscription subscription) + internal IsoTpFunctionalListener(ISubscription subscription, ProtocolActor? clock = null) { _subscription = subscription; + _clock = clock; + _time = clock?.TimeSource ?? MonotonicTimeSource.Instance; } /// @@ -68,7 +73,7 @@ public async Task> CollectAsync( if (window > MaxWindow) throw new ArgumentOutOfRangeException(nameof(window), window, "The collection window exceeds what a timer can measure (about 49 days)."); - long now = Stopwatch.GetTimestamp(); + long now = _time.GetTimestamp(); var responses = new List(_carried); _carried.Clear(); if (window <= TimeSpan.Zero) @@ -80,16 +85,15 @@ public async Task> CollectAsync( TakeBuffered(responses, now); return responses.AsReadOnly(); } - long deadline = now + (long)(window.TotalSeconds * Stopwatch.Frequency); - using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - windowCts.CancelAfter(window); + long deadline = now + (long)(window.TotalSeconds * _time.Frequency); + using var windowEnd = new FunctionalWindow(_clock, window, cancellationToken); try { // Ends when the window runs out, or when the subscription is completed underneath // -- the service disposed -- in which case what is buffered is still read, and // the next collection throws rather than return empty at once: a loop that // collects while a deadline lasts would otherwise spin on it (Bugbot on #150). - while (await _subscription.WaitToReadAsync(windowCts.Token).ConfigureAwait(false)) + while (await _subscription.WaitToReadAsync(windowEnd.Token).ConfigureAwait(false)) TakeBuffered(responses, deadline); _ended = true; } @@ -109,7 +113,7 @@ private void TakeBuffered(List responses, long deadline { while (_subscription.TryRead(out var frameEvent)) { - if (!IsoTpFunctionalClient.TryParseFunctionalResponse(frameEvent, out var response)) continue; + if (!IsoTpFunctionalClient.TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp)) continue; if (response!.HostArrivalTimestamp <= deadline) responses.Add(response); else _carried.Enqueue(response); } @@ -122,3 +126,54 @@ public void Dispose() _subscription.Dispose(); } } + +/// +/// A functional collection window: its token is cancelled by the caller's token or when the +/// window ends. Production ends it with . +/// A test that injected an actor arms the same interval on that actor's clock instead, so the +/// window follows the clock and not the runner (#171). +/// +internal sealed class FunctionalWindow : IDisposable +{ + private readonly CancellationTokenSource _linked; + private readonly IDisposable? _timer; + private readonly object _gate = new(); + private bool _disposed; + + public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken cancellationToken) + { + _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + if (actor is null) + { + _linked.CancelAfter(window); + return; + } + + _timer = actor.Schedule(window, End); + } + + public CancellationToken Token => _linked.Token; + + /// + /// The injected actor's timer. Cancelling that timer only flags it, so a callback the loop + /// has already taken can run after ; under the same lock it then finds + /// the window disposed and does nothing. Nothing here depends on the actor still running. + /// + internal void End() + { + lock (_gate) + { + if (!_disposed) _linked.Cancel(); + } + } + + public void Dispose() + { + _timer?.Dispose(); + lock (_gate) + { + _disposed = true; + _linked.Dispose(); + } + } +} diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index a285d83..b7f21ad 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs @@ -51,6 +51,7 @@ internal sealed class J1939TpChannel : IJ1939TpChannel private readonly J1939TpOptions _options; private readonly ProtocolActor _actor; + private readonly bool _ownsActor; private readonly DeadlineScheduler _deadlines; private readonly ISubscription _subscription; private readonly Task _readerTask; @@ -89,7 +90,22 @@ internal sealed class J1939TpChannel : IJ1939TpChannel /// public event EventHandler? BackgroundExceptionOccurred; - internal J1939TpChannel(ICanBusService service, byte sourceAddress, J1939TpOptions options, bool ownsService) + /// + /// Builds a channel on . A null -- the only + /// value production passes -- makes the channel create and own its own loop; an injected one + /// stays the caller's to dispose. + /// + /// + /// A seam for tests, not a feature: substituting a loop built on a hand-driven monotonic + /// source is what lets T1..T4 and the BAM packet spacing be asserted without measuring + /// wall-clock gaps on a shared runner (#171). Every one of those intervals is armed through + /// or the built on this + /// actor, and the clock is that actor's time source. The channel does not take one of its + /// own, so the loop and the timers cannot be given + /// different clocks -- the same rule J1939NodeImpl follows. + /// + internal J1939TpChannel(ICanBusService service, byte sourceAddress, J1939TpOptions options, bool ownsService, + ProtocolActor? actor = null) { _service = service ?? throw new ArgumentNullException(nameof(service)); if (sourceAddress == J1939Pgn.GlobalAddress) @@ -108,7 +124,8 @@ internal J1939TpChannel(ICanBusService service, byte sourceAddress, J1939TpOptio }; _pduInbox = Channel.CreateBounded(inboxOptions); - _actor = new ProtocolActor(); + _ownsActor = actor is null; + _actor = actor ?? new ProtocolActor(); _actor.BackgroundExceptionOccurred += OnActorBackgroundException; _deadlines = new DeadlineScheduler(_actor); @@ -142,7 +159,8 @@ internal J1939TpChannel(ICanBusService service, byte sourceAddress, J1939TpOptio } catch { - _actor.Dispose(); + _actor.BackgroundExceptionOccurred -= OnActorBackgroundException; + if (_ownsActor) _actor.Dispose(); throw; } @@ -300,41 +318,86 @@ public void Dispose() // down the reader task. _pduInbox.Writer.TryComplete(); + // The reader is joined and the subscription closed before the sessions are failed: a + // frame the reader was still handing to the actor is then queued ahead of the cleanup, + // not behind it where it could open a session nobody fails (Bugbot on #183). + try { _readerTask.Wait(TimeSpan.FromSeconds(2)); } catch (AggregateException) { /* observed via task; not fatal */ } + _subscription.Dispose(); + // Cancel every still-in-flight session on the actor so their TCSs get an - // ObjectDisposedException instead of hanging on the now-disposed inbox. - try + // ObjectDisposedException instead of hanging on the now-disposed inbox. Disposed from + // the actor's own loop, the sessions are this thread's to fail, and a post would only + // run after this call returned (Codex on #183). + var cleanup = Task.CompletedTask; + if (_actor.IsOnCurrentActor) { - _actor.Post(() => - { - foreach (var kv in _txSessions) - kv.Value.Fail(new ObjectDisposedException(nameof(J1939TpChannel))); - _txSessions.Clear(); - foreach (var kv in _txQueues) - { - foreach (var pending in kv.Value) - pending.Tcs.TrySetException(new ObjectDisposedException(nameof(J1939TpChannel))); - } - _txQueues.Clear(); - foreach (var kv in _rxSessions) - kv.Value.Cancel(); - _rxSessions.Clear(); - }); + FailSessionsOnDispose(); } - catch (ObjectDisposedException) + else { - // actor already gone; nothing more to do + try + { + cleanup = _actor.PostAsync(FailSessionsOnDispose); + } + catch (ObjectDisposedException) + { + // An injected actor its owner already disposed took its sessions with it. They + // are actor state, and its loop may still be draining, so this thread leaves + // them alone (Bugbot on #183). + } } - try { _readerTask.Wait(TimeSpan.FromSeconds(2)); } catch { /* observed via task; not fatal */ } - - _subscription.Dispose(); - _actor.Dispose(); + // An injected actor is not ours to dispose -- the caller may still be running other + // work on it -- but the handler is, so it comes off either way. + _actor.BackgroundExceptionOccurred -= OnActorBackgroundException; _readerCts.Dispose(); + if (_ownsActor) + { + // Its loop runs the cleanup as it drains. + _actor.Dispose(); + DisposeOwnedService(); + } + else if (cleanup.Wait(TimeSpan.FromSeconds(2))) + { + // A borrowed actor keeps running, possibly busy with the caller's other work, so the + // cleanup is waited for: no session is left using the service (Codex on #183). + DisposeOwnedService(); + } + else + { + // Still busy past the budget. Dispose does not block on the caller's work any longer, + // and says so where background faults go, as the actor's own Dispose does; the + // service outlives the cleanup rather than being torn down under it (Codex on #183). + RaiseBackgroundException(new TimeoutException( + "The injected actor did not run the channel's session cleanup within 2 s. Dispose returned; the sends still in flight fail, and an owned service is disposed, once the actor gets to it.")); + _ = cleanup.ContinueWith(static (_, state) => ((J1939TpChannel)state!).DisposeOwnedService(), this, + CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default); + } + } + + private void DisposeOwnedService() + { if (_ownsService) _service.Dispose(); } + private void FailSessionsOnDispose() + { + foreach (var kv in _txSessions) + kv.Value.Fail(new ObjectDisposedException(nameof(J1939TpChannel))); + _txSessions.Clear(); + foreach (var kv in _txQueues) + { + foreach (var pending in kv.Value) + pending.Tcs.TrySetException(new ObjectDisposedException(nameof(J1939TpChannel))); + } + _txQueues.Clear(); + foreach (var kv in _rxSessions) + kv.Value.Cancel(); + _rxSessions.Clear(); + } + // ----------------------------------------------------------------------------------------- // Subscription reader -- one long-lived Task per channel that pushes RX frames onto the actor. // ----------------------------------------------------------------------------------------- @@ -406,6 +469,10 @@ private async Task RunReaderAsync() private void HandleIncoming(uint pgn, byte sa, byte da, byte[] payload) { + // A frame the reader handed over is dropped once Dispose has begun, whenever the actor + // gets to it: a reader that outlived its join could otherwise post one behind the + // session cleanup and open a session nobody fails on a borrowed actor (Codex on #183). + if (Volatile.Read(ref _disposed) != 0) return; try { if (J1939Pgn.IsTransportCm(pgn)) diff --git a/src/CanKit.Pro.RawCan/CanBusService.cs b/src/CanKit.Pro.RawCan/CanBusService.cs index a1e99fe..dd4ce8e 100644 --- a/src/CanKit.Pro.RawCan/CanBusService.cs +++ b/src/CanKit.Pro.RawCan/CanBusService.cs @@ -2,7 +2,6 @@ using System.Collections.Generic; using System.Diagnostics; using System.Linq; -using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; using CanKit.Abstractions.API.Can; @@ -64,13 +63,32 @@ public sealed class CanBusService : ICanBusService private int _disposed; + // Host arrival and transmit stamps. Production reads Stopwatch; a test passes the same + // counter the protocol window is measured on, so a frame and the deadline that admits it + // are not two clocks (#171). + private readonly Func _hostTimestamp; + /// /// Creates a service that demultiplexes . Attaches to the bus's /// event immediately. /// public CanBusService(ICanBus bus) + : this(bus, hostTimestamp: null) + { + } + + /// + /// As the public constructor, but supplies the host + /// arrival and transmit stamps. Null keeps . + /// + /// + /// Internal on purpose. The stamps are facts about one clock, and the only caller that + /// needs to choose which clock is a test driving a functional window by hand (#171). + /// + internal CanBusService(ICanBus bus, Func? hostTimestamp) { _bus = bus ?? throw new ArgumentNullException(nameof(bus)); + _hostTimestamp = hostTimestamp ?? Stopwatch.GetTimestamp; _bus.FrameObserved += OnFrameObserved; _bus.FaultOccurred += OnFaultOccurred; } @@ -163,7 +181,7 @@ private void OnFrameObserved(object? sender, CanReceiveDataView e) // per-subscription buffers below can hold a frame while their reader is descheduled. // Either would make a punctual frame look late to whoever is enforcing a deadline on // it (Codex on #112). - var hostArrival = Stopwatch.GetTimestamp(); + var hostArrival = _hostTimestamp(); // Independent of subscription dispatch below: echo frames must be checked against // outstanding SendConfirmed calls regardless of whether anyone also has a @@ -321,13 +339,13 @@ private async Task SendApproximatedAsync(CanFrame frame, Cancell // runs again, and a reading taken there starts a caller's response deadline late // (Codex on #112). ExecuteSynchronously is what observes completion closest; when the // runtime declines to inline it the reading is what it would have been anyway. - var stamp = new StrongBox(); - var handoffStart = Stopwatch.GetTimestamp(); + var stamp = new HandoffStamp(_hostTimestamp); + var handoffStart = stamp.Now(); var accepted = await _bus.TransmitAsync(frame, cancellationToken) .ContinueWith( static (completed, state) => { - ((StrongBox)state!).Value = Stopwatch.GetTimestamp(); + ((HandoffStamp)state!).Value = ((HandoffStamp)state!).Now(); // GetResult rather than .Result: it surfaces a driver fault or a // cancellation as itself instead of wrapping it in an AggregateException, @@ -387,7 +405,7 @@ private async Task SendWithEchoConfirmAsync(CanFrame frame, Time RegisterPending(pending); // The other end of the driver call, inside the lock: a caller's cutoff for // what can still be a response to this frame (Codex on #147). - handoffStart = Stopwatch.GetTimestamp(); + handoffStart = _hostTimestamp(); accepted = _bus.Transmit(in frame); // Taken here, inside the lock and immediately after the driver call returns: @@ -396,7 +414,7 @@ private async Task SendWithEchoConfirmAsync(CanFrame frame, Time // for this very lock, which another send holds across its own Transmit -- and // would start a caller's response deadline while the request was still // queued behind it (Codex on #112). - handoff = Stopwatch.GetTimestamp(); + handoff = _hostTimestamp(); } } catch @@ -604,5 +622,20 @@ private void OnFaultOccurred(object? sender, Exception ex) }); } } + + /// + /// The stamp taken on the thread that completes the driver's task. A static continuation + /// cannot close over the service, so the delegate travels with the box. + /// + private sealed class HandoffStamp + { + private readonly Func _now; + + public HandoffStamp(Func now) => _now = now; + + public long Value; + + public long Now() => _now(); + } } } diff --git a/src/CanKit.Pro.Uds/SuppressedResponseWindows.cs b/src/CanKit.Pro.Uds/SuppressedResponseWindows.cs index df54f50..d600d08 100644 --- a/src/CanKit.Pro.Uds/SuppressedResponseWindows.cs +++ b/src/CanKit.Pro.Uds/SuppressedResponseWindows.cs @@ -21,7 +21,18 @@ internal sealed class SuppressedResponseWindows /// Notes that a request for went out at /// and may be answered for . public void Note(byte sid, long sentTimestamp, TimeSpan window) - => Extend(sid, sentTimestamp + (long)(window.TotalSeconds * Stopwatch.Frequency)); + => Note(sid, sentTimestamp, window, Stopwatch.Frequency); + + /// + /// As , counting in + /// so the deadline is on the same clock as + /// (#171). + /// + public void Note(byte sid, long sentTimestamp, TimeSpan window, long ticksPerSecond) + => Extend(sid, sentTimestamp + Ticks(window, ticksPerSecond)); + + internal static long Ticks(TimeSpan window, long ticksPerSecond) + => (long)(window.TotalSeconds * ticksPerSecond); /// Moves the window for out to , if later. public void Extend(byte sid, long until) @@ -76,8 +87,15 @@ public void Forget(byte sid) /// How long from now until , or zero if it has passed. public static TimeSpan Remaining(long until) + => Remaining(until, Stopwatch.GetTimestamp(), Stopwatch.Frequency); + + /// + /// As , measured from on a clock whose + /// frequency is . + /// + public static TimeSpan Remaining(long until, long now, long ticksPerSecond) { - var ticks = until - Stopwatch.GetTimestamp(); - return ticks <= 0 ? TimeSpan.Zero : TimeSpan.FromSeconds((double)ticks / Stopwatch.Frequency); + var ticks = until - now; + return ticks <= 0 ? TimeSpan.Zero : TimeSpan.FromSeconds((double)ticks / ticksPerSecond); } } diff --git a/src/CanKit.Pro.Uds/UdsFunctionalClient.cs b/src/CanKit.Pro.Uds/UdsFunctionalClient.cs index 684b4e3..d4633db 100644 --- a/src/CanKit.Pro.Uds/UdsFunctionalClient.cs +++ b/src/CanKit.Pro.Uds/UdsFunctionalClient.cs @@ -1,9 +1,9 @@ using System; using System.Collections.Generic; -using System.Diagnostics; using System.Linq; using System.Threading; using System.Threading.Tasks; +using CanKit.Pro.Actor; using CanKit.Pro.IsoTp; namespace CanKit.Pro.Uds; @@ -31,6 +31,11 @@ public sealed class UdsFunctionalClient : IDisposable private readonly bool _ownsClient; private readonly TimeSpan _responseWindow; private readonly TimeSpan _responsePendingWindow; + // Null in production: deadlines are Stopwatch readings and a listener delay is Task.Delay. + // A test injects the actor whose clock those readings and the functional collection windows + // already share, so the same advance ends both (#171). Not disposed here. + private readonly ProtocolActor? _clock; + private readonly ITimeSource _time; private readonly SuppressedResponseWindows _openWindows = new(); // Per service: the standing subscription that hears the window out, and the task reading it. private readonly Dictionary _listeners = new(); @@ -52,14 +57,21 @@ public sealed class UdsFunctionalClient : IDisposable private readonly CancellationTokenSource _lifetimeCts = new(); // Test hook: how late the listener's worker starts reading its subscription, standing in - // for a thread pool that schedules it after the window has run out (Codex on #150). - internal TimeSpan ListenerStartDelay { get; set; } + // for a thread pool that schedules it after the window has run out (Codex on #150). Only + // on an injected clock: a wall-clock start delay is the kind of sleep #171 removes. + internal void DelayListenerStart(TimeSpan delay) + => _listenerStartDelay = _clock is not null + ? delay + : throw new InvalidOperationException("A listener start delay is measured on an injected clock (#171)."); + private TimeSpan _listenerStartDelay; private int _disposed; private UdsFunctionalClient(IsoTpFunctionalClient client, bool ownsClient, TimeSpan responseWindow, - TimeSpan responsePendingWindow) + TimeSpan responsePendingWindow, ProtocolActor? clock) { _client = client ?? throw new ArgumentNullException(nameof(client)); + _clock = clock; + _time = clock?.TimeSource ?? MonotonicTimeSource.Instance; // Bounded above as a collection window is: a listener collecting for longer than a // timer measures would fault, and the window it owned go unwaited (Codex on #150). if (responseWindow <= TimeSpan.Zero || responseWindow > MaxCollectionWindow) @@ -84,8 +96,19 @@ private UdsFunctionalClient(IsoTpFunctionalClient client, bool ownsClient, TimeS /// public static UdsFunctionalClient Create(IsoTpFunctionalClient client, bool ownsClient = false, TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null) + => Create(client, clock: null, ownsClient, responseWindow, responsePendingWindow); + + /// + /// As , measuring P2 + /// and P2* on ; null is the wall clock, as the public overload uses. + /// The functional client underneath must have been opened on that same actor, and the demux + /// must stamp frames with its time source, or a deadline and an arrival are not comparable. + /// The actor is not disposed with this client. + /// + internal static UdsFunctionalClient Create(IsoTpFunctionalClient client, ProtocolActor? clock, + bool ownsClient = false, TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null) => new(client, ownsClient, responseWindow ?? UdsClientOptions.DefaultP2, - responsePendingWindow ?? UdsClientOptions.DefaultP2Star); + responsePendingWindow ?? UdsClientOptions.DefaultP2Star, clock); /// The underlying ISO-TP functional client. public IsoTpFunctionalClient Channel => _client; @@ -148,7 +171,7 @@ private async Task> SendRawLockedAsync(Read bool hadWindow; long previousUntil; lock (_listeners) hadWindow = _openWindows.TryGetDeadline(sid, out previousUntil); - StartListening(sid, Stopwatch.GetTimestamp(), inFlight: true); + StartListening(sid, Now(), inFlight: true); IsoTpTransmitStamps stamps; try { @@ -163,7 +186,7 @@ private async Task> SendRawLockedAsync(Read } catch { - StartListening(sid, Stopwatch.GetTimestamp(), inFlight: false); + StartListening(sid, Now(), inFlight: false); throw; } // Through StartListening again: the window is moved out to the transmission -- the @@ -200,7 +223,7 @@ private async Task> SendRawLockedAsync(Read bool hadWindowBefore; long previousUntilBefore; lock (_listeners) hadWindowBefore = _openWindows.TryGetDeadline(sid, out previousUntilBefore); - StartListening(sid, Stopwatch.GetTimestamp(), inFlight: true); + StartListening(sid, Now(), inFlight: true); IsoTpFunctionalCollection collected; try { @@ -219,7 +242,7 @@ private async Task> SendRawLockedAsync(Read // the request on the bus; the ECUs' P2 from the transmission is at most P2 from // now, so that window is noted before the exception leaves, and the listener, // kept through the send, has heard any 0x78 the collection lost (Codex on #150). - StartListening(sid, Stopwatch.GetTimestamp(), inFlight: false); + StartListening(sid, Now(), inFlight: false); throw; } var req = request.Span; @@ -239,7 +262,7 @@ private async Task> SendRawLockedAsync(Read long cutoff = collected.TransmitStamps.LastFrameHandoffTimestamp; long transmitted = collected.TransmitStamps.LastFrameTransmitTimestamp > 0 ? collected.TransmitStamps.LastFrameTransmitTimestamp - : Stopwatch.GetTimestamp() - Ticks(window); + : Now() - Ticks(window); lock (_listeners) NoteAnchored(sid, transmitted, cutoff); var responses = new List(raw.Count); foreach (var r in raw) @@ -286,7 +309,7 @@ private void StartListening(byte sid, long from, bool inFlight, long cutoff = 0) if (inFlight) { _inFlight[sid] = new List(); - _openWindows.Note(sid, from, _responseWindow); + _openWindows.Note(sid, from, _responseWindow, _time.Frequency); } else { @@ -296,8 +319,8 @@ private void StartListening(byte sid, long from, bool inFlight, long cutoff = 0) } } - private static long TransmittedAt(IsoTpTransmitStamps stamps) - => stamps.LastFrameTransmitTimestamp > 0 ? stamps.LastFrameTransmitTimestamp : Stopwatch.GetTimestamp(); + private long TransmittedAt(IsoTpTransmitStamps stamps) + => stamps.LastFrameTransmitTimestamp > 0 ? stamps.LastFrameTransmitTimestamp : Now(); // The functional client refuses a request before transmitting it -- oversized for a Single // Frame, an argument error -- or because it is disposed; a cancellation or a transport fault @@ -326,7 +349,7 @@ private void RestoreWindow(byte sid, bool had, long until) // out, one from after P2 does not (Codex on #150). private void NoteAnchored(byte sid, long from, long cutoff) { - _openWindows.Note(sid, from, _responseWindow); + _openWindows.Note(sid, from, _responseWindow, _time.Frequency); if (cutoff > 0) _cutoffs[sid] = cutoff; else _cutoffs.Remove(sid); if (!_inFlight.TryGetValue(sid, out var heard)) return; _inFlight.Remove(sid); @@ -348,7 +371,7 @@ private void EnsureListener(byte sid) { if (_listeners.ContainsKey(sid)) return; if (!_openWindows.TryGetDeadline(sid, out var until) - || SuppressedResponseWindows.Remaining(until) <= TimeSpan.Zero) + || RemainingUntil(until) <= TimeSpan.Zero) { _openWindows.Forget(sid); return; @@ -375,8 +398,8 @@ private async Task ListenAsync(byte sid, IsoTpFunctionalListener ears) { try { - if (ListenerStartDelay > TimeSpan.Zero) - await Task.Delay(ListenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); + if (_listenerStartDelay > TimeSpan.Zero) + await WaitOnClockAsync(_clock!, _listenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); bool drained = false; while (true) { @@ -387,7 +410,7 @@ private async Task ListenAsync(byte sid, IsoTpFunctionalListener ears) // and the entry goes with it: a note after this reads no listener and // starts one, a note before it moved the deadline this reads. if (!_openWindows.TryGetDeadline(sid, out var until) - || (remaining = SuppressedResponseWindows.Remaining(until)) <= TimeSpan.Zero) + || (remaining = RemainingUntil(until)) <= TimeSpan.Zero) { if (_inFlight.ContainsKey(sid)) { @@ -561,7 +584,23 @@ private static bool EchoMatches(ReadOnlySpan request, byte[] response, int return true; } - private static long Ticks(TimeSpan span) => (long)(span.TotalSeconds * Stopwatch.Frequency); + private long Now() => _time.GetTimestamp(); + + private TimeSpan RemainingUntil(long until) + => SuppressedResponseWindows.Remaining(until, Now(), _time.Frequency); + + private long Ticks(TimeSpan span) => SuppressedResponseWindows.Ticks(span, _time.Frequency); + + // The listener's start delay, on the injected actor's clock: a test moves it by advancing + // that clock instead of sleeping (#171). + private static async Task WaitOnClockAsync(ProtocolActor clock, TimeSpan delay, CancellationToken cancellationToken) + { + var done = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var registration = cancellationToken.Register(static state => + ((TaskCompletionSource)state!).TrySetCanceled(), done); + using var handle = clock.Schedule(delay, () => done.TrySetResult(true)); + await done.Task.ConfigureAwait(false); + } private static bool IsResponsePending(byte[] data, byte sid) => IsResponsePending(data) && data[1] == sid; diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index 35e0820..f274ac4 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Channels; using System.Threading.Tasks; @@ -29,6 +30,13 @@ internal sealed class StarvedReaderBusService : ICanBusService public void Deliver(CanFrameView frame, long hostArrivalTimestamp = 0) => _frames.Writer.TryWrite( new CanFrameEvent(frame, isEcho: false, TimeSpan.Zero, hostArrivalTimestamp)); + /// + /// Whether the subscription's Frames enumeration yields nothing until it is cancelled, + /// so what buffered is left for a TryRead drain -- the state a + /// collection window can end in when its reader was descheduled (#171). + /// + public bool HoldFrames { get; set; } + /// Lets the reader task's wait complete once; every later wait stays pending. public void WakeReader() => _wake.TrySetResult(true); @@ -55,6 +63,12 @@ public ISubscription Subscribe(CanIdFilter filter, int? bufferCapacity = null, bool includeEcho = false) => new Sub(this); + /// + /// Runs inside before it confirms: the time a driver takes to + /// confirm, modelled by moving a virtual clock (#171). + /// + public Action? OnSendConfirmed { get; set; } + /// Every frame the channel put on the wire, in order. public List Sent { get; } = new(); @@ -62,6 +76,7 @@ public Task SendConfirmed(CanFrame frame, TimeSpan? timeout = nu CancellationToken cancellationToken = default) { lock (Sent) Sent.Add(frame.Data.ToArray()); + OnSendConfirmed?.Invoke(); return Task.FromResult(new TxConfirmation { Confirmed = true }); } @@ -75,7 +90,15 @@ private sealed class Sub : ISubscription private readonly StarvedReaderBusService _owner; public Sub(StarvedReaderBusService owner) => _owner = owner; - public IAsyncEnumerable Frames => _owner._frames.Reader.ReadAllAsync(); + public IAsyncEnumerable Frames + => _owner.HoldFrames ? Held() : _owner._frames.Reader.ReadAllAsync(); + + private static async IAsyncEnumerable Held( + [EnumeratorCancellation] CancellationToken cancellationToken = default) + { + await Task.Delay(Timeout.Infinite, cancellationToken); + yield break; + } public bool TryRead(out CanFrameEvent frameEvent) { diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs index 745f257..e2e8faa 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs @@ -8,6 +8,7 @@ using CanKit.Abstractions.API.Can.Definitions; using CanKit.Abstractions.API.Common.Definitions; using CanKit.Core; +using CanKit.Pro.Actor; using CanKit.Pro.IsoTp; using CanKit.Pro.RawCan; using CanKit.Pro.Tests.Infrastructure; @@ -660,4 +661,77 @@ public async Task Functional_Listener_Disposal_Ends_A_Collection_In_Progress() Func act = () => listener.CollectAsync(TimeSpan.FromMilliseconds(10)); await act.Should().ThrowAsync(); } + + private static CanFrameView SingleFrameView(uint canId, byte[] pdu) + => new(CanFrameType.Can20, unchecked((int)canId), + IsoTpFrameCodec.BuildSingleFrame(IsoTpEndpoint.Normal(canId, 0), pdu, isCanFd: false, padding: true), + FrameFlags.None); + + // #171: on an injected clock the window ends when that clock says so, and what its reader + // had not yet taken is drained by arrival stamp against the same clock's deadline. The + // subscription yields nothing, so the drain is the only way a frame is collected; the + // three frames sit at the deadline, one tick past it, and unstamped -- stamped at the + // drain, after the clock has left the window. + [Fact] + public async Task Functional_Collect_On_An_Injected_Clock_Drains_By_That_Clocks_Deadline() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + clock.Advance(TimeSpan.FromMilliseconds(1)); // a zero stamp reads as "unstamped" + var service = new StarvedReaderBusService { HoldFrames = true }; + using var client = new IsoTpFunctionalClient(service, 0x7DF, 0x7E8, 0x7EF, FastOptions(), + ownsService: true, actor); + + var window = TimeSpan.FromMilliseconds(100); + var call = client.SendAndCollectAsync(new byte[] { 0x22, 0xF1, 0x90 }, window); + await clock.WaitUntilTimerArmedAsync(actor, window, ShortTimeout); + // The clock has not moved since the collection read it, so its deadline is this. + long deadline = actor.TimeSource.GetTimestamp() + + (long)(window.TotalSeconds * actor.TimeSource.Frequency); + service.Deliver(SingleFrameView(0x7E8, new byte[] { 0x62, 0xF1, 0x90, 0x01 }), deadline); + service.Deliver(SingleFrameView(0x7E9, new byte[] { 0x62, 0xF1, 0x90, 0x02 }), deadline + 1); + service.Deliver(SingleFrameView(0x7EA, new byte[] { 0x62, 0xF1, 0x90, 0x03 })); + await clock.AdvanceAsync(window + TimeSpan.FromMilliseconds(1)); + + var responses = await call.WaitAsync(ShortTimeout); + responses.Should().ContainSingle().Which.SourceCanId.Should().Be(0x7E8u, + "only the frame stamped inside the window, on the window's own clock, is this collection's"); + } + + // #171: a frame event built without a host stamp is stamped from the collector's clock, + // not from Stopwatch, or it could not be compared with a deadline on an injected clock. + [Fact] + public void An_Unstamped_Functional_Response_Is_Stamped_From_The_Collectors_Clock() + { + var unstamped = new CanFrameEvent(SingleFrameView(0x7E8, new byte[] { 0x7E, 0x00 }), + isEcho: false, TimeSpan.Zero); + + IsoTpFunctionalClient.TryParseFunctionalResponse(unstamped, out var response, () => 42) + .Should().BeTrue(); + response!.HostArrivalTimestamp.Should().Be(42); + } + + // Codex and Bugbot on #183: a clocked window is disposed without the actor's help, so an + // actor already torn down does not turn a collection's end into ObjectDisposedException, + // and a timer callback the loop took before the disposal finds the window gone and does + // nothing. + [Fact] + public async Task A_Clocked_Window_Ends_On_Its_Clock_And_Outlives_Its_Actor() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + + using (var ended = new FunctionalWindow(actor, TimeSpan.FromMilliseconds(10), CancellationToken.None)) + { + await clock.WaitUntilTimerArmedAsync(actor, TimeSpan.FromMilliseconds(10), ShortTimeout); + ended.Token.IsCancellationRequested.Should().BeFalse(); + await clock.AdvanceAsync(TimeSpan.FromMilliseconds(10)); + ended.Token.IsCancellationRequested.Should().BeTrue("the window ends when its clock says so"); + } + + using var late = new FunctionalWindow(actor, TimeSpan.FromMilliseconds(10), CancellationToken.None); + actor.Dispose(); + late.Invoking(w => w.Dispose()).Should().NotThrow("the actor is not needed to dispose the window"); + late.Invoking(w => w.End()).Should().NotThrow("a callback taken before the disposal does nothing"); + } } diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index c5e27f8..894b271 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -7,6 +7,7 @@ using CanKit.Abstractions.API.Can.Definitions; using CanKit.Abstractions.API.Common.Definitions; using CanKit.Core; +using CanKit.Pro.Actor; using CanKit.Pro.Addressing; using CanKit.Pro.J1939Tp; using CanKit.Pro.RawCan; @@ -2384,6 +2385,424 @@ public async Task DatagramReceived_Finds_The_Datagram_Already_Receivable() var datagram = await fromHandler.Task.WaitAsync(ShortTimeout); datagram.Payload.Should().Equal(payload); } + + // #171: the channel used to build its own actor on the wall clock, so a BAM spacing could + // only be waited out. An injected actor is the caller's, including when construction fails + // before the channel exists to be asked. + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task A_Failed_Construction_Disposes_Its_Own_Actor_But_Never_An_Injected_One(bool inject) + { + var session = NewSession(); + using var bus = Open(session, 0); + using var inner = new CanBusService(bus); + + using var clock = new VirtualClock(); + var injected = inject ? clock.NewActor() : null; + var failing = new ThrowingSubscribeService(inner, failOnCall: 1); + var loopsBefore = ProtocolActor.RunningLoopCount; + + Action construct = () => new J1939TpChannel(failing, sourceAddress: 0x10, + new J1939TpOptions(), ownsService: false, injected); + construct.Should().Throw(); + + failing.SubscribeCalls.Should().Be(1); + inner.SubscriptionCount.Should().Be(0, "the subscription that threw was never installed"); + if (injected is not null) + { + (await injected.PostAsync(() => 42).WaitAsync(ShortTimeout)).Should().Be(42, + "an injected actor belongs to the caller and must survive a construction that failed"); + } + ProtocolActor.RunningLoopCount.Should().Be(loopsBefore, + inject + ? "the channel created no actor here, so it must not have ended one either" + : "the actor the channel created for itself must be gone once construction fails"); + } + + // The spacing timer starts the send; the frame is handed to the driver on the pool. + private static async Task WaitForTransmitCount(ControllableBus bus, int count) + { + var deadline = DateTime.UtcNow + ShortTimeout; + while (bus.TransmitCount < count) + { + if (DateTime.UtcNow > deadline) + throw new TimeoutException($"TransmitCount stayed at {bus.TransmitCount}; wanted {count}."); + await Task.Yield(); + } + } + + // #171: BAM packet spacing is an actor timer. With the clock frozen, wall time past the + // spacing must not release the next TP.DT; advancing the clock must. + [Fact] + public async Task Bam_Packet_Spacing_Follows_The_Injected_Actors_Clock() + { + var spacing = TimeSpan.FromMilliseconds(50); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = new J1939TpChannel(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: spacing), ownsService: false, actor); + + var send = sender.SendBamAsync(0xFECBu, RandomPayload(9, seed: 171)); + await clock.WaitUntilTimerArmedAsync(actor, spacing, ShortTimeout); + + bus.TransmitCount.Should().Be(1, "only the BAM announce is out before the spacing elapses"); + await Task.Delay(TimeSpan.FromMilliseconds(200)); + await actor.PostAsync(() => 0); + bus.TransmitCount.Should().Be(1, + "wall time past the spacing must not release a TP.DT while the injected clock stands still"); + + await clock.AdvanceAsync(spacing); + // The timer only starts the send; the frame leaves on the pool, after the callback. + await WaitForTransmitCount(bus, 2); + bus.TransmitCount.Should().Be(2, "the first TP.DT goes out once the clock says the spacing elapsed"); + await clock.WaitUntilTimerArmedAsync(actor, spacing, ShortTimeout); + await clock.AdvanceAsync(spacing); + await send.WaitAsync(ShortTimeout); + bus.TransmitCount.Should().Be(3, "nine bytes are a BAM plus two TP.DT frames"); + + sender.Dispose(); + (await actor.PostAsync(() => 7).WaitAsync(ShortTimeout)).Should().Be(7, + "disposing the channel must not dispose the actor it was given"); + } + + private static readonly TimeSpan InFlightSpacing = TimeSpan.FromMilliseconds(50); + + private static J1939TpChannel BamSenderOn(ProtocolActor actor, ICanBusService service) + => new(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: InFlightSpacing), ownsService: false, actor); + + // A BAM on a borrowed, frozen clock: the announce is out and the first TP.DT waits on the + // spacing timer, so the send is in flight until the channel is disposed. A second BAM to the + // same global destination waits in the queue behind it. + private static async Task SendBamsInFlight(VirtualClock clock, ProtocolActor actor, J1939TpChannel sender) + { + var first = sender.SendBamAsync(0xFECBu, RandomPayload(9, seed: 183)); + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + var queued = sender.SendBamAsync(0xFECCu, RandomPayload(9, seed: 184)); + await clock.SettleAsync(); + first.IsCompleted.Should().BeFalse("the send waits on the spacing timer of a frozen clock"); + queued.IsCompleted.Should().BeFalse("the second BAM waits for the first one's session slot"); + return new[] { first, queued }; + } + + // Codex on #183: a borrowed actor is not drained by disposing it, so Dispose must wait for + // its session cleanup itself. With the actor busy on the caller's other work, a Dispose + // that only posted the cleanup would return with the send still in flight. + [Fact] + public async Task Disposing_On_A_Busy_Borrowed_Actor_Fails_The_Send_Before_It_Returns() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + var sends = await SendBamsInFlight(clock, actor, sender); + + // The caller's other work holds the actor for 100 ms, and releases it on its own: the + // assertion does not depend on how long that takes, only on Dispose's 2 s budget for + // the cleanup outlasting it. A Dispose that only posted the cleanup returns inside it. + var occupied = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + actor.Post(() => + { + occupied.SetResult(true); + Thread.Sleep(TimeSpan.FromMilliseconds(100)); + }); + await occupied.Task.WaitAsync(ShortTimeout); + var noneInFlight = await Task.Run(() => + { + sender.Dispose(); + return sends.All(t => t.IsCompleted); + }).WaitAsync(ShortTimeout); + + noneInFlight.Should().BeTrue("Dispose returned, so no send may still be in flight"); + await ShouldAllFailDisposed(sends); + } + + // Disposed from the borrowed actor's own loop: a posted cleanup could only run after this + // work item, so the sessions are failed on the spot. + [Fact] + public async Task Disposing_From_The_Borrowed_Actors_Loop_Fails_The_Send_At_Once() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + var sends = await SendBamsInFlight(clock, actor, sender); + + var failedInside = await actor.PostAsync(() => + { + sender.Dispose(); + return sends.All(t => t.IsCompleted); + }).WaitAsync(ShortTimeout); + + failedInside.Should().BeTrue("the loop that disposed the channel failed its sends itself"); + await ShouldAllFailDisposed(sends); + } + + // Bugbot on #183: a borrowed actor its owner disposed first took its sessions with it. They + // are actor state its loop may still be draining, so the channel's Dispose leaves them + // alone -- and does not throw for the actor being gone. + [Fact] + public async Task Disposing_After_The_Borrowed_Actor_Leaves_Its_Sessions_To_It() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + var sends = await SendBamsInFlight(clock, actor, sender); + + actor.Dispose(); + sender.Invoking(c => c.Dispose()).Should().NotThrow(); + + sends.Should().OnlyContain(t => !t.IsCompleted, + "nothing on this thread touched the sessions of an actor that is gone"); + } + + // Codex on #183: a borrowed actor busy past Dispose's 2 s budget. Dispose returns and says + // so on BackgroundExceptionOccurred, and the service it owns is disposed only once the + // actor has run the cleanup, not underneath it. + [Fact] + public async Task Disposing_On_A_Borrowed_Actor_Stuck_Past_The_Budget_Defers_The_Service() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new DisposalRecordingBusService(new CanBusService(bus)); // the channel owns it; disposing twice is harmless + using var sender = new J1939TpChannel(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: InFlightSpacing), ownsService: true, actor); + var sends = await SendBamsInFlight(clock, actor, sender); + var faults = new List(); + sender.BackgroundExceptionOccurred += (_, e) => { lock (faults) faults.Add(e); }; + + using var release = new ManualResetEventSlim(); + var occupied = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + actor.Post(() => + { + occupied.SetResult(true); + release.Wait(TimeSpan.FromSeconds(30)); + }); + try + { + await occupied.Task.WaitAsync(ShortTimeout); + sender.Dispose(); + + lock (faults) faults.Should().ContainSingle().Which.Should().BeOfType(); + service.Disposed.IsCompleted.Should().BeFalse("the cleanup has not run, so the service is still in use"); + sends.Should().OnlyContain(t => !t.IsCompleted); + } + finally + { + release.Set(); + } + + await ShouldAllFailDisposed(sends); + await service.Disposed.WaitAsync(ShortTimeout); + } + + // Codex on #183: a reader that outlived its join could still hand frames to a borrowed actor + // after the session cleanup. Whenever the actor gets to them, frames handed over before or + // during Dispose are dropped: here a whole BAM waits behind the caller's work while the + // channel is disposed, and must not be delivered once the actor is free. + [Fact] + public async Task Frames_The_Actor_Reaches_After_Dispose_Began_Are_Dropped() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new FrameConsumptionCountingBusService(new CanBusService(bus)); + using var receiver = new J1939TpChannel(service, sourceAddress: 0x10, new J1939TpOptions(), + ownsService: false, actor); + var delivered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + receiver.DatagramReceived += (_, _) => delivered.TrySetResult(true); + var receiving = receiver.ReceiveAsync(); + + using var release = new ManualResetEventSlim(); + var occupied = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + actor.Post(() => + { + occupied.SetResult(true); + release.Wait(TimeSpan.FromSeconds(30)); + }); + Task disposing; + try + { + await occupied.Task.WaitAsync(ShortTimeout); + + const byte peerSa = 0x20; + const uint pgn = 0xFECDu; + var payload = RandomPayload(9, seed: 185); + CanFrame FromPeer(uint framePgn, byte[] data) => CanFrame.Classic( + (int)J1939Id.ComposePgn(7, framePgn, peerSa, J1939Pgn.GlobalAddress), data, isExtendedFrame: true); + bus.RaiseObserved(FromPeer(J1939Pgn.TpCm, J1939TpFrames.BuildBam(9, 2, pgn)), isEcho: false); + bus.RaiseObserved(FromPeer(J1939Pgn.TpDt, J1939TpFrames.BuildDt(1, payload, 0)), isEcho: false); + bus.RaiseObserved(FromPeer(J1939Pgn.TpDt, J1939TpFrames.BuildDt(2, payload, 7)), isEcho: false); + // The reader asked for a fourth frame, so it has handed all three to the actor. + await service.WaitUntilConsumedAsync(3).WaitAsync(ShortTimeout); + + disposing = Task.Run(receiver.Dispose); + // The inbox is completed right after Dispose marks the channel disposed. + await receiving.Invoking(t => t).Should().ThrowAsync(); + } + finally + { + release.Set(); + } + + await disposing.WaitAsync(ShortTimeout); + await actor.PostAsync(() => 0).WaitAsync(ShortTimeout); + // The event is raised on the pool; a delivered BAM would have raised it by now. + (await Task.WhenAny(delivered.Task, Task.Delay(TimeSpan.FromMilliseconds(500)))) + .Should().NotBeSameAs(delivered.Task, "a BAM the actor reached after Dispose began is not delivered"); + } + + private static async Task ShouldAllFailDisposed(Task[] sends) + { + foreach (var send in sends) + await send.Invoking(t => t).Should().ThrowAsync(); + } +} + +/// +/// Test double: forwards to an inner service it owns, and counts the frames each subscriber has +/// finished with -- a frame counts once the subscriber asks for the next one, i.e. once the +/// J1939-TP reader has handed it to the actor (#183). +/// +internal sealed class FrameConsumptionCountingBusService : ICanBusService +{ + private readonly ICanBusService _inner; + private readonly object _gate = new(); + private int _consumed; + private readonly List<(int Count, TaskCompletionSource Reached)> _waiters = new(); + + public FrameConsumptionCountingBusService(ICanBusService inner) => _inner = inner; + + public Task WaitUntilConsumedAsync(int count) + { + lock (_gate) + { + if (_consumed >= count) return Task.CompletedTask; + var reached = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _waiters.Add((count, reached)); + return reached.Task; + } + } + + private void Consumed() + { + lock (_gate) + { + _consumed++; + foreach (var (count, reached) in _waiters) + if (_consumed >= count) reached.TrySetResult(true); + } + } + + public ICanBus Bus => _inner.Bus; + public int SubscriptionCount => _inner.SubscriptionCount; + + public event EventHandler? BackgroundExceptionOccurred + { + add => _inner.BackgroundExceptionOccurred += value; + remove => _inner.BackgroundExceptionOccurred -= value; + } + + public ISubscription Subscribe(Func? predicate = null, int? bufferCapacity = null, bool includeEcho = false) + => new Counted(this, _inner.Subscribe(predicate, bufferCapacity, includeEcho)); + + public ISubscription Subscribe(CanIdFilter filter, int? bufferCapacity = null, bool includeEcho = false) + => new Counted(this, _inner.Subscribe(filter, bufferCapacity, includeEcho)); + + public IReadOnlyList FindOverlappingFilterSubscriptions() + => _inner.FindOverlappingFilterSubscriptions(); + + public Task SendConfirmed(CanFrame frame, TimeSpan? timeout = null, + CancellationToken cancellationToken = default) + => _inner.SendConfirmed(frame, timeout, cancellationToken); + + public void Dispose() => _inner.Dispose(); + + private sealed class Counted : ISubscription + { + private readonly FrameConsumptionCountingBusService _owner; + private readonly ISubscription _inner; + + public Counted(FrameConsumptionCountingBusService owner, ISubscription inner) + { + _owner = owner; + _inner = inner; + } + + public IAsyncEnumerable Frames => Count(); + + private async IAsyncEnumerable Count( + [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default) + { + await foreach (var frameEvent in _inner.Frames.WithCancellation(cancellationToken)) + { + yield return frameEvent; + _owner.Consumed(); // runs when the subscriber asks for the next frame + } + } + + public bool TryRead(out CanFrameEvent frameEvent) => _inner.TryRead(out frameEvent); + + public ValueTask WaitToReadAsync(CancellationToken cancellationToken = default) + => _inner.WaitToReadAsync(cancellationToken); + + public void Reconfigure(CanIdFilter filter) => _inner.Reconfigure(filter); + + public void Reconfigure(Func? predicate) => _inner.Reconfigure(predicate); + + public void Dispose() => _inner.Dispose(); + } +} + +/// +/// Test double: forwards everything to an inner service it owns, and records when it is +/// disposed, so a test can tell whether the channel disposed it and when (#183). +/// +internal sealed class DisposalRecordingBusService : ICanBusService +{ + private readonly ICanBusService _inner; + private readonly TaskCompletionSource _disposed = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public DisposalRecordingBusService(ICanBusService inner) => _inner = inner; + + /// Completes when has run. + public Task Disposed => _disposed.Task; + + public ICanBus Bus => _inner.Bus; + public int SubscriptionCount => _inner.SubscriptionCount; + + public event EventHandler? BackgroundExceptionOccurred + { + add => _inner.BackgroundExceptionOccurred += value; + remove => _inner.BackgroundExceptionOccurred -= value; + } + + public ISubscription Subscribe(Func? predicate = null, int? bufferCapacity = null, bool includeEcho = false) + => _inner.Subscribe(predicate, bufferCapacity, includeEcho); + + public ISubscription Subscribe(CanIdFilter filter, int? bufferCapacity = null, bool includeEcho = false) + => _inner.Subscribe(filter, bufferCapacity, includeEcho); + + public IReadOnlyList FindOverlappingFilterSubscriptions() + => _inner.FindOverlappingFilterSubscriptions(); + + public Task SendConfirmed(CanFrame frame, TimeSpan? timeout = null, + CancellationToken cancellationToken = default) + => _inner.SendConfirmed(frame, timeout, cancellationToken); + + public void Dispose() + { + _inner.Dispose(); + _disposed.TrySetResult(true); + } } /// diff --git a/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs index 0888e88..27ca4cd 100644 --- a/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs @@ -8,6 +8,7 @@ using CanKit.Abstractions.API.Can.Definitions; using CanKit.Abstractions.API.Common.Definitions; using CanKit.Core; +using CanKit.Pro.Actor; using CanKit.Pro.IsoTp; using CanKit.Pro.RawCan; using CanKit.Pro.Tests.Infrastructure; @@ -48,6 +49,69 @@ private static ICanBus OpenClassic(string session, int channel) => CanBus.Open( private static TimeSpan Between(long earlier, long later) => TimeSpan.FromSeconds((later - earlier) / (double)Stopwatch.Frequency); + /// + /// A functional client whose P2, P2* and collection windows run on , + /// and a demux that stamps frames with that same clock. A zero stamp means "unstamped", + /// so the clock is nudged off zero before anything is sent (#171). + /// + private sealed class WindowClock : IDisposable + { + public VirtualClock Clock { get; } = new(); + public ProtocolActor Actor { get; } + + public WindowClock() + { + Actor = Clock.NewActor(); + Clock.Advance(TimeSpan.FromMilliseconds(1)); + } + + public long Now => Actor.TimeSource.GetTimestamp(); + + public void Dispose() => Clock.Dispose(); + } + + private static UdsFunctionalClient OpenOnClock(ICanBus bus, WindowClock clock, + TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null, + IsoTpFunctionalOptions? options = null) + => OpenOnClock(new CanBusService(bus, clock.Actor.TimeSource.GetTimestamp), clock, + responseWindow, responsePendingWindow, options); + + private static UdsFunctionalClient OpenOnClock(ICanBusService service, WindowClock clock, + TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null, + IsoTpFunctionalOptions? options = null) + { + var iso = new IsoTpFunctionalClient(service, FunctionalTxId, Ecu1, 0x7EF, + options ?? FastOptions(), ownsService: true, clock.Actor); + return UdsFunctionalClient.Create(iso, clock.Actor, ownsClient: true, + responseWindow: responseWindow, responsePendingWindow: responsePendingWindow); + } + + private static async Task WaitUntilArmed(WindowClock clock, TimeSpan expected, TimeSpan? lateBy = null) + => await clock.Clock.WaitUntilTimerArmedAsync(clock.Actor, expected, ShortTimeout, + lateBy ?? TimeSpan.FromMilliseconds(5)); + + // Steps the clock until the call returns. One advance of the nominal window can stop a tick + // short of a ceiling, or land before the collection timer is armed, and the call then sits + // until the wall-clock budget cancels it. + private static async Task RunCall(WindowClock clock, Task call) + { + await clock.Clock.RunUntilAsync(call, TimeSpan.FromMilliseconds(25), ShortTimeout); + return await call; + } + + // While a collection outlasts P2 the listener keeps 20 ms slices. The call's task can + // complete with that slice still armed, and it is shorter than the next window, so the + // next window is not the earliest timer until the slice has elapsed and the listener retired. + private static async Task RetireInFlightSlices(WindowClock clock) + { + for (var i = 0; i < 4; i++) + { + var delay = await clock.Actor.NextTimerDelayAsync(); + if (delay is null || delay > TimeSpan.FromMilliseconds(20)) return; + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(21)); + } + } + private static CanFrame SingleFrameFrom(uint canId, byte[] pdu) => CanFrame.Classic(unchecked((int)canId), IsoTpFrameCodec.BuildSingleFrame(IsoTpEndpoint.Normal(canId, 0), pdu, isCanFd: false, padding: true)); @@ -232,6 +296,7 @@ public async Task A_Positive_Response_For_Another_Did_Is_Not_Attributed() public async Task A_Late_Negative_Answer_To_A_Suppressed_Send_Does_Not_Land_In_The_Next_Window() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -240,19 +305,25 @@ public async Task A_Late_Negative_Answer_To_A_Suppressed_Send_Does_Not_Land_In_T busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[2] == 0x80) - _ = Task.Run(async () => { await Task.Delay(100); busEcus.Transmit(negative); }); - else + if (e.CanFrame.Data.Span[2] != 0x80) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(300)); + using var functional = OpenOnClock(busTester, clock, responseWindow: TimeSpan.FromMilliseconds(300)); using var cts = new CancellationTokenSource(ShortTimeout); await functional.TesterPresentAsync(cancellationToken: cts.Token); - var responses = await functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(300), TimeSpan.FromMilliseconds(20)); + + // 100 ms into the 300 ms window. On the wall clock a loaded runner has stretched this + // past the window and the negative landed in the next call (macOS, run 36234735257). + clock.Clock.Advance(TimeSpan.FromMilliseconds(100)); + busEcus.Transmit(negative); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(201)); + + var second = functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await WaitUntilArmed(clock, Window, TimeSpan.FromMilliseconds(20)); + var responses = await RunCall(clock, second); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the late negative answer belongs to the suppressed send and is not collected"); @@ -264,6 +335,7 @@ public async Task A_Late_Negative_Answer_To_A_Suppressed_Send_Does_Not_Land_In_T public async Task A_Late_Negative_Answer_To_A_Previous_Request_Does_Not_Land_In_The_Next_Window() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -272,20 +344,26 @@ public async Task A_Late_Negative_Answer_To_A_Previous_Request_Does_Not_Land_In_ busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[3] == 0x90) - _ = Task.Run(async () => { await Task.Delay(150); busEcus.Transmit(negative); }); - else + if (e.CanFrame.Data.Span[3] != 0x90) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(300)); + using var functional = OpenOnClock(busTester, clock, responseWindow: TimeSpan.FromMilliseconds(300)); using var cts = new CancellationTokenSource(ShortTimeout); - // A short collection window: the negative answer comes after it. - await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(30), cts.Token); - var responses = await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + // A short collection window: the negative answer comes after it, and still inside P2. + var first = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(30), cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(30)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(31)); + await first; + + clock.Clock.Advance(TimeSpan.FromMilliseconds(120)); // 151 ms from the request: after the collection, inside P2 + busEcus.Transmit(negative); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(160)); // P2 of 300 ms is over + + var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + await WaitUntilArmed(clock, Window, TimeSpan.FromMilliseconds(20)); + var responses = await RunCall(clock, second); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "F190's late negative answer is the previous request's"); @@ -296,41 +374,49 @@ public async Task A_Late_Negative_Answer_To_A_Previous_Request_Does_Not_Land_In_ public async Task A_Cancelled_Collection_Still_Leaves_Its_Window_For_The_Next_Call() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); var negative = SingleFrameFrom(Ecu1, new byte[] { 0x7F, 0x22, 0x31 }); var positive = SingleFrameFrom(Ecu1, new byte[] { 0x62, 0xF1, 0x91, 0x02 }); + var sentAt = new List(); busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[3] == 0x90) - _ = Task.Run(async () => { await Task.Delay(150); busEcus.Transmit(negative); }); - else + lock (sentAt) sentAt.Add(clock.Now); + if (e.CanFrame.Data.Span[3] != 0x90) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(300)); - - var sentAt = new List(); - busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID == unchecked((int)FunctionalTxId)) lock (sentAt) sentAt.Add(Stopwatch.GetTimestamp()); }; + using var functional = OpenOnClock(busTester, clock, responseWindow: TimeSpan.FromMilliseconds(300)); using var early = new CancellationTokenSource(TimeSpan.FromMilliseconds(30)); Func cancelled = () => functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, Window, early.Token); await cancelled.Should().ThrowAsync(); + // The cancellation is the caller's, at 30 ms of wall time, while the clock has not moved. + // The negative belongs to the cancelled request: after that cancellation, inside its P2. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(300), TimeSpan.FromMilliseconds(20)); + clock.Clock.Advance(TimeSpan.FromMilliseconds(150)); + busEcus.Transmit(negative); + using var cts = new CancellationTokenSource(ShortTimeout); - var responses = await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + // Still inside the 300 ms window: the next request must not have gone out. + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(50)); + lock (sentAt) sentAt.Should().HaveCount(1, "the cancelled request's window is still open"); + + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(110)); // past P2 + await WaitUntilArmed(clock, Window, TimeSpan.FromMilliseconds(20)); + var responses = await RunCall(clock, second); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the cancelled request's late negative answer is not the next call's"); - // And the reason it is not: the second request waited the window out. A lower bound, - // which a loaded host only raises. long gap; lock (sentAt) gap = sentAt[1] - sentAt[0]; - Between(0, gap).Should().BeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(250), - "the cancelled request's window was kept for the next call"); + TimeSpan.FromSeconds(gap / (double)clock.Actor.TimeSource.Frequency) + .Should().BeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(250), + "the cancelled request's window was kept for the next call"); } // Codex on #150: WriteMemoryByAddress echoes the addressAndLengthFormatIdentifier and the @@ -454,8 +540,8 @@ public async Task A_Window_Is_Anchored_At_The_Drivers_Acceptance_Not_Before_The_ var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // the count never decreases + bus.RaiseObserved(positive, isEcho: false); bus.DeferredEchoes.ReleaseNext(); - _ = Task.Run(async () => { await Task.Delay(20); bus.RaiseObserved(positive, isEcho: false); }); (await second).Should().ContainSingle(); long gap; @@ -493,8 +579,8 @@ public async Task A_Window_Is_Anchored_At_The_Drivers_Acceptance_Not_At_The_Conf var startedAt = Stopwatch.GetTimestamp(); var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); + bus.RaiseObserved(positive, isEcho: false); bus.DeferredEchoes.ReleaseNext(); - _ = Task.Run(async () => { await Task.Delay(20); bus.RaiseObserved(positive, isEcho: false); }); (await second).Should().ContainSingle(); long secondSentAt; @@ -549,6 +635,7 @@ public async Task A_Pending_Answer_Extends_The_Window_From_Its_Arrival_Not_From_ public async Task A_Pending_Answer_Collected_After_P2_Does_Not_Revive_The_Window() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -557,26 +644,29 @@ public async Task A_Pending_Answer_Collected_After_P2_Does_Not_Revive_The_Window busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[3] == 0x90) - _ = Task.Run(async () => { await Task.Delay(250); busEcus.Transmit(pending); }); - else busEcus.Transmit(positive); + if (e.CanFrame.Data.Span[3] != 0x90) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(2000)); + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(2000)); using var cts = new CancellationTokenSource(ShortTimeout); // P2 = 100 ms; the 0x78 at 250 ms is inside the 400 ms collection but after P2. - await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(400), cts.Token); + var first = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(400), cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(100), TimeSpan.FromMilliseconds(20)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(250)); + busEcus.Transmit(pending); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(151)); + await first; + await RetireInFlightSlices(clock); - // Revived, the window would reach 2250 ms and this call would wait most of two - // seconds; a loaded host only makes the call slower, so the bound is wide. - var sw = Stopwatch.StartNew(); - var responses = await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, TimeSpan.FromMilliseconds(50), cts.Token); - sw.Stop(); + // Revived, the window would reach 250 + 2000 ms and this call would not go out after + // a 100 ms step. It goes out after P2, which is already over. + var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, TimeSpan.FromMilliseconds(50), cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(50), TimeSpan.FromMilliseconds(20)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(51)); + var responses = await second; responses.Should().ContainSingle(); - sw.Elapsed.Should().BeLessThan(TimeSpan.FromSeconds(1), "the 0x78 from after P2 did not revive the window"); } // Codex on #150: the listener's subscription is made before the send, but its worker may @@ -586,6 +676,7 @@ public async Task A_Pending_Answer_Collected_After_P2_Does_Not_Revive_The_Window public async Task A_Listener_Starting_After_The_Window_Still_Hears_What_It_Buffered() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -595,26 +686,31 @@ public async Task A_Listener_Starting_After_The_Window_Still_Hears_What_It_Buffe busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if ((e.CanFrame.Data.Span[2] & 0x80) != 0) - { - busEcus.Transmit(pending); - _ = Task.Run(async () => { await Task.Delay(250); busEcus.Transmit(negative); }); - } + if ((e.CanFrame.Data.Span[2] & 0x80) != 0) busEcus.Transmit(pending); else busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); - functional.ListenerStartDelay = TimeSpan.FromMilliseconds(150); // past the 100 ms window + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); + functional.DelayListenerStart(TimeSpan.FromMilliseconds(150)); // past the 100 ms window using var cts = new CancellationTokenSource(ShortTimeout); await functional.SendRawAsync(new byte[] { 0x10, 0x83 }, Window, cts.Token); // suppressed: the 0x78 is buffered before the worker runs - - var responses = await functional.DiagnosticSessionControlAsync(UdsSessionType.Extended, Window, cts.Token); - // P2* = 1500 ms: the negative answer's timer must fire before the second request goes - // out, with 1250 ms to spare (macOS CI on #150 exceeded 150). The same P2* in the - // three tests below, for the same reason. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(150), TimeSpan.FromMilliseconds(20)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(150)); // the worker starts, reads the buffered 0x78 + functional.DelayListenerStart(TimeSpan.Zero); // the next call's listener is not the late one + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(1300), TimeSpan.FromMilliseconds(300)); + + // The second call waits the extended window out. The negative goes in while it is + // waiting, at 250 ms, and must not be collected as this call's answer. + var second = functional.DiagnosticSessionControlAsync(UdsSessionType.Extended, Window, cts.Token); + clock.Clock.Advance(TimeSpan.FromMilliseconds(100)); + busEcus.Transmit(negative); + // Still P2* from the buffered 0x78, not this call's own collection: a window that was + // not extended would already be arming that shorter timer. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(1200), TimeSpan.FromMilliseconds(200)); + await clock.Clock.RunUntilAsync(second, TimeSpan.FromMilliseconds(50), ShortTimeout); + var responses = await second; responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the late worker read the buffered 0x78 and kept the window to 1500 ms, past the negative at 250"); } @@ -625,6 +721,7 @@ public async Task A_Listener_Starting_After_The_Window_Still_Hears_What_It_Buffe public async Task A_Cancelled_Collection_Does_Not_Lose_The_Pending_Answer_The_Listener_Heard() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -634,24 +731,32 @@ public async Task A_Cancelled_Collection_Does_Not_Lose_The_Pending_Answer_The_Li busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[3] == 0x90) - { - busEcus.Transmit(pending); - _ = Task.Run(async () => { await Task.Delay(250); busEcus.Transmit(negative); }); - } + if (e.CanFrame.Data.Span[3] == 0x90) busEcus.Transmit(pending); else busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); using var early = new CancellationTokenSource(TimeSpan.FromMilliseconds(30)); Func cancelled = () => functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, Window, early.Token); await cancelled.Should().ThrowAsync(); + // The 0x78 went out with the request. Elapse P2 so the listener applies it, then put the + // negative inside the P2* that application opens. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(100), TimeSpan.FromMilliseconds(30)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(101)); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(1300), TimeSpan.FromMilliseconds(300)); + clock.Clock.Advance(TimeSpan.FromMilliseconds(150)); + busEcus.Transmit(negative); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(1400)); + using var cts = new CancellationTokenSource(ShortTimeout); - var responses = await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); + // P2* has elapsed, so this call's own 100 ms window is what is armed. A revived + // 1500 ms window would still be the earliest timer, and this wait would not match. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(100), TimeSpan.FromMilliseconds(40)); + var responses = await RunCall(clock, second); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the final negative answer at 250 ms belongs to the cancelled request, whose 0x78 the cancellation hid"); } @@ -662,6 +767,7 @@ public async Task A_Cancelled_Collection_Does_Not_Lose_The_Pending_Answer_The_Li public async Task A_Cancelled_Wait_Leaves_The_Listener_To_Hear_The_Pending_Answer() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -671,29 +777,31 @@ public async Task A_Cancelled_Wait_Leaves_The_Listener_To_Hear_The_Pending_Answe busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[2] == 0x80) - _ = Task.Run(async () => - { - await Task.Delay(50); busEcus.Transmit(pending); - await Task.Delay(200); busEcus.Transmit(negative); - }); - else busEcus.Transmit(positive); + if (e.CanFrame.Data.Span[2] != 0x80) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(300), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(300), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); using var cts = new CancellationTokenSource(ShortTimeout); - await functional.TesterPresentAsync(cancellationToken: cts.Token); // suppressed; 0x78 at 50 ms, negative at 250 ms, window to 650 ms + await functional.TesterPresentAsync(cancellationToken: cts.Token); // suppressed + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(300), TimeSpan.FromMilliseconds(20)); + clock.Clock.Advance(TimeSpan.FromMilliseconds(50)); + busEcus.Transmit(pending); // 0x78 inside P2; the listener applies it when the slice ends using var early = new CancellationTokenSource(TimeSpan.FromMilliseconds(100)); Func cancelled = () => functional.TesterPresentAsync(suppressPositiveResponse: false, Window, early.Token); await cancelled.Should().ThrowAsync(); // cancelled while waiting; the 0x78 it heard is lost - var responses = await functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(200)); // 250 ms: the negative, inside the extended window + busEcus.Transmit(negative); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(1400)); // P2* from the 0x78 is over + + var second = functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await WaitUntilArmed(clock, Window, TimeSpan.FromMilliseconds(20)); + var responses = await RunCall(clock, second); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( - "the negative at 250 ms belongs to the suppressed send; the window reaches to 650 ms"); + "the negative at 250 ms belongs to the suppressed send; the window reaches past it"); } // Codex on #150: between a suppressed send and the next call nobody was collecting, so a @@ -703,6 +811,7 @@ public async Task A_Cancelled_Wait_Leaves_The_Listener_To_Hear_The_Pending_Answe public async Task A_Pending_Answer_In_The_Gap_After_A_Suppressed_Send_Is_Observed() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -712,24 +821,30 @@ public async Task A_Pending_Answer_In_The_Gap_After_A_Suppressed_Send_Is_Observe busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - if (e.CanFrame.Data.Span[2] == 0x80) - _ = Task.Run(async () => - { - await Task.Delay(100); busEcus.Transmit(pending); - await Task.Delay(400); busEcus.Transmit(negative); - }); - else busEcus.Transmit(positive); + if (e.CanFrame.Data.Span[2] != 0x80) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(300), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(300), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); using var cts = new CancellationTokenSource(ShortTimeout); - await functional.TesterPresentAsync(cancellationToken: cts.Token); // suppressed; 0x78 at 100 ms, negative at 500 ms - await Task.Delay(200); // nobody collecting: the gap - - var responses = await functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await functional.TesterPresentAsync(cancellationToken: cts.Token); // suppressed + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(300), TimeSpan.FromMilliseconds(20)); + clock.Clock.Advance(TimeSpan.FromMilliseconds(100)); + busEcus.Transmit(pending); // 0x78 during the gap, inside P2 + + // The next call starts while that window is still open and blocks until the listener + // has applied the 0x78 and run P2* out. The negative at 500 ms lands in that wait. + var second = functional.TesterPresentAsync(suppressPositiveResponse: false, Window, cts.Token); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(201)); // the 300 ms slice ends; the 0x78 extends it + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(1200), TimeSpan.FromMilliseconds(200)); + clock.Clock.Advance(TimeSpan.FromMilliseconds(200)); // 500 ms: the negative, inside P2* + busEcus.Transmit(negative); + // The call is still waiting out P2*, so the earliest timer is that remainder, not the + // 300 ms collection it arms only once the extension has elapsed. + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(900), TimeSpan.FromMilliseconds(400)); + await clock.Clock.RunUntilAsync(second, TimeSpan.FromMilliseconds(50), ShortTimeout); + var responses = await second; responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the 0x78 in the gap moved the window past the negative at 500 ms"); } @@ -807,8 +922,8 @@ public async Task A_Pending_Answer_From_Before_The_Handoff_Is_Not_This_Requests( bus.DeferredEchoes.ReleaseNext(); // the other sender's confirmation await other; await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); + bus.RaiseObserved(positive, isEcho: false); // after the handoff, inside the collection bus.DeferredEchoes.ReleaseNext(); // the request's - _ = Task.Run(async () => { await Task.Delay(20); bus.RaiseObserved(positive, isEcho: false); }); var responses = await first; responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse("the 0x78 from before the handoff is not this request's"); @@ -930,8 +1045,8 @@ public async Task A_Listener_Is_Restarted_When_The_Acceptance_Outlasted_The_Wind var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, Window, cts.Token); await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); + bus.RaiseObserved(positive, isEcho: false); bus.DeferredEchoes.ReleaseNext(); - _ = Task.Run(async () => { await Task.Delay(20); bus.RaiseObserved(positive, isEcho: false); }); var responses = await second; responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( @@ -944,6 +1059,7 @@ public async Task A_Listener_Is_Restarted_When_The_Acceptance_Outlasted_The_Wind public async Task A_Collection_That_Outlasts_The_Window_Does_Not_Leave_A_Zombie_Listener() { var session = NewSession(); + using var clock = new WindowClock(); using var busTester = OpenClassic(session, 0); using var busEcus = OpenClassic(session, 1); @@ -953,22 +1069,30 @@ public async Task A_Collection_That_Outlasts_The_Window_Does_Not_Leave_A_Zombie_ busEcus.FrameObserved += (_, e) => { if (e.CanFrame.ID != unchecked((int)FunctionalTxId)) return; - int n = Interlocked.Increment(ref seen); - if (n == 2) _ = Task.Run(async () => { await Task.Delay(150); busEcus.Transmit(negative); }); - if (n == 3) busEcus.Transmit(positive); + if (Interlocked.Increment(ref seen) == 3) busEcus.Transmit(positive); }; - using var functional = UdsFunctionalClient.Create( - IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), - ownsClient: true, responseWindow: TimeSpan.FromMilliseconds(400)); + using var functional = OpenOnClock(busTester, clock, responseWindow: TimeSpan.FromMilliseconds(400)); using var cts = new CancellationTokenSource(ShortTimeout); // A collection longer than the window: the pre-send listener retires; the anchor after // it is already over and must start nothing. - await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(500), cts.Token); + var first = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(500), cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(400), TimeSpan.FromMilliseconds(20)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(501)); + await first; + await RetireInFlightSlices(clock); // A short collection; its late negative, at 150 ms, needs a living listener's window. - await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, TimeSpan.FromMilliseconds(30), cts.Token); - var responses = await functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x92 }, Window, cts.Token); + var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, TimeSpan.FromMilliseconds(30), cts.Token); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(30)); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(31)); + await second; + clock.Clock.Advance(TimeSpan.FromMilliseconds(150)); + busEcus.Transmit(negative); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(251)); // the 400 ms window is over + var third = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x92 }, Window, cts.Token); + await WaitUntilArmed(clock, Window, TimeSpan.FromMilliseconds(20)); + var responses = await RunCall(clock, third); responses.Should().ContainSingle().Which.IsNegative.Should().BeFalse( "the second request's window was honoured by a real listener, not blocked by a zombie"); @@ -1045,4 +1169,78 @@ public async Task DiagnosticSessionControl_To_Everyone_Collects_Each_Ecus_Sessio responses.Should().ContainSingle().Which.Response[1].Should().Be(0x03); } + + // #171: the listener start delay stands in for a late thread pool, and is moved by + // advancing an injected clock. On the wall clock it would be a sleep, so it is refused. + [Fact] + public void A_Listener_Start_Delay_Needs_An_Injected_Clock() + { + using var busTester = OpenClassic(NewSession(), 0); + using var functional = UdsFunctionalClient.Create( + IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), ownsClient: true); + + Action delay = () => functional.DelayListenerStart(TimeSpan.FromMilliseconds(1)); + delay.Should().Throw(); + } + + // Codex on #150: a suppressed send cancelled after the driver took the frame may still have + // put it on the bus, so the ECUs' P2 runs from the cancellation, not from before the send. + // The driver holds the echo while the clock moves 200 ms: the window ends at 500 ms, where + // the provisional one from before the send ended at 300. + [Fact] + public async Task A_Suppressed_Send_Cancelled_After_The_Driver_Took_It_Leaves_P2_From_The_Cancellation() + { + using var clock = new WindowClock(); + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + var options = new IsoTpFunctionalOptions { IsExtendedCanId = false, UseCanFd = false, UsePadding = true, NAs = ShortTimeout }; + using var functional = OpenOnClock(bus, clock, responseWindow: TimeSpan.FromMilliseconds(300), options: options); + + using var cancel = new CancellationTokenSource(); + var suppressed = functional.TesterPresentAsync(cancellationToken: cancel.Token); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(200)); + cancel.Cancel(); + Func cancelled = () => suppressed; + await cancelled.Should().ThrowAsync(); + bus.EchoMode = EchoDelivery.Synchronous; + + await NextCallWaitsUntil500(clock, functional, () => bus.TransmitCount); + } + + // A driver that reports no transmit instant: the window runs from the confirmation. It + // takes 200 ms to come, so the window ends at 500 ms, not at the provisional 300. + [Fact] + public async Task A_Suppressed_Send_Without_A_Transmit_Stamp_Is_Anchored_At_Its_Confirmation() + { + using var clock = new WindowClock(); + var service = new StarvedReaderBusService + { + OnSendConfirmed = () => clock.Clock.Advance(TimeSpan.FromMilliseconds(200)), + }; + using var functional = OpenOnClock(service, clock, responseWindow: TimeSpan.FromMilliseconds(300)); + + using var cts = new CancellationTokenSource(ShortTimeout); + await functional.TesterPresentAsync(cancellationToken: cts.Token); + service.OnSendConfirmed = null; + + await NextCallWaitsUntil500(clock, functional, () => service.Sent.Count); + } + + // At 400 ms a window ending at 500 still holds the next call back, its listener collecting + // the 100 ms left. A window that ended at 300 would have let the call out, and the earliest + // timer would be its own 30 ms collection. + private static async Task NextCallWaitsUntil500(WindowClock clock, UdsFunctionalClient functional, + Func transmitted) + { + using var cts = new CancellationTokenSource(ShortTimeout); + var next = functional.TesterPresentAsync(suppressPositiveResponse: false, + TimeSpan.FromMilliseconds(30), cts.Token); + await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(400) - clock.Clock.Elapsed + + TimeSpan.FromMilliseconds(1)); + await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(100)); + transmitted().Should().Be(1, "the next call is still waiting out the window"); + + (await RunCall(clock, next)).Should().BeEmpty(); + transmitted().Should().Be(2); + } }