Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 61 additions & 13 deletions src/CanKit.Pro.J1939Tp/J1939TpChannel.cs
Original file line number Diff line number Diff line change
Expand Up @@ -895,7 +895,13 @@ internal void RemoveQueued(TxCompletion tcs)
/// </summary>
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)
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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; }
/// <summary>The armed packet-spacing timer of a BAM; released when the session ends.</summary>
public IDisposable? SpacingTimer { get; set; }

/// <summary>
/// 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 <see cref="OnCmDtConfirmed"/>.
/// </summary>
public bool HasPendingCts { get; set; }
/// <summary>A CTS(0) hold arrived before the last DT of the block was confirmed: T4 starts when it is.</summary>
public bool HoldPending { get; set; }
public byte PendingCtsNumPackets { get; set; }
/// <summary>The stashed CTS asks for a packet already sent: applied as soon as the outstanding DT is confirmed.</summary>
public bool PendingCtsIsRetransmit { get; set; }
Expand All @@ -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();
Expand All @@ -1543,6 +1590,7 @@ public void Fail(Exception ex)
{
Deadline?.Dispose();
Deadline = null;
ReleaseSpacingTimer();
Tcs.TrySetException(ex);
}
}
Expand Down
15 changes: 15 additions & 0 deletions src/CanKit.Pro.J1939Tp/J1939TpOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,15 @@ public J1939TpOptions With(
/// </summary>
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.");
Expand All @@ -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.");
}
}
Loading
Loading