diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.CommunicationProfile.cs b/src/CanKit.Pro.CANopen/CanOpenNode.CommunicationProfile.cs index 40542f8..a5403bc 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.CommunicationProfile.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.CommunicationProfile.cs @@ -697,6 +697,9 @@ private void PerformNmtReset(bool communicationOnly) // the hook lets the application restore what the description does not hold. RaiseApplicationReset(communicationOnly ? NmtCommand.ResetCommunication : NmtCommand.ResetNode); + // A subscriber may have disposed the node; its timers are gone and nothing is announced. + if (Volatile.Read(ref _disposed) != 0) return; + _state = NmtState.PreOperational; // Boot-up (0x00) first; a heartbeat with the new state follows only when the producer is // active, in which case §7.2.8.3.2.2 regards the boot-up as its first heartbeat — and the diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.FlyingMaster.cs b/src/CanKit.Pro.CANopen/CanOpenNode.FlyingMaster.cs index 22051d8..760a58f 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.FlyingMaster.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.FlyingMaster.cs @@ -614,7 +614,7 @@ private void RaiseFlyingMaster(FlyingMasterSignal signal, byte? otherNodeId, ush var args = new FlyingMasterChangedEventArgs(signal, _flyingMasterRole, otherNodeId, otherPriority); // The dispatcher reports a subscriber throw and keeps going, so a timeout or an // EMCY already queued is still delivered. A second catch here would only repeat that. - EnqueueEvent(() => FlyingMasterChanged?.Invoke(this, args)); + EnqueueEvent(() => DeliverToSubscribers(FlyingMasterChanged, args)); } private OdWriteDecision ValidateFlyingMasterTimingWrite(byte subindex, byte[] value) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.NodeGuarding.cs b/src/CanKit.Pro.CANopen/CanOpenNode.NodeGuarding.cs index 8922aa8..456794d 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.NodeGuarding.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.NodeGuarding.cs @@ -332,7 +332,7 @@ private void RaiseLifeGuardingEvent(LifeGuardingState state) var args = new LifeGuardingEventArgs(state, _guardTime, _lifeTimeFactor); EnqueueEvent(() => { - try { LifeGuardingEvent?.Invoke(this, args); } + try { DeliverToSubscribers(LifeGuardingEvent, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }); } @@ -354,7 +354,7 @@ private void RaiseNodeGuardingReceived(byte producer, NmtState state, bool toggl var args = new NodeGuardingReceivedEventArgs(producer, state, toggle, ts); EnqueueEvent(() => { - try { NodeGuardingReceived?.Invoke(this, args); } + try { DeliverToSubscribers(NodeGuardingReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }, critical: false, key: null, emcyProducer: -1, producer); } @@ -364,7 +364,7 @@ private void RaiseNodeGuardingTimeout(byte producer, TimeSpan guardTime, byte li var args = new NodeGuardingTimeoutEventArgs(producer, guardTime, lifeTimeFactor); EnqueueEvent(() => { - try { NodeGuardingTimeout?.Invoke(this, args); } + try { DeliverToSubscribers(NodeGuardingTimeout, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }, critical: true, EventKey.NodeGuardingTimeout(producer), emcyProducer: -1, producer); } diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 4c80bf6..25e13ab 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -39,7 +39,7 @@ namespace CanKit.Pro.CANopen; /// same threading model that the J1939-TP / IsoTp / UDS clients rely on. /// /// -internal sealed partial class CanOpenNode : ICanOpenNode +internal sealed partial class CanOpenNode : ICanOpenNode, IAsyncDisposable { private readonly ICanBusService _service; private readonly bool _ownsService; @@ -138,6 +138,8 @@ internal sealed partial class CanOpenNode : ICanOpenNode private bool _nodeGuardingProducerToggle; private int _disposed; + private bool _pumpStopRequested; + private readonly TaskCompletionSource _disposeDone = new(TaskCreationOptions.RunContinuationsAsynchronously); /// public byte NodeId => _nodeId; @@ -628,68 +630,218 @@ public Task SdoDownloadAsync(byte serverNodeId, ushort index, byte subindex, /// public void Dispose() { - if (Interlocked.Exchange(ref _disposed, 1) != 0) return; - try { _readerCts.Cancel(); } catch { /* nothing else to do */ } + // From the pump, this call ends the delivery whether it wins the disposal or not: the + // other caller cannot stop the pump before its own join has run out. + if (OnEventPump) Volatile.Write(ref _pumpStopRequested, true); + if (!BeginDispose()) + { + // Another caller won. From the actor this call still has to leave the transfers ended: + // the winner's cleanup is queued behind the callback this one is made from, and the + // callback may be waiting for it. (Idempotent; the winner runs it again.) + if (_actor.IsOnCurrentActor) CleanUpOnActor(); + return; + } + // A subscriber that disposes the node runs on the thread of the task being joined: the + // reader when it reports a failed subscription, the event pump when it delivers an event. + // Neither can finish while it waits for itself, so it does not wait; the task ends when + // the subscriber returns. (The actor's own reentrant Dispose is ProtocolActor's.) + if (!OnReader) + { + try { _readerTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } + } - try + // Complete the queue so the pump exits after draining anything still queued. An event + // accepted before this point is delivered (unless the subscriber itself hangs); one that + // arrives afterwards is dropped, timeout and EMCY included. Nothing is delivered past + // Dispose. + CompleteEventQueue(); + // From the pump the events still queued are not delivered once this subscriber returns + // (the flag is set above), and the pump is not waited for. + if (!OnEventPump) { - _actor.Post(() => + try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } + StopPumpIfStillRunning(); + } + + FinishDispose(); + } + + /// + public ValueTask DisposeAsync() + { + // From inside one of the node's own callbacks (a subscriber on the pump, on the reader's + // failure report or on the actor) the caller holds the thread the waits below would need, + // and the continuation of an awaited join would run elsewhere, away from the context + // Dispose recognises. The blocking path knows how to not wait for itself, and the caller + // is on a node thread already, so nothing is lost by taking it. + if (OnEventPump || OnReader || OnDisposal || _actor.IsOnCurrentActor) + { + Dispose(); + return default; + } + + return DisposeCoreAsync(); + } + + private async ValueTask DisposeCoreAsync() + { + // A second caller awaits the disposal that is running: returning at once would tell it + // that producers are stopped and the owned service released while the first caller is + // still waiting for the reader. + if (!BeginDispose()) + { + await _disposeDone.Task.ConfigureAwait(false); + return; + } + await JoinAsync(_readerTask).ConfigureAwait(false); + CompleteEventQueue(); + await JoinAsync(_eventPumpTask).ConfigureAwait(false); + StopPumpIfStillRunning(); + FinishDispose(); + } + + private static readonly TimeSpan DisposeJoinTimeout = TimeSpan.FromSeconds(2); + + /// A subscriber that kept the pump past the join timeout does not get the events still + /// queued behind it: nothing is delivered past Dispose. + private void StopPumpIfStillRunning() + { + if (!_eventPumpTask.IsCompleted) Volatile.Write(ref _pumpStopRequested, true); + } + + /// Waits for for at most + /// without holding a thread. Neither task joined here faults (each reports its own failures + /// through ), so there is nothing to observe. + private static async Task JoinAsync(Task task) + { + using var cts = new CancellationTokenSource(); + await Task.WhenAny(task, Task.Delay(DisposeJoinTimeout, cts.Token)).ConfigureAwait(false); + cts.Cancel(); + } + + /// True on the thread that is delivering an event to a subscriber. + private bool OnEventPump => Callbacks.Contains(CallbackKind.Delivering, this); + + /// True on the thread that is finishing the disposal: a subscriber the actor calls + /// there (its shutdown-timeout report) is waited for by that disposal, so it cannot wait for it. + private bool OnDisposal => Callbacks.Contains(CallbackKind.Finishing, this); + + /// True on the thread of the reader task while it reports a failed subscription. + private bool OnReader => Callbacks.Contains(CallbackKind.Reporting, this); + + // The nodes whose callbacks the current thread is inside, per kind. A set, not one slot: a + // subscriber of A may dispose B, and B's own callback may then ask A to dispose again, with + // both of them active on the same thread at once. + private enum CallbackKind { Delivering, Reporting, Finishing } + + private static class Callbacks + { + [ThreadStatic] + private static System.Collections.Generic.List?[]? s_nodes; + + public static void Enter(CallbackKind kind, CanOpenNode node) => Of(kind).Add(node); + + // Entering and leaving are paired, and the nodes of one thread are all the same object + // when they repeat, so removing the first match is removing the right one. + public static void Leave(CallbackKind kind, CanOpenNode node) => Of(kind).Remove(node); + + public static bool Contains(CallbackKind kind, CanOpenNode node) + { + var list = s_nodes?[(int)kind]; + if (list is null) return false; + foreach (var marked in list) { - _heartbeatProducer.Dispose(); - _heartbeatConsumer.Dispose(); - _syncProducerHandle?.Dispose(); - _syncProducerHandle = null; - DisposePdoRuntime(); - _lifeGuardingDeadline?.Dispose(); - _lifeGuardingDeadline = null; - CancelFlyingMasterDeadline(); - CancelBootUp(); - - _sdoServer?.Deadline?.Dispose(); - _sdoServer = null; - foreach (var kv in _sdoClients) - { - kv.Value.Deadline?.Dispose(); - kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); - } - _sdoClients.Clear(); + if (ReferenceEquals(marked, node)) return true; + } + return false; + } - _sdoBlockServer?.Deadline?.Dispose(); - _sdoBlockServer = null; - foreach (var kv in _sdoBlockClients) - { - kv.Value.Deadline?.Dispose(); - kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); - } - _sdoBlockClients.Clear(); + private static System.Collections.Generic.List Of(CallbackKind kind) + { + var all = s_nodes ??= new System.Collections.Generic.List?[3]; + return all[(int)kind] ??= new System.Collections.Generic.List(2); + } + } - foreach (var kv in _nodeGuardingConsumers) - { - kv.Value.PollHandle?.Dispose(); - kv.Value.LifeTimeDeadline?.Dispose(); - } - _nodeGuardingConsumers.Clear(); - }); + /// Flips the disposed flag and posts the cleanup; false when already disposed. + private bool BeginDispose() + { + if (Interlocked.Exchange(ref _disposed, 1) != 0) return false; + try { _readerCts.Cancel(); } catch (AggregateException) { /* a registered callback threw; nothing else to do */ } + + // On the actor already (a subscriber disposing from ApplicationReset, say) the cleanup runs + // here: a post would wait for that subscriber to return, and the call would hand back + // with the transfers still open. + try + { + if (_actor.IsOnCurrentActor) CleanUpOnActor(); + else _actor.Post(CleanUpOnActor); } catch (ObjectDisposedException) { // actor already gone; nothing more to do } + return true; + } - try { _readerTask.Wait(TimeSpan.FromSeconds(2)); } catch { /* observed via task; not fatal */ } + /// Stops the producers and consumers and fails every open transfer. Actor-only. + private void CleanUpOnActor() + { + _heartbeatProducer.Dispose(); + _heartbeatConsumer.Dispose(); + _syncProducerHandle?.Dispose(); + _syncProducerHandle = null; + DisposePdoRuntime(); + _lifeGuardingDeadline?.Dispose(); + _lifeGuardingDeadline = null; + CancelFlyingMasterDeadline(); + CancelBootUp(); - // Complete the queue so the pump exits after draining anything still queued. An event - // accepted before this point is delivered (unless the subscriber itself hangs); one that - // arrives afterwards is dropped, timeout and EMCY included. Nothing is delivered past - // Dispose. - CompleteEventQueue(); - try { _eventPumpTask.Wait(TimeSpan.FromSeconds(2)); } catch { /* observed via task; not fatal */ } + Release(_sdoServer?.Deadline); + _sdoServer = null; + foreach (var kv in _sdoClients) + { + Release(kv.Value.Deadline); + kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); + } + _sdoClients.Clear(); + + Release(_sdoBlockServer?.Deadline); + _sdoBlockServer = null; + foreach (var kv in _sdoBlockClients) + { + Release(kv.Value.Deadline); + kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); + } + _sdoBlockClients.Clear(); + + foreach (var kv in _nodeGuardingConsumers) + { + Release(kv.Value.PollHandle); + Release(kv.Value.LifeTimeDeadline); + } + _nodeGuardingConsumers.Clear(); + } - _subscription.Dispose(); - _actor.Dispose(); - _readerCts.Dispose(); + private static void Release(IDisposable? handle) => handle?.Dispose(); - if (_ownsService) _service.Dispose(); + private void FinishDispose() + { + Callbacks.Enter(CallbackKind.Finishing, this); + try + { + _subscription.Dispose(); + _actor.Dispose(); + _readerCts.Dispose(); + + if (_ownsService) _service.Dispose(); + } + finally + { + Callbacks.Leave(CallbackKind.Finishing, this); + // Whatever a Dispose above threw, a caller waiting for this disposal is released. + _disposeDone.TrySetResult(true); + } } // ========================================================================================= @@ -725,7 +877,13 @@ private async Task RunReaderAsync() } } catch (OperationCanceledException) { /* Dispose */ } - catch (Exception ex) { RaiseBackgroundException(ex); } + catch (Exception ex) + { + // A subscriber of the report may dispose the node; see Dispose. + Callbacks.Enter(CallbackKind.Reporting, this); + try { RaiseBackgroundException(ex); } + finally { Callbacks.Leave(CallbackKind.Reporting, this); } + } } // ========================================================================================= @@ -741,17 +899,35 @@ private async Task RunEventPumpAsync() // One signal covers every event queued by the time we look. A subscriber that // throws is reported and the loop continues: a timeout or EMCY already waiting // must still be delivered. Anything outside the delegate is still a bug in the pump. - while (TryDequeueEvent() is { } raise) + // Marks the thread for the whole batch: nothing in it awaits, so a subscriber that + // disposes the node runs on this very thread (see Dispose). + Callbacks.Enter(CallbackKind.Delivering, this); + try { - try + while (!Volatile.Read(ref _pumpStopRequested) && TryDequeueEvent() is { } raise) { - raise(); + try + { + raise(); + } + catch (Exception ex) + { + RaiseBackgroundException(ex); + } } - catch (Exception ex) + // Stopped with events still queued: they are dropped, not left linked to a node + // nobody will deliver for again. + if (Volatile.Read(ref _pumpStopRequested)) { - RaiseBackgroundException(ex); + Action? skipped; + do { skipped = TryDequeueEvent(); } + while (skipped is not null); } } + finally + { + Callbacks.Leave(CallbackKind.Delivering, this); + } // The queue was just drained. Closure is the completed flag alone: an event // accepted before completion is still in the list and was delivered above, and // one accepted after completion never enters the list. @@ -2492,10 +2668,36 @@ private Task SendOrderedControlFrames(Action? onSend }); } + /// + /// Calls the subscribers of an event one after the other, ending the round when one of them + /// has disposed the node from the event pump: nothing is delivered past Dispose. + /// + private void DeliverToSubscribers(EventHandler? handler, T args) + { + if (handler is null) return; + foreach (var subscriber in handler.GetInvocationList()) + { + if (Volatile.Read(ref _pumpStopRequested)) break; + ((EventHandler)subscriber)(this, args); + } + } + private void RaiseBackgroundException(Exception ex) { - try { BackgroundExceptionOccurred?.Invoke(this, ex); } - catch { /* subscriber must not tear down the node */ } + var handler = BackgroundExceptionOccurred; + if (handler is null) return; + // Nothing is reported once the disposal has finished: a subscriber or a timed-out send that + // outlived it and fails late has nobody to tell. + if (_disposeDone.Task.IsCompleted) return; + // One subscriber that disposes the node ends the round for the rest, as for the other + // events; a node that was disposed before the report still reports it to all. + bool disposedBefore = Volatile.Read(ref _disposed) != 0; + foreach (var subscriber in handler.GetInvocationList()) + { + if (!disposedBefore && Volatile.Read(ref _disposed) != 0) break; + try { ((EventHandler)subscriber)(this, ex); } + catch { /* subscriber must not tear down the node */ } + } } private void RaiseHeartbeatReceived(byte producer, NmtState state, DateTime ts) @@ -2503,7 +2705,7 @@ private void RaiseHeartbeatReceived(byte producer, NmtState state, DateTime ts) var args = new HeartbeatReceivedEventArgs(producer, state, ts); EnqueueEvent(() => { - try { HeartbeatReceived?.Invoke(this, args); } + try { DeliverToSubscribers(HeartbeatReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }, critical: false, key: null, emcyProducer: -1, producer); } @@ -2513,7 +2715,7 @@ private void RaiseHeartbeatTimeout(byte producer, TimeSpan timeout) var args = new HeartbeatTimeoutEventArgs(producer, timeout); EnqueueEvent(() => { - try { HeartbeatTimeout?.Invoke(this, args); } + try { DeliverToSubscribers(HeartbeatTimeout, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }, critical: true, EventKey.HeartbeatTimeout(producer), emcyProducer: -1, producer); } @@ -2523,7 +2725,7 @@ private void RaiseEmcyReceived(EmcyMessage msg, DateTime ts) var args = new EmcyReceivedEventArgs(msg, ts); EnqueueEvent(() => { - try { EmcyReceived?.Invoke(this, args); } + try { DeliverToSubscribers(EmcyReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }, critical: true, EventKey.Emcy(msg), msg.ProducerNodeId, msg.ProducerNodeId); } @@ -2533,7 +2735,7 @@ private void RaiseSyncReceived(DateTime ts) var args = new SyncReceivedEventArgs(ts); EnqueueEvent(() => { - try { SyncReceived?.Invoke(this, args); } + try { DeliverToSubscribers(SyncReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }); } @@ -2543,7 +2745,7 @@ private void RaiseRpdoReceived(int pdoIndex, uint cobId, byte[] payload) var args = new RpdoReceivedEventArgs(pdoIndex, cobId, payload); EnqueueEvent(() => { - try { RpdoReceived?.Invoke(this, args); } + try { DeliverToSubscribers(RpdoReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }); } @@ -2561,6 +2763,9 @@ private void RaiseApplicationReset(NmtCommand command) // (Codex on #133). foreach (var subscriber in handler.GetInvocationList()) { + // One that disposed the node ends the round: the rest are not called on a node that + // is gone. + if (Volatile.Read(ref _disposed) != 0) break; try { ((EventHandler)subscriber)(this, args); } catch (Exception ex) { RaiseBackgroundException(ex); } } @@ -2571,7 +2776,7 @@ private void RaiseNmtCommandReceived(NmtCommand cmd, byte target) var args = new NmtCommandReceivedEventArgs(cmd, target); EnqueueEvent(() => { - try { NmtCommandReceived?.Invoke(this, args); } + try { DeliverToSubscribers(NmtCommandReceived, args); } catch (Exception ex) { RaiseBackgroundException(ex); } }); } diff --git a/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs new file mode 100644 index 0000000..c06f202 --- /dev/null +++ b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs @@ -0,0 +1,34 @@ +using System; +using System.Threading.Tasks; + +namespace CanKit.Pro.CANopen; + +/// Asynchronous disposal for . +/// +/// itself stays : adding +/// to a published interface would break everyone who implements it. +/// The nodes this library creates implement it, and reaches +/// it without a cast. It cannot be used with await using on an +/// (that needs the interface); cast to for that. +/// +public static class CanOpenNodeExtensions +{ + /// + /// Disposes without holding a thread while it waits for the node's + /// reader and event pump to finish, which does for up to two + /// seconds. Either call may be made from a subscriber of the node (an event, + /// ApplicationReset, BackgroundExceptionOccurred): the task that subscriber runs + /// on is then not waited for. A second call while the first is running returns when the + /// disposal has finished, except from a subscriber of the node, which cannot wait for a + /// disposal that is waiting for it and returns at once. A node that is not one of this + /// library's is disposed on the thread pool. + /// + /// The node to dispose. + public static ValueTask DisposeAsync(this ICanOpenNode node) + { + if (node is null) throw new ArgumentNullException(nameof(node)); + return node is IAsyncDisposable asynchronous + ? asynchronous.DisposeAsync() + : new ValueTask(Task.Run(node.Dispose)); + } +} diff --git a/src/CanKit.Pro.CANopen/ICanOpenNode.cs b/src/CanKit.Pro.CANopen/ICanOpenNode.cs index 91a7a3b..29e1a78 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -36,6 +36,18 @@ namespace CanKit.Pro.CANopen; /// the same way. A value that is not implementable is rejected before it is stored, so the /// dictionary never describes behaviour the node does not have. /// +/// +/// Disposing. fails every open transfer with +/// (on the node's actor: a subscriber that holds the actor +/// for longer than the actor's five-second shutdown timeout, which is reported through +/// , delays that until it returns), stops the producers and starts no further delivery +/// afterwards (a subscriber that is running when the disposal finishes is not interrupted). It +/// blocks for up to two seconds on the node's reader and event pump to let them finish; +/// waits for the same without +/// holding a thread. Either may be called from a subscriber of this node, whichever thread it runs +/// on (an event, ApplicationReset, BackgroundExceptionOccurred): the task that thread +/// belongs to is then not waited for, and it ends when the subscriber returns. +/// /// public interface ICanOpenNode : IDisposable { diff --git a/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt b/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt index 5585e4f..4cb6a69 100644 --- a/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt +++ b/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt @@ -75,6 +75,10 @@ namespace CanKit.Pro.CANopen public static System.Threading.Tasks.Task> ListenAsync(CanKit.Pro.RawCan.ICanBusService service, System.TimeSpan? window = default, System.Threading.CancellationToken cancellationToken = default) { } public static System.Threading.Tasks.Task> ScanAsync(CanKit.Pro.CANopen.ICanOpenNode client, System.Collections.Generic.IEnumerable? skipNodeIds = null, System.Threading.CancellationToken cancellationToken = default) { } } + public static class CanOpenNodeExtensions + { + public static System.Threading.Tasks.ValueTask DisposeAsync(this CanKit.Pro.CANopen.ICanOpenNode node) { } + } public sealed class CanOpenNodeOptions { public CanOpenNodeOptions() { } diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index e9e8e59..3a36da6 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -40,6 +40,9 @@ public void Deliver(CanFrameView frame, long hostArrivalTimestamp = 0) => _frame /// When set, the reader task's wait throws it once woken, as a subscription that failed would. public Exception? ReaderFault { get; set; } + /// When set, the Frames enumeration throws it once woken, as a subscription whose demux failed would. + public Exception? FramesFault { get; set; } + /// When set, the subscription ends on the wake: the reader's wait returns false, as when its service is disposed. public bool EndSubscriptionOnWake { get; set; } @@ -92,7 +95,21 @@ public Task SendConfirmedAsync(CanFrame frame, TimeSpan? timeout public IReadOnlyList FindOverlappingFilterSubscriptions() => Array.Empty(); - public void Dispose() => _frames.Writer.TryComplete(); + /// Whether has been called. + public bool IsDisposed { get; private set; } + + /// Runs when the subscription is disposed: a node's disposal calls it on its own thread. + public Action? OnSubscriptionDisposed { get; set; } + + /// When set, throws it after recording the call. + public Exception? DisposeFault { get; set; } + + public void Dispose() + { + IsDisposed = true; + if (DisposeFault is { } fault) throw fault; + _frames.Writer.TryComplete(); + } private sealed class Sub : ISubscription { @@ -100,7 +117,18 @@ private sealed class Sub : ISubscription public Sub(StarvedReaderBusService owner) => _owner = owner; public IAsyncEnumerable Frames - => _owner.HoldFrames ? Held() : _owner._frames.Reader.ReadAllAsync(); + => _owner.FramesFault is not null ? Faulting(_owner) + : _owner.HoldFrames ? Held() : _owner._frames.Reader.ReadAllAsync(); + + private static async IAsyncEnumerable Faulting(StarvedReaderBusService owner, + [EnumeratorCancellation] CancellationToken cancellationToken = default) + { + await owner._wake.Task.WaitAsync(cancellationToken); + throw owner.FramesFault!; +#pragma warning disable CS0162 // an iterator needs a yield to be one + yield break; +#pragma warning restore CS0162 + } private static async IAsyncEnumerable Held( [EnumeratorCancellation] CancellationToken cancellationToken = default) @@ -134,6 +162,6 @@ public void Reconfigure(CanIdFilter filter) { } public void Reconfigure(Func? predicate) { } - public void Dispose() { } + public void Dispose() => _owner.OnSubscriptionDisposed?.Invoke(); } } diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs new file mode 100644 index 0000000..ef5103c --- /dev/null +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -0,0 +1,720 @@ +using System; +using System.Diagnostics; +using System.Threading; +using System.Threading.Tasks; +using AwesomeAssertions; +using CanKit.Abstractions.API.Can; +using CanKit.Abstractions.API.Can.Definitions; +using CanKit.Abstractions.API.Common.Definitions; +using CanKit.Core; +using CanKit.Pro.CANopen; +using CanKit.Pro.CANopen.Nmt; +using CanKit.Pro.CANopen.Sdo; +using CanKit.Pro.RawCan; +using CanKit.Pro.Tests.Infrastructure; +using Xunit; + +namespace CanKit.Pro.Tests.TestCases.CANopen; + +/// +/// Disposing a node from one of its own event handlers (#251). The handler runs on the event pump, +/// and a blocking Dispose used to wait for the pump to finish, that is for itself, until its join +/// timeout ran out. +/// +public class CanOpenDisposeTests : IClassFixture +{ + private static readonly TimeSpan ShortTimeout = TimeSpan.FromSeconds(5); + + // Dispose's join timeout is two seconds. A call that waited for itself takes all of it; one + // that did not takes milliseconds. This bound sits between the two with room on both sides. + private static readonly TimeSpan NoStall = TimeSpan.FromSeconds(1.5); + + private static string NewSession() => $"canopen-dispose-{Guid.NewGuid():N}"; + + private static ICanBus Open(string session, int channel) => CanBus.Open( + $"virtual://{session}/{channel}", + cfg => cfg.SetProtocolMode(CanProtocolMode.Can20).Baud(VirtualAdapterFixture.Bitrate)); + + private static void SendHeartbeat(ICanBus bus, byte producer, byte state) + => bus.Transmit(CanFrame.Classic(unchecked((int)(CanOpenCobId.HeartbeatBase + producer)), + new[] { state }, isExtendedFrame: false)); + + public static TheoryData Contexts => new() + { + { "Dispose", "an event handler" }, + { "DisposeAsync", "an event handler" }, + { "Dispose", "ApplicationReset" }, + { "DisposeAsync", "ApplicationReset" }, + { "Dispose", "a reader failure report" }, + { "DisposeAsync", "a reader failure report" }, + }; + + // A subscriber disposes the node from each of the three threads that call subscribers: the + // event pump (an ordinary event), the actor (ApplicationReset is raised on it) and the reader + // task (it reports a failed subscription). The subscriber blocks on the call, as a handler that + // is not async has to, so a call that waits for the thread it is on takes its whole join + // timeout (two seconds; the actor's is five). This bound sits between that and milliseconds. + [Theory] + [MemberData(nameof(Contexts))] + public async Task A_Node_Can_Be_Disposed_From_Its_Own_Subscriber_Without_Waiting_For_Itself(string call, string context) + { + using var resources = new DisposeBag(); + var elapsed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var transferEnded = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + ICanOpenNode? node = null; + Task? openTransfer = null; + void DisposeNode() + { + var watch = Stopwatch.StartNew(); + if (call == "Dispose") node!.Dispose(); + else node!.DisposeAsync().AsTask().Wait(); + // What a caller may rely on when the call returns: its open transfers have ended. From + // the actor this needs the cleanup to run in place, a post waits for this very handler. + transferEnded.TrySetResult(openTransfer?.IsCompleted ?? true); + elapsed.TrySetResult(watch.Elapsed); + } + + if (context == "a reader failure report") + { + var starved = new StarvedReaderBusService { FramesFault = new InvalidOperationException("the demux broke") }; + node = resources.Add(new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: false, new ManualTimeSource())); + node.BackgroundExceptionOccurred += (_, _) => DisposeNode(); + starved.WakeReader(); + } + else + { + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + node = resources.Add(CanOpen.OpenNode(busA, nodeId: 0x01)); + if (context == "ApplicationReset") + { + PeerSdoLaboratory.Bind(node, 0x02); + openTransfer = node.SdoUploadAsync(0x02, 0x2100, 0x00); // nobody answers + node.ApplicationReset += (_, _) => DisposeNode(); + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.NmtCommand), + new byte[] { (byte)NmtCommand.ResetNode, 0x01 }, isExtendedFrame: false)); + } + else + { + node.HeartbeatReceived += (_, _) => DisposeNode(); + SendHeartbeat(rawBus, 0x11, 0x05); + } + } + + (await elapsed.Task.WithTimeoutAsync(ShortTimeout)).Should().BeLessThan(NoStall); + (await transferEnded.Task.WithTimeoutAsync(ShortTimeout)).Should().BeTrue("the open transfer ended before the call returned"); + } + + // A second DisposeAsync that arrives while the first is still waiting for the reader returns + // when the disposal has finished, not at once: the owned service is released by then. + [Fact] + public async Task A_Concurrent_DisposeAsync_Returns_When_The_Disposal_Has_Finished() + { + var starved = new StarvedReaderBusService(); + var node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + + var first = node.DisposeAsync(); + var second = node.DisposeAsync(); + await second.AsTask().WithTimeoutAsync(ShortTimeout); + + starved.IsDisposed.Should().BeTrue(); + await first.AsTask().WithTimeoutAsync(ShortTimeout); + } + + // An outside caller wins the disposal while an ApplicationReset subscriber is running on the + // actor; the winner's cleanup is queued behind that subscriber. The subscriber then disposes + // too, loses, and must still leave the open transfer ended when its call returns. + [Fact] + public async Task A_Subscriber_On_The_Actor_That_Loses_The_Disposal_Race_Still_Ends_The_Open_Transfers() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + PeerSdoLaboratory.Bind(node, 0x02); + var openTransfer = node.SdoUploadAsync(0x02, 0x2100, 0x00); // nobody answers + var inHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var ended = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var outsideHasWon = new ManualResetEventSlim(); + node.ApplicationReset += (_, _) => + { + inHandler.TrySetResult(true); + outsideHasWon.Wait(ShortTimeout); + node.DisposeAsync().AsTask().Wait(); + ended.TrySetResult(openTransfer.IsCompleted); + }; + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.NmtCommand), + new byte[] { (byte)NmtCommand.ResetNode, 0x01 }, isExtendedFrame: false)); + await inHandler.Task.WithTimeoutAsync(ShortTimeout); + + var outside = Task.Run(node.Dispose); // wins; the actor is busy, so its cleanup waits + var deadline = DateTime.UtcNow + ShortTimeout; + while (true) + { + try { _ = node.SendNmtCommandAsync(NmtCommand.Start, 0x7F); } + catch (ObjectDisposedException) { break; } // the outside call has begun + if (DateTime.UtcNow > deadline) throw new TimeoutException("The outside Dispose never began."); + await Task.Delay(5); + } + + outsideHasWon.Set(); + + (await ended.Task.WithTimeoutAsync(ShortTimeout)).Should().BeTrue("the subscriber's call ended the transfer"); + await outside.WithTimeoutAsync(ShortTimeout); + } + + // The reset goes on after its subscribers: it announces the node with a boot-up. A subscriber + // that disposed the node ends it there, as a node that is gone announces nothing. The service + // is the caller's, so a boot-up that was wrongly sent would reach the bus. (With the guard + // removed this test still passes: nothing observable happens past the subscriber today, so it + // pins the behaviour and covers the line, and proves no mutation.) + [Fact] + public async Task A_Reset_Does_Not_Announce_A_Node_That_Its_Subscriber_Disposed() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var bootups = 0; + rawBus.FrameObserved += (_, e) => + { + if ((uint)e.CanFrame.ID == CanOpenCobId.HeartbeatBase + 0x01) Interlocked.Increment(ref bootups); + }; + var service = resources.Add(new CanBusService(busA)); + var node = CanOpen.OpenNode(service, nodeId: 0x01, leaveOpen: true); + var disposed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reports = new System.Collections.Concurrent.ConcurrentQueue(); + node.BackgroundExceptionOccurred += (_, ex) => reports.Enqueue(ex); + node.ApplicationReset += (_, _) => + { + node.Dispose(); + disposed.TrySetResult(true); + }; + var atOpen = SpinWait.SpinUntil(() => Volatile.Read(ref bootups) >= 1, ShortTimeout); + atOpen.Should().BeTrue("the node announces itself when it opens"); + + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.NmtCommand), + new byte[] { (byte)NmtCommand.ResetNode, 0x01 }, isExtendedFrame: false)); + await disposed.Task.WithTimeoutAsync(ShortTimeout); + + // No signal says "nothing more is coming", so this is a negative window: it can only pass + // falsely on a slow host, never fail falsely. + await Task.Delay(TimeSpan.FromMilliseconds(300)); + Volatile.Read(ref bootups).Should().Be(1, "only the boot-up from opening; the reset's was not sent"); + reports.Should().BeEmpty("the reset stopped at the subscriber instead of arming timers on a disposed node"); + } + + // Subscribers of ApplicationReset are called one after the other; one that disposes the node + // ends the round. + [Fact] + public async Task A_Reset_Does_Not_Call_The_Next_Subscriber_After_One_Disposed_The_Node() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var first = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var secondCalled = 0; + node.ApplicationReset += (_, _) => + { + node.Dispose(); + first.TrySetResult(true); + }; + node.ApplicationReset += (_, _) => Interlocked.Increment(ref secondCalled); + + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.NmtCommand), + new byte[] { (byte)NmtCommand.ResetNode, 0x01 }, isExtendedFrame: false)); + await first.Task.WithTimeoutAsync(ShortTimeout); + + // A negative window: it can only pass falsely on a slow host, never fail falsely. + await Task.Delay(TimeSpan.FromMilliseconds(300)); + Volatile.Read(ref secondCalled).Should().Be(0); + } + + // Events still queued when a subscriber disposes the node are dropped, not delivered after + // the call has returned. The first subscriber holds the pump until the second heartbeat is + // certainly queued behind it. + [Fact] + public async Task Events_Queued_Behind_A_Subscriber_That_Disposes_The_Node_Are_Not_Delivered() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var delivered = 0; + var disposed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + node.HeartbeatReceived += (_, _) => + { + if (Interlocked.Increment(ref delivered) == 1) + { + Thread.Sleep(300); // the second frame is queued behind this event by now + node.Dispose(); + disposed.TrySetResult(true); + } + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + SendHeartbeat(rawBus, 0x12, 0x05); + await disposed.Task.WithTimeoutAsync(ShortTimeout); + + await Task.Delay(TimeSpan.FromMilliseconds(300)); + Volatile.Read(ref delivered).Should().Be(1, "the second heartbeat was queued but the node was disposed"); + ((CanOpenNode)node).QueuedEventCount.Should().Be(0, "the events skipped by the stopped pump are not left linked to the node"); + } + + // One event, several subscribers: the one that disposes the node ends the round. + [Fact] + public async Task The_Next_Subscriber_Of_An_Event_Is_Not_Called_After_One_Disposed_The_Node() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var first = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var secondCalled = 0; + node.HeartbeatReceived += (_, _) => + { + node.Dispose(); + first.TrySetResult(true); + }; + node.HeartbeatReceived += (_, _) => Interlocked.Increment(ref secondCalled); + + SendHeartbeat(rawBus, 0x11, 0x05); + await first.Task.WithTimeoutAsync(ShortTimeout); + + // A negative window: it can only pass falsely on a slow host, never fail falsely. + await Task.Delay(TimeSpan.FromMilliseconds(300)); + Volatile.Read(ref secondCalled).Should().Be(0); + } + + // An outside Dispose that gives up waiting for a subscriber that keeps the pump (after its two + // seconds) must not see the events queued behind that subscriber delivered when it returns. + [Fact] + public async Task Events_Queued_Behind_A_Subscriber_That_Outlasts_Dispose_Are_Not_Delivered() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + using var release = new ManualResetEventSlim(); + var inHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var delivered = 0; + node.HeartbeatReceived += (_, _) => + { + if (Interlocked.Increment(ref delivered) == 1) + { + inHandler.TrySetResult(true); + release.Wait(TimeSpan.FromSeconds(30)); + } + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + await inHandler.Task.WithTimeoutAsync(ShortTimeout); + SendHeartbeat(rawBus, 0x12, 0x05); // queued behind the blocked subscriber + await Task.Delay(TimeSpan.FromMilliseconds(300)); // a negative window again: the frame is queued by now + + await Task.Run(node.Dispose); // gives up on the pump after its join timeout + release.Set(); + await Task.Delay(TimeSpan.FromMilliseconds(300)); + + Volatile.Read(ref delivered).Should().Be(1, "the second heartbeat was queued but the node was disposed"); + } + + // The disposal itself calls out to the subscription (here: a hook, as the actor's shutdown-timeout + // report would be) on the thread that is finishing it. A DisposeAsync made from there cannot + // wait for the end of the disposal, which waits for it. + [Fact] + public async Task A_DisposeAsync_From_Inside_The_Disposal_Does_Not_Wait_For_Its_End() + { + var starved = new StarvedReaderBusService(); + var node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + var inner = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + starved.OnSubscriptionDisposed = () => + { + node.DisposeAsync().AsTask().Wait(); // would wait for the disposal that is running this + inner.TrySetResult(true); + }; + + await node.DisposeAsync().AsTask().WithTimeoutAsync(ShortTimeout); + + (await inner.Task.WithTimeoutAsync(ShortTimeout)).Should().BeTrue(); + } + + // A reader failure with two subscribers of BackgroundExceptionOccurred: the first disposes + // the node and the second is not called. + [Fact] + public async Task The_Next_Subscriber_Of_A_Background_Exception_Is_Not_Called_After_One_Disposed_The_Node() + { + var starved = new StarvedReaderBusService { FramesFault = new InvalidOperationException("the demux broke") }; + var node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: false, new ManualTimeSource()); + var first = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var secondCalled = 0; + node.BackgroundExceptionOccurred += (_, _) => + { + node.Dispose(); + first.TrySetResult(true); + }; + node.BackgroundExceptionOccurred += (_, _) => Interlocked.Increment(ref secondCalled); + + starved.WakeReader(); + await first.Task.WithTimeoutAsync(ShortTimeout); + + // A negative window: it can only pass falsely on a slow host, never fail falsely. + await Task.Delay(TimeSpan.FromMilliseconds(300)); + Volatile.Read(ref secondCalled).Should().Be(0); + } + + // An outside Dispose has won while the first subscriber is running on the pump; that + // subscriber's own Dispose loses, and the second subscriber is still not called. + [Fact] + public async Task The_Next_Subscriber_Is_Not_Called_When_The_Disposing_One_Lost_The_Race() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var inHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var done = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var outsideHasWon = new ManualResetEventSlim(); + var secondCalled = 0; + node.HeartbeatReceived += (_, _) => + { + inHandler.TrySetResult(true); + outsideHasWon.Wait(ShortTimeout); + node.Dispose(); // loses + done.TrySetResult(true); + }; + node.HeartbeatReceived += (_, _) => Interlocked.Increment(ref secondCalled); + + SendHeartbeat(rawBus, 0x11, 0x05); + await inHandler.Task.WithTimeoutAsync(ShortTimeout); + var outside = Task.Run(node.Dispose); // wins, and waits for the pump + var deadline = DateTime.UtcNow + ShortTimeout; + while (true) + { + try { _ = node.SendNmtCommandAsync(NmtCommand.Start, 0x7F); } + catch (ObjectDisposedException) { break; } + if (DateTime.UtcNow > deadline) throw new TimeoutException("The outside Dispose never began."); + await Task.Delay(5); + } + + outsideHasWon.Set(); + await done.Task.WithTimeoutAsync(ShortTimeout); + await outside.WithTimeoutAsync(ShortTimeout); + + Volatile.Read(ref secondCalled).Should().Be(0); + } + + // The disposal of node A calls out on its own thread (the hook), and the subscriber disposes + // node B before it asks A to dispose again. B's disposal must not take A's marker with it. + [Fact] + public async Task Disposing_Another_Node_From_Inside_A_Disposal_Keeps_The_Marker_Of_The_Outer_One() + { + var starvedA = new StarvedReaderBusService(); + var nodeA = new CanOpenNode(starvedA, 0x01, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + var starvedB = new StarvedReaderBusService(); + var nodeB = new CanOpenNode(starvedB, 0x02, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + var inner = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + starvedA.OnSubscriptionDisposed = () => + { + nodeB.Dispose(); + nodeA.DisposeAsync().AsTask().Wait(); // would wait for the disposal that is running this + inner.TrySetResult(true); + }; + + await nodeA.DisposeAsync().AsTask().WithTimeoutAsync(ShortTimeout); + + (await inner.Task.WithTimeoutAsync(ShortTimeout)).Should().BeTrue(); + starvedB.IsDisposed.Should().BeTrue(); + } + + // Both disposals are active on one thread at once: A's hook disposes B, and B's own hook asks A + // to dispose again while B is still finishing. A's marker must still be there. + [Fact] + public async Task A_Callback_Of_An_Inner_Disposal_Still_Sees_The_Outer_Disposal_As_Running() + { + var starvedA = new StarvedReaderBusService(); + var nodeA = new CanOpenNode(starvedA, 0x01, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + var starvedB = new StarvedReaderBusService(); + var nodeB = new CanOpenNode(starvedB, 0x02, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + var inner = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + starvedB.OnSubscriptionDisposed = () => + { + nodeA.DisposeAsync().AsTask().Wait(); // would wait for the disposal that is running this + inner.TrySetResult(true); + }; + starvedA.OnSubscriptionDisposed = () => nodeB.Dispose(); + + await nodeA.DisposeAsync().AsTask().WithTimeoutAsync(ShortTimeout); + + (await inner.Task.WithTimeoutAsync(ShortTimeout)).Should().BeTrue(); + } + + // A subscriber that throws is isolated: the report goes on to the others. A report made while + // the node is being disposed (the pump is still delivering) goes to all of them, since none + // of them caused the disposal. + [Fact] + public async Task A_Background_Exception_Reaches_Every_Subscriber_Even_One_That_Throws_And_During_Disposal() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var inHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var outsideHasWon = new ManualResetEventSlim(); + var reportsAfterThrow = 0; + var thrown = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + node.BackgroundExceptionOccurred += (_, _) => throw new InvalidOperationException("a subscriber that throws"); + node.BackgroundExceptionOccurred += (_, _) => + { + Interlocked.Increment(ref reportsAfterThrow); + thrown.TrySetResult(true); + }; + node.HeartbeatReceived += (_, _) => + { + inHandler.TrySetResult(true); + outsideHasWon.Wait(ShortTimeout); + throw new InvalidOperationException("reported while the node is being disposed"); + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + await inHandler.Task.WithTimeoutAsync(ShortTimeout); + var outside = Task.Run(node.Dispose); + var deadline = DateTime.UtcNow + ShortTimeout; + while (true) + { + try { _ = node.SendNmtCommandAsync(NmtCommand.Start, 0x7F); } + catch (ObjectDisposedException) { break; } + if (DateTime.UtcNow > deadline) throw new TimeoutException("The outside Dispose never began."); + await Task.Delay(5); + } + + outsideHasWon.Set(); + await thrown.Task.WithTimeoutAsync(ShortTimeout); + await outside.WithTimeoutAsync(ShortTimeout); + + Volatile.Read(ref reportsAfterThrow).Should().Be(1); + } + + // A subscriber that outlasts Dispose's join and then throws: the disposal has finished by + // then, and the failure is not reported to anyone. + [Fact] + public async Task A_Failure_After_The_Disposal_Has_Finished_Is_Not_Reported() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var busA = resources.Add(Open(session, 1)); + var rawBus = resources.Add(Open(session, 2)); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + using var release = new ManualResetEventSlim(); + var inHandler = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var thrown = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reported = 0; + node.BackgroundExceptionOccurred += (_, _) => Interlocked.Increment(ref reported); + node.HeartbeatReceived += (_, _) => + { + inHandler.TrySetResult(true); + release.Wait(TimeSpan.FromSeconds(30)); + try { throw new InvalidOperationException("late"); } + finally { thrown.TrySetResult(true); } + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + await inHandler.Task.WithTimeoutAsync(ShortTimeout); + + await Task.Run(node.Dispose); // gives up on the pump after its join timeout + release.Set(); + await thrown.Task.WithTimeoutAsync(ShortTimeout); + // The report would follow the throw within moments; a negative window, as above. + await Task.Delay(TimeSpan.FromMilliseconds(300)); + + Volatile.Read(ref reported).Should().Be(0); + } + + // The service's Dispose throws: the first call faults, and a second call that was waiting for + // that disposal is released instead of waiting for ever. + [Fact] + public async Task A_Disposal_That_Throws_Still_Releases_A_Concurrent_DisposeAsync() + { + var starved = new StarvedReaderBusService { DisposeFault = new InvalidOperationException("the service would not dispose") }; + var node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: true, new ManualTimeSource()); + + var first = node.DisposeAsync(); + var second = node.DisposeAsync(); + + await second.AsTask().WithTimeoutAsync(ShortTimeout); + await Assert.ThrowsAsync(() => first.AsTask().WithTimeoutAsync(ShortTimeout)); + } + + // Everything a node can have open when it is disposed: a segmented and a block download as a + // server, a classic and a block upload as a client, and node-guarding consumers with and + // without a reply. Each is torn down, and the client's transfers end with ObjectDisposedException. + [Fact] + public async Task Disposing_A_Node_With_Open_Sessions_And_Consumers_Ends_Them() + { + using var resources = new DisposeBag(); + var session = NewSession(); + var rawBus = resources.Add(Open(session, 5)); + var reports = new System.Collections.Concurrent.ConcurrentQueue(); + + var busOfClassicServer = resources.Add(Open(session, 1)); + using var classicServer = CanOpen.OpenNode(busOfClassicServer, nodeId: 0x01); + var busOfBlockServer = resources.Add(Open(session, 2)); + using var blockServer = CanOpen.OpenNode(busOfBlockServer, nodeId: 0x02); + var busOfClient = resources.Add(Open(session, 3)); + using var client = CanOpen.OpenNode(busOfClient, nodeId: 0x03); + foreach (var node in new[] { classicServer, blockServer, client }) + node.BackgroundExceptionOccurred += (_, ex) => reports.Enqueue(ex); + classicServer.ObjectDictionary.AddDomain(0x2100, 0x00, new byte[4]); + blockServer.ObjectDictionary.AddDomain(0x2100, 0x00, new byte[4]); + + var classicTap = new FrameTap(rawBus, CanOpenCobId.SdoTx(0x01)); + var blockTap = new FrameTap(rawBus, CanOpenCobId.SdoTx(0x02)); + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.SdoRx(0x01)), + SdoFrames.BuildDownloadInit(0x2100, 0x00, new byte[20]), isExtendedFrame: false)); + classicTap.Next(ShortTimeout); + rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.SdoRx(0x02)), + SdoBlockFrames.BuildBlockDownloadInit(0x2100, 0x00, clientCrcSupported: false, sizeIndicated: true, totalSize: 20), + isExtendedFrame: false)); + blockTap.Next(ShortTimeout); + + PeerSdoLaboratory.Bind(client, 0x11); + PeerSdoLaboratory.Bind(client, 0x12); + var classicUpload = client.SdoUploadAsync(0x11, 0x2100, 0x00); + var blockUpload = client.SdoUploadAsync(0x12, 0x2100, 0x00, SdoTransferMode.Block); + var replied = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + client.NodeGuardingReceived += (_, e) => + { + if (e.ProducerNodeId == 0x22) replied.TrySetResult(true); + }; + client.StartNodeGuardingConsumer(0x21, TimeSpan.FromSeconds(30), 3); // never answers + client.StartNodeGuardingConsumer(0x22, TimeSpan.FromSeconds(30), 3); + rawBus.Transmit(CanFrame.Classic(unchecked((int)(CanOpenCobId.HeartbeatBase + 0x22)), + new byte[] { 0x7F }, isExtendedFrame: false)); + await replied.Task.WithTimeoutAsync(ShortTimeout); // its life time is armed now + + classicServer.Dispose(); + blockServer.Dispose(); + client.Dispose(); + + await Assert.ThrowsAsync(() => classicUpload.WithTimeoutAsync(ShortTimeout)); + await Assert.ThrowsAsync(() => blockUpload.WithTimeoutAsync(ShortTimeout)); + reports.Should().BeEmpty(); + } + + private sealed class DisposeBag : IDisposable + { + private readonly System.Collections.Generic.List _items = new(); + + public T Add(T item) where T : IDisposable + { + _items.Add(item); + return item; + } + + public void Dispose() + { + for (var i = _items.Count - 1; i >= 0; i--) _items[i].Dispose(); + } + } + + [Fact] + public async Task DisposeAsync_Completes_Open_Transfers_And_Is_Idempotent() + { + var session = NewSession(); + using var busA = Open(session, 1); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + PeerSdoLaboratory.Bind(node, 0x02); + var upload = node.SdoUploadAsync(0x02, 0x2100, 0x00); // nobody answers + + await node.DisposeAsync(); + await node.DisposeAsync(); + + await Assert.ThrowsAsync(() => upload.WithTimeoutAsync(ShortTimeout)); + Assert.Throws(() => _ = node.State); + } + + [Fact] + public async Task A_Node_Can_Be_Used_With_Await_Using_Through_IAsyncDisposable() + { + var session = NewSession(); + using var busA = Open(session, 1); + ICanOpenNode captured; + await using (var node = (IAsyncDisposable)CanOpen.OpenNode(busA, nodeId: 0x01)) + { + captured = (ICanOpenNode)node; + captured.State.Should().NotBe(default); + } + + Assert.Throws(() => _ = captured.State); + } + +#if NET5_0_OR_GREATER + // An ICanOpenNode that is not this library's (an interface is public, anyone may implement it) + // has no DisposeAsync of its own; the extension disposes it on the thread pool. + [Fact] + public async Task DisposeAsync_On_A_Node_That_Is_Not_This_Librarys_Disposes_It_On_The_Thread_Pool() + { + var foreign = System.Reflection.DispatchProxy.Create(); + + await foreign.DisposeAsync(); + + ((ForeignNode)foreign).Disposed.Should().BeTrue(); + } + + private class ForeignNode : System.Reflection.DispatchProxy + { + public bool Disposed { get; private set; } + + protected override object? Invoke(System.Reflection.MethodInfo? targetMethod, object?[]? args) + { + if (targetMethod?.Name == nameof(IDisposable.Dispose)) Disposed = true; + return null; + } + } +#endif + + [Fact] + public async Task DisposeAsync_Rejects_A_Null_Node() + { + await Assert.ThrowsAsync(async () => await CanOpenNodeExtensions.DisposeAsync(null!)); + } + + private sealed class FrameTap : IDisposable + { + private readonly ICanBus _bus; + private readonly uint _cobId; + private readonly System.Collections.Concurrent.BlockingCollection _frames = new(); + + public FrameTap(ICanBus bus, uint cobId) + { + _bus = bus; + _cobId = cobId; + _bus.FrameObserved += OnFrame; + } + + private void OnFrame(object? sender, CanReceiveDataView e) + { + if ((uint)e.CanFrame.ID == _cobId) _frames.Add(e.CanFrame.Data.ToArray()); + } + + public byte[] Next(TimeSpan timeout) + => _frames.TryTake(out var frame, timeout) + ? frame + : throw new TimeoutException($"No frame on COB-ID 0x{_cobId:X3} within {timeout}."); + + public void Dispose() + { + _bus.FrameObserved -= OnFrame; + _frames.Dispose(); + } + } +} diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index e89ae3e..b8f0e7b 100644 --- a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs @@ -1095,13 +1095,21 @@ public async Task A_Failing_Subscription_Ends_The_Inbox_With_Its_Failure() using var channel = new IsoTpChannel(service, IsoTpEndpoint.Normal(0x123, 0x321), FastOptions(), ownsService: false, actor); var reported = new List(); - channel.BackgroundExceptionOccurred += (_, ex) => { lock (reported) reported.Add(ex); }; + var firstReport = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + channel.BackgroundExceptionOccurred += (_, ex) => + { + lock (reported) reported.Add(ex); + firstReport.TrySetResult(true); + }; var waiting = channel.ReceiveAsync(); service.WakeReader(); Func act = () => waiting.WaitAsync(ShortTimeout); (await act.Should().ThrowAsync()).Which.Message.Should().Be("the demux broke"); + // The reader posts the loss to the inbox and raises the report after it: the receive can + // complete first, so the report is awaited rather than read at once (#270). + await firstReport.Task.WaitAsync(ShortTimeout); lock (reported) reported.Should().ContainSingle().Which.Message.Should().Be("the demux broke"); }