From 89ce021c49e9ba1ddf793d7a89925ef9ce20c1a2 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 06:01:10 +0200 Subject: [PATCH 01/11] fix(isotp): fault a waiting receive when the service goes, guard a send after dispose, validate the protocol timers Three defects from the 2026-09-30 deep review (I6, I7, I8). A channel opened over a shared service with leaveOpen outlived the service when its owner disposed it first. The reader task ended quietly, and a ReceiveAsync waiting on the inbox waited for ever, while a send failed at once. The channel now puts the reason (an ObjectDisposedException for the bus service) into the inbox as a fault, raises it on BackgroundExceptionOccurred and completes the inbox behind it. Dispose can fall between SendAsync's disposed check and the post of the send to the actor. Its cleanup is then queued ahead of the send, finds nothing in flight, and the send that runs afterwards began on a disposed channel: a frame on the bus and, on an owned actor that is already gone, a call nobody completes. BeginSendOnLoop now has the guard HandleReceivedFrame already had. NAs, NBs and NCr are validated where LocalStMin already was: a negative one made the deadline scheduler throw after the send was on the wire or the reception begun, leaving it without a deadline. They must be greater than zero. Each fix has a regression test that fails with that fix removed. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 51 +++++++- .../IsoTp/IsoTpChannelIntegrationTests.cs | 120 ++++++++++++++++++ 2 files changed, 169 insertions(+), 2 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index d9be5b6..2ac77d1 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -162,6 +162,12 @@ internal IsoTpChannel(ICanBusService service, IsoTpEndpoint endpoint, // Encode once: EncodeStMin throws on negative values; surfacing at Open keeps the RX // path free of codec throws that ProtocolActor would only raise as BackgroundException. _localStMinRaw = IsoTpFrameCodec.EncodeStMin(_options.LocalStMin); + // The same for the protocol timers: a negative one makes the deadline scheduler throw + // after a send is already on the wire or a reception already started, which leaves it + // with no deadline at all. Zero would fire on the next loop iteration, i.e. never wait. + RequirePositive(nameof(_options.NAs), _options.NAs); + RequirePositive(nameof(_options.NBs), _options.NBs); + RequirePositive(nameof(_options.NCr), _options.NCr); var inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity)) { @@ -582,19 +588,51 @@ void FailInFlightSend() // ----------------------------------------------------------------------------------------- private async Task RunReaderAsync() { + Exception lost; try { while (await _subscription.WaitToReadAsync(_readerCts.Token).ConfigureAwait(false)) PumpSubscription(); + + // The subscription ended on its own: the service it came from was disposed under a + // channel opened with leaveOpen. + lost = new ObjectDisposedException(nameof(ICanBusService), + "The bus service was disposed while the ISO-TP channel was open."); } catch (OperationCanceledException) { - // expected on Dispose + return; // expected on Dispose } catch (Exception ex) { - RaiseBackgroundException(ex); + lost = ex; + } + + // Nothing will ever arrive again. A receiver waiting on the inbox would wait for ever, so + // it gets the reason as a fault and the inbox is completed behind it; sends already fail + // on their own, the bus refuses them. On the actor: the inbox has one writer. + try + { + _actor.Post(() => EndInboxAfterSubscriptionLoss(lost)); } + catch (ObjectDisposedException) + { + // The channel is going down too, and its disposal completes the inbox. + if (!(lost is ObjectDisposedException)) RaiseBackgroundException(lost); + } + } + + private void EndInboxAfterSubscriptionLoss(Exception lost) + { + if (Volatile.Read(ref _disposed) != 0) return; + AbortRx(lost); + _pduInbox.Writer.TryComplete(); + } + + private static void RequirePositive(string name, TimeSpan value) + { + if (value <= TimeSpan.Zero) + throw new ArgumentOutOfRangeException(name, value, name + " must be greater than zero."); } /// @@ -718,6 +756,15 @@ private void BeginSendOnLoop(byte[] pdu, TaskCompletionSource(); + channel.BackgroundExceptionOccurred += (_, ex) => { lock (reported) reported.Add(ex); }; + + var waiting = channel.ReceiveAsync(); + service.Dispose(); + + Func act = () => waiting.WaitAsync(ShortTimeout); + var thrown = (await act.Should().ThrowAsync()).Which; + thrown.ObjectName.Should().Be(nameof(ICanBusService)); + + // The inbox is complete behind the fault: the next receive ends at once instead of waiting. + Func next = () => channel.ReceiveAsync().WaitAsync(ShortTimeout); + await next.Should().ThrowAsync(); + lock (reported) reported.Should().ContainSingle().Which.Should().BeOfType(); + } + + // I7: Dispose can fall between SendAsync's disposed check and the post of the send to the + // actor. Its cleanup is then queued ahead of the send, finds nothing in flight, and the send + // that runs afterwards used to begin on a disposed channel: a frame on the bus, no one to + // complete the call. The actor double below makes Dispose run at exactly that point. + [Fact] + public async Task A_Dispose_Between_The_Check_And_The_Post_Fails_The_Send_Before_It_Reaches_The_Bus() + { + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var inner = new ProtocolActor(); + using var actor = new DisposeOnFirstPostActor(inner); + using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), + ownsService: false, actor); + actor.Channel = channel; + + Func act = () => channel.SendAsync(new byte[] { 0x01, 0x02 }).WaitAsync(ShortTimeout); + + await act.Should().ThrowAsync(); + bus.TransmitCount.Should().Be(0, "a send that began after the dispose would have put its frame on the bus"); + } + + /// + /// Runs the channel's Dispose on another thread when the first send is posted and holds that + /// post until the dispose has posted its own cleanup, so the cleanup is queued first. + /// + private sealed class DisposeOnFirstPostActor : IProtocolActor + { + private readonly IProtocolActor _inner; + private readonly ManualResetEventSlim _cleanupPosted = new(); + private int _posts; + + public DisposeOnFirstPostActor(IProtocolActor inner) => _inner = inner; + + public IDisposable? Channel { get; set; } + + public event EventHandler? BackgroundExceptionOccurred + { + add => _inner.BackgroundExceptionOccurred += value; + remove => _inner.BackgroundExceptionOccurred -= value; + } + + public void Post(Action work) + { + if (Interlocked.Increment(ref _posts) == 1) + { + var channel = Channel!; + _ = Task.Run(channel.Dispose); + if (!_cleanupPosted.Wait(TimeSpan.FromSeconds(5))) + throw new TimeoutException("Dispose did not post its cleanup."); + _inner.Post(work); + return; + } + + _inner.Post(work); + _cleanupPosted.Set(); + } + + public Task PostAsync(Action work, CancellationToken cancellationToken = default) + => _inner.PostAsync(work, cancellationToken); + + public Task PostAsync(Func work, CancellationToken cancellationToken = default) + => _inner.PostAsync(work, cancellationToken); + + public IDisposable Schedule(TimeSpan delay, Action callback) => _inner.Schedule(delay, callback); + + public void Dispose() => _cleanupPosted.Dispose(); + } + + // I8: a negative (or zero) protocol timer made the deadline scheduler throw after the send was + // on the wire or the reception begun, leaving it with no deadline; reject it where LocalStMin + // already is. + [Theory] + [InlineData(nameof(IsoTpChannelOptions.NAs), 0)] + [InlineData(nameof(IsoTpChannelOptions.NAs), -1)] + [InlineData(nameof(IsoTpChannelOptions.NBs), 0)] + [InlineData(nameof(IsoTpChannelOptions.NBs), -1)] + [InlineData(nameof(IsoTpChannelOptions.NCr), 0)] + [InlineData(nameof(IsoTpChannelOptions.NCr), -1)] + public void A_Protocol_Timer_That_Cannot_Arm_Is_Rejected_At_Open(string option, int milliseconds) + { + var session = NewSession(); + using var bus = OpenClassic(session, 0); + var value = TimeSpan.FromMilliseconds(milliseconds); + var options = option switch + { + nameof(IsoTpChannelOptions.NAs) => new IsoTpChannelOptions { NAs = value }, + nameof(IsoTpChannelOptions.NBs) => new IsoTpChannelOptions { NBs = value }, + _ => new IsoTpChannelOptions { NCr = value }, + }; + + Action act = () => IsoTpFactory.Open(bus, IsoTpEndpoint.Normal(0x250, 0x251), options); + act.Should().Throw().WithParameterName(option); + } + // -------------------------------------------------------------------------------- // FR-TP-016 — DatagramReceived event fires for a SF PDU. // -------------------------------------------------------------------------------- From 2eeef26ce25462cd2c478c64d0e028ae1399c551 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 06:13:32 +0200 Subject: [PATCH 02/11] fix(isotp): complete the inbox with the loss, not a fault item Review on #261: when the shared service went, the reason was enqueued as a fault item. In the bounded DropOldest inbox that item could push the oldest finished PDU out, and with several blocked receivers only one saw the real failure while the others got the generic completion error. The inbox is now completed with the error: buffered PDUs stay readable, then every receiver, present and later, gets the ObjectDisposedException. A reassembly under way is dropped without an item. Also takes the service in the test under a using, as the code scanner asked. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 20 ++++++-- .../IsoTp/IsoTpChannelIntegrationTests.cs | 47 +++++++++++++++---- 2 files changed, 54 insertions(+), 13 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index 2ac77d1..c968875 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -609,8 +609,9 @@ private async Task RunReaderAsync() } // Nothing will ever arrive again. A receiver waiting on the inbox would wait for ever, so - // it gets the reason as a fault and the inbox is completed behind it; sends already fail - // on their own, the bus refuses them. On the actor: the inbox has one writer. + // the inbox is completed with the reason: what is buffered stays readable, and then every + // receiver -- not just one -- gets the failure. Sends already fail on their own, the bus + // refuses them. On the actor: the inbox has one writer. try { _actor.Post(() => EndInboxAfterSubscriptionLoss(lost)); @@ -625,8 +626,19 @@ private async Task RunReaderAsync() private void EndInboxAfterSubscriptionLoss(Exception lost) { if (Volatile.Read(ref _disposed) != 0) return; - AbortRx(lost); - _pduInbox.Writer.TryComplete(); + + // A reassembly under way dies with the subscription. No fault item for it: in the bounded + // inbox it would push the oldest finished PDU out, and only one receiver would see it. + var rx = _rx; + if (rx is not null) + { + rx.CancelDeadline(); + _rx = null; + WithdrawReception(rx.Announce); + } + + _pduInbox.Writer.TryComplete(lost); + RaiseBackgroundException(lost); } private static void RequirePositive(string name, TimeSpan value) diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 044e0cf..f4f8a71 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -852,30 +852,59 @@ public async Task Dispose_Unblocks_Pending_ReceiveAsync() // The 2026-09-30 review (I6): a channel opened over a shared service (leaveOpen) outlives the // service when the owner disposes it first. The channel's reader ended quietly, and a receiver - // waiting on the inbox waited for ever; a send, by contrast, failed at once. + // waiting on the inbox waited for ever; a send, by contrast, failed at once. The inbox is + // completed with the reason: what is buffered stays readable, then every receiver gets it. [Fact] - public async Task A_Service_Disposed_Under_The_Channel_Faults_A_Waiting_Receive() + public async Task A_Service_Disposed_Under_The_Channel_Faults_Every_Waiting_Receive() { var session = NewSession(); using var bus = OpenClassic(session, 0); - var service = new CanBusService(bus); + using var service = new CanBusService(bus); using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); var reported = new List(); channel.BackgroundExceptionOccurred += (_, ex) => { lock (reported) reported.Add(ex); }; - var waiting = channel.ReceiveAsync(); + var first = channel.ReceiveAsync(); + var second = channel.ReceiveAsync(); service.Dispose(); - Func act = () => waiting.WaitAsync(ShortTimeout); - var thrown = (await act.Should().ThrowAsync()).Which; - thrown.ObjectName.Should().Be(nameof(ICanBusService)); + foreach (var waiting in new[] { first, second }) + { + Func act = () => waiting.WaitAsync(ShortTimeout); + var thrown = (await act.Should().ThrowAsync()).Which; + thrown.ObjectName.Should().Be(nameof(ICanBusService)); + } - // The inbox is complete behind the fault: the next receive ends at once instead of waiting. + // And a receive after the loss ends at once with the same reason instead of waiting. Func next = () => channel.ReceiveAsync().WaitAsync(ShortTimeout); - await next.Should().ThrowAsync(); + await next.Should().ThrowAsync(); lock (reported) reported.Should().ContainSingle().Which.Should().BeOfType(); } + // A PDU that had already arrived is not lost with the service: it stays readable, and the + // failure follows it. + [Fact] + public async Task A_PDU_Received_Before_The_Service_Goes_Is_Still_Delivered_First() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var service = new CanBusService(busA); + using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); + + // EmitPdu enqueues before it raises the event, so the event proves the PDU is in the inbox. + var arrived = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.DatagramReceived += (_, _) => arrived.TrySetResult(true); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + await arrived.Task.WaitAsync(ShortTimeout); + + service.Dispose(); + + (await channel.ReceiveAsync().WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); + Func next = () => channel.ReceiveAsync().WaitAsync(ShortTimeout); + await next.Should().ThrowAsync(); + } + // I7: Dispose can fall between SendAsync's disposed check and the post of the send to the // actor. Its cleanup is then queued ahead of the send, finds nothing in flight, and the send // that runs afterwards used to begin on a disposed channel: a frame on the bus, no one to From 9ff425f4123d4337defe874cca38edde57eb1689 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 06:19:52 +0200 Subject: [PATCH 03/11] test(isotp): keep the receive-fault assertions out of a loop that only maps its variable Code scanning (cs/linq/missed-select) on #261. Co-Authored-By: Claude Sonnet 5.5 --- .../TestCases/IsoTp/IsoTpChannelIntegrationTests.cs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index f4f8a71..ab9cb35 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -868,13 +868,16 @@ public async Task A_Service_Disposed_Under_The_Channel_Faults_Every_Waiting_Rece var second = channel.ReceiveAsync(); service.Dispose(); - foreach (var waiting in new[] { first, second }) + static async Task FaultsWithTheLossAsync(Task waiting) { Func act = () => waiting.WaitAsync(ShortTimeout); var thrown = (await act.Should().ThrowAsync()).Which; thrown.ObjectName.Should().Be(nameof(ICanBusService)); } + await FaultsWithTheLossAsync(first); + await FaultsWithTheLossAsync(second); + // And a receive after the loss ends at once with the same reason instead of waiting. Func next = () => channel.ReceiveAsync().WaitAsync(ShortTimeout); await next.Should().ThrowAsync(); From 5b498892874ebcd0d37544372357e2d6d1b52563 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 06:44:52 +0200 Subject: [PATCH 04/11] fix(isotp): cover the subscription-loss paths and drop two branches that decided nothing Codecov on #261. The reader failing, an injected actor that is already gone, and a reassembly under way when the service goes now have tests; the disposed guard in the actor callback and the special case in the catch did not change any outcome and are gone. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 7 +-- .../Infrastructure/StarvedReaderBusService.cs | 9 ++- .../IsoTp/IsoTpChannelIntegrationTests.cs | 61 +++++++++++++++++++ 3 files changed, 72 insertions(+), 5 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index c968875..e6067fb 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -618,15 +618,14 @@ private async Task RunReaderAsync() } catch (ObjectDisposedException) { - // The channel is going down too, and its disposal completes the inbox. - if (!(lost is ObjectDisposedException)) RaiseBackgroundException(lost); + // The actor is gone (an injected one its owner disposed): there is nobody to write the + // inbox for. The loss is still reported; the channel's own disposal completes the inbox. + RaiseBackgroundException(lost); } } private void EndInboxAfterSubscriptionLoss(Exception lost) { - if (Volatile.Read(ref _disposed) != 0) return; - // A reassembly under way dies with the subscription. No fault item for it: in the bounded // inbox it would push the oldest finished PDU out, and only one receiver would see it. var rx = _rx; diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index 9dfd5a2..82986bf 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -37,6 +37,9 @@ public void Deliver(CanFrameView frame, long hostArrivalTimestamp = 0) => _frame /// public bool HoldFrames { get; set; } + /// When set, the reader task's wait throws it once woken, as a subscription that failed would. + public Exception? ReaderFault { get; set; } + /// Lets the reader task's wait complete once; every later wait stays pending. public void WakeReader() => _wake.TrySetResult(true); @@ -110,7 +113,11 @@ public bool TryRead(out CanFrameEvent frameEvent) public async ValueTask WaitToReadAsync(CancellationToken cancellationToken = default) { if (Interlocked.Exchange(ref _owner._wakesServed, 1) == 0) - return await _owner._wake.Task.WaitAsync(cancellationToken); + { + var woken = await _owner._wake.Task.WaitAsync(cancellationToken); + if (_owner.ReaderFault is { } fault) throw fault; + return woken; + } await Task.Delay(Timeout.Infinite, cancellationToken); return false; } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index ab9cb35..ee88359 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -908,6 +908,67 @@ public async Task A_PDU_Received_Before_The_Service_Goes_Is_Still_Delivered_Firs await next.Should().ThrowAsync(); } + // The reader itself failing is the same loss: the inbox ends with that failure, for every + // receiver, and it is reported once. + [Fact] + public async Task A_Failing_Subscription_Ends_The_Inbox_With_Its_Failure() + { + var service = new StarvedReaderBusService { ReaderFault = new InvalidOperationException("the demux broke") }; + using var actor = new ProtocolActor(); + using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), + ownsService: false, actor); + var reported = new List(); + channel.BackgroundExceptionOccurred += (_, ex) => { lock (reported) reported.Add(ex); }; + + var waiting = channel.ReceiveAsync(); + service.WakeReader(); + + Func act = () => waiting.WaitAsync(ShortTimeout); + (await act.Should().ThrowAsync()).Which.Message.Should().Be("the demux broke"); + lock (reported) reported.Should().ContainSingle().Which.Message.Should().Be("the demux broke"); + } + + // A reassembly under way dies with the subscription: it is withdrawn, so + // GetReceptionsInProgress does not keep reporting a transfer nobody will complete. + [Fact] + public async Task A_Reception_Under_Way_Is_Withdrawn_When_The_Service_Goes() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var service = new CanBusService(busA); + using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); + + var waiting = channel.ReceiveAsync(); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x10, 0x14, 1, 2, 3, 4, 5, 6 })); // first frame of 20 bytes + (await WaitForReceptionsAsync(channel, 1)).Should().HaveCount(1); + + service.Dispose(); + + Func act = () => waiting.WaitAsync(ShortTimeout); + await act.Should().ThrowAsync(); + channel.GetReceptionsInProgress().Should().BeEmpty(); + } + + // The actor an owner disposed first cannot write the inbox, but the loss is still reported. + [Fact] + public async Task The_Loss_Is_Still_Reported_When_The_Injected_Actor_Is_Already_Gone() + { + var session = NewSession(); + using var bus = OpenClassic(session, 0); + using var service = new CanBusService(bus); + var actor = new ProtocolActor(); + using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), + ownsService: false, actor); + var reported = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.BackgroundExceptionOccurred += (_, ex) => reported.TrySetResult(ex); + + actor.Dispose(); + service.Dispose(); + + (await reported.Task.WaitAsync(ShortTimeout)).Should().BeOfType(); + } + // I7: Dispose can fall between SendAsync's disposed check and the post of the send to the // actor. Its cleanup is then queued ahead of the send, finds nothing in flight, and the send // that runs afterwards used to begin on a disposed channel: a frame on the bus, no one to From 4ae8e1fd9691193e984a8256e8e39e3165417c3c Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 06:59:56 +0200 Subject: [PATCH 05/11] fix(isotp): keep the PDUs a discard retains after the service was lost Review on #261: DiscardPendingPdus(arrivedBefore) drains the inbox and writes the items from after its stamp back; once the inbox is completed with the loss that write fails and the retained PDUs were silently dropped, although the UDS client calls it exactly then. They now move to a fresh inbox completed with the same failure. The actor in the test is a using declaration, as code scanning asked. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 34 ++++++++++++++----- .../IsoTp/IsoTpChannelIntegrationTests.cs | 30 +++++++++++++++- 2 files changed, 55 insertions(+), 9 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index e6067fb..f2c9f3f 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -67,7 +67,12 @@ internal sealed class IsoTpChannel : IIsoTpChannel // PDU or a reassembly-abort fault (N_Cr / SN mismatch / SF·FF supersede) so a blocked // ReceiveAsync completes instead of hanging (FR-TP-010) — // the FailTx analogue on the RX side. - private readonly Channel _pduInbox; + // + // Replaced only after the bus service was lost and a discard kept some items: a completed channel + // takes no writes, so the kept items move to a fresh one completed with the same failure. + private volatile Channel _pduInbox; + private readonly BoundedChannelOptions _inboxOptions; + private Exception? _inboxLost; // actor-side // Serializes SendAsync callers: one outbound PDU on the wire at a time, per ISO 15765-2's // "one N-USData at a time" model. Also avoids competition for _tx state across calls. @@ -169,13 +174,13 @@ internal IsoTpChannel(ICanBusService service, IsoTpEndpoint endpoint, RequirePositive(nameof(_options.NBs), _options.NBs); RequirePositive(nameof(_options.NCr), _options.NCr); - var inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity)) + _inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity)) { SingleReader = false, SingleWriter = true, // written only from the actor loop FullMode = BoundedChannelFullMode.DropOldest, }; - _pduInbox = Channel.CreateBounded(inboxOptions); + _pduInbox = Channel.CreateBounded(_inboxOptions); _actor = actor ?? new ProtocolActor(); _actor.BackgroundExceptionOccurred += OnActorBackgroundException; @@ -463,8 +468,21 @@ private int ClearReceptionsBefore(long stamp) else discarded++; } - foreach (var item in kept) - _pduInbox.Writer.TryWrite(item); + if (_inboxLost is { } lost && kept.Count > 0) + { + // The inbox was completed when the bus service went, and a completed channel takes no + // writes: the items from after the stamp move to a fresh one, completed the same way. + var next = Channel.CreateBounded(_inboxOptions); + foreach (var item in kept) + next.Writer.TryWrite(item); + next.Writer.TryComplete(lost); + _pduInbox = next; + } + else + { + foreach (var item in kept) + _pduInbox.Writer.TryWrite(item); + } return discarded; } @@ -475,10 +493,9 @@ public IAsyncEnumerable ReceiveAllAsync(CancellationToken cancellationTo private async IAsyncEnumerable ReadAllAsync( [EnumeratorCancellation] CancellationToken cancellationToken) { - var reader = _pduInbox.Reader; - while (await reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) + while (await _pduInbox.Reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) { - while (reader.TryRead(out var item)) + while (_pduInbox.Reader.TryRead(out var item)) yield return UnwrapInboxItem(item); } } @@ -636,6 +653,7 @@ private void EndInboxAfterSubscriptionLoss(Exception lost) WithdrawReception(rx.Announce); } + _inboxLost = lost; _pduInbox.Writer.TryComplete(lost); RaiseBackgroundException(lost); } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index ee88359..3516bc5 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -908,6 +908,34 @@ public async Task A_PDU_Received_Before_The_Service_Goes_Is_Still_Delivered_Firs await next.Should().ThrowAsync(); } + // Review on #261: a discard that keeps the items from after its stamp writes them back, and a + // completed inbox takes no writes. The PDUs that had arrived are still delivered, then the loss. + [Fact] + public async Task A_Discard_After_The_Service_Goes_Keeps_The_Pdus_From_After_Its_Stamp() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var service = new CanBusService(busA); + using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); + + var arrived = 0; + var both = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.DatagramReceived += (_, _) => { if (Interlocked.Increment(ref arrived) == 2) both.TrySetResult(true); }; + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0C, 0x0D, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + await both.Task.WaitAsync(ShortTimeout); + service.Dispose(); + + // A stamp of 1 is before everything: nothing is discarded, everything is kept. + channel.DiscardPendingPdus(arrivedBefore: 1).Should().Be(0); + + (await channel.ReceiveAsync().WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); + (await channel.ReceiveAsync().WaitAsync(ShortTimeout)).Should().Equal(0x0C, 0x0D); + Func next = () => channel.ReceiveAsync().WaitAsync(ShortTimeout); + await next.Should().ThrowAsync(); + } + // The reader itself failing is the same loss: the inbox ends with that failure, for every // receiver, and it is reported once. [Fact] @@ -957,7 +985,7 @@ public async Task The_Loss_Is_Still_Reported_When_The_Injected_Actor_Is_Already_ var session = NewSession(); using var bus = OpenClassic(session, 0); using var service = new CanBusService(bus); - var actor = new ProtocolActor(); + using var actor = new ProtocolActor(); using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), ownsService: false, actor); var reported = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); From e43b4bbf08627ad75f94621fcb3a6423ebbe826c Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 07:15:10 +0200 Subject: [PATCH 06/11] fix(isotp): publish the replacement inbox before draining the lost one Reviews on #261: a reader meeting the old, drained and failure-completed inbox in the middle of a discard threw the loss before the retained PDUs arrived. The writable replacement is now published first and completed with the failure after the retained items are in. The discard test waits for the loss to be reported before it discards: it raced the reader noticing the service went, and the interleaving where the discard came first passed without exercising the replacement at all. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 27 ++++++++++++------- .../IsoTp/IsoTpChannelIntegrationTests.cs | 6 +++++ 2 files changed, 24 insertions(+), 9 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index f2c9f3f..557cd67 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -458,7 +458,20 @@ private int ClearReceptionsBefore(long stamp) // opposite of the reset intent. Items from after it go back in their order; this runs // on the inbox's single writer. var kept = new List(); - while (_pduInbox.Reader.TryRead(out var item)) + var inbox = _pduInbox; + Channel? replacement = null; + if (_inboxLost is not null) + { + // The inbox was completed when the bus service went, and a completed channel takes no + // writes. A writable replacement is published before the old one is drained, so a + // reader never meets an empty inbox that already carries the failure while the items + // from after the stamp are still on their way over; it is completed with the same + // failure once they are in. + replacement = Channel.CreateBounded(_inboxOptions); + _pduInbox = replacement; + } + + while (inbox.Reader.TryRead(out var item)) { // An error item carries the first-frame stamp of the reception it aborted, so an // abort of a reception that began after the stamp is kept as that reception's @@ -468,20 +481,16 @@ private int ClearReceptionsBefore(long stamp) else discarded++; } - if (_inboxLost is { } lost && kept.Count > 0) + if (replacement is not null) { - // The inbox was completed when the bus service went, and a completed channel takes no - // writes: the items from after the stamp move to a fresh one, completed the same way. - var next = Channel.CreateBounded(_inboxOptions); foreach (var item in kept) - next.Writer.TryWrite(item); - next.Writer.TryComplete(lost); - _pduInbox = next; + replacement.Writer.TryWrite(item); + replacement.Writer.TryComplete(_inboxLost); } else { foreach (var item in kept) - _pduInbox.Writer.TryWrite(item); + inbox.Writer.TryWrite(item); } return discarded; } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 3516bc5..980e99f 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -925,7 +925,13 @@ public async Task A_Discard_After_The_Service_Goes_Keeps_The_Pdus_From_After_Its busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0C, 0x0D, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); await both.Task.WaitAsync(ShortTimeout); + + // The loss is reported after the inbox is completed with it: the discard below must meet + // that state, not race the reader noticing the service went. + var lost = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.BackgroundExceptionOccurred += (_, _) => lost.TrySetResult(true); service.Dispose(); + await lost.Task.WaitAsync(ShortTimeout); // A stamp of 1 is before everything: nothing is discarded, everything is kept. channel.DiscardPendingPdus(arrivedBefore: 1).Should().Be(0); From 2aa5c96ba3ed2b6937bd2977fa6bfacf7e1822d5 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 07:31:26 +0200 Subject: [PATCH 07/11] fix(isotp): keep the service loss beside the inbox instead of completing it Reviews on #261: moving the retained PDUs to a replacement inbox after a service loss needed the handoff serialised against readers and against Dispose, and each fix opened the next race. The inbox is no longer completed when the service goes (only Dispose completes it). The failure is recorded and the receivers are woken; a receive takes what is buffered first and throws the loss only on an empty inbox, for every receiver. DiscardPendingPdus is back to its original form: the inbox stays open, so what it retains is written back as before. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 98 ++++++++++--------- .../IsoTp/IsoTpChannelIntegrationTests.cs | 27 +++++ 2 files changed, 80 insertions(+), 45 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index 557cd67..c7c6e51 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -68,11 +68,14 @@ internal sealed class IsoTpChannel : IIsoTpChannel // ReceiveAsync completes instead of hanging (FR-TP-010) — // the FailTx analogue on the RX side. // - // Replaced only after the bus service was lost and a discard kept some items: a completed channel - // takes no writes, so the kept items move to a fresh one completed with the same failure. - private volatile Channel _pduInbox; - private readonly BoundedChannelOptions _inboxOptions; - private Exception? _inboxLost; // actor-side + // The inbox is not completed when the bus service is lost (only by Dispose): that failure is + // kept beside it. A receiver takes what is buffered first and meets the failure only on an + // empty inbox, which is the order the callers rely on, and a discard keeps writing back what + // it retains into an inbox that is still open. + private readonly Channel _pduInbox; + private volatile Exception? _inboxLost; + private readonly TaskCompletionSource _inboxLostSignal = + new(TaskCreationOptions.RunContinuationsAsynchronously); // Serializes SendAsync callers: one outbound PDU on the wire at a time, per ISO 15765-2's // "one N-USData at a time" model. Also avoids competition for _tx state across calls. @@ -174,13 +177,13 @@ internal IsoTpChannel(ICanBusService service, IsoTpEndpoint endpoint, RequirePositive(nameof(_options.NBs), _options.NBs); RequirePositive(nameof(_options.NCr), _options.NCr); - _inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity)) + var inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity)) { SingleReader = false, SingleWriter = true, // written only from the actor loop FullMode = BoundedChannelFullMode.DropOldest, }; - _pduInbox = Channel.CreateBounded(_inboxOptions); + _pduInbox = Channel.CreateBounded(inboxOptions); _actor = actor ?? new ProtocolActor(); _actor.BackgroundExceptionOccurred += OnActorBackgroundException; @@ -307,12 +310,10 @@ public async Task ReceiveAsync(CancellationToken cancellationToken = def public async Task ReceiveWithArrivalAsync( CancellationToken cancellationToken = default) { - while (await _pduInbox.Reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) - { - if (_pduInbox.Reader.TryRead(out var item)) - return new IsoTpReceivedPdu(UnwrapInboxItem(item), item.ArrivalTimestamp, item.FirstFrameArrivalTimestamp); - } - throw new InvalidOperationException("Channel is disposed; no more PDUs will arrive."); + var (taken, item) = await TakeNextAsync(cancellationToken).ConfigureAwait(false); + if (!taken) + throw new InvalidOperationException("Channel is disposed; no more PDUs will arrive."); + return new IsoTpReceivedPdu(UnwrapInboxItem(item), item.ArrivalTimestamp, item.FirstFrameArrivalTimestamp); } /// @@ -458,20 +459,7 @@ private int ClearReceptionsBefore(long stamp) // opposite of the reset intent. Items from after it go back in their order; this runs // on the inbox's single writer. var kept = new List(); - var inbox = _pduInbox; - Channel? replacement = null; - if (_inboxLost is not null) - { - // The inbox was completed when the bus service went, and a completed channel takes no - // writes. A writable replacement is published before the old one is drained, so a - // reader never meets an empty inbox that already carries the failure while the items - // from after the stamp are still on their way over; it is completed with the same - // failure once they are in. - replacement = Channel.CreateBounded(_inboxOptions); - _pduInbox = replacement; - } - - while (inbox.Reader.TryRead(out var item)) + while (_pduInbox.Reader.TryRead(out var item)) { // An error item carries the first-frame stamp of the reception it aborted, so an // abort of a reception that began after the stamp is kept as that reception's @@ -481,17 +469,8 @@ private int ClearReceptionsBefore(long stamp) else discarded++; } - if (replacement is not null) - { - foreach (var item in kept) - replacement.Writer.TryWrite(item); - replacement.Writer.TryComplete(_inboxLost); - } - else - { - foreach (var item in kept) - inbox.Writer.TryWrite(item); - } + foreach (var item in kept) + _pduInbox.Writer.TryWrite(item); return discarded; } @@ -502,10 +481,39 @@ public IAsyncEnumerable ReceiveAllAsync(CancellationToken cancellationTo private async IAsyncEnumerable ReadAllAsync( [EnumeratorCancellation] CancellationToken cancellationToken) { - while (await _pduInbox.Reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) + while (true) { - while (_pduInbox.Reader.TryRead(out var item)) - yield return UnwrapInboxItem(item); + var (taken, item) = await TakeNextAsync(cancellationToken).ConfigureAwait(false); + if (!taken) yield break; + yield return UnwrapInboxItem(item); + } + } + + /// + /// The next item of the inbox, waiting for one. False once the inbox is completed (the channel + /// was disposed). If the bus service was lost, what is buffered is delivered first and the loss + /// is thrown only on an empty inbox, to every receiver. + /// + private async Task<(bool Taken, RxInboxItem Item)> TakeNextAsync(CancellationToken cancellationToken) + { + while (true) + { + if (_pduInbox.Reader.TryRead(out var item)) return (true, item); + if (_inboxLost is { } lost) throw lost; + + using var release = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + var wait = _pduInbox.Reader.WaitToReadAsync(release.Token).AsTask(); + var first = await Task.WhenAny(wait, _inboxLostSignal.Task).ConfigureAwait(false); + if (first != wait) + { + // The loss came first: let go of the wait and look again, items before the loss. + release.Cancel(); + try { await wait.ConfigureAwait(false); } + catch (OperationCanceledException) { /* the wait we just released */ } + continue; + } + + if (!await wait.ConfigureAwait(false)) return (false, default); } } @@ -635,9 +643,9 @@ private async Task RunReaderAsync() } // Nothing will ever arrive again. A receiver waiting on the inbox would wait for ever, so - // the inbox is completed with the reason: what is buffered stays readable, and then every - // receiver -- not just one -- gets the failure. Sends already fail on their own, the bus - // refuses them. On the actor: the inbox has one writer. + // the reason is recorded and the receivers are woken: what is buffered is delivered first, + // and then every receiver -- not just one -- gets the failure. Sends already fail on their + // own, the bus refuses them. On the actor, like everything else that touches the inbox. try { _actor.Post(() => EndInboxAfterSubscriptionLoss(lost)); @@ -663,7 +671,7 @@ private void EndInboxAfterSubscriptionLoss(Exception lost) } _inboxLost = lost; - _pduInbox.Writer.TryComplete(lost); + _inboxLostSignal.TrySetResult(true); RaiseBackgroundException(lost); } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 980e99f..5d9cb07 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -942,6 +942,33 @@ public async Task A_Discard_After_The_Service_Goes_Keeps_The_Pdus_From_After_Its await next.Should().ThrowAsync(); } + // The streaming receive follows the same order: what arrived, then the loss. + [Fact] + public async Task ReceiveAll_Yields_The_Buffered_Pdu_And_Then_Throws_The_Loss() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var service = new CanBusService(busA); + using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); + + var arrived = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.DatagramReceived += (_, _) => arrived.TrySetResult(true); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + await arrived.Task.WaitAsync(ShortTimeout); + service.Dispose(); + + using var cts = new CancellationTokenSource(ShortTimeout); + var seen = new List(); + Func act = async () => + { + await foreach (var pdu in channel.ReceiveAllAsync(cts.Token)) + seen.Add(pdu); + }; + await act.Should().ThrowAsync(); + seen.Should().ContainSingle().Which.Should().Equal(0x0A, 0x0B); + } + // The reader itself failing is the same loss: the inbox ends with that failure, for every // receiver, and it is reported once. [Fact] From 18712a878d4c7fe401859ed7e59bf71d94178c9b Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 07:46:05 +0200 Subject: [PATCH 08/11] fix(isotp): a receive during a discard after the loss still gets the retained PDU first Review on #261: a discard takes the PDUs it retains out of the inbox and writes them back. After a service loss a receiver meeting the inbox empty in that gap threw the loss ahead of them. A flag marks the gap; a receiver waits it out and looks once more before it throws. The streaming receive's end on dispose is covered as well. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 51 ++++++++++++---- .../IsoTp/IsoTpChannelIntegrationTests.cs | 58 +++++++++++++++++++ 2 files changed, 98 insertions(+), 11 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index c7c6e51..9396056 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -74,6 +74,13 @@ internal sealed class IsoTpChannel : IIsoTpChannel // it retains into an inbox that is still open. private readonly Channel _pduInbox; private volatile Exception? _inboxLost; + // Set while a discard has taken the items it retains out of the inbox and not yet put them + // back: a receiver meeting an empty inbox then must not read that as "nothing more to come". + private int _discarding; + + /// For tests: runs between a discard taking the retained items out of the inbox and + /// putting them back, on the actor. + internal Action? DiscardGapObserver { get; set; } private readonly TaskCompletionSource _inboxLostSignal = new(TaskCreationOptions.RunContinuationsAsynchronously); @@ -459,18 +466,28 @@ private int ClearReceptionsBefore(long stamp) // opposite of the reset intent. Items from after it go back in their order; this runs // on the inbox's single writer. var kept = new List(); - while (_pduInbox.Reader.TryRead(out var item)) + Volatile.Write(ref _discarding, 1); + try { - // An error item carries the first-frame stamp of the reception it aborted, so an - // abort of a reception that began after the stamp is kept as that reception's - // outcome (Codex on #143). - if (item.FirstFrameArrivalTimestamp >= stamp) - kept.Add(item); - else - discarded++; + while (_pduInbox.Reader.TryRead(out var item)) + { + // An error item carries the first-frame stamp of the reception it aborted, so an + // abort of a reception that began after the stamp is kept as that reception's + // outcome (Codex on #143). + if (item.FirstFrameArrivalTimestamp >= stamp) + kept.Add(item); + else + discarded++; + } + + DiscardGapObserver?.Invoke(); + foreach (var item in kept) + _pduInbox.Writer.TryWrite(item); + } + finally + { + Volatile.Write(ref _discarding, 0); } - foreach (var item in kept) - _pduInbox.Writer.TryWrite(item); return discarded; } @@ -499,7 +516,19 @@ private async IAsyncEnumerable ReadAllAsync( while (true) { if (_pduInbox.Reader.TryRead(out var item)) return (true, item); - if (_inboxLost is { } lost) throw lost; + if (_inboxLost is { } lost) + { + // A discard has the retained items out of the inbox for a moment: wait it out. + if (Volatile.Read(ref _discarding) != 0) + { + await Task.Yield(); + continue; + } + + // Looked again after the flag: a discard that finished in between has put them back. + if (_pduInbox.Reader.TryRead(out item)) return (true, item); + throw lost; + } using var release = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); var wait = _pduInbox.Reader.WaitToReadAsync(release.Token).AsTask(); diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 5d9cb07..dd567cf 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -969,6 +969,64 @@ public async Task ReceiveAll_Yields_The_Buffered_Pdu_And_Then_Throws_The_Loss() seen.Should().ContainSingle().Which.Should().Equal(0x0A, 0x0B); } + // Review on #261: a discard after the loss has the retained items out of the inbox for a + // moment. A receiver arriving then must not read the empty inbox as the end and throw the + // loss ahead of the PDU that is about to be put back. + [Fact] + public async Task A_Receive_During_A_Discard_After_The_Loss_Still_Gets_The_Retained_Pdu_First() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var service = new CanBusService(busA); + using var channel = IsoTpFactory.Open(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), leaveOpen: true); + var concrete = (IsoTpChannel)channel; + + var arrived = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.DatagramReceived += (_, _) => arrived.TrySetResult(true); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + await arrived.Task.WaitAsync(ShortTimeout); + var lost = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.BackgroundExceptionOccurred += (_, _) => lost.TrySetResult(true); + service.Dispose(); + await lost.Task.WaitAsync(ShortTimeout); + + Task? receive = null; + concrete.DiscardGapObserver = () => + { + // On the actor, with the retained PDU out of the inbox: a receiver starts now and + // gets a moment to meet the empty inbox. No signal says "it looked and waited", so + // this is a negative window: it can only pass falsely on a slow host. + receive = Task.Run(() => channel.ReceiveAsync()); + Thread.Sleep(300); + receive.IsCompleted.Should().BeFalse("the retained PDU is about to be put back"); + }; + channel.DiscardPendingPdus(arrivedBefore: 1).Should().Be(0); + + (await receive!.WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); + } + + // The streaming receive ends gracefully when the channel is disposed under it. + [Fact] + public async Task ReceiveAll_Ends_When_The_Channel_Is_Disposed() + { + var session = NewSession(); + using var bus = OpenClassic(session, 0); + using var channel = IsoTpFactory.Open(bus, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions()); + + using var cts = new CancellationTokenSource(ShortTimeout); + var seen = 0; + var reading = Task.Run(async () => + { + await foreach (var _ in channel.ReceiveAllAsync(cts.Token)) seen++; + }); + await Task.Delay(50); // the enumeration is waiting on the empty inbox + channel.Dispose(); + + await reading.WaitAsync(ShortTimeout); + seen.Should().Be(0); + } + // The reader itself failing is the same loss: the inbox ends with that failure, for every // receiver, and it is reported once. [Fact] From 8b35c9a8fd887aa29fb3be36185b05ad229f6ec0 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:01:14 +0200 Subject: [PATCH 09/11] fix(isotp): a canceled token does not consume a buffered PDU Review on #261: the receive loop read the inbox before it reached any wait on the token, so a token that was already canceled still returned a buffered PDU (and the streaming receive drained them), where the wait on the token threw before. The token is checked before each read. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 3 ++ .../IsoTp/IsoTpChannelIntegrationTests.cs | 28 +++++++++++++++++++ 2 files changed, 31 insertions(+) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index 9396056..3d3c984 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -515,6 +515,9 @@ private async IAsyncEnumerable ReadAllAsync( { while (true) { + // Before any read: a canceled token must not consume a buffered PDU, as the wait on the + // token did before this loop existed (and as TryReceiveWithArrival documents). + cancellationToken.ThrowIfCancellationRequested(); if (_pduInbox.Reader.TryRead(out var item)) return (true, item); if (_inboxLost is { } lost) { diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index dd567cf..813a890 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -1027,6 +1027,34 @@ public async Task ReceiveAll_Ends_When_The_Channel_Is_Disposed() seen.Should().Be(0); } + // Review on #261: a canceled token must not consume a PDU that is already buffered. + [Fact] + public async Task A_Canceled_Token_Does_Not_Consume_A_Buffered_Pdu() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + using var channel = IsoTpFactory.Open(busA, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions()); + + var arrived = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.DatagramReceived += (_, _) => arrived.TrySetResult(true); + busB.Transmit(CanFrame.Classic(0x321, new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC })); + await arrived.Task.WaitAsync(ShortTimeout); + + using var canceled = new CancellationTokenSource(); + canceled.Cancel(); + Func receive = () => channel.ReceiveAsync(canceled.Token); + await receive.Should().ThrowAsync(); + Func all = async () => + { + await foreach (var _ in channel.ReceiveAllAsync(canceled.Token)) { } + }; + await all.Should().ThrowAsync(); + + // Still there for a receive that is not canceled. + (await channel.ReceiveAsync().WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); + } + // The reader itself failing is the same loss: the inbox ends with that failure, for every // receiver, and it is reported once. [Fact] From adf5b1d932fd58c4d18076d62fe17f6a20eb4dfe Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:16:08 +0200 Subject: [PATCH 10/11] refactor(isotp): read the inbox once more after the discard flag as a second loop pass Codecov on #261: the second read was only reachable through a race. As a second pass of the loop it gives the same guarantee and is run by every receive that ends in the loss. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index 3d3c984..fcb0134 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -513,6 +513,7 @@ private async IAsyncEnumerable ReadAllAsync( /// private async Task<(bool Taken, RxInboxItem Item)> TakeNextAsync(CancellationToken cancellationToken) { + var lossConfirmed = false; while (true) { // Before any read: a canceled token must not consume a buffered PDU, as the wait on the @@ -524,12 +525,19 @@ private async IAsyncEnumerable ReadAllAsync( // A discard has the retained items out of the inbox for a moment: wait it out. if (Volatile.Read(ref _discarding) != 0) { + lossConfirmed = false; await Task.Yield(); continue; } - // Looked again after the flag: a discard that finished in between has put them back. - if (_pduInbox.Reader.TryRead(out item)) return (true, item); + // The empty inbox was read before the flag: a discard that finished in between has + // put its items back. Read once more, after the flag, before it counts. + if (!lossConfirmed) + { + lossConfirmed = true; + continue; + } + throw lost; } From 4c3b56953a4aad3ef5d879ca2c3c293bb4516d2e Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:30:46 +0200 Subject: [PATCH 11/11] fix(isotp): post the service loss under the pump lock Review on #261: a caller pumping the subscription (GetReceptionsInProgress, a settle, a discard) holds the pump lock from taking a frame to posting it; the reader's post of the loss ran without it and could overtake that frame, so a receiver threw the loss before a PDU that had been received. The loss is posted under the same lock. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.IsoTp/IsoTpChannel.cs | 8 ++++- .../Infrastructure/StarvedReaderBusService.cs | 8 +++++ .../IsoTp/IsoTpChannelIntegrationTests.cs | 30 +++++++++++++++++++ 3 files changed, 45 insertions(+), 1 deletion(-) diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index fcb0134..dc8cf7b 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -688,7 +688,13 @@ private async Task RunReaderAsync() // own, the bus refuses them. On the actor, like everything else that touches the inbox. try { - _actor.Post(() => EndInboxAfterSubscriptionLoss(lost)); + // Under the pump lock: a caller that is pumping the subscription right now -- it took + // the last frame and has not posted it yet -- finishes first, so the frame is on the + // actor ahead of the loss and a receiver cannot meet the loss before it. + lock (_pumpGate) + { + _actor.Post(() => EndInboxAfterSubscriptionLoss(lost)); + } } catch (ObjectDisposedException) { diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index 82986bf..e9e8e59 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -40,6 +40,12 @@ public void Deliver(CanFrameView frame, long hostArrivalTimestamp = 0) => _frame /// When set, the reader task's wait throws it once woken, as a subscription that failed would. public Exception? ReaderFault { get; set; } + /// When set, the subscription ends on the wake: the reader's wait returns false, as when its service is disposed. + public bool EndSubscriptionOnWake { get; set; } + + /// When set, a TryRead that took a frame holds it until the gate opens: a caller-side pump caught between taking a frame and posting it. + public ManualResetEventSlim? TryReadGate { get; set; } + /// Lets the reader task's wait complete once; every later wait stays pending. public void WakeReader() => _wake.TrySetResult(true); @@ -107,6 +113,7 @@ public bool TryRead(out CanFrameEvent frameEvent) { bool ok = _owner._frames.Reader.TryRead(out frameEvent); if (!ok) _owner._drained.TrySetResult(true); + else _owner.TryReadGate?.Wait(); return ok; } @@ -116,6 +123,7 @@ public async ValueTask WaitToReadAsync(CancellationToken cancellationToken { var woken = await _owner._wake.Task.WaitAsync(cancellationToken); if (_owner.ReaderFault is { } fault) throw fault; + if (_owner.EndSubscriptionOnWake) return false; return woken; } await Task.Delay(Timeout.Infinite, cancellationToken); diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 813a890..c29ab05 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -1055,6 +1055,36 @@ public async Task A_Canceled_Token_Does_Not_Consume_A_Buffered_Pdu() (await channel.ReceiveAsync().WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); } + // Review on #261: a caller that pumps the subscription (GetReceptionsInProgress, a settle, a + // discard) can have taken the last frame and not yet posted it when the subscription ends. The + // loss must not overtake that frame, or a receiver throws it before a PDU that was received. + [Fact] + public async Task The_Loss_Does_Not_Overtake_A_Frame_A_Pumping_Caller_Is_About_To_Post() + { + var service = new StarvedReaderBusService { EndSubscriptionOnWake = true }; + using var actor = new ProtocolActor(); + using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), + ownsService: false, actor); + using var gate = new ManualResetEventSlim(false); + service.TryReadGate = gate; + + var receive = channel.ReceiveAsync(); + service.Deliver(new CanFrameView(CanFrameType.Can20, 0x321, + new byte[] { 0x02, 0x0A, 0x0B, 0xCC, 0xCC, 0xCC, 0xCC, 0xCC }, FrameFlags.None)); + + // A caller pumps: it takes the frame and is held before it posts it. + var pump = Task.Run(() => channel.GetReceptionsInProgress()); + await Task.Delay(100); + service.WakeReader(); // the subscription ends + // No signal says "the reader reached the post and waited for the lock", so this is a + // negative window: it can only pass falsely on a slow host. + await Task.Delay(200); + + gate.Set(); // the pump posts its frame + await pump.WaitAsync(ShortTimeout); + (await receive.WaitAsync(ShortTimeout)).Should().Equal(0x0A, 0x0B); + } + // The reader itself failing is the same loss: the inbox ends with that failure, for every // receiver, and it is reported once. [Fact]