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
183 changes: 161 additions & 22 deletions src/CanKit.Pro.IsoTp/IsoTpChannel.cs
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,22 @@ internal sealed class IsoTpChannel : IIsoTpChannel
// PDU or a reassembly-abort fault (N_Cr / SN mismatch / SF·FF supersede) so a blocked
// ReceiveAsync completes instead of hanging (FR-TP-010) —
// the FailTx analogue on the RX side.
//
// The inbox is not completed when the bus service is lost (only by Dispose): that failure is
// kept beside it. A receiver takes what is buffered first and meets the failure only on an
// empty inbox, which is the order the callers rely on, and a discard keeps writing back what
// it retains into an inbox that is still open.
private readonly Channel<RxInboxItem> _pduInbox;
private volatile Exception? _inboxLost;
// Set while a discard has taken the items it retains out of the inbox and not yet put them
// back: a receiver meeting an empty inbox then must not read that as "nothing more to come".
private int _discarding;

/// <summary>For tests: runs between a discard taking the retained items out of the inbox and
/// putting them back, on the actor.</summary>
internal Action? DiscardGapObserver { get; set; }
private readonly TaskCompletionSource<bool> _inboxLostSignal =
new(TaskCreationOptions.RunContinuationsAsynchronously);

// Serializes SendAsync callers: one outbound PDU on the wire at a time, per ISO 15765-2's
// "one N-USData at a time" model. Also avoids competition for _tx state across calls.
Expand Down Expand Up @@ -162,6 +177,12 @@ internal IsoTpChannel(ICanBusService service, IsoTpEndpoint endpoint,
// Encode once: EncodeStMin throws on negative values; surfacing at Open keeps the RX
// path free of codec throws that ProtocolActor would only raise as BackgroundException.
_localStMinRaw = IsoTpFrameCodec.EncodeStMin(_options.LocalStMin);
// The same for the protocol timers: a negative one makes the deadline scheduler throw
// after a send is already on the wire or a reception already started, which leaves it
// with no deadline at all. Zero would fire on the next loop iteration, i.e. never wait.
RequirePositive(nameof(_options.NAs), _options.NAs);
RequirePositive(nameof(_options.NBs), _options.NBs);
RequirePositive(nameof(_options.NCr), _options.NCr);

var inboxOptions = new BoundedChannelOptions(Math.Max(1, _options.ReceiveBufferCapacity))
{
Expand Down Expand Up @@ -296,12 +317,10 @@ public async Task<byte[]> ReceiveAsync(CancellationToken cancellationToken = def
public async Task<IsoTpReceivedPdu> ReceiveWithArrivalAsync(
CancellationToken cancellationToken = default)
{
while (await _pduInbox.Reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false))
{
if (_pduInbox.Reader.TryRead(out var item))
return new IsoTpReceivedPdu(UnwrapInboxItem(item), item.ArrivalTimestamp, item.FirstFrameArrivalTimestamp);
}
throw new InvalidOperationException("Channel is disposed; no more PDUs will arrive.");
var (taken, item) = await TakeNextAsync(cancellationToken).ConfigureAwait(false);
if (!taken)
throw new InvalidOperationException("Channel is disposed; no more PDUs will arrive.");
return new IsoTpReceivedPdu(UnwrapInboxItem(item), item.ArrivalTimestamp, item.FirstFrameArrivalTimestamp);
}

/// <inheritdoc />
Expand Down Expand Up @@ -447,18 +466,28 @@ private int ClearReceptionsBefore(long stamp)
// opposite of the reset intent. Items from after it go back in their order; this runs
// on the inbox's single writer.
var kept = new List<RxInboxItem>();
while (_pduInbox.Reader.TryRead(out var item))
Volatile.Write(ref _discarding, 1);
try
{
// An error item carries the first-frame stamp of the reception it aborted, so an
// abort of a reception that began after the stamp is kept as that reception's
// outcome (Codex on #143).
if (item.FirstFrameArrivalTimestamp >= stamp)
kept.Add(item);
else
discarded++;
while (_pduInbox.Reader.TryRead(out var item))
{
// An error item carries the first-frame stamp of the reception it aborted, so an
// abort of a reception that began after the stamp is kept as that reception's
// outcome (Codex on #143).
if (item.FirstFrameArrivalTimestamp >= stamp)
kept.Add(item);
else
discarded++;
}

DiscardGapObserver?.Invoke();
foreach (var item in kept)
_pduInbox.Writer.TryWrite(item);
Comment thread
dborgards marked this conversation as resolved.
}
finally
{
Volatile.Write(ref _discarding, 0);
}
foreach (var item in kept)
_pduInbox.Writer.TryWrite(item);
return discarded;
}

Expand All @@ -469,11 +498,62 @@ public IAsyncEnumerable<byte[]> ReceiveAllAsync(CancellationToken cancellationTo
private async IAsyncEnumerable<byte[]> ReadAllAsync(
[EnumeratorCancellation] CancellationToken cancellationToken)
{
var reader = _pduInbox.Reader;
while (await reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false))
while (true)
{
while (reader.TryRead(out var item))
yield return UnwrapInboxItem(item);
var (taken, item) = await TakeNextAsync(cancellationToken).ConfigureAwait(false);
if (!taken) yield break;
yield return UnwrapInboxItem(item);
}
}

/// <summary>
/// The next item of the inbox, waiting for one. False once the inbox is completed (the channel
/// was disposed). If the bus service was lost, what is buffered is delivered first and the loss
/// is thrown only on an empty inbox, to every receiver.
/// </summary>
private async Task<(bool Taken, RxInboxItem Item)> TakeNextAsync(CancellationToken cancellationToken)
{
var lossConfirmed = false;
while (true)
{
// Before any read: a canceled token must not consume a buffered PDU, as the wait on the
// token did before this loop existed (and as TryReceiveWithArrival documents).
cancellationToken.ThrowIfCancellationRequested();
if (_pduInbox.Reader.TryRead(out var item)) return (true, item);
Comment thread
dborgards marked this conversation as resolved.
if (_inboxLost is { } lost)
{
// A discard has the retained items out of the inbox for a moment: wait it out.
if (Volatile.Read(ref _discarding) != 0)
{
lossConfirmed = false;
await Task.Yield();
continue;
}

// The empty inbox was read before the flag: a discard that finished in between has
// put its items back. Read once more, after the flag, before it counts.
if (!lossConfirmed)
{
lossConfirmed = true;
continue;
}

throw lost;
}

using var release = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
var wait = _pduInbox.Reader.WaitToReadAsync(release.Token).AsTask();
var first = await Task.WhenAny(wait, _inboxLostSignal.Task).ConfigureAwait(false);
if (first != wait)
{
// The loss came first: let go of the wait and look again, items before the loss.
release.Cancel();
try { await wait.ConfigureAwait(false); }
catch (OperationCanceledException) { /* the wait we just released */ }
continue;
}

if (!await wait.ConfigureAwait(false)) return (false, default);
}
}

Expand Down Expand Up @@ -582,21 +662,71 @@ void FailInFlightSend()
// -----------------------------------------------------------------------------------------
private async Task RunReaderAsync()
{
Exception lost;
try
{
while (await _subscription.WaitToReadAsync(_readerCts.Token).ConfigureAwait(false))
PumpSubscription();

// The subscription ended on its own: the service it came from was disposed under a
// channel opened with leaveOpen.
lost = new ObjectDisposedException(nameof(ICanBusService),
"The bus service was disposed while the ISO-TP channel was open.");
}
catch (OperationCanceledException)
{
// expected on Dispose
return; // expected on Dispose
}
catch (Exception ex)
{
RaiseBackgroundException(ex);
lost = ex;
}

// Nothing will ever arrive again. A receiver waiting on the inbox would wait for ever, so
// the reason is recorded and the receivers are woken: what is buffered is delivered first,
// and then every receiver -- not just one -- gets the failure. Sends already fail on their
// own, the bus refuses them. On the actor, like everything else that touches the inbox.
try
{
// Under the pump lock: a caller that is pumping the subscription right now -- it took
// the last frame and has not posted it yet -- finishes first, so the frame is on the
// actor ahead of the loss and a receiver cannot meet the loss before it.
lock (_pumpGate)
{
_actor.Post(() => EndInboxAfterSubscriptionLoss(lost));
}
}
catch (ObjectDisposedException)
{
// The actor is gone (an injected one its owner disposed): there is nobody to write the
// inbox for. The loss is still reported; the channel's own disposal completes the inbox.
RaiseBackgroundException(lost);
}
}

private void EndInboxAfterSubscriptionLoss(Exception lost)
{
// A reassembly under way dies with the subscription. No fault item for it: in the bounded
// inbox it would push the oldest finished PDU out, and only one receiver would see it.
var rx = _rx;
if (rx is not null)
{
rx.CancelDeadline();
_rx = null;
WithdrawReception(rx.Announce);
}

_inboxLost = lost;
_inboxLostSignal.TrySetResult(true);
RaiseBackgroundException(lost);
}

private static void RequirePositive(string name, TimeSpan value)
{
if (value <= TimeSpan.Zero)
throw new ArgumentOutOfRangeException(name, value, name + " must be greater than zero.");
}

/// <summary>
/// Takes every frame the subscription has buffered and posts it to the actor, in order.
/// Run by the reader task when frames arrive, and by a caller that must see the buffered
Expand Down Expand Up @@ -718,6 +848,15 @@ private void BeginSendOnLoop(byte[] pdu, TaskCompletionSource<IsoTpTransmitStamp
if (IsSendAlreadyCanceled(tcs, ct))
return;

// Disposed between the caller's check and this callback: Dispose has already failed
// whatever was in flight, and found nothing, so a send that began now would put a frame
// on the bus after disposal and never be completed. HandleReceivedFrame has the same guard.
if (Volatile.Read(ref _disposed) != 0)
{
tcs.TrySetException(new ObjectDisposedException(nameof(IsoTpChannel)));
return;
}

if (_tx is not null)
{
// Sender gate already serializes SendAsync, so _tx must be null when we get here.
Expand Down
17 changes: 16 additions & 1 deletion tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,15 @@ public void Deliver(CanFrameView frame, long hostArrivalTimestamp = 0) => _frame
/// </summary>
public bool HoldFrames { get; set; }

/// <summary>When set, the reader task's wait throws it once woken, as a subscription that failed would.</summary>
public Exception? ReaderFault { get; set; }

/// <summary>When set, the subscription ends on the wake: the reader's wait returns false, as when its service is disposed.</summary>
public bool EndSubscriptionOnWake { get; set; }

/// <summary>When set, a <c>TryRead</c> that took a frame holds it until the gate opens: a caller-side pump caught between taking a frame and posting it.</summary>
public ManualResetEventSlim? TryReadGate { get; set; }

/// <summary>Lets the reader task's wait complete once; every later wait stays pending.</summary>
public void WakeReader() => _wake.TrySetResult(true);

Expand Down Expand Up @@ -104,13 +113,19 @@ public bool TryRead(out CanFrameEvent frameEvent)
{
bool ok = _owner._frames.Reader.TryRead(out frameEvent);
if (!ok) _owner._drained.TrySetResult(true);
else _owner.TryReadGate?.Wait();
return ok;
}

public async ValueTask<bool> WaitToReadAsync(CancellationToken cancellationToken = default)
{
if (Interlocked.Exchange(ref _owner._wakesServed, 1) == 0)
return await _owner._wake.Task.WaitAsync(cancellationToken);
{
var woken = await _owner._wake.Task.WaitAsync(cancellationToken);
if (_owner.ReaderFault is { } fault) throw fault;
if (_owner.EndSubscriptionOnWake) return false;
return woken;
}
await Task.Delay(Timeout.Infinite, cancellationToken);
return false;
}
Expand Down
Loading
Loading