diff --git a/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs b/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs index 216735d5..70922be9 100644 --- a/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/BusStateMonitorTests.cs @@ -46,15 +46,22 @@ public void CurrentState_Reflects_The_Bus_State_At_Construction() [Fact] public async Task StateChanged_Is_Not_Raised_While_The_State_Is_Unchanged() { + using var clock = new VirtualClock(); using var bus = OpenBus(); - using var actor = new ProtocolActor(); - using var monitor = new BusStateMonitor(bus, actor, TimeSpan.FromMilliseconds(20)); + var actor = clock.NewActor(); + var pollInterval = TimeSpan.FromMilliseconds(20); + using var monitor = new BusStateMonitor(bus, actor, pollInterval); var changes = 0; monitor.StateChanged += (_, _) => Interlocked.Increment(ref changes); - // Let many poll ticks run without ever changing the bus state. - await Task.Delay(TimeSpan.FromMilliseconds(200)); + // Let ten poll ticks run without ever changing the bus state, each proven armed before the + // clock moves past it. + for (var i = 0; i < 10; i++) + { + await clock.WaitUntilTimerArmedAsync(actor, pollInterval, Bounded); + await clock.AdvanceAsync(pollInterval); + } Volatile.Read(ref changes).Should().Be(0, "an unchanged state must never raise an edge-triggered event"); } @@ -105,19 +112,28 @@ public async Task StateChanged_Fires_On_Recovery_Back_Down_From_BusOff() [Fact] public async Task Dispose_Stops_Further_StateChanged_Events_And_Is_Idempotent() { + using var clock = new VirtualClock(); using var bus = OpenBus(); - using var actor = new ProtocolActor(); - var monitor = new BusStateMonitor(bus, actor, TimeSpan.FromMilliseconds(20)); + var actor = clock.NewActor(); + var pollInterval = TimeSpan.FromMilliseconds(20); + var monitor = new BusStateMonitor(bus, actor, pollInterval); var changes = 0; monitor.StateChanged += (_, _) => Interlocked.Increment(ref changes); + // Wait for the first poll to actually be armed before disposing: disposing while the + // constructor's initial RearmPoll post is still in flight would let the race resolve + // either way (a handle Dispose never gets to see and cancel is not what "Dispose stops + // the poll" is claiming). + await clock.WaitUntilTimerArmedAsync(actor, pollInterval, Bounded); + monitor.Dispose(); monitor.Dispose(); // idempotent // Change the state only after disposing: with the poll stopped, no event may arrive. bus.BusState = BusState.BusOff; - await Task.Delay(TimeSpan.FromMilliseconds(150)); // several poll intervals + await clock.AdvanceAsync(pollInterval + pollInterval + pollInterval); // several poll intervals + await clock.SettleAsync(); Volatile.Read(ref changes).Should().Be(0, "a disposed monitor must stop polling and raising events"); } diff --git a/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs b/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs index 15292e2b..6ae28568 100644 --- a/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/DeadlineTests.cs @@ -4,6 +4,7 @@ using System.Threading.Tasks; using CanKit.Pro.Actor; using CanKit.Pro.Reliability; +using CanKit.Pro.Tests.Infrastructure; using FluentAssertions; using Xunit; @@ -37,11 +38,16 @@ public async Task Deadline_Fires_OnExpired_After_The_Timeout_Elapses() [Fact] public async Task Complete_Before_Expiry_Prevents_OnExpired_And_Is_Idempotent() { - using var actor = new ProtocolActor(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); var scheduler = new DeadlineScheduler(actor); var fired = false; - var deadline = scheduler.Arm(TimeSpan.FromMilliseconds(200), () => fired = true); + var timeout = TimeSpan.FromMilliseconds(200); + var deadline = scheduler.Arm(timeout, () => fired = true); + // Arm before you advance: prove the deadline really was scheduled for the configured + // timeout before Complete races ahead of it. + await clock.WaitUntilTimerArmedAsync(actor, timeout, Bounded); deadline.Complete().Should().BeTrue("Complete wins the race well before the deadline would expire"); deadline.IsCompleted.Should().BeTrue(); @@ -53,54 +59,76 @@ public async Task Complete_Before_Expiry_Prevents_OnExpired_And_Is_Idempotent() deadline.IsCompleted.Should().BeTrue(); deadline.IsCancelled.Should().BeFalse(); - // Let the original timer's due point pass and round-trip through the loop; onExpired must - // never fire for a completed deadline. - await Task.Delay(TimeSpan.FromMilliseconds(300)); - await actor.PostAsync(() => 0); + // Move the clock past the original timer's due point and let the loop settle; onExpired + // must never fire for a completed deadline. + await clock.AdvanceAsync(timeout + TimeSpan.FromMilliseconds(100)); + await clock.SettleAsync(); fired.Should().BeFalse(); } [Fact] public async Task Cancel_Before_Expiry_Prevents_OnExpired() { - using var actor = new ProtocolActor(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); var scheduler = new DeadlineScheduler(actor); var fired = false; - var deadline = scheduler.Arm(TimeSpan.FromMilliseconds(200), () => fired = true); + var timeout = TimeSpan.FromMilliseconds(200); + var deadline = scheduler.Arm(timeout, () => fired = true); + await clock.WaitUntilTimerArmedAsync(actor, timeout, Bounded); + deadline.Dispose(); // Dispose == Cancel deadline.IsCancelled.Should().BeTrue(); - await Task.Delay(TimeSpan.FromMilliseconds(300)); - await actor.PostAsync(() => 0); + await clock.AdvanceAsync(timeout + TimeSpan.FromMilliseconds(100)); + await clock.SettleAsync(); fired.Should().BeFalse(); } [Fact] - public void Rearm_Before_Original_Expiry_Extends_The_Deadline() + public async Task Rearm_Before_Original_Expiry_Extends_The_Deadline() { - using var actor = new ProtocolActor(); - var scheduler = new DeadlineScheduler(actor); - - // #130. Task.Delay completes on the thread pool, and a saturated net48 pool injects - // threads only one or two per second. Delay(850) can still be pending when the rearmed - // 2000 ms deadline fires on the actor's dedicated thread, so WhenAny reports that fire - // and the assertion blames the superseded timer. These waits block the calling thread. - // The windows stay 50 / 850 / 2000 ms. The generation guard itself is + // Previously #130 kept this test on the wall clock: a saturated net48 thread pool injects + // threads only one or two per second, so Task.Delay(850) could still be pending when the + // rearmed 2000 ms deadline fired, and WhenAny would report that fire and blame the + // superseded timer instead. A VirtualClock sidesteps the thread pool race entirely -- the + // deadline's due points are moments this test itself schedules, so there is nothing left + // for the pool to starve. The generation guard itself is covered separately by // Rearm_Leaves_An_Already_Dispatched_Callback_Unable_To_Expire. - using var fired = new ManualResetEventSlim(false); - using var deadline = scheduler.Arm(TimeSpan.FromMilliseconds(600), () => fired.Set()); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + var scheduler = new DeadlineScheduler(actor); + var fired = false; - Thread.Sleep(TimeSpan.FromMilliseconds(50)); - deadline.Rearm(TimeSpan.FromMilliseconds(2000)).Should().BeTrue("re-arming a still-pending deadline succeeds"); + var original = TimeSpan.FromMilliseconds(600); + var extended = TimeSpan.FromMilliseconds(2000); + using var deadline = scheduler.Arm(original, () => fired = true); + await clock.WaitUntilTimerArmedAsync(actor, original, Bounded); + + deadline.Rearm(extended).Should().BeTrue("re-arming a still-pending deadline succeeds"); + // The rearmed timer must now be armed for the extended interval, not the remainder of the + // original one -- this is the generation guard's counterpart on the happy path. + await clock.WaitUntilTimerArmedAsync(actor, extended, Bounded); + + // Bracket from both sides. First: past the original due point, well short of the + // rearmed one. + await clock.AdvanceAsync(original); + await clock.SettleAsync(); + fired.Should().BeFalse("the original timeout must have been superseded by Rearm"); + deadline.IsExpired.Should().BeFalse(); - // Past the original 600 ms, and still more than a second short of the rearmed deadline. - fired.Wait(TimeSpan.FromMilliseconds(850)).Should().BeFalse( - "the original timeout must have been superseded by Rearm"); + // One tick short of the rearmed due point: still must not have fired. + var epsilon = TimeSpan.FromMilliseconds(1); + await clock.AdvanceAsync(extended - original - epsilon); + await clock.SettleAsync(); + fired.Should().BeFalse("the rearmed deadline has not reached its own due point yet"); deadline.IsExpired.Should().BeFalse(); - fired.Wait(Bounded).Should().BeTrue( - "the re-armed timeout must still fire at its new deadline"); + // The last tick reaches it. + await clock.AdvanceAsync(epsilon); + await clock.SettleAsync(); + fired.Should().BeTrue("the re-armed timeout must still fire at its new deadline"); deadline.IsExpired.Should().BeTrue(); } @@ -195,19 +223,24 @@ public async Task Exception_From_OnExpired_Surfaces_Via_The_Actor_Background_Exc [Fact] public async Task Disposing_The_Owning_Actor_While_Pending_Never_Fires_The_Deadline_And_Escapes_No_Exception() { - var actor = new ProtocolActor(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); var scheduler = new DeadlineScheduler(actor); var backgroundFaulted = false; actor.BackgroundExceptionOccurred += (_, _) => backgroundFaulted = true; var fired = false; - using var deadline = scheduler.Arm(TimeSpan.FromMilliseconds(200), () => fired = true); + var timeout = TimeSpan.FromMilliseconds(200); + using var deadline = scheduler.Arm(timeout, () => fired = true); + await clock.WaitUntilTimerArmedAsync(actor, timeout, Bounded); // The actor's FinalDrain discards not-yet-due Schedule callbacks, so a deadline that was - // still Pending simply never resolves (documented best-effort behavior). + // still Pending simply never resolves (documented best-effort behavior). actor.Dispose() + // joins the dedicated thread and only returns once FinalDrain has already run, so the + // outcome is settled the moment this call returns -- no wait, virtual or otherwise, is + // needed to observe it. actor.Dispose(); - await Task.Delay(TimeSpan.FromMilliseconds(300)); fired.Should().BeFalse("a pending deadline whose actor is disposed must never fire"); deadline.IsExpired.Should().BeFalse(); backgroundFaulted.Should().BeFalse("disposing the actor under a pending deadline must not raise any exception"); diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpBusOffTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpBusOffTests.cs index cd24f0b4..cc6a2caf 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpBusOffTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpBusOffTests.cs @@ -57,13 +57,24 @@ public async Task Active_MultiFrame_Send_Faults_With_BusOff_Instead_Of_Hanging() using var sender = IsoTpFactory.Open(bus, IsoTpEndpoint.Normal(0x300, 0x301), FastOptions()); + // #171: "give the channel a moment to register the pending FF confirmation" was a + // Task.Delay(100) guessing at a race. OnTransmitting fires synchronously from inside + // Transmit, and its own doc says CanBusService calls it "with the transmitting send + // already registered" -- i.e. once this has fired, the FF's entry is in the pending-send + // list, which is exactly the state the test needs before it drives BusOff. + var ffTransmitted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + bus.OnTransmitting = frame => + { + var payload = frame.Data.ToArray(); + if (payload.Length > 0 && (payload[0] >> 4) == 0x1) ffTransmitted.TrySetResult(true); + }; + // Multi-frame send: the FF goes out and its TX confirmation stays pending behind the // blocked echo. Driving the bus off while that confirmation is outstanding must abort // the send (L2 -> L3 propagation per FR-RAW-051), not hang. var send = sender.SendAsync(Enumerable.Range(0, 30).Select(i => (byte)i).ToArray()); - // Give the channel a moment to register the pending FF confirmation. - await Task.Delay(100); + await ffTransmitted.Task.WaitAsync(ShortTimeout); bus.BusState = BusState.BusOff; bus.RaiseFault(new InvalidOperationException("simulated bus-off")); diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpCanFdTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpCanFdTests.cs index dd376f21..ad1b88b3 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpCanFdTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpCanFdTests.cs @@ -126,6 +126,13 @@ public async Task Channel_With_UseCanFd_Emits_Only_CanFd_Frames_For_Sf_Ff_Cf_And using var receiver = IsoTpFactory.Open(busB, IsoTpEndpoint.Normal(0x7E8, 0x7E0), FastOptions(useCanFd: true)); var kinds = new List<(uint id, byte pci, CanFrameType kind)>(); + // #171: "let the hub deliver the tail frames to the sniffer" was a Task.Delay(100) + // guessing at how long the sniffer's independent subscription takes to catch up with + // the last CF/FC. sniffedAllKinds is set from inside the handler itself once every kind + // this assertion needs has actually arrived, which is what a fixed window could only + // ever approximate. + var sniffedAllKinds = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + bool sawSf = false, sawFf = false, sawCf = false, sawFc = false; snifferBus.FrameObserved += (_, view) => { var id = view.CanFrame.ID; @@ -138,6 +145,11 @@ public async Task Channel_With_UseCanFd_Emits_Only_CanFd_Frames_For_Sf_Ff_Cf_And lock (kinds) { kinds.Add(((uint)id, pci, view.CanFrame.FrameKind)); + if (id == 0x7E0 && (pci & 0xF0) == 0x00) sawSf = true; + if (id == 0x7E0 && (pci & 0xF0) == 0x10) sawFf = true; + if (id == 0x7E0 && (pci & 0xF0) == 0x20) sawCf = true; + if (id == 0x7E8 && (pci & 0xF0) == 0x30) sawFc = true; + if (sawSf && sawFf && sawCf && sawFc) sniffedAllKinds.TrySetResult(true); } }; @@ -153,7 +165,7 @@ public async Task Channel_With_UseCanFd_Emits_Only_CanFd_Frames_For_Sf_Ff_Cf_And await sender.SendAsync(sf); (await recvSf).Should().Equal(sf); - await Task.Delay(100); // let the hub deliver the tail frames to the sniffer + await sniffedAllKinds.Task.WaitAsync(ShortTimeout); List<(uint id, byte pci, CanFrameType kind)> observed; lock (kinds) diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index 2a551024..f4102b0c 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -334,6 +334,20 @@ public async Task MultiFrame_Send_Aborts_When_Peer_Exceeds_WftMax() } }; + // #171: "each Wait FC is processed within 20 ms before the next is sent" was a + // Task.Delay(20) standing in for two facts: the frame reached the sender's bus, and the + // sender's actor took it off the mailbox. fcArrived is the first (a frame on the wire, + // observed on busA where the sender listens); SettleAsync is the second (an actor + // round-trip after that frame was delivered). Sending the next Wait FC before both have + // happened would let two Waits collapse into the actor's mailbox as one delivery. + using var fcArrived = new SemaphoreSlim(0); + busA.FrameObserved += (_, e) => + { + if (e.CanFrame.ID != 0x401) return; + var payload = e.CanFrame.Data.ToArray(); + if (payload.Length > 0 && (payload[0] >> 4) == 0x3) fcArrived.Release(); + }; + byte[] pdu = Enumerable.Range(0, 30).Select(i => (byte)i).ToArray(); var sendTask = sender.SendAsync(pdu); @@ -346,7 +360,9 @@ public async Task MultiFrame_Send_Aborts_When_Peer_Exceeds_WftMax() var fc = IsoTpFrameCodec.BuildFlowControl(epBA, FlowStatus.Wait, blockSize: 0, stMinRaw: 0, isCanFd: false, padding: true); busB.Transmit(CanFrame.Classic(0x401, fc)); - await Task.Delay(20); + (await fcArrived.WaitAsync(ShortTimeout)).Should().BeTrue( + $"the sender's bus must see Wait FC #{i + 1}"); + await sender.SettleAsync().WaitAsync(ShortTimeout); } Func act = () => sendTask; @@ -432,8 +448,12 @@ public async Task Settle_Takes_A_Buffered_Single_Frame_Through_To_The_Inbox() byte[] sf = { 0x03, 0x7F, 0x3E, 0x78, 0x00, 0x00, 0x00, 0x00 }; long arrival = Stopwatch.GetTimestamp(); service.Deliver(new CanFrameView(CanFrameType.Can20, 0x7E8, sf, FrameFlags.None), arrival); - // Buffered, but nothing has taken it yet -- the reader task is starved by construction. - await Task.Delay(50); + // Buffered, but nothing has taken it yet -- the reader task is starved by construction + // (StarvedReaderBusService.WaitToReadAsync never completes). #171: that starvation is + // deterministic, not a race, so what proves "nothing has taken it yet" is not "enough + // wall time for a non-starved reader" but an actor round-trip: the actor has nothing + // queued for this frame because nothing pumped the subscription for it. + await actor.PostAsync(() => { }).WaitAsync(ShortTimeout); channel.TryReceiveWithArrival(out _).Should().BeFalse("the reader task has not run"); await channel.SettleAsync().WaitAsync(ShortTimeout); @@ -456,7 +476,10 @@ public async Task A_Discard_Given_A_Stamp_Keeps_What_Arrived_After_It() using var channel = new IsoTpChannel(service, ep, FastOptions(), ownsService: false, actor); long stamp = Stopwatch.GetTimestamp(); - await Task.Delay(5); + // #171: Task.Delay(5) was a guess at how long it takes Stopwatch to tick past `stamp`. + // What the test actually needs is a later reading, which a busy check gets deterministically + // and without depending on the runner's scheduler granularity at all. + while (Stopwatch.GetTimestamp() == stamp) { } byte[] sf = { 0x03, 0x7F, 0x3E, 0x78, 0x00, 0x00, 0x00, 0x00 }; service.Deliver(new CanFrameView(CanFrameType.Can20, 0x7E8, sf, FrameFlags.None), Stopwatch.GetTimestamp()); @@ -635,11 +658,25 @@ public async Task DiscardPendingPdus_Aborts_InFlight_Reassembly() // Mid-reassembly reset — must clear _rx so the trailing CF cannot complete a PDU. receiver.DiscardPendingPdus(); + // #171: "the trailing CF was processed within 50 ms" was a Task.Delay(50) guessing at + // cross-bus delivery. staleCfSeen is a frame on the wire (the receiver's own bus, + // busB); the SettleAsync after it is the actor round-trip that proves the receiver has + // taken it off the mailbox -- both before the fresh exchange below can be attributed to + // whichever run first. + var staleCfSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + busB.FrameObserved += (_, e) => + { + if (e.CanFrame.ID != unchecked((int)epPeer.TxCanId)) return; + var data = e.CanFrame.Data.ToArray(); + if (data.Length > 0 && (data[0] >> 4) == 0x2) staleCfSeen.TrySetResult(true); + }; + byte[] chunk = stalePayload.AsSpan(ffData).ToArray(); var cf = IsoTpFrameCodec.BuildConsecutiveFrame(epPeer, sequenceNumber: 1, chunk, isCanFd: false, padding: true); busA.Transmit(CanFrame.Classic(unchecked((int)epPeer.TxCanId), cf)); - await Task.Delay(50); + await staleCfSeen.Task.WaitAsync(ShortTimeout); + await receiver.SettleAsync().WaitAsync(ShortTimeout); // Stale multi-frame must not appear; a fresh SF must be the next ReceiveAsync result. using var sender = IsoTpFactory.Open(busA, epPeer, FastOptions()); @@ -1152,7 +1189,11 @@ public async Task Send_Cancelled_Before_Actor_Delivery_Emits_No_Frame_And_Channe Func act = () => sendTask.WaitAsync(ShortTimeout); await act.Should().ThrowAsync(); - // Give any straggling actor work time to (incorrectly) hit the wire under the bug. + // #171: "give any straggling actor work time to hit the wire" was a Task.Delay(100). A + // round-trip on the actor we parked and released proves [begin, cleanup] has run. An + // errant transmit would still leave through SendFrameOnBus's Task.Run, which the + // round-trip cannot see, so the original window follows it. + await actor.PostAsync(() => { }).WaitAsync(ShortTimeout); await Task.Delay(100); framesToPeer.Should().Be(0, "a send cancelled before the actor delivers BeginSendOnLoop must never put a frame on the bus"); @@ -1303,6 +1344,21 @@ await stillHeld.Should().ThrowAsync( // the gate is already free and this SF (DL=3) hits the bus immediately; under the fix it // must wait until we release the first confirmation. var secondSend = sender.SendAsync(new byte[] { 0x11, 0x22, 0x33 }); + + // #171: "no SF may hit the peer" was checked after a Task.Delay(100). SendAsync acquires + // _sendGate before ever posting to the actor (IsoTpChannel.cs), so while the gate is + // held secondSend cannot even reach the actor's mailbox at all -- under the fix nothing + // is scheduled, ever, and the round-trip below is a complete, deterministic proof of + // that. A regression that skipped the gate would instead post BeginSendOnLoop, whose own + // transmit is dispatched onto the thread pool (SendFrameOnBus's Task.Run) rather than + // run inline on the actor, so the round-trip alone only proves the actor's *decision*, + // not that dispatch's completion; the residual wait below covers only that one + // thread-pool hop, unchanged from the original margin, not lengthened for this. + var actorField = sender.GetType().GetField("_actor", + BindingFlags.Instance | BindingFlags.NonPublic); + actorField.Should().NotBeNull("IsoTpChannel must keep an _actor field for this race test"); + var actor = (IProtocolActor)actorField!.GetValue(sender)!; + await actor.PostAsync(() => { }).WaitAsync(ShortTimeout); await Task.Delay(100); framesToPeer.Should().Be(0, "no SF may hit the peer while the aborted send's SendConfirmed is still parked"); @@ -1377,10 +1433,27 @@ public async Task MultiFrame_Send_Accepts_FlowControl_Arriving_During_Last_Cf_Co } }; + // #171: "ensure peer FC is processed into DeferredFcs" was a Task.Delay(30). fcArrived + // is a frame on the wire (the sender's own bus, busA, where epAB.RxCanId = 0x261); the + // SettleAsync after it is the actor round-trip that proves the sender has taken it off + // the mailbox before the held CF confirm is released. + using var fcArrived = new SemaphoreSlim(0); + busA.FrameObserved += (_, e) => + { + if (e.CanFrame.ID != 0x261) return; + var payload = e.CanFrame.Data.ToArray(); + if (payload.Length > 0 && (payload[0] >> 4) == 0x3) fcArrived.Release(); + }; + // 20 bytes classic: FF(6) + CF1(7) + CF2(7). BS=1 => wait for FC after FF and after CF1. byte[] pdu = Enumerable.Range(0, 20).Select(i => (byte)(i + 1)).ToArray(); var sendTask = sender.SendAsync(pdu); + // The FC answering the FF comes first -- CF1 is not sent without it -- and is not one of + // the deferred FCs below; its permit is taken here so each wait in the loop is for the FC + // the peer sent while that CF's confirmation was parked (Codex on #185). + (await fcArrived.WaitAsync(ShortTimeout)).Should().BeTrue("the sender's bus must see the FC answering the FF"); + // Two block-ending CFs (CF1 then CF2): for each, wait until confirm is parked (FC already // sent by the peer handler above), then release so deferred FC is applied. for (int i = 0; i < 2; i++) @@ -1388,7 +1461,9 @@ public async Task MultiFrame_Send_Accepts_FlowControl_Arriving_During_Last_Cf_Co cfConfirmParked.Wait(TimeSpan.FromSeconds(3)).Should().BeTrue( $"CF confirm #{i + 1} must park after transmit so peer FC can defer"); cfConfirmParked.Reset(); - await Task.Delay(30); // ensure peer FC is processed into DeferredFcs + (await fcArrived.WaitAsync(ShortTimeout)).Should().BeTrue( + $"the sender's bus must see the FC deferred for block #{i + 1}"); + await sender.SettleAsync().WaitAsync(ShortTimeout); holdCfConfirm.Release(); } @@ -1428,6 +1503,18 @@ public async Task MultiFrame_Send_Counts_Wait_FlowControls_Deferred_During_Ff_Co ffSeen.TrySetResult(true); }; + // #171: "each Wait FC is processed within 20 ms" was a Task.Delay(20). fcArrived is a + // frame on the wire (the sender's own bus, busA, where epAB.RxCanId = 0x271); the + // SettleAsync after it is the actor round-trip that proves it has been queued as its own + // DeferredFc before the next one goes out, which is exactly what this test guards. + using var fcArrived = new SemaphoreSlim(0); + busA.FrameObserved += (_, e) => + { + if (e.CanFrame.ID != 0x271) return; + var payload = e.CanFrame.Data.ToArray(); + if (payload.Length > 0 && (payload[0] >> 4) == 0x3) fcArrived.Release(); + }; + byte[] pdu = Enumerable.Range(0, 30).Select(i => (byte)i).ToArray(); var sendTask = sender.SendAsync(pdu); @@ -1442,7 +1529,9 @@ public async Task MultiFrame_Send_Counts_Wait_FlowControls_Deferred_During_Ff_Co var fc = IsoTpFrameCodec.BuildFlowControl(epBA, FlowStatus.Wait, blockSize: 0, stMinRaw: 0, isCanFd: false, padding: true); busB.Transmit(CanFrame.Classic(0x271, fc)); - await Task.Delay(20); + (await fcArrived.WaitAsync(ShortTimeout)).Should().BeTrue( + $"the sender's bus must see Wait FC #{i + 1}"); + await sender.SettleAsync().WaitAsync(ShortTimeout); } holdConfirm.Release(); @@ -1528,6 +1617,16 @@ public async Task Rx_Classic_Rejects_CanFd_Escape_FirstFrame_Without_Allocating( Interlocked.Increment(ref anyFc); }; + // #171: "an illegal FC would have been sent within 100 ms" was a Task.Delay(100). + // hugeFfSeen is a frame on the wire (the receiver's own bus, busA); SettleAsync after it + // is the actor round-trip that proves the receiver has taken it off the mailbox and made + // its (silent) decision before anyFc is read. + var hugeFfSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + busA.FrameObserved += (_, e) => + { + if (e.CanFrame.ID == 0x7E8) hugeFfSeen.TrySetResult(true); + }; + // CAN-FD escape FF announcing 0x01000000 bytes. Classic TryParsePci must reject it // (isCanFd: false) so HandleRxFirstFrame never sees pci.Length = 16_777_216. byte[] hugeFf = @@ -1538,7 +1637,8 @@ public async Task Rx_Classic_Rejects_CanFd_Escape_FirstFrame_Without_Allocating( }; busB.Transmit(CanFrame.Classic(0x7E8, hugeFf)); - await Task.Delay(100); + await hugeFfSeen.Task.WaitAsync(ShortTimeout); + await receiver.SettleAsync().WaitAsync(ShortTimeout); anyFc.Should().Be(0, "classic channels must drop CAN-FD escape FFs without FC reply"); // Channel remains usable for a normal SF afterwards. @@ -1683,9 +1783,19 @@ public async Task A_Stale_StMin_Timer_Does_Not_Send_A_ConsecutiveFrame_Of_The_Ne ff2.Skip(2).Take(6).Should().Equal(second.Take(6)); // Now the old STmin elapses. Nothing may follow the FF: the peer has not sent Flow Control. + // #171: the trailing Task.Delay(100) ("a frame the timer released would be on the wire + // by now") added a wall-clock guess on top of an already-virtual clock. AdvanceAsync + // returns once the due timer's callback has run, and SettleAsync's round-trip proves the + // actor-side decision -- SendNextConsecutiveFrame's stale-tx check (IsoTpChannel.cs) -- + // is deterministic under the clock: a released timer for a superseded transfer decides + // not to build a frame at all, and nothing is ever handed to the wire. What neither + // proves is that a frame the timer *did* release has finished transmitting: SendFrameOnBus + // hands the actual bus write to the thread pool (Task.Run), a hop with no actor-side + // observable, so the residual wait below covers only that one hop -- unchanged from the + // original margin, not lengthened for this. await clock.AdvanceAsync(stMin); await clock.SettleAsync(); - await Task.Delay(100); // a frame the timer released would be on the wire by now + await Task.Delay(100); fromSender.Reader.TryRead(out var stray).Should().BeFalse( $"no Consecutive Frame may go out before the peer's Flow Control, but one did: {(stray is null ? "" : BitConverter.ToString(stray))}"); @@ -1741,12 +1851,17 @@ public async Task Receiver_With_LocalBlockSize_Emits_FlowControl_After_Each_Full FastOptions(localBs: 2)); var fcCount = 0; + // #171: "all three FCs have been counted within 50 ms of receive completing" was a + // Task.Delay(50). thirdFcSeen is set from inside the handler itself once the count this + // assertion needs has actually been observed on the sniffer's own (independent) + // subscription, which a fixed window could only ever approximate. + var thirdFcSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); snifferBus.FrameObserved += (_, view) => { if (view.CanFrame.ID != 0x7E8) return; var data = view.CanFrame.Data.Span; if (data.Length > 0 && (data[0] & 0xF0) == 0x30) - Interlocked.Increment(ref fcCount); + if (Interlocked.Increment(ref fcCount) == 3) thirdFcSeen.TrySetResult(true); }; // 38 bytes => FF (6) + 5 CFs (7,7,7,7,4): with BS=2 the receiver sends its initial @@ -1757,7 +1872,7 @@ public async Task Receiver_With_LocalBlockSize_Emits_FlowControl_After_Each_Full var got = await recvTask; got.Should().Equal(pdu); - await Task.Delay(50); + await thirdFcSeen.Task.WaitAsync(ShortTimeout); Volatile.Read(ref fcCount).Should().Be(3); } diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs index e2e8faaf..b9ae0a46 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpFunctionalClientTests.cs @@ -275,8 +275,29 @@ public async Task Functional_Collect_Discards_Frames_That_Arrived_Before_Send() unchecked((int)EcuResponseId), IsoTpFrameCodec.BuildSingleFrame(realEp, realPdu, isCanFd: false, padding: true)); - using var client = IsoTpFactory.OpenFunctional(busA, FunctionalTxId, RangeStart, RangeEnd, - FastOptions()); + // #171: this used to send the stale frame and sleep 50 ms, which left open whether it + // landed before the subscription (never seen) or inside it. The hook below puts it + // where the drain matters: into the collection's own subscription, after that + // subscription exists and before SendAndCollectAsync drains it and sends (Codex on #185). + // The hook returns only once busA has raised the stale frame. The service attaches to + // busA when it is created, before the handler below, so by the time that handler runs + // the service has dispatched the frame into the new subscription, where the drain has + // to find it (Codex on #185). + using var staleBuffered = new ManualResetEventSlim(); + using var service = new FrameConsumptionCountingBusService(new CanBusService(busA)) + { + OnNextSubscribe = () => + { + busB.Transmit(staleFrame); + staleBuffered.Wait(ShortTimeout).Should().BeTrue("the stale frame must reach the tester's bus"); + }, + }; + busA.FrameObserved += (_, e) => + { + if (e.CanFrame.ID == unchecked((int)StaleEcuResponseId)) staleBuffered.Set(); + }; + using var client = new IsoTpFunctionalClient(service, FunctionalTxId, RangeStart, RangeEnd, + FastOptions(), ownsService: false); busB.FrameObserved += (_, e) => { @@ -284,13 +305,6 @@ public async Task Functional_Collect_Discards_Frames_That_Arrived_Before_Send() busB.Transmit(realFrame); }; - // Blast the stale SF into the pipe. Give the virtual hub a moment to route it — - // if it lands after SendAndCollectAsync's internal Subscribe, DrainBuffered must - // drop it; if it lands before Subscribe, the subscription never sees it. Either - // way the assertion below must hold. - busB.Transmit(staleFrame); - await Task.Delay(50); - byte[] request = { 0x22, 0xF1, 0x90 }; var responses = await client.SendAndCollectAsync(request, CollectionWindow) .WaitAsync(ShortTimeout); @@ -420,8 +434,17 @@ public async Task Functional_Collect_Drains_Buffered_Frames_On_Window_Expiry() [Fact] public async Task Functional_Collect_Does_Not_Accept_Frames_After_Window_Expiry() { - // Bugbot 3604785766: after ResponseTimeout, disposing the subscription before the - // TryRead drain must prevent a post-window frame from being admitted. + // Bugbot 3604785766: a frame that arrives after the collection window must not be + // admitted, however late the collector gets to it. + // + // #171: this used to wait out the window on the wall clock and inject afterwards, by + // which time the collection had usually returned and never saw the frame. On the + // injected clock the late frame is put where the guard matters: the clock is moved past + // the deadline without waking the actor, so the window's timer has not fired yet, and + // the frame -- stamped past the deadline by the same clock -- reaches a collection that + // is still reading. The counting subscription shows the collector has taken it; only + // then is the window's timer allowed to fire. The window is 10 s so the actor's own wait + // cannot run out before that (Bugbot on #185). var session = NewSession(); using var busA = OpenClassic(session, 0); using var busB = OpenClassic(session, 1); @@ -429,31 +452,32 @@ public async Task Functional_Collect_Does_Not_Accept_Frames_After_Window_Expiry( const uint FunctionalTxId = 0x7DF; const uint EcuResponseId = 0x7E8; - using var client = IsoTpFactory.OpenFunctional(busA, FunctionalTxId, 0x7E8, 0x7EF, - FastOptions()); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var service = new FrameConsumptionCountingBusService( + new CanBusService(busA, actor.TimeSource.GetTimestamp)); + using var client = new IsoTpFunctionalClient(service, FunctionalTxId, 0x7E8, 0x7EF, + FastOptions(), ownsService: true, actor); - // Collect with a bounded window and no ECU reply during it. The window has to be - // long enough that the delay below is unambiguously past its end: CollectResponsesAsync - // arms the window on a continuation, so on a loaded runner the window can start tens of - // milliseconds after the call — with a 40 ms window and an 80 ms wait, that alone was - // enough to inject the frame while the window was still open and fail this test on - // macOS. The property under test is unchanged; only the margin is. - var window = TimeSpan.FromMilliseconds(200); + var window = TimeSpan.FromSeconds(10); var collectTask = client.CollectResponsesAsync(window); - await Task.Delay(TimeSpan.FromMilliseconds(window.TotalMilliseconds * 2)); + await clock.WaitUntilTimerArmedAsync(actor, window, ShortTimeout); + clock.Advance(window + TimeSpan.FromMilliseconds(1)); // past the deadline; the timer has not run - // Inject a late SF after the window has expired — must not appear in the result. + var taken = service.WaitUntilConsumedAsync(e => e.Frame.ID == unchecked((int)EcuResponseId)); var ep = IsoTpEndpoint.Normal(EcuResponseId, 0); byte[] latePdu = { 0x50, 0x01 }; busB.Transmit(CanFrame.Classic( unchecked((int)EcuResponseId), IsoTpFrameCodec.BuildSingleFrame(ep, latePdu, isCanFd: false, padding: true))); - await Task.Delay(60); + await taken.WaitAsync(ShortTimeout); + collectTask.IsCompleted.Should().BeFalse("the window's timer has not fired, so the collection was still reading"); + await clock.SettleAsync(); // the window's timer fires now var responses = await collectTask.WaitAsync(ShortTimeout); responses.Should().BeEmpty( - "a frame that arrives after the collection window must not be admitted " + - "(Bugbot 3604785766: dispose-before-drain)"); + "a frame stamped after the collection window must not be admitted " + + "(Bugbot 3604785766)"); } // ───────────────────────────────────────────────────────────────────────────────────────── diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs index d92c72d0..cb60dc4a 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939/J1939NodeTests.cs @@ -28,6 +28,9 @@ namespace CanKit.Pro.Tests.TestCases.J1939; public class J1939NodeTests : IClassFixture { private static readonly TimeSpan ShortTimeout = TimeSpan.FromSeconds(5); + // A frame the node's actor has sent still crosses the service and the virtual hub on their + // own threads; a negative check on the wire waits this long for one already on its way. + private static readonly TimeSpan WireWindow = TimeSpan.FromMilliseconds(200); private static string NewSession() => $"j1939-{Guid.NewGuid():N}"; @@ -1320,13 +1323,20 @@ public async Task A_Claim_During_The_Backoff_Before_A_Reclaim_Faults_And_The_Rec // Codex and Bugbot on #153: a loser that claims again before its delayed Cannot Claim went // out does not have it go out -- the new claim's announcement is the newer word on the bus. + // + // #171: the loser's actor runs on a VirtualClock so "past the backoff, with room" no longer + // trusts a wall-clock sleep to have outrun a ~153 ms pseudo-random backoff -- it moves the + // clock past it deterministically instead. [Fact] public async Task A_Delayed_Cannot_Claim_Is_Dropped_Once_A_New_Claim_Has_Started() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); using var busC = Open(session, 2); // spectator + using var serviceB = new CanBusService(busB); + var loserActor = clock.NewActor(); int cannotClaims = 0; busC.FrameObserved += (_, e) => @@ -1336,8 +1346,10 @@ public async Task A_Delayed_Cannot_Claim_Is_Dropped_Once_A_New_Claim_Has_Started if (J1939Pgn.IsAddressClaim(fields.Pgn) && fields.SourceAddress == J1939Pgn.NullAddress) Interlocked.Increment(ref cannotClaims); }; - using var winner = J1939Node.Open(busA, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); - using var loser = J1939Node.Open(busB, new J1939NodeOptions(Name(0x000158)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); // backoff 150 ms + var announce = TimeSpan.FromMilliseconds(80); + using var winner = J1939Node.Open(busA, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = announce }); + using var loser = new J1939NodeImpl(serviceB, new J1939NodeOptions(Name(0x000158)) + { ClaimAnnounceTimeout = announce }, ownsService: false, loserActor); // backoff 150 ms await winner.ClaimAddressAsync(0x60).WithTimeout(ShortTimeout); var lossSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); @@ -1345,14 +1357,25 @@ public async Task A_Delayed_Cannot_Claim_Is_Dropped_Once_A_New_Claim_Has_Started // The claim faults only once its Cannot Claim is on the bus, so the loss is awaited // through the state change here -- the exception would come too late to claim inside - // the backoff it is waiting for. + // the backoff it is waiting for. The loss itself is a real wire round trip with + // `winner` (a wall-clock node), not something blocked on the frozen clock, so it is + // awaited on the wall clock -- advancing the virtual clock here instead would race + // ahead of winner's real defensive re-announce and let the loser think it is + // uncontested. Only the second claim's own completion is blocked on the loser's clock. var lost = loser.ClaimAddressAsync(0x60); await lossSeen.Task.AsTaskWithTimeout(ShortTimeout); - await loser.ClaimAddressAsync(0x61).WithTimeout(ShortTimeout); // within the backoff, and through the new arbitration + var second = loser.ClaimAddressAsync(0x61); // within the backoff, and through the new arbitration + await clock.RunUntilAsync(second, step: TimeSpan.FromMilliseconds(20), giveUpAfter: ShortTimeout); Func awaitLost = () => lost.WithTimeout(ShortTimeout); await awaitLost.Should().ThrowAsync("dropping the Cannot Claim settles the loss that owed it"); - await Task.Delay(300); // past the backoff, with room + + // Past the backoff, with room: the dropped Cannot Claim must never appear. A frame the + // actor did send still crosses the service and the hub on their own threads, so the + // negative gets a wall window for it (Bugbot on #185). + await clock.AdvanceAsync(TimeSpan.FromMilliseconds(300)); + await clock.SettleAsync(); + await Task.Delay(WireWindow); Volatile.Read(ref cannotClaims).Should().Be(0, "the Cannot Claim was overtaken by the new claim"); loser.Address.Should().Be(0x61); } @@ -1426,21 +1449,28 @@ public async Task A_Second_Loss_Waits_Its_Own_Full_Backoff_Before_Cannot_Claim() // Codex on #153: a round waiting its backoff has announced nothing, and answers a Request // for Address Claimed by starting -- one announcement, now -- rather than by an // announcement of its own with the round's to follow. + // + // #171: node's actor runs on a VirtualClock so "past the backoff" is a deterministic advance + // rather than a 300 ms wall-clock sleep racing a ~153 ms pseudo-random backoff. [Fact] public async Task A_Request_During_The_Backoff_Starts_The_Round_With_A_Single_Announcement() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); + using var serviceA = new FrameConsumptionCountingBusService(new CanBusService(busA)); + var nodeActor = clock.NewActor(); const byte contended = 0x81; - using var winner = J1939Node.Open(busB, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); + var announce = TimeSpan.FromMilliseconds(80); + using var winner = J1939Node.Open(busB, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = announce }); await winner.ClaimAddressAsync(contended).WithTimeout(ShortTimeout); - using var node = J1939Node.Open(busA, new J1939NodeOptions(Name(0x000158)) // backoff 150 ms + using var node = new J1939NodeImpl(serviceA, new J1939NodeOptions(Name(0x000158)) // backoff 150 ms { - ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80), + ClaimAnnounceTimeout = announce, EnableArbitraryAddressClaiming = true, - }); + }, ownsService: false, nodeActor); int candidateClaims = 0; busB.FrameObserved += (_, e) => @@ -1452,15 +1482,31 @@ public async Task A_Request_During_The_Backoff_Starts_The_Round_With_A_Single_An var backingOff = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); node.AddressClaimChanged += (_, e) => { if (e.State == J1939ClaimState.Claiming && e.Address == contended + 1) backingOff.TrySetResult(true); }; + // The initial loss is a real wire round trip with `winner` (a wall-clock node), not + // something blocked on the frozen clock -- advancing the virtual clock while waiting for + // it would race ahead of winner's real defensive re-announce. Only the request-driven + // claim that follows is blocked on node's own clock. var claim = node.ClaimAddressAsync(contended); await backingOff.Task.AsTaskWithTimeout(ShortTimeout); + // The request must reach the node's actor before the clock moves, or the claim can + // finish on its original backoff and prove nothing about the request path (Bugbot on + // #185): the node's reader hands it over, and a round trip behind that post has run it. + var requestTaken = serviceA.WaitUntilConsumedAsync(e => + J1939Id.Decompose((uint)e.Frame.ID).Pgn == J1939Pgn.Request); busB.Transmit(CanFrame.Classic( (int)J1939Id.ComposePgn(6, J1939Pgn.Request, sourceAddress: 0x20, destinationAddress: J1939Pgn.GlobalAddress), new byte[] { 0x00, 0xEE, 0x00 }, isExtendedFrame: true)); + await requestTaken.WaitAsync(ShortTimeout); + await nodeActor.PostAsync(() => 0); - await claim.WithTimeout(ShortTimeout); + await clock.RunUntilAsync(claim, step: TimeSpan.FromMilliseconds(20), giveUpAfter: ShortTimeout); node.Address.Should().Be((byte)(contended + 1)); - await Task.Delay(300); // past the backoff: a round that still fired would announce again + + // Past the backoff: a round that still fired would announce again, and its frame would + // still be crossing the hub when the actor settles, hence the wall window. + await clock.AdvanceAsync(TimeSpan.FromMilliseconds(300)); + await clock.SettleAsync(); + await Task.Delay(WireWindow); Volatile.Read(ref candidateClaims).Should().Be(1, "the request started the round, whose announcement is the answer, and nothing announced twice"); } @@ -1469,33 +1515,66 @@ public async Task A_Request_During_The_Backoff_Starts_The_Round_With_A_Single_An // frame the loss owed the bus was then never sent at all. ClaimAddressAsync now faults // once the frame has gone out, so the loser here disposes as soon as it can and the // spectator still sees it. + // + // #171: the loser's actor runs on a VirtualClock (the actor is caller-injected, so disposing + // the node does not dispose it) so "a backoff that survived the dispose would have fired by + // now" is a deterministic advance instead of a 300 ms wall-clock sleep. [Fact] public async Task A_Lost_Claim_Faults_Only_Once_Its_Cannot_Claim_Is_On_The_Bus() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); using var busC = Open(session, 2); // spectator: it outlives the loser + using var serviceB = new CanBusService(busB); + var loserActor = clock.NewActor(); int cannotClaims = 0; + using var cannotClaimSeen = new SemaphoreSlim(0); busC.FrameObserved += (_, e) => { if (!e.CanFrame.IsExtendedFrame) return; var fields = J1939Id.Decompose((uint)e.CanFrame.ID); - if (J1939Pgn.IsAddressClaim(fields.Pgn) && fields.SourceAddress == J1939Pgn.NullAddress) Interlocked.Increment(ref cannotClaims); + if (J1939Pgn.IsAddressClaim(fields.Pgn) && fields.SourceAddress == J1939Pgn.NullAddress) + { + Interlocked.Increment(ref cannotClaims); + cannotClaimSeen.Release(); + } }; - using var winner = J1939Node.Open(busA, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); + var announce = TimeSpan.FromMilliseconds(80); + using var winner = J1939Node.Open(busA, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = announce }); await winner.ClaimAddressAsync(0x63).WithTimeout(ShortTimeout); - using (var loser = J1939Node.Open(busB, new J1939NodeOptions(Name(0x00015D)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) })) // backoff 153 ms - { - Func act = () => loser.ClaimAddressAsync(0x63).WithTimeout(ShortTimeout); + using (var loser = new J1939NodeImpl(serviceB, new J1939NodeOptions(Name(0x00015D)) + { ClaimAnnounceTimeout = announce }, ownsService: false, loserActor)) // backoff 153 ms + { + var lossSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + loser.AddressClaimChanged += (_, e) => { if (e.State == J1939ClaimState.CannotClaim) lossSeen.TrySetResult(true); }; + + // The loss is a real wire round trip with `winner` (a wall-clock node), not + // something blocked on the frozen clock -- advancing the virtual clock while + // waiting for it would race ahead of winner's real defensive re-announce. Only the + // Cannot Claim that follows the loss is blocked on the loser's own backoff timer. + var claim = loser.ClaimAddressAsync(0x63); + await lossSeen.Task.AsTaskWithTimeout(ShortTimeout); + Func act = () => clock.RunUntilAsync(claim, step: TimeSpan.FromMilliseconds(20), giveUpAfter: ShortTimeout); await act.Should().ThrowAsync(); } // the using a caller ends on the exception - await Task.Delay(300); // a backoff that survived the dispose would have fired by now - Volatile.Read(ref cannotClaims).Should().Be(1, + // The claim faulted once its Cannot Claim was handed to the driver; the spectator's copy + // still crosses the hub on its own thread, so it is awaited from the observer (Codex on + // #185). + (await cannotClaimSeen.WaitAsync(ShortTimeout)).Should().BeTrue( "the claim faulted only after its Cannot Claim went out, so disposing on the exception cannot suppress it"); + + // A backoff that survived the dispose would have fired by now; its frame gets the same + // hub crossing, as a wall window after the clock has moved. + await clock.AdvanceAsync(TimeSpan.FromMilliseconds(300)); + await clock.SettleAsync(); + (await cannotClaimSeen.WaitAsync(WireWindow)).Should().BeFalse( + "a disposed node's backoff must not send a second Cannot Claim"); + Volatile.Read(ref cannotClaims).Should().Be(1); } // Codex on #153: a Cannot Claim carries the null address, so two nodes answering the same @@ -1503,53 +1582,89 @@ public async Task A_Lost_Claim_Faults_Only_Once_Its_Cannot_Claim_Is_On_The_Bus() // different NAMEs on the bus. The answer a losing node owes therefore waits the same // §4.4.4.3 backoff as the Cannot Claim of the loss itself -- and a request arriving inside // that backoff is answered by that one frame, not by an immediate second copy. + // + // #171: the loser's actor runs on a VirtualClock. Rather than inferring "the request did not + // shortcut the backoff" from a Stopwatch delta after the fact, the backoff is bracketed from + // both sides directly: not yet due one tick short of it, due on the tick itself. [Fact] public async Task A_Request_During_The_Cannot_Claim_Backoff_Is_Answered_By_That_One_Frame() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); + using var serviceA = new FrameConsumptionCountingBusService(new CanBusService(busA)); + var loserActor = clock.NewActor(); const byte contended = 0x62; int cannotClaims = 0; - long firstCannotClaimAt = 0; + using var cannotClaimSeen = new SemaphoreSlim(0); busB.FrameObserved += (_, e) => { if (!e.CanFrame.IsExtendedFrame) return; var fields = J1939Id.Decompose((uint)e.CanFrame.ID); - if (!J1939Pgn.IsAddressClaim(fields.Pgn) || fields.SourceAddress != J1939Pgn.NullAddress) return; - Interlocked.CompareExchange(ref firstCannotClaimAt, Stopwatch.GetTimestamp(), 0); - Interlocked.Increment(ref cannotClaims); + if (J1939Pgn.IsAddressClaim(fields.Pgn) && fields.SourceAddress == J1939Pgn.NullAddress) + { + Interlocked.Increment(ref cannotClaims); + cannotClaimSeen.Release(); + } }; + // The actor's decision is settled by the clock, but the frame it sends still crosses the + // service and the hub on their own threads (Codex on #185): a positive is awaited from the + // observer, and each negative gets WireWindow for a frame already on its way. - using var winner = J1939Node.Open(busB, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); + var announce = TimeSpan.FromMilliseconds(80); + var backoff = TimeSpan.FromMilliseconds(153); // this NAME's §4.4.4.3 backoff + using var winner = J1939Node.Open(busB, new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = announce }); await winner.ClaimAddressAsync(contended).WithTimeout(ShortTimeout); - using var loser = J1939Node.Open(busA, new J1939NodeOptions(Name(0x00015D)) { ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80) }); // backoff 153 ms + using var loser = new J1939NodeImpl(serviceA, new J1939NodeOptions(Name(0x00015D)) + { ClaimAnnounceTimeout = announce }, ownsService: false, loserActor); - long lostAt = 0; var lossSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - loser.AddressClaimChanged += (_, e) => - { - if (e.State != J1939ClaimState.CannotClaim) return; - Interlocked.CompareExchange(ref lostAt, Stopwatch.GetTimestamp(), 0); - lossSeen.TrySetResult(true); - }; + loser.AddressClaimChanged += (_, e) => { if (e.State == J1939ClaimState.CannotClaim) lossSeen.TrySetResult(true); }; // The claim's exception arrives with the Cannot Claim itself, so the request has to be - // sent from the loss, which is the instant the backoff starts. + // sent from the loss, which is the instant the backoff starts. The loss is a real wire + // round trip with `winner` (a wall-clock node), not something blocked on the frozen + // clock -- advancing the virtual clock while waiting for it would race ahead of + // winner's real defensive re-announce. var lost = loser.ClaimAddressAsync(contended); await lossSeen.Task.AsTaskWithTimeout(ShortTimeout); + // Arm before you advance: prove the Cannot Claim really is scheduled for this NAME's + // backoff before the request is allowed to race it. + await clock.WaitUntilTimerArmedAsync(loserActor, backoff, ShortTimeout); + + // The request must be on the loser's actor before the clock moves, or it could be + // handled after the Cannot Claim and prove nothing (Bugbot on #185): the node's reader + // hands it over, which the counting subscription sees, and a round trip behind that + // post has run it. + var requestTaken = serviceA.WaitUntilConsumedAsync(e => + J1939Id.Decompose((uint)e.Frame.ID).Pgn == J1939Pgn.Request); busB.Transmit(CanFrame.Classic( (int)J1939Id.ComposePgn(6, J1939Pgn.Request, sourceAddress: 0x20, destinationAddress: J1939Pgn.GlobalAddress), new byte[] { 0x00, 0xEE, 0x00 }, isExtendedFrame: true)); + await requestTaken.WaitAsync(ShortTimeout); + await loserActor.PostAsync(() => 0); + + // Bracket the backoff from both sides: the request must not shortcut it. + var epsilon = TimeSpan.FromMilliseconds(1); + await clock.AdvanceAsync(backoff - epsilon); + await clock.SettleAsync(); + (await cannotClaimSeen.WaitAsync(WireWindow)).Should().BeFalse("the request did not shortcut the backoff"); + + await clock.AdvanceAsync(epsilon); + await clock.SettleAsync(); + (await cannotClaimSeen.WaitAsync(ShortTimeout)).Should().BeTrue("the backoff has now elapsed and the answer is on the bus"); + Volatile.Read(ref cannotClaims).Should().Be(1); Func awaitLost = () => lost.WithTimeout(ShortTimeout); await awaitLost.Should().ThrowAsync(); - await Task.Delay(500); // past the backoff, with room for a second copy to show up - Volatile.Read(ref cannotClaims).Should().Be(1, "the answer already waiting is the answer to the request"); - var waited = TimeSpan.FromSeconds((Interlocked.Read(ref firstCannotClaimAt) - Interlocked.Read(ref lostAt)) / (double)Stopwatch.Frequency); - waited.Should().BeGreaterThanOrEqualTo(TimeSpan.FromMilliseconds(100), - "the request did not shortcut the backoff -- a lower bound on the 153 ms a loaded host only lengthens"); + + // With room for a second copy to show up: the request must not have queued a duplicate. + await clock.AdvanceAsync(TimeSpan.FromMilliseconds(200)); + await clock.SettleAsync(); + (await cannotClaimSeen.WaitAsync(WireWindow)).Should().BeFalse("the answer already waiting is the answer to the request"); + Volatile.Read(ref cannotClaims).Should().Be(1); } // #58: a second ClaimAddressAsync while one is in arbitration faults instead of silently @@ -2638,21 +2753,27 @@ await peerTp.SendCmAsync(pgn: 0xEF00u, destinationAddress: 0xA0, payload) // Verified by restoring the regression (dropping both the teardown post and the // already-completed guard in OnClaimAnnounceElapsed): this test fails, as the version // before the timing rework did. + // + // #171: node's actor runs on a VirtualClock, so "wait past the original arbitration window" + // is a deterministic advance rather than a 700 ms wall-clock sleep racing a 500 ms window. // --------------------------------------------------------------------------------------- [Fact] public async Task ClaimAddressAsync_CancelDuringArbitration_TearsDownPendingClaim() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); // spectator: sees the announcement on the wire + using var service = new CanBusService(busA); + var actor = clock.NewActor(); - // Long arbitration window so the test can cancel comfortably in the middle. 500 ms is - // well above the actor scheduling jitter we need to observe. + // Long arbitration window so the test can cancel comfortably in the middle. + var announceTimeout = TimeSpan.FromMilliseconds(500); var opts = new J1939NodeOptions(Name(1)) { - ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(500), + ClaimAnnounceTimeout = announceTimeout, }; - using var node = J1939Node.Open(busA, opts); + using var node = new J1939NodeImpl(service, opts, ownsService: false, actor); // The announcement leaving the bus is the signal that the arbitration window is about // to be armed -- see the note above on why the Claiming transition is too early and a @@ -2672,6 +2793,9 @@ public async Task ClaimAddressAsync_CancelDuringArbitration_TearsDownPendingClai var claimTask = node.ClaimAddressAsync(0x33, cts.Token); await announced.Task.WithTimeout(ShortTimeout); + // Arm before you advance: prove the arbitration deadline really is scheduled before + // the cancel is allowed to race it. + await clock.WaitUntilTimerArmedAsync(actor, announceTimeout, ShortTimeout); cts.Cancel(); @@ -2681,9 +2805,10 @@ public async Task ClaimAddressAsync_CancelDuringArbitration_TearsDownPendingClai Func awaitClaim = () => claimTask.WithTimeout(ShortTimeout); await awaitClaim.Should().ThrowAsync(); - // Wait past the original arbitration window so any surviving timer would have - // fired. - await Task.Delay(700); + // Move the clock past the original arbitration window so any surviving timer would + // have fired. + await clock.AdvanceAsync(announceTimeout + TimeSpan.FromMilliseconds(200)); + await clock.SettleAsync(); // NotClaimed specifically, not merely "not Claimed": a missing teardown leaves the // node wedged in Claiming, which NotBe(Claimed) would wave through. This is the @@ -2692,7 +2817,8 @@ public async Task ClaimAddressAsync_CancelDuringArbitration_TearsDownPendingClai node.Address.Should().BeNull(); // A fresh claim must still work (i.e. teardown left the state machine consistent). - await node.ClaimAddressAsync(0x44).WithTimeout(ShortTimeout); + var fresh = node.ClaimAddressAsync(0x44); + await clock.RunUntilAsync(fresh, step: TimeSpan.FromMilliseconds(20), giveUpAfter: ShortTimeout); node.ClaimState.Should().Be(J1939ClaimState.Claimed); node.Address.Should().Be((byte)0x44); } @@ -3104,164 +3230,129 @@ await awaitSend.Should().ThrowAsync( // FR-J1939-007: periodic single-frame PGN send. Every periodic PGN flows through the // node's SendAsync / actor loop (L2 scheduling) — the previous dual `IPeriodicTx` path // was collapsed to a single implementation (PR #33) so error handling and claim-gate - // semantics are uniform across payload sizes. The test collects a run of frames on a - // spectator bus and asserts the two rate properties that survive a loaded runner. - // - // What this test deliberately does NOT assert is grid alignment, and the reason is - // measurement resolution rather than modesty. PeriodicSchedule skips a tick whose - // previous emission is still in flight and coalesces the ticks that fell behind by - // advancing the anchor whole periods at a time, so under load it emits at 2 x period, or - // 3 x, *by design*. Checking that those gaps land on the grid means resolving them to - // better than half a period -- and these stamps are taken on a spectator bus, where an - // emission observed late next to one observed on time already moves a gap by ~50 ms. - // Against a 120 ms period that is 40 % of the grid spacing, so the check cannot separate - // a drifting scheduler from a busy host no matter how the tolerance is set. Measured, not - // assumed: under an 8x CPU overload the observed gaps scatter across 130..410 ms. - // - // The anti-drift property is real and is tested -- by the multi-frame sibling below, - // which earns the resolution by deriving its period from a measured send time so that one - // period is several times the jitter. Asserting it twice, once where it cannot be - // measured, bought a standing red leg on macOS (#92) and no coverage. - // - // So, the two properties that hold regardless of host load: + // semantics are uniform across payload sizes. // - // * No bursting -- every gap rounds to at least one slot. Coalescing must advance the - // anchor, never queue the ticks it skipped and release them back to back. Jitter - // cannot fake this: half a period separates a real emission from slot zero. - // * Never faster than requested -- the mean gap is at least one period. Dropped ticks - // only ever make it longer, so this bound is one-sided in the direction load pushes - // and cannot be tripped by a slow runner. - // - // The collection loop above closes the other side: requiring requiredSamples emissions - // inside ShortTimeout bounds the mean gap from above too, loosely but honestly. + // #171 (audit "still tight" budget): this test used to run on the wall clock, deliberately + // NOT asserting grid alignment because the required measurement resolution was not + // achievable there -- see StartPeriodicSend_MultiFrame_Emits_On_An_Exact_Grid_On_A_Clock_ + // The_Test_Drives's doc comment for why. The exact-grid pattern that sibling test earns + // through a measured per-frame send cost applies just as well to a single frame, whose + // emission is one CAN frame rather than a TP session: on a clock this test drives, "on the + // grid" is exact, so this is now the same, stronger assertion the multi-frame sibling + // makes, rather than a weaker mean-gap heuristic sized against a 10.56 s wall-clock budget + // a 3x-loaded runner could still exhaust (the property the audit flagged). // --------------------------------------------------------------------------------------- [Fact] public async Task StartPeriodicSend_SingleFrame_FiresAtConfiguredPeriod() { - var session = NewSession(); - using var busA = Open(session, 0); - using var busB = Open(session, 1); // spectator: samples arrival times + var period = TimeSpan.FromMilliseconds(120); + const int requiredEmissions = 6; + + using var clock = new VirtualClock(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + var senderActor = clock.NewActor(); + using var sender = new J1939NodeImpl(service, new J1939NodeOptions(Name(1)), ownsService: false, senderActor); + + // The claim waits out its contention window on the same clock, so it needs the clock + // moved before it can succeed -- it is a precondition here, not the subject. + await clock.RunUntilAsync(sender.ClaimAddressAsync(0xC1), + step: TimeSpan.FromMilliseconds(50), giveUpAfter: ShortTimeout); + + // A send costs virtual time from here on, as in the multi-frame sibling: on a clock + // where sending is free, a send-then-delay loop lands on the same grid as an anchored + // schedule and the slot assertions below could not tell them apart (Bugbot on #185). + var perFrameCost = TimeSpan.FromMilliseconds(1); + bus.OnTransmitting = _ => clock.Advance(perFrameCost); - using var sender = J1939Node.Open(busA, new J1939NodeOptions(Name(1))); - await sender.ClaimAddressAsync(0xC1).WithTimeout(ShortTimeout); - - // The stamp collection is protected by its own lock; the FrameObserved handler runs - // on the bus's dispatch thread and multiple readers might in principle observe the - // frame concurrently on some adapters. - // Stopwatch ticks, not DateTime.UtcNow: these samples are only ever subtracted from - // each other, and a wall clock can step under them mid-run. Same monotonic basis the - // actor measures its own deadlines on. - var stamps = new List(); - var stampsLock = new object(); const uint targetPgn = 0xFEE5u; // PDU2, PS=0xE5 (arbitrary), well-known-ish - busB.FrameObserved += (_, e) => + var stamps = new List(); + var stampsLock = new object(); + bus.FrameObserved += (_, e) => { if (!e.CanFrame.IsExtendedFrame) return; var fields = J1939Id.Decompose((uint)e.CanFrame.ID); if (fields.SourceAddress != 0xC1) return; if (fields.Pgn != targetPgn) return; - lock (stampsLock) stamps.Add(Stopwatch.GetTimestamp()); + lock (stampsLock) stamps.Add(clock.Elapsed); }; - // 120 ms period is comfortably above both the ~1 ms virtual-loopback latency and the - // ~15.6 ms default timer granularity on Windows, and short enough to gather the - // samples in under two seconds. - var period = TimeSpan.FromMilliseconds(120); - // 22, not 10. The bound below is undercut by however late the run's first emission was, - // divided by the number of *measured* gaps -- see the assertion for why that is the error - // term that matters. Note the two subtractions: 22 emissions give 21 gaps, and trimming - // the warm-up leaves 20. Ten samples would divide a 240 ms cold start by nine and lose - // 27 ms of a 120 ms period; twenty measured gaps divide it by twenty, which is what the - // 10 % allowance below is worth. This is the knob that makes the assertion sound, so it - // is not a free parameter. - const int requiredSamples = 22; + int Count() { lock (stampsLock) return stamps.Count; } + var payload = new byte[] { 0x11, 0x22, 0x33, 0x44 }; var message = new J1939Message(targetPgn, payload, priority: 6, destinationAddress: 0xFF); + // The virtual resolution the period is bracketed to. + var Step = TimeSpan.FromMilliseconds(1); + + var startedAt = clock.Elapsed; using (var handle = sender.StartPeriodicSend(message, period)) { handle.Should().NotBeNull(); - // Collect until we have enough samples for a stable mean, or bail out with a - // clear failure message if the schedule never fires. - // Not ShortTimeout: 21 samples at 120 ms need 2.5 s even when nothing is dropped, and - // a loaded runner coalescing to 2x or 3x the period needs several times that. Four - // times the nominal run is generous enough not to fail for being slow, and still - // bounds the rate from above -- the one direction the assertion below does not cover. - var collectBudget = TimeSpan.FromMilliseconds(period.TotalMilliseconds * requiredSamples * 4); - var deadline = Stopwatch.GetTimestamp() + (long)(collectBudget.TotalSeconds * Stopwatch.Frequency); - while (true) + for (var slot = 1; slot <= requiredEmissions; slot++) { - int count; - lock (stampsLock) count = stamps.Count; - if (count >= requiredSamples) break; - if (Stopwatch.GetTimestamp() >= deadline) - throw new TimeoutException( - $"Expected at least {requiredSamples} periodic emissions within " + - $"{collectBudget.TotalSeconds:F1}s; observed {count}."); - await Task.Delay(20); - } - } + var slotPoint = startedAt + TimeSpan.FromTicks(period.Ticks * slot); - // Post-Dispose: no additional frames should arrive after a settle window. - int countAtDispose; - lock (stampsLock) countAtDispose = stamps.Count; - await Task.Delay(period + period); // wait 2 periods - int countAfterSettle; - lock (stampsLock) countAfterSettle = stamps.Count; + // Which period the schedule is actually on is decided here, by asking the actor + // which instant its next tick is armed for -- see the multi-frame sibling's doc + // comment for why the wire cannot answer it and a jump straight to the slot is + // not sufficient on its own. + var remaining = slotPoint - clock.Elapsed; + await clock.WaitUntilTimerArmedAsync(senderActor, remaining, ShortTimeout, Step); + var armed = await senderActor.NextTimerDelayAsync(); + var lateBy = armed is { } delay && delay > remaining ? delay - remaining : TimeSpan.Zero; - countAfterSettle.Should().BeLessOrEqualTo(countAtDispose + 1, - "disposing the handle must stop the periodic loop so at most an already-in-flight " + - "SendAsync may still land after Dispose returns"); + // One tick short of the slot: corroboration on the wire that the tick armed + // above has not fired early. It is the barrier, not this, that pins the period. + await clock.AdvanceToAsync(slotPoint - Step); + await clock.SettleAsync(); + Count().Should().Be(slot - 1, + "the clock is one tick short of slot {0}, so that emission is not due yet", + slot); - List snapshot; - lock (stampsLock) snapshot = new List(stamps); - snapshot.Count.Should().BeGreaterOrEqualTo(requiredSamples); + await clock.AdvanceToAsync(slotPoint); + if (lateBy > TimeSpan.Zero) + await clock.AdvanceAsync(lateBy); + await WaitForAnnouncesAsync(Count, slot); - double targetMs = period.TotalMilliseconds; - var gaps = new List(snapshot.Count - 1); - for (int i = 1; i < snapshot.Count; i++) - gaps.Add((snapshot[i] - snapshot[i - 1]) * 1000d / Stopwatch.Frequency); + // The frame on the wire is not the end of the emission: OnTick drops a tick + // whose predecessor is still in flight, and that state clears on the node's own + // loop. Waiting for the node to say so is what keeps the next slot's assertion + // about the grid rather than about a race. + await WaitForAnnouncesAsync(() => sender.PeriodicEmissionsCompleted, slot); + } + } - // Never faster than requested, as the mean gap over a warm-up-trimmed sample. - // - // Three statistics have been tried here and the first two were chosen by intuition; this - // one is chosen by its error term, which is the only way to size it honestly. - // - // Every emission lands on its own grid slot, late by however long the host stalled: - // t(i) = slot(i) * period + late(i). Summing the gaps telescopes, so - // - // mean gap = (slots spanned / gaps) * period + (late(last) - late(first)) / gaps - // - // The first term is at least the period, because slots are distinct. The second is the - // whole problem, and it is bounded by the *number of gaps* -- nothing else. That rules - // out both earlier attempts: - // - // * The plain mean over nine gaps divides a cold start by nine. A first tick 239 ms - // late leaves a mean of 106.8 ms against this bound, from a scheduler doing exactly - // what it documents. - // * The median has no such error term to shrink, which looked like an advantage and is - // not: it is sensitive to the shape of the jitter instead of its size. On a real - // macOS runner the gaps came out 83, 132, 73, 173, 67, 188, 62, 105, 185 -- mean - // 118.7 ms, so the rate was right to within 1 % -- and the median was 105 ms, because - // an odd number of alternating gaps has one more short than long. It failed a - // perfectly good run. - // - // So: trim the first gap, which is the only one measured from a cold schedule, and take - // the mean of the rest. Twenty measured gaps hold the residual endpoint term under a - // tenth of a period for any swing below 20 x 12 ms = 240 ms between the second emission's - // lateness and the last one's -- which is exactly the worst cold start observed here, and - // an order of magnitude beyond the jitter a loaded runner has otherwise produced. - // Oscillation cancels in a mean by construction, so the case above passes. - var measured = gaps.GetRange(1, gaps.Count - 1); - var meanGap = measured.Sum() / measured.Count; - - // One-sided on purpose: load can only lengthen gaps, so there is no honest upper bound to - // pair with this one. The collection loop bounds the rate from above instead. - meanGap.Should().BeGreaterOrEqualTo(targetMs * 0.9, - $"mean gap over {measured.Count} samples ({meanGap:F0} ms, first gap discarded as " - + $"warm-up) must not undercut the configured period ({targetMs:F0} ms); " - + "observed gaps: " + string.Join(", ", gaps.ConvertAll(g => $"{g:F0}"))); + // Post-dispose. Every emission above was waited out to completion, so nothing is in + // flight: the schedule must have left no tick armed, and period after period -- one at a + // time, since a single jump would let a live schedule coalesce its ticks into one -- + // must add no frame (Codex on #185). + var countAtDispose = Count(); + (await senderActor.NextTimerDelayAsync()).Should().BeNull( + "disposing the handle cancels the schedule's tick"); + for (var quiet = 0; quiet < 3; quiet++) + { + await clock.AdvanceAsync(period); + await clock.SettleAsync(); + } + sender.PeriodicEmissionsCompleted.Should().Be(requiredEmissions, "no tick started a send after Dispose"); + Count().Should().Be(countAtDispose, "disposing the handle must stop the periodic loop"); + + List snapshot; + lock (stampsLock) snapshot = new List(stamps); + snapshot.Count.Should().BeGreaterOrEqualTo(requiredEmissions); + + // Every emission sits on its slot -- no earlier, and not on a later one either (which a + // schedule that restarted its period after each send would have drifted onto). + for (var i = 0; i < requiredEmissions; i++) + { + var slotPoint = startedAt + TimeSpan.FromTicks(period.Ticks * (i + 1)); + snapshot[i].Should().BeGreaterThanOrEqualTo(slotPoint, + "emission {0} is triggered by the clock reaching its slot", i + 1); + snapshot[i].Should().BeLessThan(slotPoint + period, + "emission {0} belongs to slot {0} and not to a later one", i + 1); + } } /// @@ -3975,30 +4066,33 @@ public async Task StartPeriodicSend_SingleFrame_BeforeClaim_ThrowsNoAddress() // SendAsync's pre-flight claim gate throws J1939NoAddressException on every subsequent // tick, so no CAN frame is emitted while the node is un-claimed. Assert the wire goes // quiet after unseating. + // + // #171: owner's actor runs on a VirtualClock, so both "the schedule really was running" + // and "~2 periods of quiet after the loss" are deterministic advances instead of + // wall-clock polling and a 130 ms sleep. [Fact] public async Task StartPeriodicSend_SingleFrame_StopsAfterAddressLoss() { + using var clock = new VirtualClock(); var session = NewSession(); using var busA = Open(session, 0); using var busB = Open(session, 1); using var busC = Open(session, 2); // spectator: counts periodic emissions + using var serviceA = new CanBusService(busA); + var ownerActor = clock.NewActor(); // Owner has a HIGHER numeric NAME → lower priority → will be unseated when the // peer with a lower NAME claims the same SA per SAE J1939-81 §4.4.3.2. - var ownerOpts = new J1939NodeOptions(Name(0x000200)) - { - ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80), - }; - var peerOpts = new J1939NodeOptions(Name(0x000010)) - { - ClaimAnnounceTimeout = TimeSpan.FromMilliseconds(80), - }; + var announce = TimeSpan.FromMilliseconds(80); + var ownerOpts = new J1939NodeOptions(Name(0x000200)) { ClaimAnnounceTimeout = announce }; + var peerOpts = new J1939NodeOptions(Name(0x000010)) { ClaimAnnounceTimeout = announce }; - using var owner = J1939Node.Open(busA, ownerOpts); + using var owner = new J1939NodeImpl(serviceA, ownerOpts, ownsService: false, ownerActor); using var peer = J1939Node.Open(busB, peerOpts); const byte contendedSa = 0x50; - await owner.ClaimAddressAsync(contendedSa).WithTimeout(ShortTimeout); + await clock.RunUntilAsync(owner.ClaimAddressAsync(contendedSa), + step: TimeSpan.FromMilliseconds(20), giveUpAfter: ShortTimeout); owner.ClaimState.Should().Be(J1939ClaimState.Claimed); owner.Address.Should().Be(contendedSa); @@ -4006,15 +4100,14 @@ public async Task StartPeriodicSend_SingleFrame_StopsAfterAddressLoss() // assertion is independent of the owner node's internal state and matches what // downstream ECUs actually observe. const uint targetPgn = 0xFEE7u; - var stamps = new List(); - var stampsLock = new object(); + var count = 0; busC.FrameObserved += (_, e) => { if (!e.CanFrame.IsExtendedFrame) return; var fields = J1939Id.Decompose((uint)e.CanFrame.ID); if (fields.Pgn != targetPgn) return; if (fields.SourceAddress != contendedSa) return; - lock (stampsLock) stamps.Add(DateTime.UtcNow); + Interlocked.Increment(ref count); }; var period = TimeSpan.FromMilliseconds(40); @@ -4022,17 +4115,13 @@ public async Task StartPeriodicSend_SingleFrame_StopsAfterAddressLoss() destinationAddress: 0xFF); using var handle = owner.StartPeriodicSend(message, period); - // Wait until the schedule has actually put a few frames on the wire so the - // "stop" assertion below is meaningful (the schedule really was running). - var readyDeadline = DateTime.UtcNow + ShortTimeout; - while (true) + // Move the clock through three ticks so the "stop" assertion below is meaningful (the + // schedule really was running), each one proven armed before the clock advances past it. + for (var i = 0; i < 3; i++) { - int c; - lock (stampsLock) c = stamps.Count; - if (c >= 3) break; - if (DateTime.UtcNow >= readyDeadline) - throw new TimeoutException("Expected ≥3 periodic frames from owner before contest."); - await Task.Delay(10); + await clock.WaitUntilTimerArmedAsync(ownerActor, period, ShortTimeout, TimeSpan.FromMilliseconds(2)); + await clock.AdvanceAsync(period); + await WaitForAnnouncesAsync(() => Volatile.Read(ref count), i + 1); } // Peer with lower NAME claims the same SA. HandleIncomingAddressClaim's @@ -4042,23 +4131,40 @@ public async Task StartPeriodicSend_SingleFrame_StopsAfterAddressLoss() await peer.ClaimAddressAsync(contendedSa).WithTimeout(ShortTimeout); peer.Address.Should().Be(contendedSa); - // Wait for the owner's state machine to observe the contest. + // Wait for the owner's state machine to observe the contest. This is a real wire round + // trip with `peer` (a wall-clock node) landing on the owner's actor loop, not something + // blocked on the frozen clock, so it stays a category-1 poll on the wall clock. var lossDeadline = DateTime.UtcNow + ShortTimeout; while (owner.ClaimState == J1939ClaimState.Claimed && DateTime.UtcNow < lossDeadline) await Task.Delay(10); owner.ClaimState.Should().NotBe(J1939ClaimState.Claimed); owner.Address.Should().BeNull(); - // Give the schedule ~2 periods to observe the state transition and let the - // in-flight SendAsync (if any) drain. Peer traffic on `contendedSa` is filtered by - // NAME (owner's Name(0x200) ≠ peer's Name(0x010)), so any frames on `contendedSa` - // that arrive here originate from the owner's periodic loop *not yet stopping* — - // that is exactly the bug we are guarding against. - int countAfterLoss; - lock (stampsLock) countAfterLoss = stamps.Count; - await Task.Delay(period + period + TimeSpan.FromMilliseconds(50)); - int countAfterQuiet; - lock (stampsLock) countAfterQuiet = stamps.Count; + // Give the schedule a few more ticks to observe the state transition and let each + // dispatched SendAsync (successful or gated) drain. Peer traffic on `contendedSa` is + // filtered by NAME (owner's Name(0x200) ≠ peer's Name(0x010)), so any frames on + // `contendedSa` that arrive here originate from the owner's periodic loop *not yet + // stopping* -- that is exactly the bug we are guarding against. + // + // PeriodicEmissionsCompleted increments whether SendAsync succeeded or was gated (it + // is bumped in OnTick's finally), so waiting for it here -- rather than merely settling + // the actor -- is what proves each tick's send has actually run its course on the + // thread pool before the wire count below is sampled. + // Not armed-checked: the loss also schedules its own Cannot Claim backoff, which can + // legitimately be the actor's nearer timer now, so "the periodic tick is the next + // thing due" is no longer a safe assumption to check for -- only that enough ticks of + // it have gone by. + var countAfterLoss = Volatile.Read(ref count); + var completedBeforeQuiet = owner.PeriodicEmissionsCompleted; + for (var i = 0; i < 3; i++) + { + await clock.AdvanceAsync(period); + await WaitForAnnouncesAsync(() => owner.PeriodicEmissionsCompleted, completedBeforeQuiet + i + 1); + } + // A completed send has handed its frame to the driver; the hub still delivers it on its + // own thread, so the count gets a wall window before it is read (Bugbot on #185). + await Task.Delay(WireWindow); + var countAfterQuiet = Volatile.Read(ref count); // We tolerate at most one already-in-flight emission slipping past the state // transition. Anything more means the loop kept sending under a stale SA. diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index 894b271c..c9479c45 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -537,6 +537,13 @@ async Task DriveAsync(ICanBus bus, byte peerSa, uint pgn, byte[] pdu) // Bugbot 3596183535: BAM announce TX rejection must fail SendBamAsync (not only raise // BackgroundExceptionOccurred) and must not proceed to TP.DT after Th. + // + // #171: "give Th a chance to fire" used to be a wall-clock sleep. The actor is now the + // caller's and its clock is frozen: instead of hoping 80 ms was enough (or, on a live actor, + // long enough that a wrongly-armed spacing timer could have already fired and been trimmed + // before the check ran), the test reads the actor's own timer list once the send has + // faulted, on a clock nothing can advance out from under it -- a rejected announce must have + // armed nothing at all, and a clock that never moves cannot hide a timer that was. [Fact] public async Task Bam_AnnounceTxRejected_FailsSendAndDoesNotEmitDt() { @@ -547,7 +554,9 @@ public async Task Bam_AnnounceTxRejected_FailsSendAndDoesNotEmitDt() using var inner = new CanBusService(busA); using var rejecting = new RejectTpCmBusService(inner); var opts = new J1939TpOptions().With(bamPacketSpacing: TimeSpan.FromMilliseconds(5)); - using var sender = J1939TpFactory.Open(rejecting, sourceAddress: 0x51, options: opts, leaveOpen: true); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var sender = new J1939TpChannel(rejecting, sourceAddress: 0x51, opts, ownsService: false, actor); var dtSeen = 0; busB.FrameObserved += (_, e) => @@ -565,13 +574,23 @@ public async Task Bam_AnnounceTxRejected_FailsSendAndDoesNotEmitDt() .WithTimeout(ShortTimeout); await act.Should().ThrowAsync(); - // Give Th a chance to fire if DT were incorrectly scheduled after a rejected BAM. - await Task.Delay(80); + // The send has already faulted on the actor loop; settle it once more so any work the + // fault path itself posted has also run, then read the timer list directly instead of + // waiting out a spacing period on the wall clock. + await actor.PostAsync(() => 0); + (await actor.NextTimerDelayAsync()).Should().BeNull( + "a rejected BAM announce must leave no timer armed -- not the spacing timer, not anything else"); Volatile.Read(ref dtSeen).Should().Be(0, "rejected BAM announce must not schedule TP.DT"); Volatile.Read(ref bgSeen).Should().Be(0, "CM TX failure must fail the send TCS, not only BackgroundExceptionOccurred"); } // Bugbot 3596025915: canceling before BeginTxOnLoop runs must not emit TP.CM/TP.DT. + // + // #171: SendCmAsync registers the cancellation and posts BeginTxOnLoop to the actor before + // it ever returns, both on this thread and in that order, so a token already canceled means + // the work item is already sitting in the actor's mailbox by the time the send throws. One + // round trip on the same, injected actor is therefore an exact barrier -- not a guess at how + // long draining it might take. [Fact] public async Task SendCm_CanceledBeforeStart_DoesNotTransmit() { @@ -579,7 +598,9 @@ public async Task SendCm_CanceledBeforeStart_DoesNotTransmit() using var busA = Open(session, 0); using var busB = Open(session, 1); - using var sender = J1939TpFactory.Open(busA, sourceAddress: 0x61); + using var actor = new ProtocolActor(); + using var sender = new J1939TpChannel(new CanBusService(busA), sourceAddress: 0x61, + new J1939TpOptions(), ownsService: true, actor); using var _ = J1939TpFactory.Open(busB, sourceAddress: 0x62); var seen = 0; @@ -599,7 +620,12 @@ public async Task SendCm_CanceledBeforeStart_DoesNotTransmit() RandomPayload(50, seed: 61), cts.Token); await act.Should().ThrowAsync(); - // Give the actor a beat to drain any incorrectly queued BeginTx work. + // BeginTxOnLoop was already enqueued behind us; this drains it -- the actor's decision + // not to start a session is now certain. What is not covered by the round trip is the + // wire itself: a start it should not have made still transmits through a fire-and-forget + // Task.Run outside the actor's mailbox (SendControlFrame), so the original 100 ms window + // follows the deterministic half; shortening it would only weaken the negative. + await actor.PostAsync(() => 0); await Task.Delay(100); Volatile.Read(ref seen).Should().Be(0, "canceled send must not emit TP.CM/TP.DT"); } @@ -994,6 +1020,11 @@ public async Task Cm_Sender_PrematureEom_FailsSend() // Bugbot 3596489078: EOM totals that disagree with the session must fail SendCmAsync // (not complete successfully with only a BackgroundExceptionOccurred). + // + // #171: the frame that reaches the peer's wire is not proof OnCmDtConfirmed has run on the + // sender's own actor -- that confirmation comes back through the sender's own TX-confirm + // path, on the sender's own loop, and is what arms WaitEom's T3. Waiting for that timer to + // be armed replaces the 20 ms guess. [Fact] public async Task Cm_Sender_EomSizeMismatch_FailsSend() { @@ -1006,7 +1037,11 @@ public async Task Cm_Sender_EomSizeMismatch_FailsSend() const uint pgn = 0xFE93u; var payload = RandomPayload(14, seed: 93); // 2 packets - using var sender = J1939TpFactory.Open(senderBus, sourceAddress: senderSa); + var options = new J1939TpOptions(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var sender = new J1939TpChannel(new CanBusService(senderBus), sourceAddress: senderSa, + options, ownsService: true, actor); var rtsSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var lastDtSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); @@ -1033,8 +1068,9 @@ public async Task Cm_Sender_EomSizeMismatch_FailsSend() peerBus.Transmit(CanFrame.Classic((int)cmId, cts, isExtendedFrame: true)); await lastDtSeen.Task.AsTaskWithTimeout(ShortTimeout); - // Small settle so OnCmDtConfirmed arms WaitEom before we inject the bad EOM. - await Task.Delay(20); + // OnCmDtConfirmed arms T3 on the move to WaitEom; wait for that arm instead of guessing + // how long the sender's own confirmation takes to come back. + await clock.WaitUntilTimerArmedAsync(actor, options.T3, ShortTimeout); var badEom = J1939TpFrames.BuildEomAck(payload.Length, J1939TpFrames.TotalPackets(payload.Length), pgn); badEom[1] = (byte)(payload.Length + 1); // mismatch totals vs session @@ -1906,6 +1942,11 @@ public async Task Cm_Sender_CtsForPacketZero_AbortsAsBadSequenceNumber() // #31 — the window from the receiver's CTS to the first TP.DT is T2 (1250 ms), not Tr // (200 ms): Tr is the time a node has to *send* a response it owes. A conforming but slow // originator that needs 300 ms to get its first DT out must not be rejected. + // + // #171: 300 ms used to be a wall-clock lower bound -- a loaded host only ever widened the + // gap, up to T2's margin. On the receiver's own injected, virtual clock the gap is exact: + // arm-before-advance (wait for T2 itself to be armed, not just for the CTS to be on the + // wire) and then advance by exactly 300 ms, strictly inside T2's 1250 ms. // ----------------------------------------------------------------------------------- [Fact] public async Task Cm_Receiver_Survives_A_First_Dt_That_Arrives_300ms_After_Cts() @@ -1918,10 +1959,14 @@ public async Task Cm_Receiver_Survives_A_First_Dt_That_Arrives_300ms_After_Cts() const byte peerSa = 0xB2; const uint pgn = 0xFEB1u; var payload = RandomPayload(14, seed: 177); // 2 packets + var gap = TimeSpan.FromMilliseconds(300); - // Defaults: T2 = 1250 ms is the timer under test. The 300 ms below is a lower bound on - // the delay, so a loaded host only widens the gap it must survive, up to T2's margin. - using var receiver = J1939TpFactory.Open(receiverBus, sourceAddress: receiverSa); + // Defaults: T2 = 1250 ms is the timer under test; 300 ms < T2 is the bracket. + var options = new J1939TpOptions(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var receiver = new J1939TpChannel(new CanBusService(receiverBus), sourceAddress: receiverSa, + options, ownsService: true, actor); var ctsSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); peerBus.FrameObserved += (_, e) => @@ -1940,7 +1985,10 @@ public async Task Cm_Receiver_Survives_A_First_Dt_That_Arrives_300ms_After_Cts() peerBus.Transmit(CanFrame.Classic((int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, receiverSa), rts, isExtendedFrame: true)); await ctsSeen.Task.AsTaskWithTimeout(ShortTimeout); - await Task.Delay(300); + // The CTS on the wire is not proof T2 is armed yet -- ArmT2 runs after the CTS is sent, + // on the receiver's own actor. Wait for the arm itself before moving the clock. + await clock.WaitUntilTimerArmedAsync(actor, options.T2, ShortTimeout); + await clock.AdvanceAsync(gap); var dtId = (int)J1939Id.ComposePgn(7, J1939Pgn.TpDt, peerSa, receiverSa); peerBus.Transmit(CanFrame.Classic(dtId, J1939TpFrames.BuildDt(sn: 1, pdu: payload, offset: 0), isExtendedFrame: true)); @@ -2081,6 +2129,11 @@ public async Task A_Cts_With_The_Destination_In_The_Pgns_Low_Byte_Reaches_The_Se // #58: a PDU1 PGN with its low byte set is not a PGN -- the destination is the address // argument -- and is refused before anything goes out, as J1939Id.ComposePgn refuses it // (#55); the session is keyed on what the peer names in its CTS. + // + // #171: the 50 ms "chance for a wrongly-scheduled transmit to show up" was never needed. + // ValidateSendPayload runs synchronously, before either method ever posts to the actor, so + // by the time both awaits below have returned there is no work in flight to wait out -- + // nothing was ever queued that could still transmit. [Fact] public async Task A_Pdu1_Pgn_With_A_Low_Byte_Is_Refused_Before_Anything_Goes_Out() { @@ -2094,8 +2147,8 @@ public async Task A_Pdu1_Pgn_With_A_Low_Byte_Is_Refused_Before_Anything_Goes_Out await cm.Should().ThrowAsync().WithParameterName("pgn"); Func bam = () => sender.SendBamAsync(0xEE8Du, RandomPayload(14, seed: 1)); await bam.Should().ThrowAsync().WithParameterName("pgn"); - await Task.Delay(50); - transmitted.Should().Be(0, "nothing was transmitted"); + transmitted.Should().Be(0, + "validation throws before either call ever posts to the actor, so nothing was ever queued to transmit"); } // #58: an RTS that allows no packet per CTS can never be served -- every CTS would be a @@ -2186,15 +2239,21 @@ public async Task An_End_Of_Message_After_A_Partial_Retransmit_Completes_The_Sen // Codex on #152: a retransmit request that arrives while a block is still draining takes // effect as soon as the outstanding DT is confirmed -- the receiver is missing a packet, // and every later one it gets meanwhile is out of sequence to it -- not after the block. + // + // #171: "the CTS is on the actor before the confirmation is released" was a 50 ms sleep. + // The CTS reaches the actor through the reader task's own post, so the test waits for the + // reader to have handed that frame over, then for one round trip on the injected actor. [Fact] public async Task A_Retransmit_Request_Mid_Block_Takes_Effect_After_The_Outstanding_Packet() { using var bus = ControllableBus.DeferredEchoCapable(NewSession()); - using var service = new CanBusService(bus); + using var service = new FrameConsumptionCountingBusService(new CanBusService(bus)); const byte subjectSa = 0x10, peerSa = 0x20; const uint pgn = 0xFEC6u; var payload = RandomPayload(21, seed: 7); // three packets - using var sender = J1939TpFactory.Open(service, sourceAddress: subjectSa); + using var actor = new ProtocolActor(); + using var sender = new J1939TpChannel(service, sourceAddress: subjectSa, new J1939TpOptions(), + ownsService: false, actor); var dtSns = new List(); bus.OnTransmitting = f => @@ -2212,9 +2271,15 @@ static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) bus.RaiseObserved(PeerCm(peerSa, subjectSa, J1939TpFrames.BuildCts(numPackets: 3, nextPacketSn: 1, dataPgn: pgn)), isEcho: false); await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // DT 1, its confirmation held - // Packet 1 asked for again while DT 1 is still outstanding. - bus.RaiseObserved(PeerCm(peerSa, subjectSa, J1939TpFrames.BuildCts(numPackets: 1, nextPacketSn: 1, dataPgn: pgn)), isEcho: false); - await Task.Delay(50); // the CTS is on the actor before the confirmation is released + // Packet 1 asked for again while DT 1 is still outstanding. The CTS must be on the actor + // before the confirmation is released: the reader hands it over (the counting + // subscription sees the reader ask for the next frame), and a round trip behind that post + // has run it. + var retransmit = J1939TpFrames.BuildCts(numPackets: 1, nextPacketSn: 1, dataPgn: pgn); + var handedOver = service.WaitUntilConsumedAsync(e => e.Frame.Data.ToArray().SequenceEqual(retransmit)); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, retransmit), isEcho: false); + await handedOver.WaitAsync(ShortTimeout); + await actor.PostAsync(() => 0); bus.DeferredEchoes.ReleaseNext(); await bus.DeferredEchoes.WaitForEnqueuedAsync(3, ShortTimeout); // the next DT @@ -2231,6 +2296,14 @@ static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) // Codex on #152: while a block drains, "already sent" reaches only up to the outstanding // packet; a CTS for a later packet of the grant would skip the ones between, and is a // sequence error (table 7, code 7), not a retransmit. + // + // #171: the trailing 50 ms was doing two jobs at once -- giving the one expected abort time + // to actually reach the wire (AbortTx sets the fault and fires the transmit through a + // fire-and-forget Task.Run, outside the actor's mailbox, so the fault and the wire are not + // the same instant) and giving a wrongly-triggered second abort a window to show up after + // the now-orphaned DT 1 confirmation is released. Split into what each half actually needs: + // a wait for the abort that must arrive, then a settle plus the original window for the + // one that must not. [Fact] public async Task A_Cts_For_An_Unsent_Packet_Of_The_Block_Is_A_Sequence_Error_Not_A_Retransmit() { @@ -2239,14 +2312,21 @@ public async Task A_Cts_For_An_Unsent_Packet_Of_The_Block_Is_A_Sequence_Error_No const byte subjectSa = 0x10, peerSa = 0x20; const uint pgn = 0xFEC7u; var payload = RandomPayload(21, seed: 8); // three packets - using var sender = J1939TpFactory.Open(service, sourceAddress: subjectSa); + using var actor = new ProtocolActor(); + using var sender = new J1939TpChannel(service, sourceAddress: subjectSa, new J1939TpOptions(), + ownsService: false, actor); var aborts = new List(); + var firstAbort = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); bus.OnTransmitting = f => { var fields = J1939Id.Decompose((uint)f.ID); if (J1939Pgn.IsTransportCm(fields.Pgn) && f.Data.Span[0] == J1939TpFrames.ControlAbort) - lock (aborts) aborts.Add(f.Data.ToArray()); + { + var frame = f.Data.ToArray(); + lock (aborts) aborts.Add(frame); + firstAbort.TrySetResult(frame); + } }; static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) => CanFrame.Classic((int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, subjectSa), data, isExtendedFrame: true); @@ -2263,7 +2343,13 @@ static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) Func failed = () => send; var ex = await failed.Should().ThrowAsync(); ex.Which.Reason.Should().Be(J1939TpAbortReason.BadSequenceNumber); + await firstAbort.Task.AsTaskWithTimeout(ShortTimeout); // the one expected abort, actually on the wire + bus.DeferredEchoes.ReleaseNext(); + // The stale DT 1 confirmation is now on the actor; settle it, then the original 50 ms + // window covers whatever it might still fire onto the wire through the same + // fire-and-forget path, which the settle cannot see. + await actor.PostAsync(() => 0); await Task.Delay(50); lock (aborts) aborts.Should().ContainSingle().Which[1].Should().Be((byte)J1939TpAbortReason.BadSequenceNumber); } @@ -2678,9 +2764,23 @@ internal sealed class FrameConsumptionCountingBusService : ICanBusService private readonly object _gate = new(); private int _consumed; private readonly List<(int Count, TaskCompletionSource Reached)> _waiters = new(); + private readonly List<(Func Match, TaskCompletionSource Seen)> _matchers = new(); public FrameConsumptionCountingBusService(ICanBusService inner) => _inner = inner; + /// Runs once, right after the next subscription is created and before it is returned. + public Action? OnNextSubscribe { get; set; } + + private ISubscription Subscribed(ISubscription inner) + { + var hook = Interlocked.Exchange(ref _onNextSubscribeTaken, 1) == 0 ? OnNextSubscribe : null; + var counted = new Counted(this, inner); + hook?.Invoke(); + return counted; + } + + private int _onNextSubscribeTaken; + public Task WaitUntilConsumedAsync(int count) { lock (_gate) @@ -2692,13 +2792,26 @@ public Task WaitUntilConsumedAsync(int count) } } - private void Consumed() + /// + /// Completes once a subscriber has finished with a frame matching + /// -- for the J1939-TP reader, once it has handed that frame to the actor. + /// + public Task WaitUntilConsumedAsync(Func match) + { + var seen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + lock (_gate) _matchers.Add((match, seen)); + return seen.Task; + } + + private void Consumed(CanFrameEvent frameEvent) { lock (_gate) { _consumed++; foreach (var (count, reached) in _waiters) if (_consumed >= count) reached.TrySetResult(true); + foreach (var (match, seen) in _matchers) + if (match(frameEvent)) seen.TrySetResult(true); } } @@ -2712,10 +2825,10 @@ public event EventHandler? BackgroundExceptionOccurred } public ISubscription Subscribe(Func? predicate = null, int? bufferCapacity = null, bool includeEcho = false) - => new Counted(this, _inner.Subscribe(predicate, bufferCapacity, includeEcho)); + => Subscribed(_inner.Subscribe(predicate, bufferCapacity, includeEcho)); public ISubscription Subscribe(CanIdFilter filter, int? bufferCapacity = null, bool includeEcho = false) - => new Counted(this, _inner.Subscribe(filter, bufferCapacity, includeEcho)); + => Subscribed(_inner.Subscribe(filter, bufferCapacity, includeEcho)); public IReadOnlyList FindOverlappingFilterSubscriptions() => _inner.FindOverlappingFilterSubscriptions(); @@ -2745,7 +2858,7 @@ private async IAsyncEnumerable Count( await foreach (var frameEvent in _inner.Frames.WithCancellation(cancellationToken)) { yield return frameEvent; - _owner.Consumed(); // runs when the subscriber asks for the next frame + _owner.Consumed(frameEvent); // runs when the subscriber asks for the next frame } } diff --git a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs index 829f96f1..5a8233f7 100644 --- a/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/ProtocolActorTests.cs @@ -6,6 +6,7 @@ using System.Threading; using System.Threading.Tasks; using CanKit.Pro.Actor; +using CanKit.Pro.Tests.Infrastructure; using FluentAssertions; using Xunit; @@ -130,16 +131,18 @@ public async Task Schedule_Fires_Callback_After_The_Configured_Delay() [Fact] public async Task Disposing_The_Schedule_Handle_Before_Due_Prevents_The_Callback_From_Firing() { - using var actor = new ProtocolActor(); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); var fired = false; - var handle = actor.Schedule(TimeSpan.FromMilliseconds(100), () => fired = true); + var timeout = TimeSpan.FromMilliseconds(100); + var handle = actor.Schedule(timeout, () => fired = true); + await clock.WaitUntilTimerArmedAsync(actor, timeout, TimeSpan.FromSeconds(5)); handle.Dispose(); - await Task.Delay(TimeSpan.FromMilliseconds(300)); - // Round-trip through the actor once more so we know the loop has definitely passed the - // point where the (cancelled) timer would have fired. - await actor.PostAsync(() => 0); + // Move the clock well past the (cancelled) timer's due point and let the loop settle. + await clock.AdvanceAsync(timeout + TimeSpan.FromMilliseconds(200)); + await clock.SettleAsync(); fired.Should().BeFalse(); } @@ -191,8 +194,12 @@ public async Task Dispose_Called_Reentrantly_From_A_Posted_Callback_Does_Not_Dea (await Task.WhenAny(completed.Task, Task.Delay(TimeSpan.FromSeconds(5)))).Should().Be(completed.Task, "a self-disposing callback must return promptly instead of deadlocking on its own loop"); - // The loop must still actually finish tearing itself down shortly afterward. - await Task.Delay(TimeSpan.FromMilliseconds(200)); + // No further wait is needed: Dispose flips _disposedFlag inside the same _disposeGate lock + // that Post checks, and it does so as the very first thing it does -- even on the reentrant + // path that returns immediately without joining the loop thread. Since actor.Dispose() is + // called (and therefore has already flipped the flag) strictly before completed.TrySetResult + // runs inside that same synchronous callback, the flag is already set by the time + // completed.Task above resolves. Action postAfter = () => actor.Post(() => { }); postAfter.Should().Throw(); }