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
49 changes: 40 additions & 9 deletions src/CanKit.Pro.Uds/UdsClientImpl.cs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,37 @@ internal sealed class UdsClientImpl : IUdsClient
private readonly SemaphoreSlim _requestLock = new(1, 1);
private readonly CancellationTokenSource _lifetimeCts = new();

// 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
// subscribes to learn the wait has begun rather than guessing how long the holder needs.
internal event Action? RequestLockContended;

// Test hook: fires once _requestLock has actually been taken -- contended or not -- so a
// test driving a holder that then blocks on something else (a gated channel double) knows
// the lock is held without guessing how long the acquisition itself takes (#171).
internal event Action? RequestLockAcquired;

/// <summary>
/// Acquires <see cref="_requestLock"/>, raising <see cref="RequestLockContended"/> first if
/// it is already held. The non-blocking probe
/// (<see cref="SemaphoreSlim.Wait(int, CancellationToken)"/> with a zero timeout) keeps the
/// common uncontended path free of the event's cost and, more to the point, is what makes
/// the signal mean "a holder is in the way" rather than "a wait was requested" -- the latter
/// would fire on every call, contended or not (#171). The probe takes the token too, so a
/// call cancelled before it starts throws rather than taking a free lock, as the plain
/// <see cref="SemaphoreSlim.WaitAsync(CancellationToken)"/> it replaced did.
/// </summary>
private async Task AcquireRequestLockAsync(CancellationToken cancellationToken)
{
if (!_requestLock.Wait(0, cancellationToken))
{
RequestLockContended?.Invoke();
await _requestLock.WaitAsync(cancellationToken).ConfigureAwait(false);
}
RequestLockAcquired?.Invoke();
}

private byte _currentSession = (byte)UdsSessionType.Default;
private TesterPresentKeepAlive? _keepAlive;
private int _disposed;
Expand Down Expand Up @@ -322,7 +353,7 @@ public async Task SecurityAccessAsync(byte requestSeedLevel, Func<byte[], byte[]
// Hold the request lock across seed + sendKey so TesterPresent keep-alive (or another
// UDS call) cannot interleave and break the ISO 14229-1 security-access sequence
// (NRC requestSequenceError on real ECUs).
await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
var seedRequest = new byte[] { (byte)UdsServiceId.SecurityAccess, requestSeedLevel };
Expand Down Expand Up @@ -437,7 +468,7 @@ private async Task SendWithoutResponseAsync(byte[] request, CancellationToken ca
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
// The stale-reply discard the answered path runs before its send, here too: a late
Expand Down Expand Up @@ -611,7 +642,7 @@ public async Task<UdsDownloadResponse> RequestDownloadAsync(
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
return await RequestTransferSetupCoreAsync(
Expand Down Expand Up @@ -641,7 +672,7 @@ public async Task<UdsUploadResponse> RequestUploadAsync(
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
return await RequestTransferSetupCoreAsync(
Expand Down Expand Up @@ -729,7 +760,7 @@ public async Task<byte[]> TransferDataAsync(byte blockSequenceCounter,
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
return await TransferDataCoreAsync(blockSequenceCounter, data, linkedToken)
Expand Down Expand Up @@ -779,7 +810,7 @@ public async Task RequestTransferExitAsync(
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
await RequestTransferExitCoreAsync(transferRequestParameterRecord, linkedToken)
Expand Down Expand Up @@ -864,7 +895,7 @@ public async Task DownloadAsync(
// SecurityAccessAsync: acquire the lock once, then call the *Core helpers that assume
// the lock is held. The public RequestDownload/TransferData/RequestTransferExit APIs
// remain unchanged for single-step callers.
await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
var download = await RequestTransferSetupCoreAsync(
Expand Down Expand Up @@ -932,7 +963,7 @@ public async Task<byte[]> UploadAsync(

// Same lock discipline as DownloadAsync: one continuous 0x35 → 0x36…0x36 → 0x37
// sequence, so keep-alive traffic cannot desynchronise the ECU's block-sequence counter.
await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
_ = await RequestTransferSetupCoreAsync(
Expand Down Expand Up @@ -989,7 +1020,7 @@ private async Task<byte[]> ExecuteAsync(UdsServiceId serviceId, byte[] request,
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
return await ExecuteCoreAsync(serviceId, request, linkedToken).ConfigureAwait(false);
Expand Down
17 changes: 16 additions & 1 deletion src/CanKit.Pro.Uds/UdsFunctionalClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,21 @@ internal void DelayListenerStart(TimeSpan delay)
private TimeSpan _listenerStartDelay;
private int _disposed;

// Test hook: fires whenever a call finds _requestLock already held and starts waiting on
// it -- the observable a call queued behind another is waiting on, standing in for a
// wall-clock sleep timed to land while the earlier call's window is still open (#171).
// No-op in production.
internal event Action? RequestLockContended;

private async Task AcquireRequestLockAsync(CancellationToken cancellationToken)
{
if (!_requestLock.Wait(0, cancellationToken))
{
RequestLockContended?.Invoke();
await _requestLock.WaitAsync(cancellationToken).ConfigureAwait(false);
}
}

private UdsFunctionalClient(IsoTpFunctionalClient client, bool ownsClient, TimeSpan responseWindow,
TimeSpan responsePendingWindow, ProtocolActor? clock)
{
Expand Down Expand Up @@ -140,7 +155,7 @@ public async Task<IReadOnlyList<UdsFunctionalResponse>> SendRawAsync(ReadOnlyMem
cancellationToken, _lifetimeCts.Token);
var linkedToken = linked.Token;

await _requestLock.WaitAsync(linkedToken).ConfigureAwait(false);
await AcquireRequestLockAsync(linkedToken).ConfigureAwait(false);
try
{
// Disposed while queued behind another call: the lock is released, not used.
Expand Down
67 changes: 62 additions & 5 deletions tests/CanKit.Pro.Tests/TestCases/Uds/UdsClientTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -490,6 +490,22 @@ public async Task SendRaw_With_The_Suppress_Bit_Does_Not_Wait_For_A_Response()
}
}

// A contended lock with nobody subscribed to RequestLockContended is the production shape
// of every request lock acquisition: the event is a test hook and normally has no listener
// (#171).
[Fact]
public async Task A_Contended_Lock_Needs_No_Subscriber()
{
var (client, _, dispose) = BuildPair(e => e.On(0x22, req => new byte[] { req[1], req[2] }));
using (dispose)
{
using var cts = new CancellationTokenSource(ShortTimeout);
var first = client.ReadDataByIdentifierAsync(0xF190, cts.Token);
var second = client.ReadDataByIdentifierAsync(0xF191, cts.Token);
await Task.WhenAll(first, second);
}
}

// Codex on #150: a suppressed send may still draw a negative response, up to P2 after it.
// The next request for the same service waits that window out rather than taking the
// negative response as its own.
Expand Down Expand Up @@ -1190,6 +1206,16 @@ public async Task SecurityAccess_Holds_Lock_Across_Seed_And_Key()

using (dispose)
{
var impl = (UdsClientImpl)client;
// The property is that a keep-alive tick queued behind SecurityAccess does not
// transmit until the lock is released -- proven by observing at least one tick
// actually contend for the lock while computeKey blocks, not by guessing how many
// 30 ms periods a fixed wait covers (#171). Stronger than the counted-periods guess
// it replaces: that could pass with zero ticks ever firing.
var keepAliveContended = new TaskCompletionSource<bool>(
TaskCreationOptions.RunContinuationsAsynchronously);
impl.RequestLockContended += () => keepAliveContended.TrySetResult(true);

using var keepAlive = client.StartTesterPresentKeepAlive(TimeSpan.FromMilliseconds(30));

var unlock = client.SecurityAccessAsync(
Expand All @@ -1198,14 +1224,14 @@ public async Task SecurityAccess_Holds_Lock_Across_Seed_And_Key()
{
keyStarted.TrySetResult(true);
// Block inside computeKey (still under the request lock) long enough that
// several keep-alive ticks fire; they must not transmit until unlock ends.
// a keep-alive tick contends for it; it must not transmit until unlock ends.
releaseKey.Task.Wait(ShortTimeout);
return s.Select(b => (byte)(b ^ 0x55)).ToArray();
},
cancellationToken: new CancellationTokenSource(ShortTimeout).Token);

await keyStarted.Task.WaitAsync(ShortTimeout);
await Task.Delay(120); // several keep-alive periods while lock is held
await keepAliveContended.Task.WaitAsync(ShortTimeout); // a keep-alive tick queued behind SecurityAccess
releaseKey.TrySetResult(true);
await unlock;

Expand Down Expand Up @@ -1677,9 +1703,15 @@ public async Task Dispose_During_InFlight_Request_Does_Not_Race_RequestLock()
P2StarClientMax = TimeSpan.FromSeconds(5),
});
using var teardown = dispose;
var impl = (UdsClientImpl)client;
// The ECU never answers (EcuSilent), so once the lock is held the request stays parked
// in the receive: the observable is the lock, not a guess at how long entering the
// receive on top of it takes (#171).
var lockHeld = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
impl.RequestLockAcquired += () => lockHeld.TrySetResult(true);
var inFlight = client.ReadDataByIdentifierAsync(0xF190,
new CancellationTokenSource(ShortTimeout).Token);
await Task.Delay(50); // enter ReceiveWithTimeout under the request lock
await lockHeld.Task; // holds _requestLock, parked in ReceiveWithTimeout

Action act = () => client.Dispose();
act.Should().NotThrow(
Expand All @@ -1705,15 +1737,20 @@ public async Task Dispose_Cancels_Suppress_TesterPresent_Blocked_On_RequestLock(
P2StarClientMax = TimeSpan.FromSeconds(5),
});
using var teardown = dispose;
var impl = (UdsClientImpl)client;
// Hold the request lock with a silent ECU read so suppress TesterPresent blocks
// in WaitAsync rather than racing through Send.
var readHoldsLock = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
impl.RequestLockAcquired += () => readHoldsLock.TrySetResult(true);
var inFlight = client.ReadDataByIdentifierAsync(0xF190,
new CancellationTokenSource(ShortTimeout).Token);
await Task.Delay(50);
await readHoldsLock.Task;

var testerPresentQueued = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
impl.RequestLockContended += () => testerPresentQueued.TrySetResult(true);
var testerPresent = client.TesterPresentAsync(suppressPositiveResponse: true,
new CancellationTokenSource(ShortTimeout).Token);
await Task.Delay(30); // park on _requestLock.WaitAsync
await testerPresentQueued.Task; // parked on _requestLock.WaitAsync

Action act = () => client.Dispose();
act.Should().NotThrow();
Expand All @@ -1726,6 +1763,26 @@ await waitTp.Should().ThrowAsync<OperationCanceledException>(
await waitRead.Should().ThrowAsync<OperationCanceledException>();
}

// -----------------------------------------------------------------------------------
// #171 — the uncontended fast path of the request lock must honour a token that is
// already cancelled, as the plain WaitAsync it replaced did: a call cancelled before it
// starts never takes the lock.
// -----------------------------------------------------------------------------------
[Fact]
public async Task An_Already_Cancelled_Call_Does_Not_Take_A_Free_Request_Lock()
{
var (client, _, dispose) = BuildPair(e => e.On(0x22, _ => new byte[] { 0xF1, 0x90, 0x01 }));
using var teardown = dispose;
var impl = (UdsClientImpl)client;
var acquired = false;
impl.RequestLockAcquired += () => acquired = true;

Func<Task> act = () => client.ReadDataByIdentifierAsync(0xF190, new CancellationToken(canceled: true));

await act.Should().ThrowAsync<OperationCanceledException>();
acquired.Should().BeFalse("a call cancelled before it starts must not take the request lock");
}

/// <summary>
/// Records <see cref="TxConfirmation.HostTransmitTimestamp"/> of the first
/// <see cref="ICanBusService.SendConfirmed"/> and completes that task only afterwards,
Expand Down
12 changes: 9 additions & 3 deletions tests/CanKit.Pro.Tests/TestCases/Uds/UdsExpiredDeadlineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -368,10 +368,16 @@ public async Task N_Dispose_Leaves_The_Lock_To_A_Holder_That_Outlasts_The_Wait()
stampArrivalAtDelivery: true)
{ Gate = gate };
using var client = NewClient(channel); // a second Dispose is idempotent
((UdsClientImpl)client).DisposeLockTimeout = TimeSpan.FromMilliseconds(100);

var impl = (UdsClientImpl)client;
impl.DisposeLockTimeout = TimeSpan.FromMilliseconds(100);

// The request is gated indefinitely on `gate` (set below), so once it holds the lock it
// stays blocked in the receive: an observable that the lock was taken is enough, with
// no need to also observe the receive itself (#171).
var lockHeld = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
impl.RequestLockAcquired += () => lockHeld.TrySetResult(true);
var inFlight = client.ReadDataByIdentifierAsync(0xF190, CancellationToken.None);
await Task.Delay(50); // let the request take the lock and block in the receive
await lockHeld.Task; // the request holds _requestLock and is blocked in the receive

client.Dispose(); // returns after its 100 ms wait, the holder still inside

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -225,10 +225,16 @@ public async Task A_Call_Queued_Behind_Another_Does_Not_Send_After_Dispose()
using var functional = UdsFunctionalClient.Create( // a second Dispose is idempotent
IsoTpFactory.OpenFunctional(busTester, FunctionalTxId, Ecu1, 0x7EF, FastOptions()), ownsClient: true);

// The observable a queued call is waiting on: the second SendRawAsync finds
// _requestLock already held by the first (its collection window still open) and starts
// waiting on it, rather than a guess at how long the first's window needs (#171).
var secondQueued = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
functional.RequestLockContended += () => secondQueued.TrySetResult(true);

using var cts = new CancellationTokenSource(ShortTimeout);
var first = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x90 }, TimeSpan.FromMilliseconds(300), cts.Token);
var second = functional.SendRawAsync(new byte[] { 0x22, 0xF1, 0x91 }, TimeSpan.FromMilliseconds(300), cts.Token);
await Task.Delay(50); // the first is in its window, the second queued behind it
await secondQueued.Task; // the first is in its window, the second queued behind it

functional.Dispose();

Expand Down Expand Up @@ -976,6 +982,8 @@ public async Task An_Invalid_Collection_Window_Transmits_Nothing(long windowMill
var window = windowMilliseconds == long.MaxValue ? TimeSpan.MaxValue : TimeSpan.FromMilliseconds(windowMilliseconds);
Func<Task> act = () => functional.DiagnosticSessionControlAsync(UdsSessionType.Extended, window);
await act.Should().ThrowAsync<ArgumentOutOfRangeException>();
// A negative check on another bus: a frame sent by a regression that validated after
// the send would reach busEcus asynchronously, so this keeps its wall window (#171).
await Task.Delay(50);
seen.Should().Be(0, "nothing was transmitted");
}
Expand Down
Loading
Loading