From 7059d99cdc352b7320bbcd9e292100cfa56c7fed Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 18:23:58 +0200 Subject: [PATCH 01/17] feat(canopen): add DisposeAsync, and let a node be disposed from its own event handler ICanOpenNode now also implements IAsyncDisposable. DisposeAsync waits for the reader and the event pump without holding a thread. Dispose and DisposeAsync no longer wait for the event pump when they are called from a subscriber, which runs on it: Dispose used to stall for its two-second join timeout waiting for itself (#251, the CANopen part). Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 103 ++++++++++++--- src/CanKit.Pro.CANopen/ICanOpenNode.cs | 11 +- .../CanKit.Pro.CANopen.approved.txt | 2 +- .../TestCases/CANopen/CanOpenDisposeTests.cs | 117 ++++++++++++++++++ 4 files changed, 214 insertions(+), 19 deletions(-) create mode 100644 tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 4c80bf6..2333a8c 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -628,7 +628,71 @@ public Task SdoDownloadAsync(byte serverNodeId, ushort index, byte subindex, /// public void Dispose() { - if (Interlocked.Exchange(ref _disposed, 1) != 0) return; + if (!BeginDispose()) return; + try { _readerTask.Wait(DisposeJoinTimeout); } catch { /* observed via task; not fatal */ } + + // 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(); + // Called from a subscriber, this is the pump's own thread: it cannot finish while it is + // waiting for itself, so it does not wait, and the pump ends when the subscriber returns. + if (!OnEventPump) + { + try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch { /* observed via task; not fatal */ } + } + + FinishDispose(); + } + + /// + public async ValueTask DisposeAsync() + { + // Read before the first await: after it the continuation is on another thread, and a + // subscriber that blocks on this call is still holding the pump's. + bool onPump = OnEventPump; + if (!BeginDispose()) return; + await JoinQuietlyAsync(_readerTask).ConfigureAwait(false); + CompleteEventQueue(); + if (!onPump) + await JoinQuietlyAsync(_eventPumpTask).ConfigureAwait(false); + FinishDispose(); + } + + private static readonly TimeSpan DisposeJoinTimeout = TimeSpan.FromSeconds(2); + + /// Waits for for at most + /// without holding a thread; a fault is observed and dropped, as the blocking join does. + private static async Task JoinQuietlyAsync(Task task) + { + using var cts = new CancellationTokenSource(); + var timeout = Task.Delay(DisposeJoinTimeout, cts.Token); + try + { + if (await Task.WhenAny(task, timeout).ConfigureAwait(false) == task) + await task.ConfigureAwait(false); + } + catch + { + // observed via task; not fatal + } + finally + { + cts.Cancel(); + } + } + + /// True on the thread that is delivering an event to a subscriber. + private bool OnEventPump => ReferenceEquals(t_deliveringFor, this); + + [ThreadStatic] + private static CanOpenNode? t_deliveringFor; + + /// 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 { /* nothing else to do */ } try @@ -675,16 +739,11 @@ public void Dispose() { // actor already gone; nothing more to do } + return true; + } - try { _readerTask.Wait(TimeSpan.FromSeconds(2)); } catch { /* observed via task; not fatal */ } - - // 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 */ } - + private void FinishDispose() + { _subscription.Dispose(); _actor.Dispose(); _readerCts.Dispose(); @@ -741,17 +800,27 @@ 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). + t_deliveringFor = this; + try { - try - { - raise(); - } - catch (Exception ex) + while (TryDequeueEvent() is { } raise) { - RaiseBackgroundException(ex); + try + { + raise(); + } + catch (Exception ex) + { + RaiseBackgroundException(ex); + } } } + finally + { + t_deliveringFor = null; + } // 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. diff --git a/src/CanKit.Pro.CANopen/ICanOpenNode.cs b/src/CanKit.Pro.CANopen/ICanOpenNode.cs index 91a7a3b..32a292e 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -36,8 +36,17 @@ 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. and +/// both fail every open transfer with , stop the producers +/// and deliver nothing afterwards. Dispose blocks for up to two seconds on the node's +/// reader and event pump to let them finish; DisposeAsync waits for the same without +/// holding a thread, and is the one to use on a thread that must not block. Either may be called +/// from an event handler of this node: the handler runs on the event pump, which is then not waited +/// for, and it ends when the handler returns. +/// /// -public interface ICanOpenNode : IDisposable +public interface ICanOpenNode : IDisposable, IAsyncDisposable { /// Node identifier (1..127) this instance answers as on the bus. byte NodeId { get; } 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..dc59128 100644 --- a/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt +++ b/tests/CanKit.Pro.Tests/ApiApprovals/CanKit.Pro.CANopen.approved.txt @@ -203,7 +203,7 @@ namespace CanKit.Pro.CANopen public byte ProducerNodeId { get; } public System.TimeSpan Timeout { get; } } - public interface ICanOpenNode : System.IDisposable + public interface ICanOpenNode : System.IAsyncDisposable, System.IDisposable { byte? ActiveFlyingMasterNodeId { get; } ushort? ActiveFlyingMasterPriority { get; } 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..b40400e --- /dev/null +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -0,0 +1,117 @@ +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.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)); + + [Fact] + public async Task Dispose_From_An_Event_Handler_Does_Not_Wait_For_The_Pump_It_Runs_On() + { + var session = NewSession(); + using var busA = Open(session, 1); + using var rawBus = Open(session, 2); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var elapsed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + node.HeartbeatReceived += (_, _) => + { + var watch = Stopwatch.StartNew(); + node.Dispose(); + elapsed.TrySetResult(watch.Elapsed); + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + + (await elapsed.Task.WithTimeoutAsync(ShortTimeout)).Should().BeLessThan(NoStall); + } + + // Whether the continuation after DisposeAsync's first await changes thread depends on whether + // the reader task had finished by then, so one run does not always reach the case that matters: + // the pump's own thread has to be recognised before that await. Repeated, a mistake in that + // shows up; a correct implementation passes every time. + [Fact] + public async Task DisposeAsync_From_An_Event_Handler_Does_Not_Wait_For_The_Pump_It_Runs_On() + { + for (var attempt = 0; attempt < 10; attempt++) + { + var session = NewSession(); + using var busA = Open(session, 1); + using var rawBus = Open(session, 2); + var node = CanOpen.OpenNode(busA, nodeId: 0x01); + var elapsed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + node.HeartbeatReceived += (_, _) => + { + var watch = Stopwatch.StartNew(); + node.DisposeAsync().AsTask().Wait(); + elapsed.TrySetResult(watch.Elapsed); + }; + + SendHeartbeat(rawBus, 0x11, 0x05); + + (await elapsed.Task.WithTimeoutAsync(ShortTimeout)).Should().BeLessThan(NoStall); + } + } + + [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() + { + var session = NewSession(); + using var busA = Open(session, 1); + ICanOpenNode captured; + await using (var node = CanOpen.OpenNode(busA, nodeId: 0x01)) + { + captured = node; + node.State.Should().NotBe(default); + } + + Assert.Throws(() => _ = captured.State); + } +} From 4991838ffc6d960f39599716464056fddefba380 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 18:40:26 +0200 Subject: [PATCH 02/17] fix(canopen): dispose from the actor and the reader's failure report as well, not only from the event pump DisposeAsync called from inside one of the node's own callbacks now takes the blocking path, which knows which task it must not wait for (the reader on its failure report, the event pump on an event); the actor's reentrant Dispose stays on the thread it recognises. The join no longer swallows exceptions it cannot see: the two tasks do not fault. Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 84 ++++++++++------- src/CanKit.Pro.CANopen/ICanOpenNode.cs | 5 +- .../Infrastructure/StarvedReaderBusService.cs | 16 +++- .../TestCases/CANopen/CanOpenDisposeTests.cs | 91 ++++++++++++------- 4 files changed, 128 insertions(+), 68 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 2333a8c..f5848cf 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -629,71 +629,87 @@ public Task SdoDownloadAsync(byte serverNodeId, ushort index, byte subindex, public void Dispose() { if (!BeginDispose()) return; - try { _readerTask.Wait(DisposeJoinTimeout); } catch { /* observed via task; not fatal */ } + // 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 */ } + } // 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(); - // Called from a subscriber, this is the pump's own thread: it cannot finish while it is - // waiting for itself, so it does not wait, and the pump ends when the subscriber returns. if (!OnEventPump) { - try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch { /* observed via task; not fatal */ } + try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } } FinishDispose(); } /// - public async ValueTask DisposeAsync() + 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 || _actor.IsOnCurrentActor) + { + Dispose(); + return default; + } + + return DisposeCoreAsync(); + } + + private async ValueTask DisposeCoreAsync() { - // Read before the first await: after it the continuation is on another thread, and a - // subscriber that blocks on this call is still holding the pump's. - bool onPump = OnEventPump; if (!BeginDispose()) return; - await JoinQuietlyAsync(_readerTask).ConfigureAwait(false); + await JoinAsync(_readerTask).ConfigureAwait(false); CompleteEventQueue(); - if (!onPump) - await JoinQuietlyAsync(_eventPumpTask).ConfigureAwait(false); + await JoinAsync(_eventPumpTask).ConfigureAwait(false); FinishDispose(); } private static readonly TimeSpan DisposeJoinTimeout = TimeSpan.FromSeconds(2); /// Waits for for at most - /// without holding a thread; a fault is observed and dropped, as the blocking join does. - private static async Task JoinQuietlyAsync(Task task) + /// 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(); - var timeout = Task.Delay(DisposeJoinTimeout, cts.Token); - try - { - if (await Task.WhenAny(task, timeout).ConfigureAwait(false) == task) - await task.ConfigureAwait(false); - } - catch - { - // observed via task; not fatal - } - finally - { - cts.Cancel(); - } + 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 => ReferenceEquals(t_deliveringFor, this); + /// True on the thread of the reader task while it reports a failed subscription. + private bool OnReader => ReferenceEquals(t_reportingFor, this); + [ThreadStatic] private static CanOpenNode? t_deliveringFor; + [ThreadStatic] + private static CanOpenNode? t_reportingFor; + + private static void MarkDelivering(CanOpenNode? node) => t_deliveringFor = node; + + private static void MarkReporting(CanOpenNode? node) => t_reportingFor = node; + /// 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 { /* nothing else to do */ } + try { _readerCts.Cancel(); } catch (AggregateException) { /* a registered callback threw; nothing else to do */ } try { @@ -784,7 +800,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. + MarkReporting(this); + try { RaiseBackgroundException(ex); } + finally { MarkReporting(null); } + } } // ========================================================================================= @@ -802,7 +824,7 @@ private async Task RunEventPumpAsync() // must still be delivered. Anything outside the delegate is still a bug in the pump. // 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). - t_deliveringFor = this; + MarkDelivering(this); try { while (TryDequeueEvent() is { } raise) @@ -819,7 +841,7 @@ private async Task RunEventPumpAsync() } finally { - t_deliveringFor = null; + MarkDelivering(null); } // 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 diff --git a/src/CanKit.Pro.CANopen/ICanOpenNode.cs b/src/CanKit.Pro.CANopen/ICanOpenNode.cs index 32a292e..283bad7 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -42,8 +42,9 @@ namespace CanKit.Pro.CANopen; /// and deliver nothing afterwards. Dispose blocks for up to two seconds on the node's /// reader and event pump to let them finish; DisposeAsync waits for the same without /// holding a thread, and is the one to use on a thread that must not block. Either may be called -/// from an event handler of this node: the handler runs on the event pump, which is then not waited -/// for, and it ends when the handler returns. +/// 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, IAsyncDisposable diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index e9e8e59..e6be05e 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; } @@ -100,7 +103,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) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index b40400e..8c81948 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -8,6 +8,7 @@ using CanKit.Abstractions.API.Common.Definitions; using CanKit.Core; using CanKit.Pro.CANopen; +using CanKit.Pro.CANopen.Nmt; using CanKit.Pro.RawCan; using CanKit.Pro.Tests.Infrastructure; using Xunit; @@ -37,51 +38,73 @@ private static void SendHeartbeat(ICanBus bus, byte producer, byte state) => bus.Transmit(CanFrame.Classic(unchecked((int)(CanOpenCobId.HeartbeatBase + producer)), new[] { state }, isExtendedFrame: false)); - [Fact] - public async Task Dispose_From_An_Event_Handler_Does_Not_Wait_For_The_Pump_It_Runs_On() + 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) { - var session = NewSession(); - using var busA = Open(session, 1); - using var rawBus = Open(session, 2); - var node = CanOpen.OpenNode(busA, nodeId: 0x01); var elapsed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - node.HeartbeatReceived += (_, _) => + ICanOpenNode? node = null; + void DisposeNode() { var watch = Stopwatch.StartNew(); - node.Dispose(); + if (call == "Dispose") node!.Dispose(); + else node!.DisposeAsync().AsTask().Wait(); elapsed.TrySetResult(watch.Elapsed); - }; - - SendHeartbeat(rawBus, 0x11, 0x05); - - (await elapsed.Task.WithTimeoutAsync(ShortTimeout)).Should().BeLessThan(NoStall); - } + } - // Whether the continuation after DisposeAsync's first await changes thread depends on whether - // the reader task had finished by then, so one run does not always reach the case that matters: - // the pump's own thread has to be recognised before that await. Repeated, a mistake in that - // shows up; a correct implementation passes every time. - [Fact] - public async Task DisposeAsync_From_An_Event_Handler_Does_Not_Wait_For_The_Pump_It_Runs_On() - { - for (var attempt = 0; attempt < 10; attempt++) + ICanBus? busA = null, rawBus = null; + StarvedReaderBusService? starved = null; + try { - var session = NewSession(); - using var busA = Open(session, 1); - using var rawBus = Open(session, 2); - var node = CanOpen.OpenNode(busA, nodeId: 0x01); - var elapsed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - node.HeartbeatReceived += (_, _) => + if (context == "a reader failure report") { - var watch = Stopwatch.StartNew(); - node.DisposeAsync().AsTask().Wait(); - elapsed.TrySetResult(watch.Elapsed); - }; - - SendHeartbeat(rawBus, 0x11, 0x05); + starved = new StarvedReaderBusService { FramesFault = new InvalidOperationException("the demux broke") }; + node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: false, new ManualTimeSource()); + node.BackgroundExceptionOccurred += (_, _) => DisposeNode(); + starved.WakeReader(); + } + else + { + var session = NewSession(); + busA = Open(session, 1); + rawBus = Open(session, 2); + node = CanOpen.OpenNode(busA, nodeId: 0x01); + if (context == "ApplicationReset") + { + 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); } + finally + { + node?.Dispose(); + busA?.Dispose(); + rawBus?.Dispose(); + } } [Fact] From 9792880fc74e06402d0f5de7f603e71f5a50bc29 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 18:55:21 +0200 Subject: [PATCH 03/17] fix(canopen): run the dispose cleanup in place on the actor, make a concurrent DisposeAsync wait for the disposal Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 92 +++++++++++-------- .../Infrastructure/StarvedReaderBusService.cs | 9 +- .../TestCases/CANopen/CanOpenDisposeTests.cs | 86 +++++++++++------ 3 files changed, 120 insertions(+), 67 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index f5848cf..aaac4cc 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -138,6 +138,7 @@ internal sealed partial class CanOpenNode : ICanOpenNode private bool _nodeGuardingProducerToggle; private int _disposed; + private readonly TaskCompletionSource _disposeDone = new(TaskCreationOptions.RunContinuationsAsynchronously); /// public byte NodeId => _nodeId; @@ -670,7 +671,14 @@ public ValueTask DisposeAsync() private async ValueTask DisposeCoreAsync() { - if (!BeginDispose()) return; + // 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); @@ -711,45 +719,13 @@ 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 { - _actor.Post(() => - { - _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(); - - _sdoBlockServer?.Deadline?.Dispose(); - _sdoBlockServer = null; - foreach (var kv in _sdoBlockClients) - { - kv.Value.Deadline?.Dispose(); - kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); - } - _sdoBlockClients.Clear(); - - foreach (var kv in _nodeGuardingConsumers) - { - kv.Value.PollHandle?.Dispose(); - kv.Value.LifeTimeDeadline?.Dispose(); - } - _nodeGuardingConsumers.Clear(); - }); + if (_actor.IsOnCurrentActor) CleanUpOnActor(); + else _actor.Post(CleanUpOnActor); } catch (ObjectDisposedException) { @@ -758,6 +734,45 @@ private bool BeginDispose() return true; } + /// 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(); + + _sdoServer?.Deadline?.Dispose(); + _sdoServer = null; + foreach (var kv in _sdoClients) + { + kv.Value.Deadline?.Dispose(); + kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); + } + _sdoClients.Clear(); + + _sdoBlockServer?.Deadline?.Dispose(); + _sdoBlockServer = null; + foreach (var kv in _sdoBlockClients) + { + kv.Value.Deadline?.Dispose(); + kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); + } + _sdoBlockClients.Clear(); + + foreach (var kv in _nodeGuardingConsumers) + { + kv.Value.PollHandle?.Dispose(); + kv.Value.LifeTimeDeadline?.Dispose(); + } + _nodeGuardingConsumers.Clear(); + } + private void FinishDispose() { _subscription.Dispose(); @@ -765,6 +780,7 @@ private void FinishDispose() _readerCts.Dispose(); if (_ownsService) _service.Dispose(); + _disposeDone.TrySetResult(true); } // ========================================================================================= diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index e6be05e..94a087e 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -95,7 +95,14 @@ 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; } + + public void Dispose() + { + IsDisposed = true; + _frames.Writer.TryComplete(); + } private sealed class Sub : ISubscription { diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 8c81948..24b3051 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -57,53 +57,83 @@ private static void SendHeartbeat(ICanBus bus, byte producer, byte state) [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); } - ICanBus? busA = null, rawBus = null; - StarvedReaderBusService? starved = null; - try + if (context == "a reader failure report") { - 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") { - starved = new StarvedReaderBusService { FramesFault = new InvalidOperationException("the demux broke") }; - node = new CanOpenNode(starved, 0x01, new CanOpenNodeOptions(), ownsService: false, new ManualTimeSource()); - node.BackgroundExceptionOccurred += (_, _) => DisposeNode(); - starved.WakeReader(); + 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 { - var session = NewSession(); - busA = Open(session, 1); - rawBus = Open(session, 2); - node = CanOpen.OpenNode(busA, nodeId: 0x01); - if (context == "ApplicationReset") - { - 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); - } + 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()); - (await elapsed.Task.WithTimeoutAsync(ShortTimeout)).Should().BeLessThan(NoStall); + var first = node.DisposeAsync(); + var second = node.DisposeAsync(); + await second.AsTask().WithTimeoutAsync(ShortTimeout); + + starved.IsDisposed.Should().BeTrue(); + await first.AsTask().WithTimeoutAsync(ShortTimeout); + } + + 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; } - finally + + public void Dispose() { - node?.Dispose(); - busA?.Dispose(); - rawBus?.Dispose(); + for (var i = _items.Count - 1; i >= 0; i--) _items[i].Dispose(); } } From 47145c42e8ad37f304cb6740dc26b364a42c8612 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:10:53 +0200 Subject: [PATCH 04/17] fix(canopen): release a waiting DisposeAsync when the disposal throws; cover the cleanup with open sessions Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 17 +++- .../Infrastructure/StarvedReaderBusService.cs | 4 + .../TestCases/CANopen/CanOpenDisposeTests.cs | 99 +++++++++++++++++++ 3 files changed, 115 insertions(+), 5 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index aaac4cc..b8d6c62 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -775,12 +775,19 @@ private void CleanUpOnActor() private void FinishDispose() { - _subscription.Dispose(); - _actor.Dispose(); - _readerCts.Dispose(); + try + { + _subscription.Dispose(); + _actor.Dispose(); + _readerCts.Dispose(); - if (_ownsService) _service.Dispose(); - _disposeDone.TrySetResult(true); + if (_ownsService) _service.Dispose(); + } + finally + { + // Whatever a Dispose above threw, a caller waiting for this disposal is released. + _disposeDone.TrySetResult(true); + } } // ========================================================================================= diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index 94a087e..9d7f439 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -98,9 +98,13 @@ public IReadOnlyList FindOverlappingFilterSubscriptions() /// Whether has been called. public bool IsDisposed { get; private 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(); } diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 24b3051..6ef8eb2 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -9,6 +9,7 @@ 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; @@ -121,6 +122,74 @@ public async Task A_Concurrent_DisposeAsync_Returns_When_The_Disposal_Has_Finish await first.AsTask().WithTimeoutAsync(ShortTimeout); } + // 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 classicServer = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 1)), nodeId: 0x01)); + var blockServer = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 2)), nodeId: 0x02)); + var client = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 3)), 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(); @@ -167,4 +236,34 @@ public async Task A_Node_Can_Be_Used_With_Await_Using() Assert.Throws(() => _ = captured.State); } + + 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(); + } + } } From c636a17ceb3b71bd15d6158d32a81b4ef92e8c5a Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:22:08 +0200 Subject: [PATCH 05/17] feat(canopen): offer DisposeAsync as an extension instead of widening ICanOpenNode ICanOpenNode stays IDisposable: adding IAsyncDisposable to a published interface breaks every third-party implementer and would need a major release. The nodes this library creates implement IAsyncDisposable, and CanOpenNodeExtensions.DisposeAsync(ICanOpenNode) reaches it (a node that is not this library's is disposed on the thread pool). Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 2 +- .../CanOpenNodeExtensions.cs | 33 ++++++++++++++++ src/CanKit.Pro.CANopen/ICanOpenNode.cs | 17 ++++---- .../CanKit.Pro.CANopen.approved.txt | 6 ++- .../TestCases/CANopen/CanOpenDisposeTests.cs | 39 +++++++++++++++++-- 5 files changed, 82 insertions(+), 15 deletions(-) create mode 100644 src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index b8d6c62..00947da 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; diff --git a/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs new file mode 100644 index 0000000..2602c31 --- /dev/null +++ b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs @@ -0,0 +1,33 @@ +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. 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 283bad7..0636f19 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -37,17 +37,16 @@ namespace CanKit.Pro.CANopen; /// dictionary never describes behaviour the node does not have. /// /// -/// Disposing. and -/// both fail every open transfer with , stop the producers -/// and deliver nothing afterwards. Dispose blocks for up to two seconds on the node's -/// reader and event pump to let them finish; DisposeAsync waits for the same without -/// holding a thread, and is the one to use on a thread that must not block. 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. +/// Disposing. fails every open transfer with +/// , stops the producers and delivers nothing afterwards. 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, IAsyncDisposable +public interface ICanOpenNode : IDisposable { /// Node identifier (1..127) this instance answers as on the bus. byte NodeId { get; } 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 dc59128..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() { } @@ -203,7 +207,7 @@ namespace CanKit.Pro.CANopen public byte ProducerNodeId { get; } public System.TimeSpan Timeout { get; } } - public interface ICanOpenNode : System.IAsyncDisposable, System.IDisposable + public interface ICanOpenNode : System.IDisposable { byte? ActiveFlyingMasterNodeId { get; } ushort? ActiveFlyingMasterPriority { get; } diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 6ef8eb2..a239864 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -223,20 +223,51 @@ public async Task DisposeAsync_Completes_Open_Transfers_And_Is_Idempotent() } [Fact] - public async Task A_Node_Can_Be_Used_With_Await_Using() + 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 = CanOpen.OpenNode(busA, nodeId: 0x01)) + await using (var node = (IAsyncDisposable)CanOpen.OpenNode(busA, nodeId: 0x01)) { - captured = node; - node.State.Should().NotBe(default); + 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)(object)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; From 22dca28d04e8ab244ca4a22ac2c6cc048c86a465 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:36:45 +0200 Subject: [PATCH 06/17] fix(canopen): a subscriber that loses the disposal race still ends the transfers; bound the wait of a losing DisposeAsync Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 29 ++++++---- .../TestCases/CANopen/CanOpenDisposeTests.cs | 54 +++++++++++++++++-- 2 files changed, 70 insertions(+), 13 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 00947da..7da7626 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -629,7 +629,14 @@ public Task SdoDownloadAsync(byte serverNodeId, ushort index, byte subindex, /// public void Dispose() { - if (!BeginDispose()) return; + 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 @@ -673,10 +680,12 @@ 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. + // still waiting for the reader. The wait is bounded: a subscriber that the first + // disposal itself calls (the actor reports a shutdown timeout on the disposing thread) + // may land here, and the disposal cannot finish before it returns. if (!BeginDispose()) { - await _disposeDone.Task.ConfigureAwait(false); + await JoinAsync(_disposeDone.Task).ConfigureAwait(false); return; } await JoinAsync(_readerTask).ConfigureAwait(false); @@ -747,32 +756,34 @@ private void CleanUpOnActor() CancelFlyingMasterDeadline(); CancelBootUp(); - _sdoServer?.Deadline?.Dispose(); + Release(_sdoServer?.Deadline); _sdoServer = null; foreach (var kv in _sdoClients) { - kv.Value.Deadline?.Dispose(); + Release(kv.Value.Deadline); kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); } _sdoClients.Clear(); - _sdoBlockServer?.Deadline?.Dispose(); + Release(_sdoBlockServer?.Deadline); _sdoBlockServer = null; foreach (var kv in _sdoBlockClients) { - kv.Value.Deadline?.Dispose(); + Release(kv.Value.Deadline); kv.Value.Tcs.TrySetException(new ObjectDisposedException(nameof(CanOpenNode))); } _sdoBlockClients.Clear(); foreach (var kv in _nodeGuardingConsumers) { - kv.Value.PollHandle?.Dispose(); - kv.Value.LifeTimeDeadline?.Dispose(); + Release(kv.Value.PollHandle); + Release(kv.Value.LifeTimeDeadline); } _nodeGuardingConsumers.Clear(); } + private static void Release(IDisposable? handle) => handle?.Dispose(); + private void FinishDispose() { try diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index a239864..042313d 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -122,6 +122,49 @@ public async Task A_Concurrent_DisposeAsync_Returns_When_The_Disposal_Has_Finish 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 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] @@ -148,9 +191,12 @@ public async Task Disposing_A_Node_With_Open_Sessions_And_Consumers_Ends_Them() var rawBus = resources.Add(Open(session, 5)); var reports = new System.Collections.Concurrent.ConcurrentQueue(); - var classicServer = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 1)), nodeId: 0x01)); - var blockServer = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 2)), nodeId: 0x02)); - var client = resources.Add(CanOpen.OpenNode(resources.Add(Open(session, 3)), nodeId: 0x03)); + var busOfClassicServer = resources.Add(Open(session, 1)); + var classicServer = resources.Add(CanOpen.OpenNode(busOfClassicServer, nodeId: 0x01)); + var busOfBlockServer = resources.Add(Open(session, 2)); + var blockServer = resources.Add(CanOpen.OpenNode(busOfBlockServer, nodeId: 0x02)); + var busOfClient = resources.Add(Open(session, 3)); + var client = resources.Add(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]); @@ -247,7 +293,7 @@ public async Task DisposeAsync_On_A_Node_That_Is_Not_This_Librarys_Disposes_It_O await foreign.DisposeAsync(); - ((ForeignNode)(object)foreign).Disposed.Should().BeTrue(); + ((ForeignNode)foreign).Disposed.Should().BeTrue(); } private class ForeignNode : System.Reflection.DispatchProxy From db290f5b90ddd0c4d0c763d08640354e0ba5501d Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:52:52 +0200 Subject: [PATCH 07/17] fix(canopen): stop a reset at a subscriber that disposed the node; wait for a running disposal for as long as it can take Co-Authored-By: Claude Sonnet 5.5 --- .../CanOpenNode.CommunicationProfile.cs | 3 ++ src/CanKit.Pro.CANopen/CanOpenNode.cs | 15 ++++-- .../TestCases/CANopen/CanOpenDisposeTests.cs | 47 +++++++++++++++++-- 3 files changed, 57 insertions(+), 8 deletions(-) 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.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 7da7626..650856d 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -680,12 +680,13 @@ 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. The wait is bounded: a subscriber that the first + // still waiting for the reader. The wait is bounded, by more than the first disposal can + // take (reader, pump and the actor's own shutdown timeout): a subscriber that the first // disposal itself calls (the actor reports a shutdown timeout on the disposing thread) // may land here, and the disposal cannot finish before it returns. if (!BeginDispose()) { - await JoinAsync(_disposeDone.Task).ConfigureAwait(false); + await JoinAsync(_disposeDone.Task, ConcurrentDisposeWait).ConfigureAwait(false); return; } await JoinAsync(_readerTask).ConfigureAwait(false); @@ -696,13 +697,17 @@ private async ValueTask DisposeCoreAsync() private static readonly TimeSpan DisposeJoinTimeout = TimeSpan.FromSeconds(2); - /// Waits for for at most + // Reader join, pump join and the actor's five-second shutdown timeout, with room. + private static readonly TimeSpan ConcurrentDisposeWait = TimeSpan.FromSeconds(12); + + /// Waits for for at most (default + /// ) /// 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) + private static async Task JoinAsync(Task task, TimeSpan? timeout = null) { using var cts = new CancellationTokenSource(); - await Task.WhenAny(task, Task.Delay(DisposeJoinTimeout, cts.Token)).ConfigureAwait(false); + await Task.WhenAny(task, Task.Delay(timeout ?? DisposeJoinTimeout, cts.Token)).ConfigureAwait(false); cts.Cancel(); } diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 042313d..79a5b3b 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -165,6 +165,47 @@ public async Task A_Subscriber_On_The_Actor_That_Loses_The_Disposal_Race_Still_E 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"); + } + // 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] @@ -192,11 +233,11 @@ public async Task Disposing_A_Node_With_Open_Sessions_And_Consumers_Ends_Them() var reports = new System.Collections.Concurrent.ConcurrentQueue(); var busOfClassicServer = resources.Add(Open(session, 1)); - var classicServer = resources.Add(CanOpen.OpenNode(busOfClassicServer, nodeId: 0x01)); + using var classicServer = CanOpen.OpenNode(busOfClassicServer, nodeId: 0x01); var busOfBlockServer = resources.Add(Open(session, 2)); - var blockServer = resources.Add(CanOpen.OpenNode(busOfBlockServer, nodeId: 0x02)); + using var blockServer = CanOpen.OpenNode(busOfBlockServer, nodeId: 0x02); var busOfClient = resources.Add(Open(session, 3)); - var client = resources.Add(CanOpen.OpenNode(busOfClient, nodeId: 0x03)); + 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]); From 5292c1adcbfa46a2d5245c3779bcad0c09cfb33d Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 20:07:32 +0200 Subject: [PATCH 08/17] fix(canopen): drop queued events and stop the reset round when a subscriber disposes the node Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 14 ++++- .../CanOpenNodeExtensions.cs | 5 +- .../TestCases/CANopen/CanOpenDisposeTests.cs | 59 +++++++++++++++++++ 3 files changed, 74 insertions(+), 4 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 650856d..e72443d 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -138,6 +138,7 @@ internal sealed partial class CanOpenNode : ICanOpenNode, IAsyncDisposable private bool _nodeGuardingProducerToggle; private int _disposed; + private bool _pumpStopRequested; private readonly TaskCompletionSource _disposeDone = new(TaskCreationOptions.RunContinuationsAsynchronously); /// @@ -651,7 +652,13 @@ public void Dispose() // arrives afterwards is dropped, timeout and EMCY included. Nothing is delivered past // Dispose. CompleteEventQueue(); - if (!OnEventPump) + if (OnEventPump) + { + // The events still queued are not delivered once this subscriber returns: nothing is + // delivered past Dispose. + Volatile.Write(ref _pumpStopRequested, true); + } + else { try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } } @@ -866,7 +873,7 @@ private async Task RunEventPumpAsync() MarkDelivering(this); try { - while (TryDequeueEvent() is { } raise) + while (!Volatile.Read(ref _pumpStopRequested) && TryDequeueEvent() is { } raise) { try { @@ -2691,6 +2698,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); } } diff --git a/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs index 2602c31..c06f202 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNodeExtensions.cs @@ -19,8 +19,9 @@ public static class CanOpenNodeExtensions /// 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. A node that is not one of this library's is disposed on the thread - /// pool. + /// 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) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 79a5b3b..fd7fe32 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -206,6 +206,65 @@ public async Task A_Reset_Does_Not_Announce_A_Node_That_Its_Subscriber_Disposed( 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"); + } + // 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] From d9c4eae1fa44e377bf72be097dafa4265b2bb34b Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 20:22:27 +0200 Subject: [PATCH 09/17] fix(canopen): end an event's round at the subscriber that disposed the node, stop the pump when its join times out, wait for a running disposal without a bound Co-Authored-By: Claude Sonnet 5.5 --- .../CanOpenNode.FlyingMaster.cs | 2 +- .../CanOpenNode.NodeGuarding.cs | 6 +- src/CanKit.Pro.CANopen/CanOpenNode.cs | 63 +++++++++++++------ .../TestCases/CANopen/CanOpenDisposeTests.cs | 60 ++++++++++++++++++ 4 files changed, 109 insertions(+), 22 deletions(-) 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 e72443d..9527094 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -661,6 +661,7 @@ public void Dispose() else { try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } + StopPumpIfStillRunning(); } FinishDispose(); @@ -674,7 +675,7 @@ public ValueTask DisposeAsync() // 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 || _actor.IsOnCurrentActor) + if (OnEventPump || OnReader || OnDisposal || _actor.IsOnCurrentActor) { Dispose(); return default; @@ -687,40 +688,45 @@ 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. The wait is bounded, by more than the first disposal can - // take (reader, pump and the actor's own shutdown timeout): a subscriber that the first - // disposal itself calls (the actor reports a shutdown timeout on the disposing thread) - // may land here, and the disposal cannot finish before it returns. + // still waiting for the reader. if (!BeginDispose()) { - await JoinAsync(_disposeDone.Task, ConcurrentDisposeWait).ConfigureAwait(false); + 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); - // Reader join, pump join and the actor's five-second shutdown timeout, with room. - private static readonly TimeSpan ConcurrentDisposeWait = TimeSpan.FromSeconds(12); + /// 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 (default - /// ) + /// 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, TimeSpan? timeout = null) + private static async Task JoinAsync(Task task) { using var cts = new CancellationTokenSource(); - await Task.WhenAny(task, Task.Delay(timeout ?? DisposeJoinTimeout, cts.Token)).ConfigureAwait(false); + 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 => ReferenceEquals(t_deliveringFor, 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 => ReferenceEquals(t_finishingFor, this); + /// True on the thread of the reader task while it reports a failed subscription. private bool OnReader => ReferenceEquals(t_reportingFor, this); @@ -730,6 +736,11 @@ private static async Task JoinAsync(Task task, TimeSpan? timeout = null) [ThreadStatic] private static CanOpenNode? t_reportingFor; + [ThreadStatic] + private static CanOpenNode? t_finishingFor; + + private static void MarkFinishing(CanOpenNode? node) => t_finishingFor = node; + private static void MarkDelivering(CanOpenNode? node) => t_deliveringFor = node; private static void MarkReporting(CanOpenNode? node) => t_reportingFor = node; @@ -798,6 +809,7 @@ private void CleanUpOnActor() private void FinishDispose() { + MarkFinishing(this); try { _subscription.Dispose(); @@ -808,6 +820,7 @@ private void FinishDispose() } finally { + MarkFinishing(null); // Whatever a Dispose above threw, a caller waiting for this disposal is released. _disposeDone.TrySetResult(true); } @@ -2629,6 +2642,20 @@ 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); } @@ -2640,7 +2667,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); } @@ -2650,7 +2677,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); } @@ -2660,7 +2687,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); } @@ -2670,7 +2697,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); } }); } @@ -2680,7 +2707,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); } }); } @@ -2711,7 +2738,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/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index fd7fe32..d6d82aa 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -265,6 +265,66 @@ public async Task Events_Queued_Behind_A_Subscriber_That_Disposes_The_Node_Are_N Volatile.Read(ref delivered).Should().Be(1, "the second heartbeat was queued but the node was disposed"); } + // 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 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] From bd88c1251a5ef62eb2efc218067db31b6fdf763f Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 20:38:55 +0200 Subject: [PATCH 10/17] fix(canopen): a losing Dispose from the pump stops delivery too, BackgroundExceptionOccurred rounds end at a disposing subscriber Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 26 ++++-- .../Infrastructure/StarvedReaderBusService.cs | 5 +- .../TestCases/CANopen/CanOpenDisposeTests.cs | 86 +++++++++++++++++++ 3 files changed, 107 insertions(+), 10 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 9527094..d0e152b 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -630,6 +630,9 @@ public Task SdoDownloadAsync(byte serverNodeId, ushort index, byte subindex, /// public void Dispose() { + // 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: @@ -652,13 +655,9 @@ public void Dispose() // arrives afterwards is dropped, timeout and EMCY included. Nothing is delivered past // Dispose. CompleteEventQueue(); - if (OnEventPump) - { - // The events still queued are not delivered once this subscriber returns: nothing is - // delivered past Dispose. - Volatile.Write(ref _pumpStopRequested, true); - } - else + // 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) { try { _eventPumpTask.Wait(DisposeJoinTimeout); } catch (AggregateException) { /* observed via task; not fatal */ } StopPumpIfStillRunning(); @@ -2658,8 +2657,17 @@ private void DeliverToSubscribers(EventHandler? handler, T 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; + // 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) diff --git a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs index 9d7f439..3a36da6 100644 --- a/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs +++ b/tests/CanKit.Pro.Tests/Infrastructure/StarvedReaderBusService.cs @@ -98,6 +98,9 @@ public IReadOnlyList FindOverlappingFilterSubscriptions() /// 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; } @@ -159,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 index d6d82aa..6f7cf2a 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -325,6 +325,92 @@ public async Task Events_Queued_Behind_A_Subscriber_That_Outlasts_Dispose_Are_No 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 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] From c20bef4356ef38520d6673069834b2f34302ebfa Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 20:54:23 +0200 Subject: [PATCH 11/17] fix(canopen): nested disposals restore the outer thread markers; cover the background-exception round Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 36 +++++++--- .../TestCases/CANopen/CanOpenDisposeTests.cs | 70 +++++++++++++++++++ 2 files changed, 97 insertions(+), 9 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index d0e152b..44bdbc2 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -738,11 +738,29 @@ private static async Task JoinAsync(Task task) [ThreadStatic] private static CanOpenNode? t_finishingFor; - private static void MarkFinishing(CanOpenNode? node) => t_finishingFor = node; + private static CanOpenNode? MarkFinishing(CanOpenNode? node) + { + var previous = t_finishingFor; + t_finishingFor = node; + return previous; + } - private static void MarkDelivering(CanOpenNode? node) => t_deliveringFor = node; + // Each marker returns what it replaced, and the caller puts that back: the thread may be + // inside another node's callback (a subscriber of A disposing B), and B's markers must not + // erase A's. + private static CanOpenNode? MarkDelivering(CanOpenNode? node) + { + var previous = t_deliveringFor; + t_deliveringFor = node; + return previous; + } - private static void MarkReporting(CanOpenNode? node) => t_reportingFor = node; + private static CanOpenNode? MarkReporting(CanOpenNode? node) + { + var previous = t_reportingFor; + t_reportingFor = node; + return previous; + } /// Flips the disposed flag and posts the cleanup; false when already disposed. private bool BeginDispose() @@ -808,7 +826,7 @@ private void CleanUpOnActor() private void FinishDispose() { - MarkFinishing(this); + var outerFinishing = MarkFinishing(this); try { _subscription.Dispose(); @@ -819,7 +837,7 @@ private void FinishDispose() } finally { - MarkFinishing(null); + MarkFinishing(outerFinishing); // Whatever a Dispose above threw, a caller waiting for this disposal is released. _disposeDone.TrySetResult(true); } @@ -861,9 +879,9 @@ private async Task RunReaderAsync() catch (Exception ex) { // A subscriber of the report may dispose the node; see Dispose. - MarkReporting(this); + var outer = MarkReporting(this); try { RaiseBackgroundException(ex); } - finally { MarkReporting(null); } + finally { MarkReporting(outer); } } } @@ -882,7 +900,7 @@ private async Task RunEventPumpAsync() // must still be delivered. Anything outside the delegate is still a bug in the pump. // 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). - MarkDelivering(this); + var outerDelivering = MarkDelivering(this); try { while (!Volatile.Read(ref _pumpStopRequested) && TryDequeueEvent() is { } raise) @@ -899,7 +917,7 @@ private async Task RunEventPumpAsync() } finally { - MarkDelivering(null); + MarkDelivering(outerDelivering); } // 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 diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 6f7cf2a..2627069 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -411,6 +411,76 @@ public async Task The_Next_Subscriber_Is_Not_Called_When_The_Disposing_One_Lost_ 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(); + } + + // 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); + } + // 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] From decf158fead730da21cb11b0e63604767eed35a5 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:09:45 +0200 Subject: [PATCH 12/17] fix(canopen): report nothing once the disposal has finished Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 3 ++ .../TestCases/CANopen/CanOpenDisposeTests.cs | 35 +++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 44bdbc2..6279f9d 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -2677,6 +2677,9 @@ private void RaiseBackgroundException(Exception ex) { 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; diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 2627069..191f94a 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -481,6 +481,41 @@ public async Task A_Background_Exception_Reaches_Every_Subscriber_Even_One_That_ 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] From 65830c9b4c87ece29c7fb411a011a907e7a22f34 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:25:07 +0200 Subject: [PATCH 13/17] fix(canopen): drop the events a stopped pump skipped; state the delivery guarantee as 'starts no further delivery' Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 8 ++++++++ src/CanKit.Pro.CANopen/ICanOpenNode.cs | 3 ++- .../TestCases/CANopen/CanOpenDisposeTests.cs | 1 + 3 files changed, 11 insertions(+), 1 deletion(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 6279f9d..111d4ed 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -914,6 +914,14 @@ private async Task RunEventPumpAsync() RaiseBackgroundException(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)) + { + while (TryDequeueEvent() is not null) + { + } + } } finally { diff --git a/src/CanKit.Pro.CANopen/ICanOpenNode.cs b/src/CanKit.Pro.CANopen/ICanOpenNode.cs index 0636f19..e358630 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -38,7 +38,8 @@ namespace CanKit.Pro.CANopen; /// /// /// Disposing. fails every open transfer with -/// , stops the producers and delivers nothing afterwards. It +/// , 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 diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index 191f94a..ef516c2 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -263,6 +263,7 @@ public async Task Events_Queued_Behind_A_Subscriber_That_Disposes_The_Node_Are_N 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. From b443730a1b29dbcea978355fac4a6e220a18c3f5 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:40:04 +0200 Subject: [PATCH 14/17] refactor(canopen): no empty loop body when dropping skipped events; document the actor-timeout limit of the transfer cleanup Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 6 +++--- src/CanKit.Pro.CANopen/ICanOpenNode.cs | 4 +++- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 111d4ed..0d78094 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -918,9 +918,9 @@ private async Task RunEventPumpAsync() // nobody will deliver for again. if (Volatile.Read(ref _pumpStopRequested)) { - while (TryDequeueEvent() is not null) - { - } + Action? skipped; + do { skipped = TryDequeueEvent(); } + while (skipped is not null); } } finally diff --git a/src/CanKit.Pro.CANopen/ICanOpenNode.cs b/src/CanKit.Pro.CANopen/ICanOpenNode.cs index e358630..29e1a78 100644 --- a/src/CanKit.Pro.CANopen/ICanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/ICanOpenNode.cs @@ -38,7 +38,9 @@ namespace CanKit.Pro.CANopen; /// /// /// Disposing. fails every open transfer with -/// , stops the producers and starts no further delivery +/// (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 From 70db3e89e493fcfac52629111061587fb7a7c7e7 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:53:21 +0200 Subject: [PATCH 15/17] test(isotp): await the failure report instead of reading it as soon as the receive faults The reader posts the loss to the inbox before it raises the report, so the receive can complete first. With the report delayed by 200 ms the old assertion fails every time and this one does not (#270). Co-Authored-By: Claude Sonnet 5.5 --- .../TestCases/IsoTp/IsoTpChannelIntegrationTests.cs | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/IsoTp/IsoTpChannelIntegrationTests.cs index c29ab05..ae1a0df 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"); } From 476630f73dc75bd1fce3156a36ca6ec5414118d2 Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sun, 4 Oct 2026 07:27:31 +0200 Subject: [PATCH 16/17] fix(canopen): keep every active disposal marker of a thread, not only the innermost Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 73 ++++++++++--------- .../TestCases/CANopen/CanOpenDisposeTests.cs | 22 ++++++ 2 files changed, 60 insertions(+), 35 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 0d78094..71e124e 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -720,47 +720,50 @@ private static async Task JoinAsync(Task task) } /// True on the thread that is delivering an event to a subscriber. - private bool OnEventPump => ReferenceEquals(t_deliveringFor, this); + private bool OnEventPump => t_delivering?.Contains(this) == true; /// 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 => ReferenceEquals(t_finishingFor, this); + private bool OnDisposal => t_finishing?.Contains(this) == true; /// True on the thread of the reader task while it reports a failed subscription. - private bool OnReader => ReferenceEquals(t_reportingFor, this); + private bool OnReader => t_reporting?.Contains(this) == true; - [ThreadStatic] - private static CanOpenNode? t_deliveringFor; + // 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 sealed class NodeMarks + { + private readonly System.Collections.Generic.List _nodes = new(2); - [ThreadStatic] - private static CanOpenNode? t_reportingFor; + public void Enter(CanOpenNode node) => _nodes.Add(node); - [ThreadStatic] - private static CanOpenNode? t_finishingFor; + public void Leave(CanOpenNode node) + { + var at = _nodes.LastIndexOf(node); + if (at >= 0) _nodes.RemoveAt(at); + } - private static CanOpenNode? MarkFinishing(CanOpenNode? node) - { - var previous = t_finishingFor; - t_finishingFor = node; - return previous; + public bool Contains(CanOpenNode node) + { + foreach (var marked in _nodes) + { + if (ReferenceEquals(marked, node)) return true; + } + return false; + } } - // Each marker returns what it replaced, and the caller puts that back: the thread may be - // inside another node's callback (a subscriber of A disposing B), and B's markers must not - // erase A's. - private static CanOpenNode? MarkDelivering(CanOpenNode? node) - { - var previous = t_deliveringFor; - t_deliveringFor = node; - return previous; - } + [ThreadStatic] + private static NodeMarks? t_delivering; - private static CanOpenNode? MarkReporting(CanOpenNode? node) - { - var previous = t_reportingFor; - t_reportingFor = node; - return previous; - } + [ThreadStatic] + private static NodeMarks? t_reporting; + + [ThreadStatic] + private static NodeMarks? t_finishing; + + private static NodeMarks Marks(ref NodeMarks? slot) => slot ??= new NodeMarks(); /// Flips the disposed flag and posts the cleanup; false when already disposed. private bool BeginDispose() @@ -826,7 +829,7 @@ private void CleanUpOnActor() private void FinishDispose() { - var outerFinishing = MarkFinishing(this); + Marks(ref t_finishing).Enter(this); try { _subscription.Dispose(); @@ -837,7 +840,7 @@ private void FinishDispose() } finally { - MarkFinishing(outerFinishing); + Marks(ref t_finishing).Leave(this); // Whatever a Dispose above threw, a caller waiting for this disposal is released. _disposeDone.TrySetResult(true); } @@ -879,9 +882,9 @@ private async Task RunReaderAsync() catch (Exception ex) { // A subscriber of the report may dispose the node; see Dispose. - var outer = MarkReporting(this); + Marks(ref t_reporting).Enter(this); try { RaiseBackgroundException(ex); } - finally { MarkReporting(outer); } + finally { Marks(ref t_reporting).Leave(this); } } } @@ -900,7 +903,7 @@ private async Task RunEventPumpAsync() // must still be delivered. Anything outside the delegate is still a bug in the pump. // 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). - var outerDelivering = MarkDelivering(this); + Marks(ref t_delivering).Enter(this); try { while (!Volatile.Read(ref _pumpStopRequested) && TryDequeueEvent() is { } raise) @@ -925,7 +928,7 @@ private async Task RunEventPumpAsync() } finally { - MarkDelivering(outerDelivering); + Marks(ref t_delivering).Leave(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 diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs index ef516c2..ef5103c 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDisposeTests.cs @@ -435,6 +435,28 @@ public async Task Disposing_Another_Node_From_Inside_A_Disposal_Keeps_The_Marker 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. From 3f632e87144c65f3874da61d3c5f86b943d0eb8d Mon Sep 17 00:00:00 2001 From: Dietmar Borgards <2646931+dborgards@users.noreply.github.com> Date: Sun, 4 Oct 2026 08:00:33 +0200 Subject: [PATCH 17/17] refactor(canopen): keep the per-thread callback sets in a static class Co-Authored-By: Claude Sonnet 5.5 --- src/CanKit.Pro.CANopen/CanOpenNode.cs | 58 +++++++++++++-------------- 1 file changed, 28 insertions(+), 30 deletions(-) diff --git a/src/CanKit.Pro.CANopen/CanOpenNode.cs b/src/CanKit.Pro.CANopen/CanOpenNode.cs index 71e124e..25e13ab 100644 --- a/src/CanKit.Pro.CANopen/CanOpenNode.cs +++ b/src/CanKit.Pro.CANopen/CanOpenNode.cs @@ -720,50 +720,48 @@ private static async Task JoinAsync(Task task) } /// True on the thread that is delivering an event to a subscriber. - private bool OnEventPump => t_delivering?.Contains(this) == true; + 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 => t_finishing?.Contains(this) == true; + 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 => t_reporting?.Contains(this) == true; + 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 sealed class NodeMarks + private enum CallbackKind { Delivering, Reporting, Finishing } + + private static class Callbacks { - private readonly System.Collections.Generic.List _nodes = new(2); + [ThreadStatic] + private static System.Collections.Generic.List?[]? s_nodes; - public void Enter(CanOpenNode node) => _nodes.Add(node); + public static void Enter(CallbackKind kind, CanOpenNode node) => Of(kind).Add(node); - public void Leave(CanOpenNode node) - { - var at = _nodes.LastIndexOf(node); - if (at >= 0) _nodes.RemoveAt(at); - } + // 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 bool Contains(CanOpenNode node) + public static bool Contains(CallbackKind kind, CanOpenNode node) { - foreach (var marked in _nodes) + var list = s_nodes?[(int)kind]; + if (list is null) return false; + foreach (var marked in list) { if (ReferenceEquals(marked, node)) return true; } return false; } - } - - [ThreadStatic] - private static NodeMarks? t_delivering; - [ThreadStatic] - private static NodeMarks? t_reporting; - - [ThreadStatic] - private static NodeMarks? t_finishing; - - private static NodeMarks Marks(ref NodeMarks? slot) => slot ??= new NodeMarks(); + 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); + } + } /// Flips the disposed flag and posts the cleanup; false when already disposed. private bool BeginDispose() @@ -829,7 +827,7 @@ private void CleanUpOnActor() private void FinishDispose() { - Marks(ref t_finishing).Enter(this); + Callbacks.Enter(CallbackKind.Finishing, this); try { _subscription.Dispose(); @@ -840,7 +838,7 @@ private void FinishDispose() } finally { - Marks(ref t_finishing).Leave(this); + Callbacks.Leave(CallbackKind.Finishing, this); // Whatever a Dispose above threw, a caller waiting for this disposal is released. _disposeDone.TrySetResult(true); } @@ -882,9 +880,9 @@ private async Task RunReaderAsync() catch (Exception ex) { // A subscriber of the report may dispose the node; see Dispose. - Marks(ref t_reporting).Enter(this); + Callbacks.Enter(CallbackKind.Reporting, this); try { RaiseBackgroundException(ex); } - finally { Marks(ref t_reporting).Leave(this); } + finally { Callbacks.Leave(CallbackKind.Reporting, this); } } } @@ -903,7 +901,7 @@ private async Task RunEventPumpAsync() // must still be delivered. Anything outside the delegate is still a bug in the pump. // 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). - Marks(ref t_delivering).Enter(this); + Callbacks.Enter(CallbackKind.Delivering, this); try { while (!Volatile.Read(ref _pumpStopRequested) && TryDequeueEvent() is { } raise) @@ -928,7 +926,7 @@ private async Task RunEventPumpAsync() } finally { - Marks(ref t_delivering).Leave(this); + 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