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
8 changes: 8 additions & 0 deletions src/CanKit.Pro.Actor/CanKit.Pro.Actor.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,14 @@
<AssemblyAttribute Include="System.Runtime.CompilerServices.InternalsVisibleTo">
<_Parameter1>CanKit.Pro.CANopen</_Parameter1>
</AssemblyAttribute>
<!-- Functional ISO-TP collection windows and the UDS client that sits on them measure P2
and P2* on the same source, so a test can end the window by moving the clock (#171). -->
<AssemblyAttribute Include="System.Runtime.CompilerServices.InternalsVisibleTo">
<_Parameter1>CanKit.Pro.IsoTp</_Parameter1>
</AssemblyAttribute>
<AssemblyAttribute Include="System.Runtime.CompilerServices.InternalsVisibleTo">
<_Parameter1>CanKit.Pro.Uds</_Parameter1>
</AssemblyAttribute>
</ItemGroup>

</Project>
39 changes: 27 additions & 12 deletions src/CanKit.Pro.IsoTp/IsoTpFunctionalClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
using System.Threading.Tasks;
using CanKit.Abstractions.API.Can.Definitions;
using CanKit.Abstractions.API.Common.Definitions;
using CanKit.Pro.Actor;
using CanKit.Pro.RawCan;

namespace CanKit.Pro.IsoTp;
Expand Down Expand Up @@ -48,6 +49,8 @@ public sealed class IsoTpFunctionalClient : IDisposable
private readonly ICanBusService _service;
private readonly bool _ownsService;
private readonly IsoTpFunctionalOptions _options;
private readonly ProtocolActor? _clock;
private readonly ITimeSource _time;

// Endpoint used only for building the outbound SF payload (Normal addressing, no AE byte).
private readonly IsoTpEndpoint _txEndpoint;
Expand All @@ -63,13 +66,21 @@ public sealed class IsoTpFunctionalClient : IDisposable

private int _disposed;

/// <summary>
/// A null <paramref name="clock"/> is production: windows end with
/// <see cref="CancellationTokenSource.CancelAfter(TimeSpan)"/> and deadlines are
/// <see cref="Stopwatch"/> readings. An injected actor is the test seam (#171). Its timers
/// and this client's deadlines share that actor's clock, so a test ends a window by
/// advancing the clock instead of sleeping. The actor stays the caller's to dispose.
/// </summary>
internal IsoTpFunctionalClient(
ICanBusService service,
uint functionalTxCanId,
uint responseRxCanIdRangeStart,
uint responseRxCanIdRangeEnd,
IsoTpFunctionalOptions options,
bool ownsService)
bool ownsService,
ProtocolActor? clock = null)
{
if (responseRxCanIdRangeEnd < responseRxCanIdRangeStart)
throw new ArgumentOutOfRangeException(nameof(responseRxCanIdRangeEnd),
Expand All @@ -78,6 +89,8 @@ internal IsoTpFunctionalClient(
_service = service ?? throw new ArgumentNullException(nameof(service));
_options = options ?? throw new ArgumentNullException(nameof(options));
_ownsService = ownsService;
_clock = clock;
_time = clock?.TimeSource ?? MonotonicTimeSource.Instance;

_txEndpoint = IsoTpEndpoint.Normal(functionalTxCanId, 0, options.IsExtendedCanId);

Expand Down Expand Up @@ -240,7 +253,7 @@ public async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectResponsesAsync(
public IsoTpFunctionalListener Listen()
{
ThrowIfDisposed();
return new IsoTpFunctionalListener(_service.Subscribe(_responseFilter, includeEcho: true));
return new IsoTpFunctionalListener(_service.Subscribe(_responseFilter, includeEcho: true), _clock);
}

/// <inheritdoc/>
Expand Down Expand Up @@ -302,7 +315,7 @@ private async Task<IsoTpTransmitStamps> SendSingleFrameAsync(ReadOnlyMemory<byte
return new IsoTpTransmitStamps(confirmation.HostHandoffTimestamp, confirmation.HostTransmitTimestamp);
}

private static async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectFromSubscriptionAsync(
private async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectFromSubscriptionAsync(
ISubscription sub, TimeSpan window, CancellationToken cancellationToken)
{
var responses = new List<IsoTpFunctionalResponse>();
Expand All @@ -311,16 +324,16 @@ private static async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectFromSub
// The window's end is also held as an arrival stamp: the timer's callback and this
// method's continuations are scheduling, and a frame that arrived after the deadline
// but before they ran is not the window's (Codex on #150).
long deadline = Stopwatch.GetTimestamp() + (long)(window.TotalSeconds * Stopwatch.Frequency);
using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
windowCts.CancelAfter(window);
var windowToken = windowCts.Token;
long deadline = _time.GetTimestamp() + (long)(window.TotalSeconds * _time.Frequency);
using var windowEnd = new FunctionalWindow(_clock, window, cancellationToken);
var windowToken = windowEnd.Token;

try
{
await foreach (var frameEvent in sub.Frames.WithCancellation(windowToken).ConfigureAwait(false))
{
if (TryParseFunctionalResponse(frameEvent, out var response) && response!.HostArrivalTimestamp <= deadline)
if (TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp)
&& response!.HostArrivalTimestamp <= deadline)
responses.Add(response);
}
}
Expand All @@ -340,7 +353,8 @@ private static async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectFromSub
sub.Dispose();
while (sub.TryRead(out var frameEvent))
{
if (TryParseFunctionalResponse(frameEvent, out var response) && response!.HostArrivalTimestamp <= deadline)
if (TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp)
&& response!.HostArrivalTimestamp <= deadline)
responses.Add(response);
}
}
Expand All @@ -359,7 +373,7 @@ private static void DrainBuffered(ISubscription sub)
}

internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent,
out IsoTpFunctionalResponse? response)
out IsoTpFunctionalResponse? response, Func<long> now)
{
var frame = frameEvent.Frame;
var payload = frame.Data.ToArray();
Expand Down Expand Up @@ -389,8 +403,9 @@ internal static bool TryParseFunctionalResponse(in CanFrameEvent frameEvent,

var pdu = new byte[pci.Length];
Array.Copy(payload, pci.DataOffset, pdu, 0, pci.Length);
// Stamped by the demux at arrival; "now" only for an event built without a stamp.
var arrival = frameEvent.HostArrivalTimestamp > 0 ? frameEvent.HostArrivalTimestamp : Stopwatch.GetTimestamp();
// Stamped by the demux at arrival; "now", on the collector's clock, only for an event
// built without a stamp.
var arrival = frameEvent.HostArrivalTimestamp > 0 ? frameEvent.HostArrivalTimestamp : now();
response = new IsoTpFunctionalResponse((uint)frame.ID, pdu, arrival);
return true;
}
Expand Down
69 changes: 62 additions & 7 deletions src/CanKit.Pro.IsoTp/IsoTpFunctionalListener.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;
using CanKit.Pro.Actor;
using CanKit.Pro.RawCan;

namespace CanKit.Pro.IsoTp;
Expand All @@ -27,15 +28,19 @@ public sealed class IsoTpFunctionalListener : IDisposable
private static readonly TimeSpan MaxWindow = TimeSpan.FromMilliseconds(uint.MaxValue - 1);

private readonly ISubscription _subscription;
private readonly ProtocolActor? _clock;
private readonly ITimeSource _time;
// Arrived after a collection's deadline but before its drain read the buffer: kept for the
// next collection, whose window they are in (Codex on #150).
private readonly Queue<IsoTpFunctionalResponse> _carried = new();
private bool _ended;
private int _disposed;

internal IsoTpFunctionalListener(ISubscription subscription)
internal IsoTpFunctionalListener(ISubscription subscription, ProtocolActor? clock = null)
{
_subscription = subscription;
_clock = clock;
_time = clock?.TimeSource ?? MonotonicTimeSource.Instance;
}

/// <summary>
Expand Down Expand Up @@ -68,7 +73,7 @@ public async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectAsync(
if (window > MaxWindow)
throw new ArgumentOutOfRangeException(nameof(window), window,
"The collection window exceeds what a timer can measure (about 49 days).");
long now = Stopwatch.GetTimestamp();
long now = _time.GetTimestamp();
var responses = new List<IsoTpFunctionalResponse>(_carried);
_carried.Clear();
if (window <= TimeSpan.Zero)
Expand All @@ -80,16 +85,15 @@ public async Task<IReadOnlyList<IsoTpFunctionalResponse>> CollectAsync(
TakeBuffered(responses, now);
return responses.AsReadOnly();
}
long deadline = now + (long)(window.TotalSeconds * Stopwatch.Frequency);
using var windowCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
windowCts.CancelAfter(window);
long deadline = now + (long)(window.TotalSeconds * _time.Frequency);
using var windowEnd = new FunctionalWindow(_clock, window, cancellationToken);
try
{
// Ends when the window runs out, or when the subscription is completed underneath
// -- the service disposed -- in which case what is buffered is still read, and
// the next collection throws rather than return empty at once: a loop that
// collects while a deadline lasts would otherwise spin on it (Bugbot on #150).
while (await _subscription.WaitToReadAsync(windowCts.Token).ConfigureAwait(false))
while (await _subscription.WaitToReadAsync(windowEnd.Token).ConfigureAwait(false))
TakeBuffered(responses, deadline);
_ended = true;
}
Expand All @@ -109,7 +113,7 @@ private void TakeBuffered(List<IsoTpFunctionalResponse> responses, long deadline
{
while (_subscription.TryRead(out var frameEvent))
{
if (!IsoTpFunctionalClient.TryParseFunctionalResponse(frameEvent, out var response)) continue;
if (!IsoTpFunctionalClient.TryParseFunctionalResponse(frameEvent, out var response, _time.GetTimestamp)) continue;
if (response!.HostArrivalTimestamp <= deadline) responses.Add(response);
else _carried.Enqueue(response);
}
Expand All @@ -122,3 +126,54 @@ public void Dispose()
_subscription.Dispose();
}
}

/// <summary>
/// A functional collection window: its token is cancelled by the caller's token or when the
/// window ends. Production ends it with <see cref="CancellationTokenSource.CancelAfter(TimeSpan)"/>.
/// A test that injected an actor arms the same interval on that actor's clock instead, so the
/// window follows the clock and not the runner (#171).
/// </summary>
internal sealed class FunctionalWindow : IDisposable
{
private readonly CancellationTokenSource _linked;
private readonly IDisposable? _timer;
private readonly object _gate = new();
private bool _disposed;

public FunctionalWindow(ProtocolActor? actor, TimeSpan window, CancellationToken cancellationToken)
{
_linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
if (actor is null)
{
_linked.CancelAfter(window);
return;
}

_timer = actor.Schedule(window, End);
}

public CancellationToken Token => _linked.Token;

/// <summary>
/// The injected actor's timer. Cancelling that timer only flags it, so a callback the loop
/// has already taken can run after <see cref="Dispose"/>; under the same lock it then finds
/// the window disposed and does nothing. Nothing here depends on the actor still running.
/// </summary>
internal void End()
{
lock (_gate)
{
if (!_disposed) _linked.Cancel();
}
}

public void Dispose()
{
_timer?.Dispose();
lock (_gate)
{
_disposed = true;
_linked.Dispose();
}
}
}
Loading
Loading