Skip to content
20 changes: 19 additions & 1 deletion src/CanKit.Pro.Actor/ProtocolActor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -345,7 +345,25 @@ public IDisposable Schedule(TimeSpan delay, Action callback)
if (callback is null) throw new ArgumentNullException(nameof(callback));
if (delay < TimeSpan.Zero) throw new ArgumentOutOfRangeException(nameof(delay), "Delay must not be negative.");

var entry = new TimerEntry(DueTimestamp(delay), callback);
return Insert(new TimerEntry(DueTimestamp(delay), callback));
}

/// <summary>
/// As <see cref="Schedule"/>, due at <paramref name="dueTimestamp"/> on
/// <see cref="TimeSource"/> rather than a delay from the reading this call takes. For a
/// caller whose deadline is fixed already: a delay computed from its own earlier reading
/// lands late by however far the clock moved in between, and on a clock a test moves in
/// steps that can be a whole step -- a timer that then never fires (#171). A due instant
/// already past fires on the loop's next pass.
/// </summary>
internal IDisposable ScheduleAt(long dueTimestamp, Action callback)
{
if (callback is null) throw new ArgumentNullException(nameof(callback));
return Insert(new TimerEntry(dueTimestamp, callback));
}

private IDisposable Insert(TimerEntry entry)
{
lock (_disposeGate)
{
ThrowIfDisposed();
Expand Down
19 changes: 19 additions & 0 deletions src/CanKit.Pro.Uds/UdsClient.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System;
using CanKit.Pro.Actor;
using CanKit.Pro.IsoTp;

namespace CanKit.Pro.Uds;
Expand Down Expand Up @@ -34,4 +35,22 @@ public static IUdsClient Create(
if (channel is null) throw new ArgumentNullException(nameof(channel));
return new UdsClientImpl(channel, options ?? new UdsClientOptions(), ownsChannel: !leaveOpen);
}

/// <summary>
/// As <see cref="Create(IIsoTpChannel, UdsClientOptions?, bool)"/>, measuring and waiting
/// out P2, P2* and the suppressed-response windows on <paramref name="clock"/>; null is the
/// wall clock, as
/// the public overload uses. <paramref name="channel"/> must have been opened on that same
/// actor, and its demux must stamp frames with its time source, or a deadline and an
/// arrival are not comparable (#171).
/// </summary>
internal static IUdsClient Create(
IIsoTpChannel channel,
ProtocolActor? clock,
UdsClientOptions? options = null,
bool leaveOpen = true)
{
if (channel is null) throw new ArgumentNullException(nameof(channel));
return new UdsClientImpl(channel, options ?? new UdsClientOptions(), ownsChannel: !leaveOpen, clock);
}
}
135 changes: 115 additions & 20 deletions src/CanKit.Pro.Uds/UdsClientImpl.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using CanKit.Pro.Actor;
using CanKit.Pro.IsoTp;

namespace CanKit.Pro.Uds;
Expand Down Expand Up @@ -53,6 +54,15 @@ internal sealed class UdsClientImpl : IUdsClient
private readonly SemaphoreSlim _requestLock = new(1, 1);
private readonly CancellationTokenSource _lifetimeCts = new();

// Null in production: P2/P2*, the suppressed-response windows and the busy-repeat delay are
// measured on Stopwatch.GetTimestamp() and waited out on real timers. A test injects the
// actor whose clock the channel it opened shares; then every timestamp this class compares
// against a deadline *and* every wait for one is on that clock, so a test advances it
// instead of sleeping, and "did the answer come before the window closed" is decided by
// the test's order of events rather than by the host's scheduling (#171). Not disposed here.
private readonly ProtocolActor? _clock;
private readonly ITimeSource _time;

// Test hook: fires whenever a caller finds _requestLock already held and starts waiting on
// it -- the observable a queued call is waiting on, standing in for a wall-clock sleep
// timed to land while an earlier call holds the lock (#171). No-op in production; a test
Expand Down Expand Up @@ -89,10 +99,26 @@ private async Task AcquireRequestLockAsync(CancellationToken cancellationToken)
private int _disposed;

public UdsClientImpl(IIsoTpChannel channel, UdsClientOptions options, bool ownsChannel)
: this(channel, options, ownsChannel, clock: null)
{
}

/// <summary>
/// As the public constructor, measuring and waiting out P2, P2*, the suppressed-response
/// windows and the busy-repeat delay on <paramref name="clock"/>; null (the public
/// constructor's choice) is the wall clock. The
/// channel underneath must have been opened on that same actor, and its demux must stamp
/// frames with its time source, or a deadline and an arrival are not comparable (#171). The
/// actor is not disposed with this client.
/// </summary>
internal UdsClientImpl(IIsoTpChannel channel, UdsClientOptions options, bool ownsChannel,
ProtocolActor? clock)
{
_channel = channel;
_options = options;
_ownsChannel = ownsChannel;
_clock = clock;
_time = clock?.TimeSource ?? MonotonicTimeSource.Instance;

// Every duration here runs a timer, and a timer measures about 49 days at most: one
// beyond that would throw when it is armed, after the request went out (Codex on #150).
Expand Down Expand Up @@ -482,7 +508,7 @@ private async Task SendWithoutResponseAsync(byte[] request, CancellationToken ca
// confirmation, the frame is on the bus and may still be answered (Codex on #150).
// Moved out to the transmit stamp afterwards.
bool hadWindow = _suppressedWindows.TryGetDeadline(request[0], out var previousUntil);
_suppressedWindows.Note(request[0], Stopwatch.GetTimestamp(), _options.P2ClientMax);
_suppressedWindows.Note(request[0], Now(), _options.P2ClientMax, _time.Frequency);
IsoTpTransmitStamps stamps;
try
{
Expand All @@ -499,11 +525,11 @@ private async Task SendWithoutResponseAsync(byte[] request, CancellationToken ca
// A send that leaves by exception -- cancelled, or a transport fault -- may
// have put the frame on the bus after the provisional window ran out; its P2
// from the transmission is at most P2 from now (Codex on #150).
_suppressedWindows.Note(request[0], Stopwatch.GetTimestamp(), _options.P2ClientMax);
_suppressedWindows.Note(request[0], Now(), _options.P2ClientMax, _time.Frequency);
throw;
}
var sent = stamps.LastFrameTransmitTimestamp > 0 ? stamps.LastFrameTransmitTimestamp : Stopwatch.GetTimestamp();
_suppressedWindows.Note(request[0], sent, _options.P2ClientMax);
var sent = stamps.LastFrameTransmitTimestamp > 0 ? stamps.LastFrameTransmitTimestamp : Now();
_suppressedWindows.Note(request[0], sent, _options.P2ClientMax, _time.Frequency);
}
finally
{
Expand Down Expand Up @@ -531,7 +557,7 @@ private async Task WaitOutSuppressedResponseWindowAsync(UdsServiceId serviceId,
{
while (true)
{
var remaining = SuppressedResponseWindows.Remaining(until);
var remaining = SuppressedResponseWindows.Remaining(until, Now(), _time.Frequency);
if (remaining <= TimeSpan.Zero)
{
// The window is over as measured now -- but a 0x78 may be queued already,
Expand All @@ -542,7 +568,7 @@ private async Task WaitOutSuppressedResponseWindowAsync(UdsServiceId serviceId,
if (DrainExtends(sid, ref until)) continue;
break;
}
using var slice = new CancellationTokenSource(remaining);
using var slice = CancelAt(until);
using var combined = CancellationTokenSource.CreateLinkedTokenSource(linkedToken, slice.Token);
IsoTpReceivedPdu pdu;
try
Expand Down Expand Up @@ -602,7 +628,7 @@ private bool ExtendOnPending(byte sid, in IsoTpReceivedPdu pdu, ref long until)
var data = pdu.Pdu;
if (data.Length < 3 || data[0] != NegativeResponseSid || data[2] != NrcResponsePending)
return false;
var extendedUntil = pdu.FirstFrameArrivalTimestamp + (long)(_options.P2StarClientMax.TotalSeconds * Stopwatch.Frequency);
var extendedUntil = pdu.FirstFrameArrivalTimestamp + Ticks(_options.P2StarClientMax);
if (data[1] != sid)
{
// Another service's: only a window still open when the 0x78 arrived (Codex on #150).
Expand Down Expand Up @@ -1052,7 +1078,7 @@ private async Task<byte[]> ExecuteCoreAsync(UdsServiceId serviceId, byte[] reque
when (ex.Code == NrcBusyRepeatRequest && repeats < _options.MaxBusyRepeatRequests)
{
if (_options.BusyRepeatRequestDelay > TimeSpan.Zero)
await Task.Delay(_options.BusyRepeatRequestDelay, linkedToken).ConfigureAwait(false);
await DelayAsync(_options.BusyRepeatRequestDelay, linkedToken).ConfigureAwait(false);
}
}
}
Expand Down Expand Up @@ -1082,15 +1108,15 @@ private async Task<byte[]> ExchangeOnceAsync(UdsServiceId serviceId, byte[] requ
// Read before the request is handed to the channel, as the fallback for a channel that
// reports no handoff instant: nothing that reached the wire after this reading can be
// an earlier request's response.
var requestStarted = Stopwatch.GetTimestamp();
var requestStarted = Now();
var stamps = await _channel.SendWithTransmitStampAsync(request, linkedToken)
.ConfigureAwait(false);
var transmitStamp = stamps.LastFrameTransmitTimestamp;

// Zero means the channel reported no transmit instant. Falling back to now is the old
// behaviour, which is worse but not broken; treating zero as a timestamp would read as
// infinitely long ago and time out every request.
var budgetStart = transmitStamp > 0 ? transmitStamp : Stopwatch.GetTimestamp();
var budgetStart = transmitStamp > 0 ? transmitStamp : Now();
// A response whose first frame arrived before this is an earlier request's (Codex on
// #143). The bound is the channel's handoff of the request's *last* frame, taken just
// before the driver call: a peer answers only a complete request, so nothing on the
Expand Down Expand Up @@ -1234,7 +1260,7 @@ private static bool IsAllZero(byte[] data, int offset, int count)
// among it is routed to its service's window rather than dropped unseen (Codex on #150).
private async Task DiscardStalePdusAsync()
{
long arrivedBefore = Stopwatch.GetTimestamp();
long arrivedBefore = Now();
await SettleAsync().ConfigureAwait(false);
DiscardStalePdus(arrivedBefore);
}
Expand Down Expand Up @@ -1292,7 +1318,7 @@ private void RouteStrayPending(in IsoTpReceivedPdu pdu)
// Only a window still open when the 0x78 arrived: one that had run out is not revived
// for a full P2* by a late frame (Codex on #150).
_suppressedWindows.ExtendIfOpenAt(data[1], pdu.FirstFrameArrivalTimestamp,
pdu.FirstFrameArrivalTimestamp + (long)(_options.P2StarClientMax.TotalSeconds * Stopwatch.Frequency));
pdu.FirstFrameArrivalTimestamp + Ticks(_options.P2StarClientMax));
}

/// <summary>
Expand Down Expand Up @@ -1324,7 +1350,7 @@ private async Task<IsoTpReceivedPdu> ReceiveWithTimeoutAsync(UdsServiceId servic
notBefore, linkedToken).ConfigureAwait(false);
}

using var timeoutCts = new CancellationTokenSource(remaining);
using var timeoutCts = CancelAt(budgetStart + Ticks(budget));
using var combined = CancellationTokenSource.CreateLinkedTokenSource(
linkedToken, timeoutCts.Token);

Expand All @@ -1347,7 +1373,9 @@ private async Task<IsoTpReceivedPdu> ReceiveWithTimeoutAsync(UdsServiceId servic
// is still there. The channel publishes a First Frame when it is read off the bus and
// withdraws it if the actor then refuses the frame -- with nothing put in the inbox -- so
// a wait on it must not be unbounded. The re-check costs nothing when the PDU arrives:
// completion or abort puts an item in the inbox and the wait returns at once.
// completion or abort puts an item in the inbox and the wait returns at once. A real timer
// even on an injected clock: it polls the channel's state, it is no protocol deadline, and a
// test would otherwise have to advance its clock for a reception it is merely waiting on.
private static readonly TimeSpan InProgressRecheck = TimeSpan.FromMilliseconds(50);

/// <summary>
Expand Down Expand Up @@ -1424,17 +1452,84 @@ private bool ResponseBeganInTime(UdsServiceId serviceId, TimeSpan budget, long b
return false;
}

// Now, on _time: Stopwatch.GetTimestamp() in production, the injected actor's clock in a
// test (#171).
private long Now() => _time.GetTimestamp();

// window in ticks of _time.Frequency, so a deadline noted from Now() and one computed here
// are on the same clock (#171).
private long Ticks(TimeSpan window) => SuppressedResponseWindows.Ticks(window, _time.Frequency);

/// <summary>
/// A token source cancelled once <see cref="_time"/> reaches <paramref name="deadline"/>: a
/// real <see cref="CancellationTokenSource"/> timer in production, the injected actor's
/// timer in a test, so a deadline and the wait bounded by it are on one clock (#171). On the
/// actor the deadline is armed as the instant it is, not as a delay from a fresh reading --
/// a test that moves the clock between this client's reading and the arming would otherwise
/// push the timer past the deadline, and it would never fire. The actor's timer fires on its
/// loop; the cancellation is handed to the thread pool rather than run there, because
/// cancelling runs the channel's registrations, and whatever they resume must not run on --
/// and block -- the loop the channel itself needs.
/// </summary>
private ClockTimeout CancelAt(long deadline)
=> _clock is null
? new ClockTimeout(SuppressedResponseWindows.Remaining(deadline, Now(), _time.Frequency))
: new ClockTimeout(_clock, deadline);

// A wait of delay on _time, as CancelAt measures it.
private Task DelayAsync(TimeSpan delay, CancellationToken cancellationToken)
=> _clock is null ? Task.Delay(delay, cancellationToken) : WaitOnClockAsync(_clock, delay, cancellationToken);

private static async Task WaitOnClockAsync(ProtocolActor clock, TimeSpan delay, CancellationToken cancellationToken)
{
var done = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
using var registration = cancellationToken.Register(static state =>
((TaskCompletionSource<bool>)state!).TrySetCanceled(), done);
using var handle = clock.Schedule(delay, () => done.TrySetResult(true));
await done.Task.ConfigureAwait(false);
}

private sealed class ClockTimeout : IDisposable
{
private readonly CancellationTokenSource _cts;
private readonly IDisposable? _timer;

public ClockTimeout(TimeSpan remaining)
=> _cts = new CancellationTokenSource(remaining);

public ClockTimeout(ProtocolActor clock, long deadline)
{
_cts = new CancellationTokenSource();
_timer = clock.ScheduleAt(deadline, () => ThreadPool.QueueUserWorkItem(static state =>
((CancellationTokenSource)state!).Cancel(), _cts));
}

public CancellationToken Token => _cts.Token;

public bool IsCancellationRequested => _cts.IsCancellationRequested;

// On the actor's clock only the timer is released: a source without a timer of its own
// holds nothing that needs it, and one whose cancellation the timer has already handed
// to the thread pool must stay cancellable rather than throw there.
public void Dispose()
{
if (_timer is null) _cts.Dispose();
else _timer.Dispose();
}
}

/// <summary>
/// Elapsed time between two <see cref="Stopwatch.GetTimestamp"/> readings, defaulting the
/// second to now. Kept in one place so the pre-check (how much budget is left) and the
/// post-check (was this PDU inside it) can never drift onto different clocks.
/// Elapsed time between two <see cref="Stopwatch.GetTimestamp"/> readings (or the injected
/// clock's equivalent), defaulting the second to now. Kept in one place so the pre-check
/// (how much budget is left) and the post-check (was this PDU inside it) can never drift
/// onto different clocks.
/// </summary>
private static TimeSpan ElapsedSince(long startTimestamp, long? endTimestamp = null)
private TimeSpan ElapsedSince(long startTimestamp, long? endTimestamp = null)
{
var end = endTimestamp ?? Stopwatch.GetTimestamp();
var end = endTimestamp ?? Now();
var ticks = end - startTimestamp;
if (ticks <= 0) return TimeSpan.Zero;
return TimeSpan.FromSeconds((double)ticks / Stopwatch.Frequency);
return TimeSpan.FromSeconds((double)ticks / _time.Frequency);
}

private void ThrowIfDisposed()
Expand Down
20 changes: 19 additions & 1 deletion tests/CanKit.Pro.Tests/TestCases/J1939TpTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2862,7 +2862,25 @@ private async IAsyncEnumerable<CanFrameEvent> Count(
}
}

public bool TryRead(out CanFrameEvent frameEvent) => _inner.TryRead(out frameEvent);
// The same meaning for a reader that drains with TryRead, as the ISO-TP channel's pump
// does: a frame is finished with once the reader asks for the next one -- by then it has
// been handed to the actor. Only one reader drains a subscription at a time (the
// channel's pump lock), so the frame in hand needs no lock of its own.
private CanFrameEvent _taken;
private bool _hasTaken;

public bool TryRead(out CanFrameEvent frameEvent)
{
if (_hasTaken)
{
_hasTaken = false;
_owner.Consumed(_taken);
}
if (!_inner.TryRead(out frameEvent)) return false;
_taken = frameEvent;
_hasTaken = true;
return true;
}

public ValueTask<bool> WaitToReadAsync(CancellationToken cancellationToken = default)
=> _inner.WaitToReadAsync(cancellationToken);
Expand Down
Loading
Loading