diff --git a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs index b9871ef..ed28716 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpChannel.cs @@ -895,7 +895,13 @@ internal void RemoveQueued(TxCompletion tcs) /// private void EndTx(TxSessionKey key) { - _txSessions.Remove(key); + // The spacing timer of a BAM holds the session -- and its payload -- until it fires, which + // with a long spacing outlives a cancel by that long. + if (_txSessions.TryGetValue(key, out var ended)) + { + ended.ReleaseSpacingTimer(); + _txSessions.Remove(key); + } if (_disposed != 0 || !_txQueues.TryGetValue(key.DestinationAddress, out var queue)) return; PendingTx? start = null; while (start is null && queue.Count > 0) @@ -935,17 +941,24 @@ private void StartTx(TxSessionKey key, byte[] pdu, TxCompletion tcs, bool isCm) // Do not schedule TP.DT until the BAM announce is TX-confirmed. Otherwise a rejected // BAM can still complete SendBamAsync after the packet spacing once DTs finish. SendTpCm(bam, destinationAddress: J1939TpFrames.GlobalDestinationAddress, session, - onConfirmed: () => OnBamAnnounceConfirmed(key)); + onConfirmed: () => OnBamAnnounceConfirmed(session)); } } - private void OnBamAnnounceConfirmed(TxSessionKey key) + // A BAM chain is driven by timers and send confirmations that outlive a cancel: the caller can + // cancel and send the same PGN again within the packet spacing, and the same key then names a + // different session. Every step therefore carries the session it belongs to and acts only + // while that very instance is the registered one -- the key alone cannot tell them apart. + private bool IsCurrentBamSession(TxSession session) + => _txSessions.TryGetValue(session.Key, out var current) && ReferenceEquals(current, session); + + private void OnBamAnnounceConfirmed(TxSession session) { - if (!_txSessions.TryGetValue(key, out var session) || session.IsCm) return; + if (!IsCurrentBamSession(session)) return; // BAM sender: the packet spacing between BAM and first DT, then between subsequent DTs. session.State = TxStage.SendingDt; session.NextSn = 1; - _actor.Schedule(_options.BamPacketSpacing, () => TrySendNextBamDt(key)); + session.SpacingTimer = _actor.Schedule(_options.BamPacketSpacing, () => TrySendNextBamDt(session)); } private void HandleRxTxSideResponse(byte sa, uint dataPgn, byte[] payload) @@ -966,11 +979,27 @@ private void HandleRxTxSideResponse(byte sa, uint dataPgn, byte[] payload) // Stay in WaitCts until a non-zero CTS resumes the transfer. if (numPackets == 0) { + // Only a sender waiting for a CTS has anything to hold. After the last DT (T3 + // for the EndOfMsgAck) or mid-block (T3 follows the block) the running timer + // stays: swapping it for a T4 whose expiry is ignored in those states left the + // send without any timer at all. + if (session.State != TxStage.WaitCts) + { + // A fast peer can hold right after the last DT of a block, before our + // confirmation of that DT has run -- the same early-response race a + // non-zero CTS is stashed for. The hold is remembered and takes effect when + // the block ends: T4 then, not T3. After the message's last packet the + // confirmation waits for the EndOfMsgAck under T3 and never reads it. + if (session.State == TxStage.SendingDt && session.BlockRemaining == 1) + session.HoldPending = true; + return; + } session.Deadline?.Dispose(); session.Deadline = _deadlines.Arm(_options.T4, () => OnTxT4Expired(key)); return; } + session.HoldPending = false; // a CTS that grants packets ends any hold // What has gone out: the highest packet ever confirmed (HighestSentSn -- an int, so a // 255-packet message's last packet counts, where the byte NextSn wraps to 0), and, // while a block drains, the one outstanding (NextSn, unconfirmed). A CTS for a @@ -1083,19 +1112,19 @@ private void HandleRxTxSideResponse(byte sa, uint dataPgn, byte[] payload) } } - private void TrySendNextBamDt(TxSessionKey key) + private void TrySendNextBamDt(TxSession session) { - if (!_txSessions.TryGetValue(key, out var session) || session.IsCm) return; + if (!IsCurrentBamSession(session)) return; int offset = (session.NextSn - 1) * J1939TpFrames.DtDataBytes; byte sn = session.NextSn; var dt = J1939TpFrames.BuildDt(sn, session.Pdu, offset); SendControlFrame(J1939Pgn.TpDt, dt, destinationAddress: J1939TpFrames.GlobalDestinationAddress, - session, onConfirmed: () => OnBamDtConfirmed(key, sn)); + session, onConfirmed: () => OnBamDtConfirmed(session, sn)); } - private void OnBamDtConfirmed(TxSessionKey key, byte confirmedSn) + private void OnBamDtConfirmed(TxSession session, byte confirmedSn) { - if (!_txSessions.TryGetValue(key, out var session) || session.IsCm) return; + if (!IsCurrentBamSession(session)) return; if (session.NextSn != confirmedSn) return; // stale confirmation // Compute the next SN as an int first: for a maximum-length PDU (TotalPackets=255) the @@ -1106,13 +1135,13 @@ private void OnBamDtConfirmed(TxSessionKey key, byte confirmedSn) { // BAM has no ack -- complete once every DT has been transmitted. session.Tcs.TrySetResult(null); - EndTx(key); + EndTx(session.Key); return; } session.NextSn = (byte)nextSn; // The spacing between two consecutive BAM DTs (J1939-21 §5.10.3, 50..200 ms). - _actor.Schedule(_options.BamPacketSpacing, () => TrySendNextBamDt(key)); + session.SpacingTimer = _actor.Schedule(_options.BamPacketSpacing, () => TrySendNextBamDt(session)); } private void TrySendNextCmDt(TxSessionKey key) @@ -1173,7 +1202,15 @@ private void OnCmDtConfirmed(TxSessionKey key, byte confirmedSn) session.State = TxStage.WaitCts; session.Deadline?.Dispose(); - session.Deadline = _deadlines.Arm(_options.T3, () => OnTxT3Expired(key)); + if (session.HoldPending) + { + session.HoldPending = false; + session.Deadline = _deadlines.Arm(_options.T4, () => OnTxT4Expired(key)); + } + else + { + session.Deadline = _deadlines.Arm(_options.T3, () => OnTxT3Expired(key)); + } return; } @@ -1503,12 +1540,16 @@ public TxSession(TxSessionKey key, byte[] pdu, int totalPackets, TxCompletion tc public byte NextSn { get; set; } public int BlockRemaining { get; set; } public IDeadline? Deadline { get; set; } + /// The armed packet-spacing timer of a BAM; released when the session ends. + public IDisposable? SpacingTimer { get; set; } /// /// Set when a CTS for the next block arrives before the last DT of the current block has /// been confirmed (Virtual-loopback race). Applied in . /// public bool HasPendingCts { get; set; } + /// A CTS(0) hold arrived before the last DT of the block was confirmed: T4 starts when it is. + public bool HoldPending { get; set; } public byte PendingCtsNumPackets { get; set; } /// The stashed CTS asks for a packet already sent: applied as soon as the outstanding DT is confirmed. public bool PendingCtsIsRetransmit { get; set; } @@ -1532,6 +1573,12 @@ public TxSession(TxSessionKey key, byte[] pdu, int totalPackets, TxCompletion tc public bool LastDtQueued { get; set; } public byte PendingCtsNextSn { get; set; } + public void ReleaseSpacingTimer() + { + SpacingTimer?.Dispose(); + SpacingTimer = null; + } + public void Cancel() { Deadline?.Dispose(); @@ -1543,6 +1590,7 @@ public void Fail(Exception ex) { Deadline?.Dispose(); Deadline = null; + ReleaseSpacingTimer(); Tcs.TrySetException(ex); } } diff --git a/src/CanKit.Pro.J1939Tp/J1939TpOptions.cs b/src/CanKit.Pro.J1939Tp/J1939TpOptions.cs index d2643fa..022e71e 100644 --- a/src/CanKit.Pro.J1939Tp/J1939TpOptions.cs +++ b/src/CanKit.Pro.J1939Tp/J1939TpOptions.cs @@ -152,6 +152,15 @@ public J1939TpOptions With( /// internal void Validate() { + // A negative timer makes Arm throw after the session is registered and the RTS is on the + // wire, which leaves a send with no deadline and its destination blocked; reject it here. + RequirePositive(nameof(T1), T1); + RequirePositive(nameof(T2), T2); + RequirePositive(nameof(T3), T3); + RequirePositive(nameof(T4), T4); + if (BamPacketSpacing < TimeSpan.Zero) + throw new ArgumentOutOfRangeException(nameof(BamPacketSpacing), BamPacketSpacing, + "BamPacketSpacing must not be negative."); if (MaxPacketsPerCts == 0) throw new ArgumentOutOfRangeException(nameof(MaxPacketsPerCts), MaxPacketsPerCts, "MaxPacketsPerCts must be in [1, 255]; 0 is not a valid CTS grant size."); @@ -168,4 +177,10 @@ internal void Validate() throw new ArgumentOutOfRangeException(nameof(MaxQueuedSendsPerDestination), MaxQueuedSendsPerDestination, "MaxQueuedSendsPerDestination must be >= 0 (0 admits no waiting send)."); } + + private static void RequirePositive(string name, TimeSpan value) + { + if (value <= TimeSpan.Zero) + throw new ArgumentOutOfRangeException(name, value, name + " must be greater than zero."); + } } diff --git a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs index 235fb01..c21fdbc 100644 --- a/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs @@ -1462,6 +1462,190 @@ public async Task Cm_Sender_T4Timeout_WhenPeerHoldsWithCtsZero() ex.Message.Should().Contain("T4"); } + // T4 is the hold timer after a CTS with numPackets = 0, and it only means something while the + // sender waits for a CTS. A hold that arrives after the last DT must leave the running T3 for + // the EndOfMsgAck alone: swapping it for a T4 whose expiry is ignored in that state left the + // send without any timer, and every later send to that destination queued behind it. + [Fact] + public async Task A_Cts_Hold_After_The_Last_Dt_Leaves_The_EndOfMsgAck_Timer_Running() + { + var session = NewSession(); + using var senderBus = Open(session, 0); + using var peerBus = Open(session, 1); + + const byte senderSa = 0x05; + const byte peerSa = 0x06; + const uint pgn = 0xFF05u; + const uint probePgn = 0xFF06u; + var payload = RandomPayload(21, seed: 41); // 3 packets + + // Distinct values, so the armed timer names which one is running. + var options = new J1939TpOptions().With( + t2: TimeSpan.FromSeconds(5), + t3: TimeSpan.FromSeconds(2), + t4: TimeSpan.FromMilliseconds(1050)); + 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); + var probeCtsSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + peerBus.FrameObserved += (_, e) => + { + if (!e.CanFrame.IsExtendedFrame) return; + var fields = J1939Id.Decompose((uint)e.CanFrame.ID); + if (fields.SourceAddress != senderSa) return; + var data = e.CanFrame.Data.Span; + if (J1939Pgn.IsTransportCm(fields.Pgn) && data.Length >= 8) + { + if (data[0] == J1939TpFrames.ControlRts && J1939TpFrames.ReadDataPgn(data) == pgn) + rtsSeen.TrySetResult(null); + if (data[0] == J1939TpFrames.ControlCts && J1939TpFrames.ReadDataPgn(data) == probePgn) + probeCtsSeen.TrySetResult(null); + } + else if (fields.Pgn == J1939Pgn.TpDt && data.Length >= 1 && data[0] == 3) + { + lastDtSeen.TrySetResult(null); + } + }; + + var send = sender.SendCmAsync(pgn, destinationAddress: peerSa, payload); + await rtsSeen.Task.AsTaskWithTimeout(ShortTimeout); + + var cmId = (int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, senderSa); + peerBus.Transmit(CanFrame.Classic(cmId, J1939TpFrames.BuildCts(numPackets: 3, nextPacketSn: 1, dataPgn: pgn), + isExtendedFrame: true)); + await lastDtSeen.Task.AsTaskWithTimeout(ShortTimeout); + + // The last DT on the wire is not proof the timer for the EndOfMsgAck is armed yet. + await clock.WaitUntilTimerArmedAsync(actor, options.T3, ShortTimeout); + + var hold = new byte[8]; + hold[0] = J1939TpFrames.ControlCts; + hold[1] = 0x00; + hold[2] = 0x04; + hold[3] = 0xFF; + hold[4] = 0xFF; + hold[5] = (byte)(pgn & 0xFF); + hold[6] = (byte)((pgn >> 8) & 0xFF); + hold[7] = (byte)((pgn >> 16) & 0xFF); + peerBus.Transmit(CanFrame.Classic(cmId, hold, isExtendedFrame: true)); + + // The channel reads its frames in order: its answer to this RTS is proof the hold above + // has been handled, without any wall time standing in for it. + var probe = J1939TpFrames.BuildRts(totalBytes: 14, totalPackets: 2, maxPacketsPerCts: 0xFF, dataPgn: probePgn); + peerBus.Transmit(CanFrame.Classic(cmId, probe, isExtendedFrame: true)); + await probeCtsSeen.Task.AsTaskWithTimeout(ShortTimeout); + + await clock.AdvanceAsync(options.T3); + + Func act = async () => await send.WithTimeout(ShortTimeout); + var ex = (await act.Should().ThrowAsync()).Which; + ex.Reason.Should().Be(J1939TpAbortReason.Timeout); + ex.Message.Should().Contain("EndOfMsgAck"); + } + + // A hold can arrive right after the last DT of a block, before this side has seen the + // confirmation of that DT -- a fast peer on a loopback bus. It is the same early response a + // non-zero CTS is stashed for, so the hold takes effect when the block ends: T4, not T3. + [Fact] + public async Task A_Cts_Hold_Before_The_Block_End_Confirmation_Starts_T4_Not_T3() + { + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + using var service = new FrameConsumptionCountingBusService(new CanBusService(bus)); + const byte subjectSa = 0x10, peerSa = 0x20; + const uint pgn = 0xFEC7u; + var payload = RandomPayload(21, seed: 8); // three packets, granted one at a time + var options = new J1939TpOptions().With( + t2: TimeSpan.FromSeconds(5), t3: TimeSpan.FromSeconds(2), t4: TimeSpan.FromMilliseconds(1050)); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var sender = new J1939TpChannel(service, sourceAddress: subjectSa, options, + ownsService: false, actor); + static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) + => CanFrame.Classic((int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, subjectSa), data, isExtendedFrame: true); + + var send = sender.SendCmAsync(pgn, peerSa, payload); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); // the RTS + bus.DeferredEchoes.ReleaseNext(); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, + J1939TpFrames.BuildCts(numPackets: 1, nextPacketSn: 1, dataPgn: pgn)), isEcho: false); + await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // DT 1, its confirmation held + + var hold = new byte[8]; + hold[0] = J1939TpFrames.ControlCts; + hold[2] = 0x02; + hold[3] = hold[4] = 0xFF; + hold[5] = (byte)(pgn & 0xFF); + hold[6] = (byte)((pgn >> 8) & 0xFF); + hold[7] = (byte)((pgn >> 16) & 0xFF); + var handedOver = service.WaitUntilConsumedAsync(e => e.Frame.Data.ToArray().SequenceEqual(hold)); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, hold), isEcho: false); + await handedOver.WaitAsync(ShortTimeout); + await actor.PostAsync(() => 0); + + bus.DeferredEchoes.ReleaseNext(); // the block ends now + await clock.WaitUntilTimerArmedAsync(actor, options.T4, ShortTimeout); + await clock.AdvanceAsync(options.T4); + + Func act = async () => await send.WithTimeout(ShortTimeout); + var ex = (await act.Should().ThrowAsync()).Which; + ex.Message.Should().Contain("T4"); + } + + // The same hold in the middle of a block has nothing to hold yet: the block runs on, and the + // timer after it is T3 as usual. + [Fact] + public async Task A_Cts_Hold_In_The_Middle_Of_A_Block_Is_Ignored() + { + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + using var service = new FrameConsumptionCountingBusService(new CanBusService(bus)); + const byte subjectSa = 0x10, peerSa = 0x20; + const uint pgn = 0xFEC5u; + var payload = RandomPayload(21, seed: 9); // three packets, granted in one block + var options = new J1939TpOptions().With( + t2: TimeSpan.FromSeconds(5), t3: TimeSpan.FromSeconds(2), t4: TimeSpan.FromMilliseconds(1050)); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var sender = new J1939TpChannel(service, sourceAddress: subjectSa, options, + ownsService: false, actor); + static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) + => CanFrame.Classic((int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, subjectSa), data, isExtendedFrame: true); + + var send = sender.SendCmAsync(pgn, peerSa, payload); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); + 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 + + var hold = new byte[8]; + hold[0] = J1939TpFrames.ControlCts; + hold[2] = 0x02; + hold[3] = hold[4] = 0xFF; + hold[5] = (byte)(pgn & 0xFF); + hold[6] = (byte)((pgn >> 8) & 0xFF); + hold[7] = (byte)((pgn >> 16) & 0xFF); + var handedOver = service.WaitUntilConsumedAsync(e => e.Frame.Data.ToArray().SequenceEqual(hold)); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, hold), isEcho: false); + await handedOver.WaitAsync(ShortTimeout); + await actor.PostAsync(() => 0); + + bus.DeferredEchoes.ReleaseNext(); // DT 1 + await bus.DeferredEchoes.WaitForEnqueuedAsync(3, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); // DT 2 + await bus.DeferredEchoes.WaitForEnqueuedAsync(4, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); // DT 3: the last packet, so T3 for the EndOfMsgAck + await clock.WaitUntilTimerArmedAsync(actor, options.T3, ShortTimeout); + await clock.AdvanceAsync(options.T3); + + Func act = async () => await send.WithTimeout(ShortTimeout); + var ex = (await act.Should().ThrowAsync()).Which; + ex.Message.Should().Contain("EndOfMsgAck"); + } + // FR-TP-031: the receiver must cap its CTS grant at the originator's RTS-advertised // maximum (here 2), even though its own MaxPacketsPerCts (16) is larger — and keep the // cap on every subsequent block's CTS. @@ -3099,6 +3283,320 @@ public void Options_Default_The_Queue_Limit_To_Eight_And_Reject_A_Negative_One() new J1939TpOptions().With(maxQueuedSendsPerDestination: 0).MaxQueuedSendsPerDestination.Should().Be(0); } + // The BAM chain is driven by timers and confirmations that outlive a cancel. Cancelling and + // sending the same PGN again inside the packet spacing puts a second session under the same + // key; nothing of the first one -- its timer, or a confirmation still in flight -- may drive + // the new session. Two lines of defence hold that (the timer is released with the session, and + // every step checks it belongs to the registered instance); the test fails only without both. + [Fact] + public async Task A_Cancelled_Bam_Does_Not_Drive_The_Next_Bam_Of_The_Same_Pgn() + { + const uint pgn = 0xFECBu; + var payload = RandomPayload(21, seed: 248); // 3 TP.DT + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + var dtSequence = new List(); + bus.OnTransmitting = frame => + { + var fields = J1939Id.Decompose((uint)frame.ID); + if (fields.Pgn != J1939Pgn.TpDt) return; + lock (dtSequence) dtSequence.Add(frame.Data.Span[0]); + }; + var background = new List(); + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + sender.BackgroundExceptionOccurred += (_, ex) => { lock (background) background.Add(ex); }; + + using var cts = new CancellationTokenSource(); + var first = sender.SendBamAsync(pgn, payload, cts.Token); + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + + cts.Cancel(); + Func cancelled = async () => await first.WithTimeout(ShortTimeout); + await cancelled.Should().ThrowAsync(); + + var second = sender.SendBamAsync(pgn, payload); + // The second announce is out; its spacing timer is armed once that is confirmed. + await WaitForTransmitCount(bus, 2); + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + + await DrainAsync(clock, actor, second); + + lock (dtSequence) + dtSequence.Should().Equal(new byte[] { 1, 2, 3 }, + "the stale timer of the cancelled BAM must not send a TP.DT for the new one"); + lock (background) background.Should().BeEmpty(); + } + + // The spacing timer of a BAM holds its session and payload until it fires. Cancelling must + // release it: with a long spacing, every cancelled BAM would otherwise stay reachable for + // that long. + [Fact] + public async Task Cancelling_A_Bam_Releases_Its_Spacing_Timer() + { + var spacing = TimeSpan.FromMinutes(10); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = new J1939TpChannel(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: spacing), ownsService: false, actor); + + using var cts = new CancellationTokenSource(); + var send = sender.SendBamAsync(0xFECBu, RandomPayload(21, seed: 249), cts.Token); + await clock.WaitUntilTimerArmedAsync(actor, spacing, ShortTimeout); + + cts.Cancel(); + Func cancelled = async () => await send.WithTimeout(ShortTimeout); + await cancelled.Should().ThrowAsync(); + + (await actor.NextTimerDelayAsync()).Should().BeNull( + "the cancelled BAM's spacing timer is released, not left to fire in ten minutes"); + } + + // The other half of the same race: the announce of a cancelled BAM is still unconfirmed when + // the same PGN is sent again, and its confirmation arrives after the new session is + // registered under the same key. It belongs to the old session and must not start the new + // session's chain a second time. + [Fact] + public async Task A_Late_Announce_Confirmation_Of_A_Cancelled_Bam_Does_Not_Start_The_Next_Bam() + { + const uint pgn = 0xFECAu; + var payload = RandomPayload(21, seed: 250); // 3 TP.DT + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + var dtSequence = new List(); + bus.OnTransmitting = frame => + { + if (J1939Id.Decompose((uint)frame.ID).Pgn != J1939Pgn.TpDt) return; + lock (dtSequence) dtSequence.Add(frame.Data.Span[0]); + }; + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + + using var cts = new CancellationTokenSource(); + var first = sender.SendBamAsync(pgn, payload, cts.Token); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); // its announce, confirmation held + cts.Cancel(); + Func cancelled = async () => await first.WithTimeout(ShortTimeout); + await cancelled.Should().ThrowAsync(); + + var second = sender.SendBamAsync(pgn, payload); + await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // the new announce + + bus.DeferredEchoes.ReleaseNext(); // the cancelled BAM's announce + await actor.PostAsync(() => 0); + bus.DeferredEchoes.ReleaseNext(); // the new one + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + + // Every TP.DT confirmation is held too; release them as they appear. + for (int sent = 0; sent < 3; sent++) + { + await clock.AdvanceAsync(InFlightSpacing); + await bus.DeferredEchoes.WaitForEnqueuedAsync(3 + sent, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); + if (sent < 2) await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + } + await second.WaitAsync(ShortTimeout); + + lock (dtSequence) + dtSequence.Should().Equal(new byte[] { 1, 2, 3 }); + } + + // A hold while the last packet of the message is outstanding has nothing to hold: the + // EndOfMsgAck is what the sender waits for next, under T3. + [Fact] + public async Task A_Cts_Hold_While_The_Last_Packet_Is_Outstanding_Leaves_T3_For_The_EndOfMsgAck() + { + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + using var service = new FrameConsumptionCountingBusService(new CanBusService(bus)); + const byte subjectSa = 0x10, peerSa = 0x20; + const uint pgn = 0xFEC4u; + var payload = RandomPayload(14, seed: 11); // two packets, one block + var options = new J1939TpOptions().With( + t2: TimeSpan.FromSeconds(5), t3: TimeSpan.FromSeconds(2), t4: TimeSpan.FromMilliseconds(1050)); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var sender = new J1939TpChannel(service, sourceAddress: subjectSa, options, + ownsService: false, actor); + static CanFrame PeerCm(byte peerSa, byte subjectSa, byte[] data) + => CanFrame.Classic((int)J1939Id.ComposePgn(7, J1939Pgn.TpCm, peerSa, subjectSa), data, isExtendedFrame: true); + + var send = sender.SendCmAsync(pgn, peerSa, payload); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, + J1939TpFrames.BuildCts(numPackets: 2, nextPacketSn: 1, dataPgn: pgn)), isEcho: false); + await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // DT 1 + bus.DeferredEchoes.ReleaseNext(); + await bus.DeferredEchoes.WaitForEnqueuedAsync(3, ShortTimeout); // DT 2, the last, held + + var hold = new byte[8]; + hold[0] = J1939TpFrames.ControlCts; + hold[2] = 0x03; + hold[3] = hold[4] = 0xFF; + hold[5] = (byte)(pgn & 0xFF); + hold[6] = (byte)((pgn >> 8) & 0xFF); + hold[7] = (byte)((pgn >> 16) & 0xFF); + var handedOver = service.WaitUntilConsumedAsync(e => e.Frame.Data.ToArray().SequenceEqual(hold)); + bus.RaiseObserved(PeerCm(peerSa, subjectSa, hold), isEcho: false); + await handedOver.WaitAsync(ShortTimeout); + await actor.PostAsync(() => 0); + + bus.DeferredEchoes.ReleaseNext(); + await clock.WaitUntilTimerArmedAsync(actor, options.T3, ShortTimeout); + await clock.AdvanceAsync(options.T3); + + Func act = async () => await send.WithTimeout(ShortTimeout); + var ex = (await act.Should().ThrowAsync()).Which; + ex.Message.Should().Contain("EndOfMsgAck"); + } + + // The same release when the whole channel goes: a BAM in flight on a borrowed actor must not + // leave its spacing timer -- and so its session and payload -- behind. + [Fact] + public async Task Disposing_The_Channel_Releases_The_Spacing_Timer_Of_A_Bam_In_Flight() + { + var spacing = TimeSpan.FromMinutes(10); + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.EchoCapable(NewSession()); + using var service = new CanBusService(bus); + using var sender = new J1939TpChannel(service, sourceAddress: 0x10, + new J1939TpOptions().With(bamPacketSpacing: spacing), ownsService: false, actor); + + var send = sender.SendBamAsync(0xFECBu, RandomPayload(21, seed: 251)); + await clock.WaitUntilTimerArmedAsync(actor, spacing, ShortTimeout); + + sender.Dispose(); + Func failed = async () => await send.WithTimeout(ShortTimeout); + await failed.Should().ThrowAsync(); + + (await actor.NextTimerDelayAsync()).Should().BeNull( + "the spacing timer is released with the session, not left to fire in ten minutes"); + } + + // And the TP.DT confirmation: a DT of a cancelled BAM can still be in flight when the same PGN + // is sent again, and its confirmation, arriving afterwards, belongs to the old session. The + // sequence-number check in the handler also holds this one, so the test pins the outcome, not + // the instance check alone (the announce test above does that). + [Fact] + public async Task A_Late_Dt_Confirmation_Of_A_Cancelled_Bam_Does_Not_Advance_The_Next_Bam() + { + const uint pgn = 0xFEC9u; + var payload = RandomPayload(21, seed: 252); // 3 TP.DT + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + var dtSequence = new List(); + bus.OnTransmitting = frame => + { + if (J1939Id.Decompose((uint)frame.ID).Pgn != J1939Pgn.TpDt) return; + lock (dtSequence) dtSequence.Add(frame.Data.Span[0]); + }; + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + + using var cts = new CancellationTokenSource(); + var first = sender.SendBamAsync(pgn, payload, cts.Token); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); // announce + bus.DeferredEchoes.ReleaseNext(); + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + await clock.AdvanceAsync(InFlightSpacing); + await bus.DeferredEchoes.WaitForEnqueuedAsync(2, ShortTimeout); // DT 1, its confirmation held + + cts.Cancel(); + Func cancelled = async () => await first.WithTimeout(ShortTimeout); + await cancelled.Should().ThrowAsync(); + var second = sender.SendBamAsync(pgn, payload); + await bus.DeferredEchoes.WaitForEnqueuedAsync(3, ShortTimeout); // the new announce + + bus.DeferredEchoes.ReleaseNext(); // the cancelled BAM's DT 1 + await actor.PostAsync(() => 0); + bus.DeferredEchoes.ReleaseNext(); // the new announce + for (int sent = 0; sent < 3; sent++) + { + await clock.WaitUntilTimerArmedAsync(actor, InFlightSpacing, ShortTimeout); + await clock.AdvanceAsync(InFlightSpacing); + await bus.DeferredEchoes.WaitForEnqueuedAsync(4 + sent, ShortTimeout); + bus.DeferredEchoes.ReleaseNext(); + } + await second.WaitAsync(ShortTimeout); + + lock (dtSequence) + dtSequence.Should().Equal(new byte[] { 1, 1, 2, 3 }, + "the old session's DT 1, then the new session's own three"); + } + + // The same late confirmation with no successor: nothing is registered under the key any more, + // and the cancelled BAM's announce confirmation must not start a chain for it. + [Fact] + public async Task A_Late_Announce_Confirmation_Of_A_Cancelled_Bam_With_No_Successor_Sends_Nothing() + { + using var clock = new VirtualClock(); + var actor = clock.NewActor(); + using var bus = ControllableBus.DeferredEchoCapable(NewSession()); + var dtCount = 0; + bus.OnTransmitting = frame => + { + if (J1939Id.Decompose((uint)frame.ID).Pgn == J1939Pgn.TpDt) Interlocked.Increment(ref dtCount); + }; + using var service = new CanBusService(bus); + using var sender = BamSenderOn(actor, service); + + using var cts = new CancellationTokenSource(); + var send = sender.SendBamAsync(0xFECAu, RandomPayload(21, seed: 253), cts.Token); + await bus.DeferredEchoes.WaitForEnqueuedAsync(1, ShortTimeout); // its announce, confirmation held + cts.Cancel(); + Func cancelled = async () => await send.WithTimeout(ShortTimeout); + await cancelled.Should().ThrowAsync(); + + bus.DeferredEchoes.ReleaseNext(); + await actor.PostAsync(() => 0); + await clock.AdvanceAsync(InFlightSpacing); + + (await actor.NextTimerDelayAsync()).Should().BeNull("no chain was started for the cancelled BAM"); + Volatile.Read(ref dtCount).Should().Be(0); + } + + [Theory] + [InlineData(nameof(J1939TpOptions.T1), 0)] + [InlineData(nameof(J1939TpOptions.T1), -1)] + [InlineData(nameof(J1939TpOptions.T2), 0)] + [InlineData(nameof(J1939TpOptions.T2), -1)] + [InlineData(nameof(J1939TpOptions.T3), 0)] + [InlineData(nameof(J1939TpOptions.T3), -1)] + [InlineData(nameof(J1939TpOptions.T4), 0)] + [InlineData(nameof(J1939TpOptions.T4), -1)] + [InlineData(nameof(J1939TpOptions.BamPacketSpacing), -1)] + public void A_Timer_Option_That_Cannot_Arm_Is_Rejected_At_Open(string option, int milliseconds) + { + var value = TimeSpan.FromMilliseconds(milliseconds); + var options = option switch + { + nameof(J1939TpOptions.T1) => new J1939TpOptions { T1 = value }, + nameof(J1939TpOptions.T2) => new J1939TpOptions { T2 = value }, + nameof(J1939TpOptions.T3) => new J1939TpOptions { T3 = value }, + nameof(J1939TpOptions.T4) => new J1939TpOptions { T4 = value }, + _ => new J1939TpOptions { BamPacketSpacing = value }, + }; + + using var bus = Open(NewSession(), 0); + Action open = () => J1939TpFactory.Open(bus, sourceAddress: 0x85, options: options); + open.Should().Throw().Which.ParamName.Should().Be(option); + } + + [Fact] + public void A_Zero_Bam_Packet_Spacing_Is_Still_Accepted() + { + using var bus = Open(NewSession(), 0); + using var channel = J1939TpFactory.Open(bus, sourceAddress: 0x86, + options: new J1939TpOptions { BamPacketSpacing = TimeSpan.Zero }); + channel.Should().NotBeNull(); + } + private static J1939TpChannel BamSenderOn(ProtocolActor actor, ICanBusService service) => new(service, sourceAddress: 0x10, new J1939TpOptions().With(bamPacketSpacing: InFlightSpacing), ownsService: false, actor);