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