diff --git a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs index d9be5b6..dc8cf7b 100644 --- a/src/CanKit.Pro.IsoTp/IsoTpChannel.cs +++ b/src/CanKit.Pro.IsoTp/IsoTpChannel.cs @@ -67,7 +67,22 @@ 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. + // + // 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; + // 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); // 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. @@ -162,6 +177,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)) { @@ -296,12 +317,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); } /// @@ -447,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; } @@ -469,11 +498,62 @@ public IAsyncEnumerable ReceiveAllAsync(CancellationToken cancellationTo private async IAsyncEnumerable ReadAllAsync( [EnumeratorCancellation] CancellationToken cancellationToken) { - var reader = _pduInbox.Reader; - while (await reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) + while (true) { - while (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) + { + var lossConfirmed = false; + 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) + { + // 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; + } + + // 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; + } + + 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); } } @@ -582,21 +662,71 @@ 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 + // 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 + { + // 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) + { + // 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) + { + // 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); + } + + _inboxLost = lost; + _inboxLostSignal.TrySetResult(true); + RaiseBackgroundException(lost); + } + + private static void RequirePositive(string name, TimeSpan value) + { + if (value <= TimeSpan.Zero) + throw new ArgumentOutOfRangeException(name, value, name + " must be greater than zero."); + } + /// /// Takes every frame the subscription has buffered and posts it to the actor, in order. /// Run by the reader task when frames arrive, and by a caller that must see the buffered @@ -718,6 +848,15 @@ private void BeginSendOnLoop(byte[] pdu, TaskCompletionSource _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; } + + /// 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); @@ -104,13 +113,19 @@ 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; } 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; + if (_owner.EndSubscriptionOnWake) return false; + 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 6afeb76..c29ab05 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -850,6 +850,396 @@ public async Task Dispose_Unblocks_Pending_ReceiveAsync() thrown.Message.Should().Contain("disposed"); } + // 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. 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_Every_Waiting_Receive() + { + var session = NewSession(); + using var bus = OpenClassic(session, 0); + 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 first = channel.ReceiveAsync(); + var second = channel.ReceiveAsync(); + service.Dispose(); + + 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(); + 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(); + } + + // 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); + + // 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); + + (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 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); + } + + // 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); + } + + // 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); + } + + // 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] + 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); + 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); + 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 + // 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. // --------------------------------------------------------------------------------