diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 3cf9186..98af19b 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -94,6 +94,11 @@ Every behavioural change needs a test, and tests here are expected to be determi - **No hardware, no timing luck.** Use the `virtual://` loopback adapter for anything about real bus behaviour, and `ControllableBus` (`tests/CanKit.Pro.Tests/Infrastructure/`) when the test needs to control what the bus does — echo frames, bus state, whether a transmit is accepted. + `ControllableBus.DeferredEchoCapable(...)` parks each TX echo in a `DeferredEchoQueue` instead + of raising it inside `Transmit`, which is the only way to have two sends pending at once: a + synchronous echo re-enters `CanBusService`'s pending-send lock on the transmitting thread, so + the pending list never holds more than that thread's own entry. Reach for it whenever the + behaviour under test is about how several in-flight sends relate to each other. - **Do not test through an adapter's internals.** If a test needs reflection into another package's private state, it is testing that package, not ours; drive the scenario through the double instead. diff --git a/tests/CanKit.Pro.Tests/Infrastructure/ControllableBus.cs b/tests/CanKit.Pro.Tests/Infrastructure/ControllableBus.cs index 0d60f8a..28fd20c 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/ControllableBus.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/ControllableBus.cs @@ -9,6 +9,27 @@ namespace CanKit.Pro.Tests.Infrastructure; +/// +/// When the TX echo of an accepted transmit reaches . +/// +public enum EchoDelivery +{ + /// + /// Raised from inside — on the + /// transmitting thread, inside whatever lock the caller holds while transmitting. What a real + /// echo-mode adapter (CanKit.Adapter.Virtual in ChannelWorkMode.Echo) does, and what + /// makes the reentrancy in CanBusService.SendWithEchoConfirmAsync observable. + /// + Synchronous, + + /// + /// Parked in instead of being raised. The test + /// chooses when each echo is delivered, which lets more than one pending send exist at once — + /// see for why that is not achievable synchronously. + /// + Deferred, +} + /// /// An the test drives directly: it decides whether a transmit is accepted, /// whether (and when) a TX echo comes back, and what the controller reports. @@ -40,10 +61,12 @@ public sealed class ControllableBus : ICanBus private readonly IBusRTOptionsConfigurator _options; private int _disposed; - private ControllableBus(ICanBus configurationSource) + private ControllableBus(ICanBus configurationSource, EchoDelivery echoDelivery) { _configurationSource = configurationSource; _options = new EchoCapableOptions(configurationSource.Options); + EchoMode = echoDelivery; + DeferredEchoes = new DeferredEchoQueue(frame => RaiseObserved(frame, isEcho: true)); // What a healthy CAN controller reports; tests move it from here. BusState = BusState.ErrActive; } @@ -51,12 +74,26 @@ private ControllableBus(ICanBus configurationSource) /// /// Creates a double whose report ChannelWorkMode.Echo and the /// CanFeature.Echo capability — the combination that makes SendConfirmed take - /// the real-echo-matching path (FR-RAW-031). + /// the real-echo-matching path (FR-RAW-031) — and that echoes synchronously from inside + /// , exactly as a real echo-mode adapter does. /// public static ControllableBus EchoCapable(string session) - => new(VirtualAdapterFixture.Open(session, 0, ChannelWorkMode.Echo)); - + => new(VirtualAdapterFixture.Open(session, 0, ChannelWorkMode.Echo), EchoDelivery.Synchronous); + /// + /// Same echo-capable configuration as , but every accepted transmit's + /// echo is parked in until the test releases it. + /// + /// + /// Use this whenever the behaviour under test needs two or more sends to be pending at the + /// same time. A synchronous echo makes that impossible — it re-enters + /// CanBusService's pending-send lock on the transmitting thread before that thread ever + /// leaves Transmit, so the pending list only ever holds the entry belonging to the + /// thread currently inside it, no matter how many callers race. See + /// . + /// + public static ControllableBus DeferredEchoCapable(string session) + => new(VirtualAdapterFixture.Open(session, 0, ChannelWorkMode.Echo), EchoDelivery.Deferred); /// Whether reports the frame as accepted. public bool AcceptTransmit { get; set; } = true; @@ -64,6 +101,20 @@ public static ControllableBus EchoCapable(string session) /// Whether an accepted frame is echoed back through . public bool EchoAcceptedFrames { get; set; } = true; + /// + /// Whether an accepted frame's echo is raised inside or + /// parked in . Settable mid-test so a scenario can, for example, + /// let the first send confirm normally and only defer the ones it needs to overlap. + /// + public EchoDelivery EchoMode { get; set; } + + /// + /// Echoes parked by mode, and the handle that releases + /// them. Always present; stays empty while is + /// . + /// + public DeferredEchoQueue DeferredEchoes { get; } + /// Number of frames handed to . public int TransmitCount => Volatile.Read(ref _transmitCount); @@ -92,9 +143,17 @@ public int Transmit(in CanFrame frame) Interlocked.Increment(ref _transmitCount); if (!AcceptTransmit) return 0; - // A real echo-mode adapter delivers the echo synchronously from inside Transmit; matching - // that is what makes the reentrancy in CanBusService.SendWithEchoConfirmAsync observable. - if (EchoAcceptedFrames) RaiseObserved(frame, isEcho: true); + if (EchoAcceptedFrames) + { + // Synchronous is the default because that is what a real echo-mode adapter does, and + // matching it is what makes the reentrancy in CanBusService.SendWithEchoConfirmAsync + // observable. Deferred parks the echo instead, so the caller leaves Transmit — and + // releases the service's pending-send lock — with its entry still pending; see + // DeferredEchoQueue for why some FR-RAW-031 behaviour is only reachable that way. + if (EchoMode == EchoDelivery.Deferred) DeferredEchoes.Park(frame); + else RaiseObserved(frame, isEcho: true); + } + return 1; } diff --git a/tests/CanKit.Pro.Tests/Infrastructure/DeferredEchoQueue.cs b/tests/CanKit.Pro.Tests/Infrastructure/DeferredEchoQueue.cs new file mode 100644 index 0000000..2f334f6 --- /dev/null +++ b/tests/CanKit.Pro.Tests/Infrastructure/DeferredEchoQueue.cs @@ -0,0 +1,180 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using CanKit.Abstractions.API.Can.Definitions; + +namespace CanKit.Pro.Tests.Infrastructure; + +/// +/// The parking lot for TX echoes of a running in +/// mode: Transmit hands the frame here instead of +/// echoing it, and the test decides when — and in which order — each echo reaches +/// FrameObserved. +/// +/// +/// Why this exists: a real echo-mode adapter delivers the echo synchronously from inside +/// Transmit, and CanBusService.SendWithEchoConfirmAsync transmits while holding its +/// pending-send lock. A synchronous echo therefore re-enters that lock on the transmitting +/// thread, so the pending list can never hold more than the one entry that thread just +/// registered. Every property that is only observable with two or more entries queued for the same +/// key — FIFO matching of byte-identical concurrent sends (SRS FR-RAW-031), an expired pending +/// send poisoning the FIFO for a later one — is therefore untestable against a synchronous echo: +/// the assertion holds no matter what the matching code does. Deferring the echo is what puts the +/// second entry in the list. +/// +/// +/// +/// Everything here is deliberately explicit rather than time-based: +/// is the only wait, and it waits on a transmit actually having happened rather than on a delay +/// that "should be long enough". Nothing in this class starts a timer, a thread, or a task — an +/// echo moves only when the test says so, on the test's own thread. +/// +/// +/// +/// Reusable beyond FR-RAW-031: any scenario that needs a pending send to still be pending while a +/// second one is registered (a late echo arriving after its own send already timed out, an echo +/// that never arrives at all while later ones do — see ) is expressed by +/// parking, then releasing or discarding, in whatever order the scenario calls for. +/// +public sealed class DeferredEchoQueue +{ + private readonly Action _deliver; + + private readonly object _gate = new(); + private readonly List _parked = new(); + private readonly List<(int Threshold, TaskCompletionSource Tcs)> _waiters = new(); + + // Monotonic: counts every frame ever parked, so a waiter's threshold cannot be un-met by a + // subsequent release. "Two sends have been transmitted" must stay true once it is true. + private int _enqueued; + + internal DeferredEchoQueue(Action deliver) => _deliver = deliver; + + /// Echoes parked and not yet released or discarded, oldest first. + public int Count + { + get { lock (_gate) return _parked.Count; } + } + + /// + /// Total number of echoes ever parked — i.e. accepted transmits observed while in + /// mode. Never decreases. + /// + public int Enqueued + { + get { lock (_gate) return _enqueued; } + } + + /// + /// Completes once has reached — the + /// deterministic replacement for "sleep a bit and hope the sends got that far". + /// + /// + /// A transmit is parked from inside ICanBus.Transmit, which + /// CanBusService.SendWithEchoConfirmAsync calls after it has registered its pending + /// entry and while still holding the pending-send lock. So "n echoes parked" is a hard + /// guarantee that n pending sends are registered, not an approximation of it. + /// + /// Number of parked echoes to wait for. + /// How long to wait before failing; a bound against a hang, not a + /// scheduling assumption. + /// Fewer than echoes were parked + /// within . + public async Task WaitForEnqueuedAsync(int count, TimeSpan timeout) + { + if (count <= 0) throw new ArgumentOutOfRangeException(nameof(count), count, "Count must be positive."); + + Task wait; + lock (_gate) + { + if (_enqueued >= count) return; + var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _waiters.Add((count, tcs)); + wait = tcs.Task; + } + + try + { + await wait.WaitAsync(timeout).ConfigureAwait(false); + } + catch (TimeoutException) + { + throw new TimeoutException( + $"Only {Enqueued} of {count} expected echoes were parked within {timeout}."); + } + } + + /// + /// Delivers the oldest parked echo through the bus's FrameObserved event, on the + /// calling thread. Returns false when nothing is parked. + /// + public bool ReleaseNext() + { + CanFrame frame; + lock (_gate) + { + if (_parked.Count == 0) return false; + frame = _parked[0]; + _parked.RemoveAt(0); + } + + // Delivered outside the lock: the echo runs the service's whole match-and-complete path + // (and any continuation it resumes) on this thread, and none of that may be serialized + // against a concurrent Transmit parking the next echo. + _deliver(frame); + return true; + } + + /// + /// Delivers every parked echo, oldest first. Returns how many were delivered. + /// + public int ReleaseAll() + { + var released = 0; + while (ReleaseNext()) released++; + return released; + } + + /// + /// Drops the oldest parked echo without ever delivering it — the frame reached the wire but + /// its echo is lost, while later echoes still arrive normally. Returns false when + /// nothing is parked. + /// + public bool DiscardNext() + { + lock (_gate) + { + if (_parked.Count == 0) return false; + _parked.RemoveAt(0); + return true; + } + } + + /// Parks ; called by . + internal void Park(in CanFrame frame) + { + (int Threshold, TaskCompletionSource Tcs)[]? satisfied = null; + lock (_gate) + { + _parked.Add(frame); + _enqueued++; + + if (_waiters.Count > 0) + { + var met = _waiters.FindAll(w => w.Threshold <= _enqueued); + if (met.Count > 0) + { + satisfied = met.ToArray(); + _waiters.RemoveAll(w => w.Threshold <= _enqueued); + } + } + } + + // Never complete a TCS under the lock: the waiter's continuation may call straight back + // into ReleaseNext/Count, and Park runs inside the service's pending-send lock. + if (satisfied is null) return; + foreach (var (_, tcs) in satisfied) + tcs.TrySetResult(); + } +} diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs index 9dede2e..9e5051f 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs @@ -166,12 +166,20 @@ await master.SdoDownloadAsync(serverNodeId: 0x11, index: 0x2100, subindex: 0x00, // applied to the abandoned buffer or committed to the OD. Prior to the fix the // expedited path did not clear _sdoServer, so the stale download session would still // accept and commit those segment frames. + // + // Part (2) is observed on the wire rather than by waiting: "0x3000 is still all zeros" is + // equally true of a server that rejected the stray segments and of one that never received + // them, so the fixed Task.Delay this used to sit on decided which behaviour was being + // asserted. A third bus watches the slave's SDO TX, and the test waits for the server's own + // responses — the session ack, then one CommandSpecifierInvalid abort per stray segment. + // That the aborts arrive is the proof the segments were delivered and refused. [Fact] public async Task Sdo_ExpeditedInitiate_ClearsStaleSegmentedServerSession() { var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); + using var busObserver = Open(session, 2); using var master = CanOpen.OpenNode(busA, nodeId: 0x01); using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); @@ -204,6 +212,36 @@ public async Task Sdo_ExpeditedInitiate_ClearsStaleSegmentedServerSession() // are rejected with SdoAbortCode.CommandSpecifierInvalid instead. slave.ObjectDictionary.AddDomain(0x3000, 0x00, new byte[8]); + // Watch what the slave itself puts on the wire (COB-ID 0x580 + 0x11 = 0x591) from a + // third bus, so the slave's responses are distinguishable from the master's requests. + var sessionInstalled = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var straysRejected = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var rejections = new List(); + busObserver.FrameObserved += (_, e) => + { + var frame = e.CanFrame; + if (frame.IsExtendedFrame || (uint)frame.ID != 0x580u + 0x11u) return; + var data = frame.Data.ToArray(); + if (data.Length < 8) return; + + // scs=0x60 for (0x3000, 0x00): the server acknowledged the segmented initiate, so + // the session this test needs to be superseded is now genuinely installed. + if (data[0] == 0x60 && data[1] == 0x00 && data[2] == 0x30 && data[3] == 0x00) + sessionInstalled.TrySetResult(data); + + if (data[0] != 0x80) return; // not an abort + uint code = (uint)(data[4] | (data[5] << 8) | (data[6] << 16) | (data[7] << 24)); + // Only CommandSpecifierInvalid: the expedited initiate below also emits a supersede + // abort for 0x3000 (SdoAbortCode.General, see Sdo_ServerSupersede_EmitsWireAbort_ + // ForPriorTransfer), and that one says nothing about the stray segments. + if (code != (uint)SdoAbortCode.CommandSpecifierInvalid) return; + lock (rejections) + { + rejections.Add(data); + if (rejections.Count == 2) straysRejected.TrySetResult(null); + } + }; + // Segmented download initiate (cs=0x21) for 0x3000:00, declared length 8. var initFrame = new byte[8] { @@ -213,8 +251,9 @@ public async Task Sdo_ExpeditedInitiate_ClearsStaleSegmentedServerSession() }; busA.Transmit(CanFrame.Classic(0x600 + 0x11, initFrame, isExtendedFrame: false)); - // Give the actor loop a moment to install the segmented session for 0x3000. - await Task.Delay(50); + // Wait for the server's own ack rather than a delay: the segmented session for 0x3000 + // is installed exactly when that ack goes out. + await sessionInstalled.Task.WithTimeoutAsync(ShortTimeout); // Expedited SDO download to 0x2001 (an unrelated U16 slot). With the fix, this // supersedes the still-open 0x3000 segmented session and clears _sdoServer. @@ -232,8 +271,15 @@ public async Task Sdo_ExpeditedInitiate_ClearsStaleSegmentedServerSession() busA.Transmit(CanFrame.Classic(0x600 + 0x11, seg1Frame, isExtendedFrame: false)); busA.Transmit(CanFrame.Classic(0x600 + 0x11, seg2Frame, isExtendedFrame: false)); - // Wait long enough for both segment frames to be processed on the actor loop. - await Task.Delay(100); + // Both segment frames reached the server and were refused as "no such transfer is + // open". This is the load-bearing wait: without it, every assertion below would also + // hold for a run in which the strays never arrived at all. + await straysRejected.Task.WithTimeoutAsync(ShortTimeout); + lock (rejections) + { + rejections.Should().HaveCount(2, + "each stray segment must be individually rejected, not silently absorbed"); + } // Verification: // * With the fix: 0x3000:00 stays untouched (all zeros) because the expedited @@ -649,6 +695,12 @@ public async Task Nmt_Broadcast_TransitionsAllNodes() } // FR-CO-007: Stop then EnterPre-Op state-machine coverage. + // + // Each transition waits for the heartbeat the slave emits to announce its new state, not for + // a fixed 30 ms. The slave applies the command and sends that heartbeat in the same actor-loop + // work item, so the heartbeat's arrival is proof the transition happened — whereas 30 ms was + // only ever a guess that happened to hold on an unloaded machine, and three of them in a row + // gave the test three chances to fail under CI load for reasons unrelated to NMT. [Fact] public async Task Nmt_StopAndPreOp_TransitionsWork() { @@ -659,17 +711,28 @@ public async Task Nmt_StopAndPreOp_TransitionsWork() using var master = CanOpen.OpenNode(busA, nodeId: 0x01); using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); - await master.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); - await Task.Delay(30); - slave.State.Should().Be(NmtState.Operational); + // One handler for the whole test: each step arms the state it is waiting for. The slave + // also emits an unsolicited boot-up heartbeat, which no step waits for and which the + // state match therefore ignores. + TaskCompletionSource? pending = null; + NmtState expected = default; + master.HeartbeatReceived += (_, e) => + { + if (e.ProducerNodeId == 0x11 && e.State == expected) pending?.TrySetResult(e.State); + }; - await master.SendNmtCommandAsync(NmtCommand.Stop, targetNodeId: 0x11); - await Task.Delay(30); - slave.State.Should().Be(NmtState.Stopped); + async Task CommandAndAwaitState(NmtCommand command, NmtState state) + { + expected = state; + pending = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + await master.SendNmtCommandAsync(command, targetNodeId: 0x11); + await pending.Task.WithTimeoutAsync(ShortTimeout); + slave.State.Should().Be(state); + } - await master.SendNmtCommandAsync(NmtCommand.EnterPreOperational, targetNodeId: 0x11); - await Task.Delay(30); - slave.State.Should().Be(NmtState.PreOperational); + await CommandAndAwaitState(NmtCommand.Start, NmtState.Operational); + await CommandAndAwaitState(NmtCommand.Stop, NmtState.Stopped); + await CommandAndAwaitState(NmtCommand.EnterPreOperational, NmtState.PreOperational); } // FR-CO-007: reset node causes a bootup frame (0x00 on 0x700+id) to be re-emitted. diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 4f67e8b..895960a 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -417,9 +417,19 @@ public async Task Dispose_Unblocks_Pending_ReceiveAsync() channel.Dispose(); channel.Dispose(); - // ReceiveAsync should now throw (channel disposed) rather than hang. - Func act = () => recvTask; - await act.Should().ThrowAsync(); + // Pinned to the exact failure ReceiveAsync documents for a disposed channel. "Any + // exception" would also have accepted the two outcomes this test exists to rule out: a + // ChannelClosedException leaking the inbox implementation to the caller, and an + // OperationCanceledException from the reader's own CTS, which callers would reasonably + // treat as "my token was cancelled, retry" rather than "this channel is finished". + // WaitAsync bounds the wait so a Dispose that fails to unblock the reader fails this + // test instead of hanging the run. + Func act = () => recvTask.WaitAsync(ShortTimeout); + var thrown = (await act.Should().ThrowAsync()).Which; + // BeOfType, not the ThrowAsync above: ObjectDisposedException derives from + // InvalidOperationException and would satisfy it. + thrown.Should().BeOfType(); + thrown.Message.Should().Contain("disposed"); } // -------------------------------------------------------------------------------- @@ -705,14 +715,20 @@ public async Task MultiFrame_Receive_Ignores_Empty_ConsecutiveFrame_Without_Adva } // -------------------------------------------------------------------------------- - // Bugbot 3594960783 (HIGH) — a codec throw inside BeginSendOnLoop (e.g. > 4095 bytes - // on classic CAN triggers ArgumentOutOfRangeException from BuildFirstFrame) must - // (1) fault the awaiting SendAsync with the codec exception, (2) release the send-gate, - // and (3) leave the channel usable for subsequent sends -- rather than leaking _tx and - // hanging every future SendAsync forever behind the gate. + // A PDU longer than this channel's frame kind can address is rejected by SendAsync's own + // pre-check, before the send-gate is taken and before anything reaches the actor. That is + // the entire behaviour here, so it is asserted exactly: the specific exception, the + // parameter it names, and that nothing was put on the wire. + // + // This test used to be named for the codec throw inside BeginSendOnLoop (Bugbot 3594960783) + // and accepted any of three exception types. It never reached that code: SendAsync's + // pre-check enforces the same limit the codec does (MaxClassicFirstFrameLength = 4095), so + // an oversized PDU is refused two layers above the actor, and "any of three exception types" + // could not distinguish the layer it came from. The actor-side failure contract that Bugbot + // finding is about is asserted by the next test, over a failure that is actually reachable. // -------------------------------------------------------------------------------- [Fact] - public async Task Send_Faults_On_Codec_Throw_And_Channel_Remains_Usable() + public async Task SendAsync_Rejects_Oversized_Pdu_Before_Anything_Reaches_The_Bus() { var session = NewSession(); using var busA = OpenClassic(session, 0); @@ -724,26 +740,61 @@ public async Task Send_Faults_On_Codec_Throw_And_Channel_Remains_Usable() using var sender = IsoTpFactory.Open(busA, epAB, FastOptions()); using var receiver = IsoTpFactory.Open(busB, epBA, FastOptions()); - // >4095 bytes on classic-CAN forces BuildFirstFrame to throw ArgumentOutOfRangeException - // synchronously on the actor loop -- the exact "codec throws inside BeginSendOnLoop" path - // Bugbot flagged. - byte[] oversized = new byte[4096]; + var framesOnWire = 0; + busB.FrameObserved += (_, _) => Interlocked.Increment(ref framesOnWire); + + // One byte more than the 12-bit classic First-Frame length field can announce. + byte[] oversized = new byte[IsoTpFrameCodec.MaxClassicFirstFrameLength + 1]; - // WaitAsync bounds the wait: under the bug this SendAsync would hang forever because - // the actor's synchronous BuildFirstFrame throw is swallowed by - // BackgroundExceptionOccurred without ever completing the TCS. Func act = () => sender.SendAsync(oversized).WaitAsync(ShortTimeout); - var caught = (await act.Should().ThrowAsync()).Which; - // Codec-thrown ArgumentOutOfRangeException surfaces directly (wrapped only in the actor's - // synchronous invocation path; unwrapped as-is by FailTx -> TCS -> await). - caught.Should().Match(e => - e is ArgumentOutOfRangeException - || e is IsoTpException - || e is InvalidOperationException, - "codec throw must fault the awaiting SendAsync, not hang it"); - - // The gate MUST be released and _tx cleared. A normal send after the failure must - // succeed within the same short timeout. + var thrown = (await act.Should().ThrowAsync()).Which; + thrown.ParamName.Should().Be("pdu"); + + Volatile.Read(ref framesOnWire).Should().Be(0, + "the length check runs before the send-gate, so no frame is ever built or transmitted"); + + // The rejected call must not have consumed the send-gate: a normal send still works. + var recvTask = receiver.ReceiveAsync(new CancellationTokenSource(ShortTimeout).Token); + byte[] normal = { 0x11, 0x22, 0x33 }; + await sender.SendAsync(normal).WaitAsync(ShortTimeout); + (await recvTask).Should().Equal(normal); + } + + // -------------------------------------------------------------------------------- + // Bugbot 3594960783 (HIGH), the reachable half — when a send fails *after* the actor has + // taken ownership of the PDU (here: the bus layer throws out of SendConfirmed), the channel + // must (1) fault the awaiting SendAsync with that exact exception, (2) release the send-gate + // and clear _tx, and (3) stay usable — rather than reporting the failure only through + // BackgroundExceptionOccurred and hanging every future SendAsync behind the gate. + // + // The failure is injected at the bus layer because that is the only way in: SendAsync's + // pre-check duplicates the codec's own length limit, so no PDU that reaches BeginSendOnLoop + // can make the codec throw. Asserting the exact exception is what gives this test teeth -- + // a channel that swallowed it, rewrote it as a generic IsoTpException, or completed the send + // anyway all fail here. + // -------------------------------------------------------------------------------- + [Fact] + public async Task Send_Faults_With_The_Bus_Layer_Exception_And_Channel_Remains_Usable() + { + var session = NewSession(); + using var busA = OpenClassic(session, 0); + using var busB = OpenClassic(session, 1); + + var epAB = IsoTpEndpoint.Normal(0x210, 0x211); + var epBA = IsoTpEndpoint.Normal(0x211, 0x210); + + using var svcA = new CanBusService(busA); + var failing = new ThrowOnFirstConfirmService(svcA, new InvalidOperationException("bus-layer boom")); + + using var sender = IsoTpFactory.Open(failing, epAB, FastOptions(), leaveOpen: true); + using var receiver = IsoTpFactory.Open(busB, epBA, FastOptions()); + + Func act = () => sender.SendAsync(new byte[] { 1, 2, 3 }).WaitAsync(ShortTimeout); + var thrown = (await act.Should().ThrowAsync()).Which; + thrown.Message.Should().Be("bus-layer boom", + "the failure the bus layer reported must reach the caller unrewritten"); + + // Gate free and _tx cleared: the very next send goes through end to end. var recvTask = receiver.ReceiveAsync(new CancellationTokenSource(ShortTimeout).Token); byte[] normal = { 0x11, 0x22, 0x33 }; await sender.SendAsync(normal).WaitAsync(ShortTimeout); @@ -1269,6 +1320,53 @@ public async Task Receiver_With_LocalBlockSize_Emits_FlowControl_After_Each_Full Volatile.Read(ref fcCount).Should().Be(3); } + /// + /// Test double: the first call throws the supplied + /// exception instead of transmitting; every later call is forwarded to the inner service + /// untouched. Models the L2/driver layer failing outright — as opposed to reporting a + /// that says the send failed — which is the one failure a + /// channel's send can hit *after* the actor already owns the PDU. + /// + private sealed class ThrowOnFirstConfirmService : ICanBusService + { + private readonly ICanBusService _inner; + private readonly Exception _failure; + private int _calls; + + public ThrowOnFirstConfirmService(ICanBusService inner, Exception failure) + { + _inner = inner; + _failure = failure; + } + + 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) + => _inner.Subscribe(predicate, bufferCapacity); + + public ISubscription Subscribe(CanIdFilter filter, int? bufferCapacity = null) + => _inner.Subscribe(filter, bufferCapacity); + + public IReadOnlyList<(ISubscription First, ISubscription Second)> FindOverlappingFilterSubscriptions() + => _inner.FindOverlappingFilterSubscriptions(); + + public Task SendConfirmed(CanFrame frame, TimeSpan? timeout = null, + CancellationToken cancellationToken = default) + { + if (Interlocked.Increment(ref _calls) == 1) throw _failure; + return _inner.SendConfirmed(frame, timeout, cancellationToken); + } + + public void Dispose() { /* wrapper: the test owns and disposes the inner service */ } + } + /// /// Test double: puts frames on the wire for real but never delivers an echo, letting the /// requested confirm timeout (the channel's N_As) expire as a Timeout failure — the exact diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs index 7700702..f6618f8 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs @@ -791,16 +791,51 @@ public async Task ReClaim_RejectsSendUntilNewClaimSucceeds() // // The pre-fix window between Open and Task.Run(Dispose) is sub-millisecond on the // virtual bus, so this test cannot deterministically reproduce the race on every run; - // it pins the correctness invariant (received <= sent) across many claim/re-claim - // cycles under continuous broadcast BAM traffic. The synchronous-dispose fix makes the - // invariant hold by construction. + // it pins the correctness invariant across many claim/re-claim cycles under continuous + // broadcast BAM traffic. The synchronous-dispose fix makes the invariant hold by + // construction. + // + // Two separate invariants, each carried by its own traffic, because they need opposite + // things from the test: + // + // * "never twice" needs *concurrent* traffic — a BAM must be in flight while a rebind + // happens for the two-live-channels window to be hit at all. That is the background + // stream below. It is free-running, and nothing is asserted about how much of it + // arrives; only that no sequence number arrives more than once. + // + // * "still delivering" needs *quiet* traffic — one BAM sent after a rebind has finished, + // with no further rebind due until the test issues the next claim, so it cannot be + // straddled and its delivery is guaranteed rather than probable. That is the probe. + // + // Trying to get both from one free-running stream is what made this test fail on Windows CI + // and pass everywhere else. It asserted a delivery floor of sent/2 over the background + // stream, and how much of that stream survives is a ratio of two unrelated clocks: the + // rebind cadence (ClaimAnnounceTimeout, 40 ms) against how long one BAM occupies the wire + // (Th between DT frames). A BAM that straddles a rebind is dropped by the disposed channel + // — that is by design, reassembly state is not carried across a rebind — so as the BAM + // period approaches the rebind spacing, *every* BAM straddles one and delivery collapses. + // Measured on this suite by stretching Th, which is what a runner with coarse timer + // granularity does to the peer's requested 2 ms: + // + // Th = 2 ms -> 118 sent, 101 delivered, losses in 15 runs of at most 2 + // Th = 16 ms -> 17 sent, 2 delivered, losses in 2 runs, longest 12 + // Th = 32 ms -> 8 sent, 0 delivered, one run of 8 + // + // No duplicate was ever observed, in any of those. The floor was measuring the runner, not + // the node — so it is gone, replaced by the probe, which is exact (all 8 arrive) and holds + // however slow the machine is. // --------------------------------------------------------------------------------------- [Fact] public async Task RebindTransport_DoesNotDeliverBamMoreThanOncePerRebind() { + const uint backgroundPgn = 0xFED1u; + const uint probePgn = 0xFED2u; + const int claims = 8; + var session = NewSession(); using var busPeer = Open(session, 0); using var busNode = Open(session, 1); + using var busProbe = Open(session, 2); // Short arbitration window so many rebinds happen while peer traffic is in flight; // small Th so a single BAM takes a couple of ms end-to-end. @@ -811,16 +846,35 @@ public async Task RebindTransport_DoesNotDeliverBamMoreThanOncePerRebind() }; using var node = J1939Node.Open(busNode, opts); - int received = 0; + // Every delivered BAM is recorded as (PGN, sequence number), so a duplicate is + // identifiable as such instead of only showing up as "one more than expected". + var delivered = new List<(uint Pgn, int Seq)>(); + var deliveredGate = new object(); + TaskCompletionSource? probeArrived = null; node.MessageReceived += (_, m) => { - if (m.Pgn == 0xFED1u) Interlocked.Increment(ref received); + if (m.Pgn != backgroundPgn && m.Pgn != probePgn) return; + var seq = BitConverter.ToInt32(m.Payload.Span.Slice(0, 4)); + lock (deliveredGate) delivered.Add((m.Pgn, seq)); + if (m.Pgn == probePgn) probeArrived?.TrySetResult(seq); }; + static byte[] Datagram(int seq) + { + // 12 bytes, so still a genuine multi-frame BAM; the first four carry the sequence + // number and the rest is the same filler as before. + var payload = new byte[12]; + BitConverter.TryWriteBytes(payload.AsSpan(0, 4), seq); + for (int b = 4; b < payload.Length; b++) payload[b] = (byte)(0xE0 + b); + return payload; + } + using var peerTp = CanKit.Pro.J1939Tp.J1939Tp.Open(busPeer, sourceAddress: 0x77, new J1939TpOptions().With(th: TimeSpan.FromMilliseconds(2))); - var payload = new byte[12]; - for (int b = 0; b < payload.Length; b++) payload[b] = (byte)(0xE0 + b); + // A second source address for the probes: one SA may only run one BAM session at a + // time, and the probe must not have to queue behind the background stream. + using var probeTp = CanKit.Pro.J1939Tp.J1939Tp.Open(busProbe, sourceAddress: 0x78, + new J1939TpOptions().With(th: TimeSpan.FromMilliseconds(2))); int sent = 0; using var peerCts = new CancellationTokenSource(); @@ -830,7 +884,7 @@ public async Task RebindTransport_DoesNotDeliverBamMoreThanOncePerRebind() { while (!peerCts.IsCancellationRequested) { - await peerTp.SendBamAsync(pgn: 0xFED1u, payload, peerCts.Token) + await peerTp.SendBamAsync(backgroundPgn, Datagram(Volatile.Read(ref sent)), peerCts.Token) .ConfigureAwait(false); Interlocked.Increment(ref sent); } @@ -842,36 +896,71 @@ await peerTp.SendBamAsync(pgn: 0xFED1u, payload, peerCts.Token) // Cycle re-claims to a new SA every iteration. Each ClaimAddressAsync triggers two // RebindTransportOnLoop calls (unbind to 0xFE, then rebind to the new SA) — that is // where the old/new channel overlap window lived pre-fix. - for (int i = 0; i < 8; i++) + for (int i = 0; i < claims; i++) { byte sa = (byte)(0x30 + i); await node.ClaimAddressAsync(sa).WithTimeout(ShortTimeout); node.ClaimState.Should().Be(J1939ClaimState.Claimed); - await Task.Delay(30); + + // The probe, sent the moment the rebind is done. Two things ride on it: the node + // must still receive broadcasts through the channel it just opened (asserted right + // here, so a build that stops delivering fails at the first claim rather than in an + // aggregate at the end), and this is also the instant a fire-and-forget-disposed + // predecessor would still be subscribed — so a duplicate is most likely exactly + // here, where the uniqueness check below will see it. + probeArrived = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + await probeTp.SendBamAsync(probePgn, Datagram(i)).WithTimeout(ShortTimeout); + (await probeArrived.Task.AsTaskWithTimeout(ShortTimeout)).Should().Be(i, + "the node must still receive broadcast BAM after rebinding to a new source address"); } peerCts.Cancel(); try { await peerTask.WithTimeout(ShortTimeout); } catch { /* peer cancel/dispose */ } - // Let any in-flight reassembly surface before the final count check. + // Let any in-flight reassembly surface before the final duplicate check. This is a + // settle for a negative assertion — there is no event that says "no duplicate is + // coming" — so it is a delay by necessity, not by convenience. await Task.Delay(150); int finalSent = Volatile.Read(ref sent); - int finalReceived = Volatile.Read(ref received); - - // The received count must never exceed sent: any excess means a rebind delivered - // the same broadcast BAM through two overlapping node-side transports. (Received < - // sent is expected — BAMs whose DT frames land during the ~ms rebind window get - // aborted / dropped by the disposed channel and never reassembled by the new one. - // The bug we are guarding against is duplicate delivery, not loss.) - finalReceived.Should().BeLessOrEqualTo(finalSent, - "no broadcast TP.BAM may be surfaced twice — before the fix, the fire-and-" + - "forget Dispose of the previous channel overlapped a freshly-opened channel " + - "and both subscriptions delivered the same reassembled datagram"); - // Sanity: this test is only meaningful if the peer actually managed to run many - // BAMs across the rebind cycles. + List<(uint Pgn, int Seq)> got; + lock (deliveredGate) got = new List<(uint, int)>(delivered); + + // Sanity: the background stream exists to put a BAM in flight across the rebinds, so + // the uniqueness assertion below is only meaningful if it actually produced traffic. + // The slowest configuration measured above still managed 8, and Windows CI 29. finalSent.Should().BeGreaterThan(5, - "the peer must generate enough BAM traffic to exercise the rebind window"); + "the background stream must generate BAM traffic across the rebind windows"); + + // The actual invariant. Before the fix, the fire-and-forget Dispose of the previous + // channel overlapped a freshly-opened one and both subscriptions delivered the same + // reassembled datagram — which shows up here as the same (PGN, sequence) twice. + got.Should().OnlyHaveUniqueItems( + "no broadcast TP.BAM may be surfaced twice, and each carries its own sequence number"); + // NotContain rather than OnlyContain: how much of the background stream survives is + // exactly what this test refuses to assert, and OnlyContain fails on an empty + // collection — which would smuggle "at least one background BAM arrived" back in as a + // hidden throughput assumption. Stated as a negative, it holds at any delivery rate + // including zero, while still catching a corrupted or invented sequence number. + // + // The bound is <= finalSent, not <: `sent` is incremented only after SendBamAsync + // completes, so cancelling the peer leaves the in-flight BAM uncounted while its DT + // frames may already be on the wire. The settle above exists precisely so that + // datagram can still reassemble, and its sequence is then equal to finalSent — a + // cancelled send that got through, not a corrupt payload. (Carried over from + // 2b13918, which found this on the previous form of the assertion; it applies + // unchanged here.) + got.Where(d => d.Pgn == backgroundPgn).Should() + .NotContain(d => d.Seq < 0 || d.Seq > finalSent, + "every delivered datagram must be one the peer actually sent, reassembled intact"); + + // Exact, and independent of how fast the runner is: one probe per claim, all delivered. + // Loss in the background stream is expected and deliberately not asserted on — a BAM + // straddling a rebind is dropped by design — but a node that stops receiving after a + // rebind cannot get all eight probes through. + got.Where(d => d.Pgn == probePgn).Select(d => d.Seq).Should() + .Equal(Enumerable.Range(0, claims), + "each rebind must be followed by a delivered probe, exactly once, in order"); } // --------------------------------------------------------------------------------------- diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index 1bec3ff..e2d16f4 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -508,6 +508,7 @@ public async Task SendCm_CancelInFlight_SendsConnectionAbort() using var sender = J1939TpFactory.Open(senderBus, sourceAddress: senderSa); + var rtsSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var abortSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); peerBus.FrameObserved += (_, e) => { @@ -515,17 +516,21 @@ public async Task SendCm_CancelInFlight_SendsConnectionAbort() var fields = J1939Id.Decompose((uint)e.CanFrame.ID); if (fields.SourceAddress != senderSa || !J1939Pgn.IsTransportCm(fields.Pgn)) return; var data = e.CanFrame.Data.ToArray(); - if (data.Length >= 8 && data[0] == J1939TpFrames.ControlAbort - && J1939TpFrames.ReadDataPgn(data) == pgn) - abortSeen.TrySetResult(data); + if (data.Length < 8 || J1939TpFrames.ReadDataPgn(data) != pgn) return; + if (data[0] == J1939TpFrames.ControlRts) rtsSeen.TrySetResult(data); + if (data[0] == J1939TpFrames.ControlAbort) abortSeen.TrySetResult(data); }; // Peer never replies with CTS, so the session stays open after RTS until we cancel. using var cts = new CancellationTokenSource(); var send = sender.SendCmAsync(pgn, destinationAddress: peerSa, RandomPayload(50, seed: 71), cts.Token); - // Wait until RTS has hit the wire (actor has registered the TX session). - await Task.Delay(50); + // Cancel only once the RTS is actually on the wire, i.e. the actor has registered the TX + // session there is something to abort. Waiting for the frame instead of 50 ms removes + // both failure modes of the delay: cancelling too early on a loaded machine (nothing + // registered yet, so no Connection Abort is due and the test fails for the wrong reason), + // and spending 50 ms per run when the RTS is out in microseconds. + await rtsSeen.Task.AsTaskWithTimeout(ShortTimeout); cts.Cancel(); Func act = async () => await send.WithTimeout(ShortTimeout); diff --git a/tests/CanKit.Pro.Tests/TestCases/RawCanSubscriptionTests.cs b/tests/CanKit.Pro.Tests/TestCases/RawCanSubscriptionTests.cs index db783e3..da9aeb2 100644 --- a/tests/CanKit.Pro.Tests/TestCases/RawCanSubscriptionTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/RawCanSubscriptionTests.cs @@ -53,6 +53,17 @@ private static async Task> Drain(ISubscription sub, int count return result; } + // Empties a subscription's buffer through the synchronous TryRead path and returns the IDs in + // arrival order. Used by the ControllableBus-driven tests, where delivery already happened + // synchronously inside RaiseObserved, so there is nothing left to wait for and Drain's + // timeout would only add a way for the test to pass by accident. + private static List DrainIds(ISubscription sub) + { + var ids = new List(); + while (sub.TryRead(out var frame)) ids.Add(frame.ID); + return ids; + } + private static readonly TimeSpan ShortTimeout = TimeSpan.FromSeconds(2); // FR-RAW-010: two subscriptions with disjoint ID filters on the same bus each receive only @@ -349,6 +360,112 @@ public async Task Reconfigure_Filter_At_Runtime_Subsequent_Frames_Follow_New_Cri after.Select(f => f.ID).Should().Equal(0x201); } + // FR-RAW-014, predicate overload: the same runtime-reconfiguration guarantee as for the + // ID-filter overload, plus the two things only this overload can express -- replacing an + // ID filter with a predicate, and passing null to go back to accepting everything. + // Driven through ControllableBus so each RaiseObserved is delivered synchronously on this + // thread: "which frames arrived after the reconfigure" is then a fact, not a drain race. + [Fact] + public void Reconfigure_Predicate_At_Runtime_Subsequent_Frames_Follow_New_Criterion() + { + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + + using var sub = service.Subscribe(CanIdFilter.Range(0x100, 0x1FF, CanFilterIDType.Standard)); + + bus.RaiseObserved(CanFrame.Classic(0x100, new byte[] { 1 }), isEcho: false); + bus.RaiseObserved(CanFrame.Classic(0x200, new byte[] { 2 }), isEcho: false); + DrainIds(sub).Should().Equal(0x100); + + // Predicate replaces the ID filter outright: 0x200 now matches and 0x100 no longer does, + // which an "and-ed on top of the old filter" implementation could not produce. + sub.Reconfigure(f => f.ID >= 0x200); + + bus.RaiseObserved(CanFrame.Classic(0x100, new byte[] { 3 }), isEcho: false); + bus.RaiseObserved(CanFrame.Classic(0x201, new byte[] { 4 }), isEcho: false); + DrainIds(sub).Should().Equal(0x201); + + // null is documented as "accept all", so both of the above must now arrive. + sub.Reconfigure((Func?)null); + + bus.RaiseObserved(CanFrame.Classic(0x100, new byte[] { 5 }), isEcho: false); + bus.RaiseObserved(CanFrame.Classic(0x201, new byte[] { 6 }), isEcho: false); + DrainIds(sub).Should().Equal(0x100, 0x201); + } + + // FR-RAW-012/014: a disposed subscription is not a silently inert one. Reconfigure after + // Dispose must throw ObjectDisposedException -- reconfiguring a subscription that will never + // deliver another frame is a caller bug, and swallowing it would hide it. + [Fact] + public void Reconfigure_After_Dispose_Throws_On_Both_Overloads() + { + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + + var sub = service.Subscribe(); + sub.Dispose(); + + var byFilter = () => sub.Reconfigure(CanIdFilter.Range(0x100, 0x1FF, CanFilterIDType.Standard)); + var byPredicate = () => sub.Reconfigure(f => f.ID == 0x100); + + byFilter.Should().Throw(); + byPredicate.Should().Throw(); + } + + // FR-RAW-011: the per-subscription buffer is bounded *and* drop-oldest. Bounded alone is not + // the requirement -- a drop-newest buffer is also bounded, and also never blocks dispatch, so + // only asserting "the consumer is not blocked" leaves the discard policy untested. What must + // hold is that an undrained subscription keeps the most recent `capacity` frames: monitoring + // consumers want current traffic, not a snapshot frozen at the moment they fell behind. + [Fact] + public void Full_Subscription_Buffer_Drops_The_Oldest_Frames_Not_The_Newest() + { + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + + const int capacity = 4; + const int sent = 7; + using var sub = service.Subscribe(bufferCapacity: capacity); + + // Never drained while sending: every frame past the fourth has to displace one. + for (var i = 0; i < sent; i++) + bus.RaiseObserved(CanFrame.Classic(0x100 + i, new byte[] { (byte)i }), isEcho: false); + + var buffered = DrainIds(sub); + + buffered.Should().HaveCount(capacity, "the buffer is bounded at its configured capacity"); + buffered.Should().Equal(new[] { 0x103, 0x104, 0x105, 0x106 }, + "a full drop-oldest buffer discards the three oldest frames and keeps the newest four, " + + "still in arrival order"); + } + + // TryRead is the non-blocking counterpart to Frames: it removes what is already buffered and + // reports emptiness rather than waiting. Both halves matter -- a TryRead that awaited arrival + // would deadlock the request/reply clients that use it to drain stale chatter before issuing + // a request, and one that returned a stale frame twice would double-deliver it. + [Fact] + public void TryRead_Removes_One_Buffered_Frame_And_Reports_An_Empty_Buffer() + { + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sub = service.Subscribe(); + + sub.TryRead(out var nothing).Should().BeFalse("nothing has been delivered yet"); + nothing.Should().Be(default(CanFrameView)); + + bus.RaiseObserved(CanFrame.Classic(0x100, new byte[] { 1 }), isEcho: false); + bus.RaiseObserved(CanFrame.Classic(0x101, new byte[] { 2 }), isEcho: false); + + sub.TryRead(out var first).Should().BeTrue(); + first.ID.Should().Be(0x100); + first.Data.ToArray().Should().Equal(new byte[] { 1 }); + + sub.TryRead(out var second).Should().BeTrue("TryRead consumes, so the next call sees the next frame"); + second.ID.Should().Be(0x101); + + sub.TryRead(out _).Should().BeFalse("both buffered frames have been consumed"); + } + // FR-RAW-013: the ID-range/mask fast path matches and excludes correctly (unit-level, no bus). [Fact] public void IdFilter_FastPath_Matches_And_Excludes() diff --git a/tests/CanKit.Pro.Tests/TestCases/TxConfirmTests.cs b/tests/CanKit.Pro.Tests/TestCases/TxConfirmTests.cs index 2ed34e3..d8e6da0 100644 --- a/tests/CanKit.Pro.Tests/TestCases/TxConfirmTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/TxConfirmTests.cs @@ -29,6 +29,12 @@ public class TxConfirmTests : IClassFixture private static ControllableBus OpenEcho() => ControllableBus.EchoCapable(VirtualAdapterFixture.NewSession("txconfirm")); + private static ControllableBus OpenDeferredEcho() + => ControllableBus.DeferredEchoCapable(VirtualAdapterFixture.NewSession("txconfirm")); + + // Only ever a bound against a hang: every wait below is on an event the test itself caused. + private static readonly TimeSpan ShortTimeout = TimeSpan.FromSeconds(5); + // FR-RAW-030/032: without echo, SendConfirmed resolves as soon as the driver accepts the // frame, explicitly marked as an approximation. [Fact] @@ -59,14 +65,71 @@ public async Task Echo_Bus_Confirms_Via_Real_Echo_Match() result.FailureReason.Should().Be(TxConfirmFailureReason.None); } - // FR-RAW-031: concurrent, byte-identical sends must each be matched to their own confirmation - // -- no cross-matching, no crash. This is the exact class of bug the review flagged for the - // ISO-TP prototype's deadline queue crashing on identical in-flight frames. Launched via - // Task.Run so they can genuinely interleave across real threads; the echo is delivered - // synchronously inside Transmit (as a real echo-mode adapter does), so true overlap of pending - // registrations isn't guaranteed on every single run, but this still exercises the exact - // thread-safety-sensitive paths (concurrent register/match/remove under the service's pending - // registry lock) end to end. + // FR-RAW-031, the actual FIFO assertion: with several byte-identical sends pending at once, + // the n-th echo must confirm the n-th *transmitted* send. Nothing else can tell the callers + // apart -- the frames are identical on the wire and the TxConfirmations they receive are + // identical too -- so the only observable of "matched FIFO" is which caller's task completes + // when, and that is exactly what this asserts. + // + // This needs ControllableBus in Deferred echo mode. With the synchronous echo a real adapter + // delivers, the echo re-enters CanBusService's pending-send lock on the transmitting thread + // before that thread leaves Transmit, so the pending list holds exactly one entry -- that + // thread's own -- for the entire match. Every FIFO ordering rule is then vacuously satisfied + // and the test cannot fail however the matching code is written (see DeferredEchoQueue). + [Fact] + public async Task Echo_Bus_Matches_Identical_Pending_Sends_In_Fifo_Order() + { + using var sender = OpenDeferredEcho(); + using var service = new CanBusService(sender); + + const int n = 4; + var frame = CanFrame.Classic(0x500, new byte[] { 42 }); + + // Long per-call timeout: the sends must still be pending when the last one registers, and + // the only thing that may ever complete one is an echo this test releases. + var sends = new Task[n]; + for (var i = 0; i < n; i++) + { + sends[i] = service.SendConfirmed(frame, TimeSpan.FromSeconds(30)); + // A parked echo means Transmit ran, which CanBusService does after registering the + // pending entry -- so this waits on registration order, not on wall-clock luck. + await sender.DeferredEchoes.WaitForEnqueuedAsync(i + 1, ShortTimeout); + } + + sender.DeferredEchoes.Count.Should().Be(n, "no echo has been released yet"); + sends.Should().OnlyContain(t => !t.IsCompleted, + "a send may only resolve once its own echo comes back"); + + for (var i = 0; i < n; i++) + { + sender.DeferredEchoes.ReleaseNext().Should().BeTrue(); + + // WhenAny over everything still outstanding, rather than awaiting sends[i] directly: + // a LIFO (or arbitrary) match resolves the *wrong* caller, and this reports that as + // "the wrong send was confirmed" immediately instead of as a timeout minutes later. + var completed = await Task.WhenAny(sends.Skip(i)).WaitAsync(ShortTimeout); + completed.Should().BeSameAs(sends[i], + "the {0}. echo must confirm the {0}. transmitted send, not a later one", i + 1); + + var confirmation = await completed; + confirmation.Confirmed.Should().BeTrue(); + confirmation.IsApproximated.Should().BeFalse(); + confirmation.FailureReason.Should().Be(TxConfirmFailureReason.None); + + for (var later = i + 1; later < n; later++) + sends[later].IsCompleted.Should().BeFalse( + "one echo confirms exactly one send, so send {0} must still be pending", later); + } + } + + // FR-RAW-031 under real concurrency: byte-identical sends racing across threads must each get + // their own confirmation -- no cross-matching, no crash. This is the exact class of bug the + // review flagged for the ISO-TP prototype's deadline queue crashing on identical in-flight + // frames. Deliberately kept on the *synchronous* echo bus: that is the reentrant + // register-transmit-match-remove path a real echo-mode adapter drives, and this test exists to + // hammer it from many threads at once. It says nothing about FIFO ordering -- with a + // synchronous echo there is never more than one pending entry to order. The test above is + // where ordering is proven. [Fact] public async Task Echo_Bus_Matches_Concurrent_Identical_Frames_Individually_Without_Crashing() {