From a986d4654a3860fbd23aea693e5f4defa9e2f102 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 27 Sep 2026 11:42:30 +0000 Subject: [PATCH 01/11] refactor(j1939tp): accept an injected actor so timers share the caller's clock The channel always built its own ProtocolActor, so BAM spacing and the other TP deadlines could only be waited out on the wall clock. An optional actor follows the node: the channel owns the loop only when it created it, and production still constructs one when the caller passes none. Co-authored-by: Dietmar Borgards --- src/CanKit.Pro.J1939Tp/J1939TpChannel.cs | 29 ++++++- .../TestCases/J1939TpTests.cs | 83 +++++++++++++++++++ 2 files changed, 108 insertions(+), 4 deletions(-) diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index a285d83..da1a3c6 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; } @@ -328,7 +346,10 @@ public void Dispose() 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; + if (_ownsActor) _actor.Dispose(); _readerCts.Dispose(); if (_ownsService) diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index c5e27f8..e67b5e5 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,88 @@ 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"); + } } /// From 6554993b0d3d5acf0729ca89e0c0d6c46037b5ed Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 27 Sep 2026 11:42:30 +0000 Subject: [PATCH 02/11] test(uds): drive functional response windows from an injected clock Functional P2, P2* and the ISO-TP collection window were Stopwatch plus CancelAfter, which is how a late negative on a busy macOS runner landed in the next window. Tests that build the stack can pass one actor, and the demux stamps frames from that same clock. A null clock keeps the wall-clock path. Co-authored-by: Dietmar Borgards --- src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj | 8 + src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs | 35 +- .../IsoTpFunctionalListener.cs | 38 +- src/CanKit.Pro.RawCan/CanBusService.cs | 47 ++- .../SuppressedResponseWindows.cs | 24 +- src/CanKit.Pro.Uds/UdsFunctionalClient.cs | 73 +++- .../TestCases/Uds/UdsFunctionalClientTests.cs | 329 ++++++++++++------ 7 files changed, 409 insertions(+), 145 deletions(-) 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..6df4634 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,17 @@ 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); + long deadline = _time.GetTimestamp() + (long)(window.TotalSeconds * _time.Frequency); using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - windowCts.CancelAfter(window); + using var windowTimer = FunctionalWindow.Arm(_clock, window, windowCts); var windowToken = windowCts.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 +354,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 +374,7 @@ private static void DrainBuffered(ISubscription sub) } internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent, - out IsoTpFunctionalResponse? response) + out IsoTpFunctionalResponse? response, Func? now = null) { var frame = frameEvent.Frame; var payload = frame.Data.ToArray(); @@ -390,7 +405,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(); + var arrival = frameEvent.HostArrivalTimestamp > 0 + ? frameEvent.HostArrivalTimestamp + : (now ?? Stopwatch.GetTimestamp)(); 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..f15ed5f 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,9 +85,9 @@ public async Task> CollectAsync( TakeBuffered(responses, now); return responses.AsReadOnly(); } - long deadline = now + (long)(window.TotalSeconds * Stopwatch.Frequency); + long deadline = now + (long)(window.TotalSeconds * _time.Frequency); using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - windowCts.CancelAfter(window); + using var windowTimer = FunctionalWindow.Arm(_clock, window, windowCts); try { // Ends when the window runs out, or when the subscription is completed underneath @@ -109,7 +114,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 +127,26 @@ public void Dispose() _subscription.Dispose(); } } + +/// +/// Ends a functional collection window. Production uses . +/// 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 static class FunctionalWindow +{ + public static IDisposable? Arm(ProtocolActor? actor, TimeSpan window, CancellationTokenSource windowCts) + { + if (actor is null) + { + windowCts.CancelAfter(window); + return null; + } + + return actor.Schedule(window, () => + { + try { windowCts.Cancel(); } + catch (ObjectDisposedException) { } + }); + } +} 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..ef951bd 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(); @@ -57,9 +62,11 @@ public sealed class UdsFunctionalClient : IDisposable 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) @@ -85,7 +92,18 @@ private UdsFunctionalClient(IsoTpFunctionalClient client, bool ownsClient, TimeS public static UdsFunctionalClient Create(IsoTpFunctionalClient client, bool ownsClient = false, TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null) => new(client, ownsClient, responseWindow ?? UdsClientOptions.DefaultP2, - responsePendingWindow ?? UdsClientOptions.DefaultP2Star); + responsePendingWindow ?? UdsClientOptions.DefaultP2Star, clock: null); + + /// + /// As , measuring P2 + /// and P2* on . 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, clock); /// The underlying ISO-TP functional client. public IsoTpFunctionalClient Channel => _client; @@ -148,7 +166,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 +181,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 +218,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 +237,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 +257,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 +304,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 +314,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 +344,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 +366,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; @@ -376,7 +394,7 @@ private async Task ListenAsync(byte sid, IsoTpFunctionalListener ears) try { if (ListenerStartDelay > TimeSpan.Zero) - await Task.Delay(ListenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); + await WaitOnClockAsync(ListenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); bool drained = false; while (true) { @@ -387,7 +405,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 +579,30 @@ 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); + + // Production waits on the wall clock. An injected actor waits on its own, which is what + // lets a test move ListenerStartDelay without sleeping (#171). + private async Task WaitOnClockAsync(TimeSpan delay, CancellationToken cancellationToken) + { + if (_clock is null) + { + await Task.Delay(delay, cancellationToken).ConfigureAwait(false); + return; + } + + if (delay <= TimeSpan.Zero) return; + 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/TestCases/Uds/UdsFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs index 0888e88..47358fc 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,64 @@ 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) + { + var service = new CanBusService(bus, clock.Actor.TimeSource.GetTimestamp); + 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 +291,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 +300,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 +330,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 +339,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 +369,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 +535,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 +574,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 +630,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 +639,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 +671,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 +681,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)); + using var functional = OpenOnClock(busTester, clock, + responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); functional.ListenerStartDelay = 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.ListenerStartDelay = 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 +716,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 +726,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 +762,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 +772,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 +806,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 +816,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 +917,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 +1040,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 +1054,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 +1064,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"); From 8274eadb83c16c25f03514174fa3c5123e0fd9e0 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards Date: Sun, 27 Sep 2026 14:45:54 +0200 Subject: [PATCH 03/11] Potential fix for pull request finding 'CodeQL / Poor error handling: empty catch block' Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com> --- src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs index f15ed5f..21a8ecd 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs @@ -146,7 +146,11 @@ internal static class FunctionalWindow return actor.Schedule(window, () => { try { windowCts.Cancel(); } - catch (ObjectDisposedException) { } + catch (ObjectDisposedException) + { + // Expected race: the collection may have been disposed before this scheduled callback fires. + // Cancellation is best-effort here, so a disposed CTS can be safely ignored. + } }); } } From 33382e72c187520d33a6a8f0f2db5c7d22e8ef27 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 13:21:10 +0000 Subject: [PATCH 04/11] test(uds): cover the functional clock seam's remaining branches Codecov reported 12 uncovered or partial lines in the functional window seam of this branch. Each is now either exercised by a test that fails when the line is broken, or gone: - FunctionalWindow cancels a source of its own that is never disposed, so a timer callback racing the window's disposal cannot meet a disposed source. That removes the ObjectDisposedException catch the CodeQL autofix left behind, rather than papering over the race. - TryParseFunctionalResponse takes the collector's clock as required: every caller passes one, so the Stopwatch fallback was dead. - The public UdsFunctionalClient.Create delegates to the internal overload with a null clock, so the defaults live in one place. - The listener start delay is a method that refuses a wall clock; the wall-clock Task.Delay branch it replaced had no remaining caller. - New tests: the post-window drain on an injected clock, the unstamped arrival fallback, a suppressed send cancelled after the driver took it, a suppressed send whose driver reports no transmit stamp, and the start-delay guard. Each was checked by mutating its line. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs | 14 ++- .../IsoTpFunctionalListener.cs | 43 ++++++---- src/CanKit.Pro.Uds/UdsFunctionalClient.cs | 40 +++++---- .../Infrastructure/StarvedReaderBusService.cs | 25 +++++- .../IsoTp/IsoTpFunctionalClientTests.cs | 50 +++++++++++ .../TestCases/Uds/UdsFunctionalClientTests.cs | 85 ++++++++++++++++++- 6 files changed, 207 insertions(+), 50 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs index 6df4634..cc39ddd 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs @@ -325,9 +325,8 @@ private async Task> CollectFromSubscripti // 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 = _time.GetTimestamp() + (long)(window.TotalSeconds * _time.Frequency); - using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - using var windowTimer = FunctionalWindow.Arm(_clock, window, windowCts); - var windowToken = windowCts.Token; + using var windowEnd = new FunctionalWindow(_clock, window, cancellationToken); + var windowToken = windowEnd.Token; try { @@ -374,7 +373,7 @@ private static void DrainBuffered(ISubscription sub) } internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent, - out IsoTpFunctionalResponse? response, Func? now = null) + out IsoTpFunctionalResponse? response, Func now) { var frame = frameEvent.Frame; var payload = frame.Data.ToArray(); @@ -404,10 +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 - : (now ?? 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 21a8ecd..9c2acb9 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs @@ -86,15 +86,14 @@ public async Task> CollectAsync( return responses.AsReadOnly(); } long deadline = now + (long)(window.TotalSeconds * _time.Frequency); - using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - using var windowTimer = FunctionalWindow.Arm(_clock, window, windowCts); + 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; } @@ -129,28 +128,38 @@ public void Dispose() } /// -/// Ends a functional collection window. Production uses . +/// 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 static class FunctionalWindow +internal sealed class FunctionalWindow : IDisposable { - public static IDisposable? Arm(ProtocolActor? actor, TimeSpan window, CancellationTokenSource windowCts) + private readonly CancellationTokenSource _linked; + private readonly IDisposable? _timer; + + public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken cancellationToken) { if (actor is null) { - windowCts.CancelAfter(window); - return null; + _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + _linked.CancelAfter(window); + return; } - return actor.Schedule(window, () => - { - try { windowCts.Cancel(); } - catch (ObjectDisposedException) - { - // Expected race: the collection may have been disposed before this scheduled callback fires. - // Cancellation is best-effort here, so a disposed CTS can be safely ignored. - } - }); + // The actor's callback cancels a source of its own, which is never disposed -- it holds + // no timer and no wait handle -- so a callback already running when the window is + // disposed cannot meet a disposed source. Disposing the linked source unhooks it. + var end = new CancellationTokenSource(); + _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, end.Token); + _timer = actor.Schedule(window, end.Cancel); + } + + public CancellationToken Token => _linked.Token; + + public void Dispose() + { + _timer?.Dispose(); + _linked.Dispose(); } } diff --git a/src/CanKit.Pro.Uds/UdsFunctionalClient.cs b/src/CanKit.Pro.Uds/UdsFunctionalClient.cs index ef951bd..d4633db 100644 --- a/src/CanKit.Pro.Uds/UdsFunctionalClient.cs +++ b/src/CanKit.Pro.Uds/UdsFunctionalClient.cs @@ -57,8 +57,13 @@ 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, @@ -91,16 +96,16 @@ private UdsFunctionalClient(IsoTpFunctionalClient client, bool ownsClient, TimeS /// public static UdsFunctionalClient Create(IsoTpFunctionalClient client, bool ownsClient = false, TimeSpan? responseWindow = null, TimeSpan? responsePendingWindow = null) - => new(client, ownsClient, responseWindow ?? UdsClientOptions.DefaultP2, - responsePendingWindow ?? UdsClientOptions.DefaultP2Star, clock: null); + => Create(client, clock: null, ownsClient, responseWindow, responsePendingWindow); /// /// As , measuring P2 - /// and P2* on . 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. + /// 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, + 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, clock); @@ -393,8 +398,8 @@ private async Task ListenAsync(byte sid, IsoTpFunctionalListener ears) { try { - if (ListenerStartDelay > TimeSpan.Zero) - await WaitOnClockAsync(ListenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); + if (_listenerStartDelay > TimeSpan.Zero) + await WaitOnClockAsync(_clock!, _listenerStartDelay, _lifetimeCts.Token).ConfigureAwait(false); bool drained = false; while (true) { @@ -586,21 +591,14 @@ private TimeSpan RemainingUntil(long until) private long Ticks(TimeSpan span) => SuppressedResponseWindows.Ticks(span, _time.Frequency); - // Production waits on the wall clock. An injected actor waits on its own, which is what - // lets a test move ListenerStartDelay without sleeping (#171). - private async Task WaitOnClockAsync(TimeSpan delay, CancellationToken cancellationToken) + // 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) { - if (_clock is null) - { - await Task.Delay(delay, cancellationToken).ConfigureAwait(false); - return; - } - - if (delay <= TimeSpan.Zero) return; 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)); + using var handle = clock.Schedule(delay, () => done.TrySetResult(true)); await done.Task.ConfigureAwait(false); } 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..a5baedb 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,53 @@ 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); + } } diff --git a/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs index 47358fc..27ca4cd 100644 --- a/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/Uds/UdsFunctionalClientTests.cs @@ -73,8 +73,13 @@ public WindowClock() 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 service = new CanBusService(bus, clock.Actor.TimeSource.GetTimestamp); var iso = new IsoTpFunctionalClient(service, FunctionalTxId, Ecu1, 0x7EF, options ?? FastOptions(), ownsService: true, clock.Actor); return UdsFunctionalClient.Create(iso, clock.Actor, ownsClient: true, @@ -687,13 +692,13 @@ public async Task A_Listener_Starting_After_The_Window_Still_Hears_What_It_Buffe using var functional = OpenOnClock(busTester, clock, responseWindow: TimeSpan.FromMilliseconds(100), responsePendingWindow: TimeSpan.FromMilliseconds(1500)); - functional.ListenerStartDelay = TimeSpan.FromMilliseconds(150); // past the 100 ms window + 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 await WaitUntilArmed(clock, TimeSpan.FromMilliseconds(150), TimeSpan.FromMilliseconds(20)); await clock.Clock.AdvanceAsync(TimeSpan.FromMilliseconds(150)); // the worker starts, reads the buffered 0x78 - functional.ListenerStartDelay = TimeSpan.Zero; // the next call's listener is not the late one + 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 @@ -1164,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); + } } From 482f7d1f1d2f25f66480d86805cb7b005331e8c9 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 13:42:30 +0000 Subject: [PATCH 05/11] refactor(isotp): dispose a clocked window's source on its actor CodeQL flagged the source the clocked FunctionalWindow never disposed. Disposing it directly would race a timer callback the loop has already taken, since cancelling the timer only flags it. The window now has one linked source again: the timer cancels it on the actor, and Dispose posts the disposal to that same actor, so the two cannot overlap and nothing is left undisposed. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../IsoTpFunctionalListener.cs | 25 ++++++++++++------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs index 9c2acb9..851a95e 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs @@ -136,30 +136,37 @@ public void Dispose() internal sealed class FunctionalWindow : IDisposable { private readonly CancellationTokenSource _linked; + private readonly ProtocolActor? _actor; private readonly IDisposable? _timer; public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken cancellationToken) { + _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); if (actor is null) { - _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); _linked.CancelAfter(window); return; } - // The actor's callback cancels a source of its own, which is never disposed -- it holds - // no timer and no wait handle -- so a callback already running when the window is - // disposed cannot meet a disposed source. Disposing the linked source unhooks it. - var end = new CancellationTokenSource(); - _linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, end.Token); - _timer = actor.Schedule(window, end.Cancel); + _actor = actor; + _timer = actor.Schedule(window, _linked.Cancel); } public CancellationToken Token => _linked.Token; public void Dispose() { - _timer?.Dispose(); - _linked.Dispose(); + if (_actor is null) + { + _linked.Dispose(); + return; + } + + // Cancelling the timer only flags it: a callback the loop has already taken still runs. + // It runs on the actor, so the source is disposed there too, after it -- a window + // disposed while its timer fires never cancels a disposed source. The actor is the + // caller's and outlives the collections on it. + _timer!.Dispose(); + _actor.Post(_linked.Dispose); } } From cfe03a48b33e7cedf0e33417102e560c8e23959c Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:03:23 +0000 Subject: [PATCH 06/11] refactor(isotp): dispose a clocked window without the actor's help Posting the linked source's disposal to the injected actor threw ObjectDisposedException once that actor was gone. Client disposal does not wait for listener tasks, so a collection can unwind after the actor has been disposed (Codex and Bugbot on #183). The window now takes a lock around its timer callback and its disposal. A callback the loop picked up before the timer was cancelled finds the window disposed and does nothing, and disposal no longer needs the actor to be running. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../IsoTpFunctionalListener.cs | 31 ++++++++++++------- .../IsoTp/IsoTpFunctionalClientTests.cs | 24 ++++++++++++++ 2 files changed, 43 insertions(+), 12 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs index 851a95e..4a8305c 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs @@ -136,8 +136,9 @@ public void Dispose() internal sealed class FunctionalWindow : IDisposable { private readonly CancellationTokenSource _linked; - private readonly ProtocolActor? _actor; private readonly IDisposable? _timer; + private readonly object _gate = new(); + private bool _disposed; public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken cancellationToken) { @@ -148,25 +149,31 @@ public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken return; } - _actor = actor; - _timer = actor.Schedule(window, _linked.Cancel); + _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() { - if (_actor is null) + _timer?.Dispose(); + lock (_gate) { + _disposed = true; _linked.Dispose(); - return; } - - // Cancelling the timer only flags it: a callback the loop has already taken still runs. - // It runs on the actor, so the source is disposed there too, after it -- a window - // disposed while its timer fires never cancels a disposed source. The actor is the - // caller's and outlives the collections on it. - _timer!.Dispose(); - _actor.Post(_linked.Dispose); } } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs index a5baedb..6477a84 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs @@ -710,4 +710,28 @@ public void An_Unstamped_Functional_Response_Is_Stamped_From_The_Collectors_Cloc .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"); + } + + 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"); + } } From 16c2c5c314eb5278c995a94555a1641234658316 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:28:09 +0000 Subject: [PATCH 07/11] refactor(j1939tp): wait for a borrowed actor's session cleanup on dispose An owned actor runs the posted session cleanup as it is disposed, but an injected one keeps running. If it was busy with the caller's other work, Dispose returned with sends still in flight and the service was torn down underneath them (Codex on #183). Dispose now waits up to 2 s for the cleanup on a borrowed actor. It fails the sessions inline when it runs on the actor's own loop, where waiting would deadlock, and when the actor is already disposed, where no loop is left to do it. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- src/CanKit.Pro.J1939Tp/J1939TpChannel.cs | 56 +++++++---- .../TestCases/J1939TpTests.cs | 95 +++++++++++++++++++ 2 files changed, 132 insertions(+), 19 deletions(-) diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index da1a3c6..e86ec0d 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs @@ -319,32 +319,34 @@ public void Dispose() _pduInbox.Writer.TryComplete(); // 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; an actor already disposed runs no loop that could race + // this thread for them (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) + { + FailSessionsOnDispose(); + } } try { _readerTask.Wait(TimeSpan.FromSeconds(2)); } catch { /* observed via task; not fatal */ } + // An owned actor runs the cleanup as it is disposed below. A borrowed one keeps running, + // possibly busy with the caller's other work, so the cleanup is waited for here: once + // this returns no session may still be using the service (Codex on #183). + if (!_ownsActor) cleanup.Wait(TimeSpan.FromSeconds(2)); + _subscription.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. @@ -356,6 +358,22 @@ public void Dispose() _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. // ----------------------------------------------------------------------------------------- diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index e67b5e5..a7a8d23 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -2467,6 +2467,101 @@ public async Task Bam_Packet_Spacing_Follows_The_Injected_Actors_Clock() (await actor.PostAsync(() => 7).WaitAsync(ShortTimeout)).Should().Be(7, "disposing the channel must not dispose the actor it was given"); } + + // 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<(J1939TpChannel Sender, Task[] Sends)> SendBamInFlight( + VirtualClock clock, ProtocolActor actor, ICanBusService service) + { + var spacing = TimeSpan.FromMilliseconds(50); + var sender = new J1939TpChannel(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: spacing), ownsService: false, actor); + var first = sender.SendBamAsync(0xFECBu, RandomPayload(9, seed: 183)); + await clock.WaitUntilTimerArmedAsync(actor, spacing, 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 (sender, 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); + var (sender, sends) = await SendBamInFlight(clock, actor, service); + + // 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); + var (sender, sends) = await SendBamInFlight(clock, actor, service); + + 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); + } + + // A borrowed actor disposed before the channel runs no loop any more: the channel fails its + // sessions itself rather than leaving the send hanging on a timer that will never fire. + [Fact] + public async Task Disposing_After_The_Borrowed_Actor_Still_Fails_The_Send() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + var (sender, sends) = await SendBamInFlight(clock, actor, service); + + actor.Dispose(); + sender.Invoking(c => c.Dispose()).Should().NotThrow(); + + sends.Should().OnlyContain(t => t.IsCompleted); + await ShouldAllFailDisposed(sends); + } + + private static async Task ShouldAllFailDisposed(Task[] sends) + { + foreach (var send in sends) + await send.Invoking(t => t).Should().ThrowAsync(); + } } /// From 3a53266d1e8eef85d67c1dad96988ffc2f898a29 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:33:19 +0000 Subject: [PATCH 08/11] test(j1939tp): dispose the windows and senders the new tests create CodeQL alert 405: the late FunctionalWindow in the clocked-window test was not disposed if an assertion threw first. The J1939-TP dispose tests had the same shape, with a helper that built and returned the sender. Each test now owns what it creates with `using`, and both Dispose methods are idempotent. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../IsoTp/IsoTpFunctionalClientTests.cs | 2 +- .../TestCases/J1939TpTests.cs | 25 +++++++++++-------- 2 files changed, 16 insertions(+), 11 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs index 6477a84..e2e8faa 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs @@ -729,7 +729,7 @@ public async Task A_Clocked_Window_Ends_On_Its_Clock_And_Outlives_Its_Actor() ended.Token.IsCancellationRequested.Should().BeTrue("the window ends when its clock says so"); } - var late = new FunctionalWindow(actor, TimeSpan.FromMilliseconds(10), CancellationToken.None); + 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 a7a8d23..b140407 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -2468,22 +2468,24 @@ public async Task Bam_Packet_Spacing_Follows_The_Injected_Actors_Clock() "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<(J1939TpChannel Sender, Task[] Sends)> SendBamInFlight( - VirtualClock clock, ProtocolActor actor, ICanBusService service) + private static async Task SendBamsInFlight(VirtualClock clock, ProtocolActor actor, J1939TpChannel sender) { - var spacing = TimeSpan.FromMilliseconds(50); - var sender = new J1939TpChannel(service, sourceAddress: 0x10, - new J1939TpOptions().With(bamPacketSpacing: spacing), ownsService: false, actor); var first = sender.SendBamAsync(0xFECBu, RandomPayload(9, seed: 183)); - await clock.WaitUntilTimerArmedAsync(actor, spacing, ShortTimeout); + 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 (sender, new[] { first, queued }); + return new[] { first, queued }; } // Codex on #183: a borrowed actor is not drained by disposing it, so Dispose must wait for @@ -2496,7 +2498,8 @@ public async Task Disposing_On_A_Busy_Borrowed_Actor_Fails_The_Send_Before_It_Re var actor = clock.NewActor(); using var bus = ControllableBus.EchoCapable(NewSession()); using var service = new CanBusService(bus); - var (sender, sends) = await SendBamInFlight(clock, actor, service); + 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 @@ -2527,7 +2530,8 @@ public async Task Disposing_From_The_Borrowed_Actors_Loop_Fails_The_Send_At_Once var actor = clock.NewActor(); using var bus = ControllableBus.EchoCapable(NewSession()); using var service = new CanBusService(bus); - var (sender, sends) = await SendBamInFlight(clock, actor, service); + using var sender = BamSenderOn(actor, service); + var sends = await SendBamsInFlight(clock, actor, sender); var failedInside = await actor.PostAsync(() => { @@ -2548,7 +2552,8 @@ public async Task Disposing_After_The_Borrowed_Actor_Still_Fails_The_Send() var actor = clock.NewActor(); using var bus = ControllableBus.EchoCapable(NewSession()); using var service = new CanBusService(bus); - var (sender, sends) = await SendBamInFlight(clock, actor, service); + using var sender = BamSenderOn(actor, service); + var sends = await SendBamsInFlight(clock, actor, sender); actor.Dispose(); sender.Invoking(c => c.Dispose()).Should().NotThrow(); From e9496145343fe28b803bf8d0ba0acc7b7ea93a16 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:44:17 +0000 Subject: [PATCH 09/11] refactor(j1939tp): order a borrowed actor's dispose and bound its wait Three findings on 16c2c5c (Codex and Bugbot on #183): - The reader is now joined and the subscription closed before the session cleanup is posted. A frame the reader was still handing to the actor is queued ahead of the cleanup, not behind it. - When the injected actor is already disposed, the channel no longer fails its sessions from the calling thread. They are actor state, and its loop may still be draining. - When the borrowed actor is still busy after 2 s, Dispose reports a TimeoutException on BackgroundExceptionOccurred. It disposes an owned service only once the cleanup has run, not underneath it. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- src/CanKit.Pro.J1939Tp/J1939TpChannel.cs | 48 +++++++--- .../TestCases/J1939TpTests.cs | 93 ++++++++++++++++++- 2 files changed, 125 insertions(+), 16 deletions(-) diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index e86ec0d..15503db 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs @@ -318,11 +318,16 @@ 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 { /* 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. 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; an actor already disposed runs no loop that could race - // this thread for them (Codex on #183). + // run after this call returned (Codex on #183). var cleanup = Task.CompletedTask; if (_actor.IsOnCurrentActor) { @@ -336,24 +341,43 @@ public void Dispose() } catch (ObjectDisposedException) { - FailSessionsOnDispose(); + // 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 */ } - - // An owned actor runs the cleanup as it is disposed below. A borrowed one keeps running, - // possibly busy with the caller's other work, so the cleanup is waited for here: once - // this returns no session may still be using the service (Codex on #183). - if (!_ownsActor) cleanup.Wait(TimeSpan.FromSeconds(2)); - - _subscription.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; - if (_ownsActor) _actor.Dispose(); _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(); } diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index b140407..321a1d3 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -2543,10 +2543,11 @@ public async Task Disposing_From_The_Borrowed_Actors_Loop_Fails_The_Send_At_Once await ShouldAllFailDisposed(sends); } - // A borrowed actor disposed before the channel runs no loop any more: the channel fails its - // sessions itself rather than leaving the send hanging on a timer that will never fire. + // 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_Still_Fails_The_Send() + public async Task Disposing_After_The_Borrowed_Actor_Leaves_Its_Sessions_To_It() { using var clock = new VirtualClock(); var actor = clock.NewActor(); @@ -2558,8 +2559,49 @@ public async Task Disposing_After_The_Borrowed_Actor_Still_Fails_The_Send() actor.Dispose(); sender.Invoking(c => c.Dispose()).Should().NotThrow(); - sends.Should().OnlyContain(t => t.IsCompleted); + 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()); + var service = new DisposalRecordingBusService(new CanBusService(bus)); + 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); } private static async Task ShouldAllFailDisposed(Task[] sends) @@ -2569,6 +2611,49 @@ private static async Task ShouldAllFailDisposed(Task[] sends) } } +/// +/// 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); + } +} + /// /// Test double: rejects every TP.CM frame at SendConfirmed, forwards everything else. /// Used to prove BAM/RTS TX failure fails the send TCS (Bugbot 3596183535). From ddae84dc25f0b86eb812a60db26b65ac1f2747be Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 14:46:19 +0000 Subject: [PATCH 10/11] test(j1939tp): own the recording service in the stuck-actor test The channel disposes the service it owns, and the test's `using` covers an assertion that throws before that. CanBusService.Dispose is idempotent, so disposing twice is harmless. This is the same shape CodeQL alert 406 flagged in the earlier helper. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index 321a1d3..a6cdb42 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -2572,7 +2572,7 @@ public async Task Disposing_On_A_Borrowed_Actor_Stuck_Past_The_Budget_Defers_The using var clock = new VirtualClock(); var actor = clock.NewActor(); using var bus = ControllableBus.EchoCapable(NewSession()); - var service = new DisposalRecordingBusService(new CanBusService(bus)); + 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); From 5f84e17f54fd41b751d47853ec9c734065d0e94b Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 15:15:33 +0000 Subject: [PATCH 11/11] refactor(j1939tp): drop frames the actor reaches after dispose began A reader that outlived its 2 s join could still post HandleIncoming to a borrowed actor after the session cleanup had run, and open a session nobody fails (Codex on #183). HandleIncoming now returns immediately once Dispose has begun, whenever the actor gets to the frame. The reader-join catch is narrowed to AggregateException, which is what Task.Wait throws for a faulted or cancelled reader (CodeQL alert 407). Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- src/CanKit.Pro.J1939Tp/J1939TpChannel.cs | 6 +- .../TestCases/J1939TpTests.cs | 151 ++++++++++++++++++ 2 files changed, 156 insertions(+), 1 deletion(-) diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index 15503db..b7f21ad 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs @@ -321,7 +321,7 @@ public void Dispose() // 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 { /* observed via task; not fatal */ } + 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 @@ -469,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/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index a6cdb42..894b271 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -2604,6 +2604,62 @@ public async Task Disposing_On_A_Borrowed_Actor_Stuck_Past_The_Budget_Defers_The 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) @@ -2611,6 +2667,101 @@ private static async Task ShouldAllFailDisposed(Task[] sends) } } +/// +/// 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).