From 31e664c3a4ca21166a27ac1fd5904b6710aac17a Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 19:41:10 +0000 Subject: [PATCH 1/6] refactor(canopen): add a write-gate-waiters test seam to ObjectDictionary WriteUnsigned and Add already serialize on the private _writeGate mutex, but nothing let a test observe that a concurrent writer had genuinely reached it. Three CanOpenCommunicationProfileTests relied on a fixed Thread.Sleep instead, assuming a concurrent "hammer" writer had reached the gate within the sleep. Internal WriteGateWaiters (incremented before lock (_writeGate), decremented once acquired) gives tests a real signal to poll instead. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- src/CanKit.Pro.CANopen/ObjectDictionary.cs | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/src/CanKit.Pro.CANopen/ObjectDictionary.cs b/src/CanKit.Pro.CANopen/ObjectDictionary.cs index bba0054..f95b106 100644 --- a/src/CanKit.Pro.CANopen/ObjectDictionary.cs +++ b/src/CanKit.Pro.CANopen/ObjectDictionary.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Threading; using CanKit.Pro.CANopen.Sdo; namespace CanKit.Pro.CANopen; @@ -43,6 +44,17 @@ public sealed class ObjectDictionary // may. Readers never take this lock, so a validator reading under _sync cannot deadlock. private readonly object _writeGate = new(); + // Test seam (#171): the number of callers that have reached the write gate and are not yet + // inside it — WriteUnsigned's and Add's own lock (_writeGate), the only two entry + // points a concurrent test "hammer" reaches. While a test holds the gate, every caller + // counted here is blocked on it; tests poll it as that signal instead of a fixed sleep that + // only ever guessed how long reaching the gate takes. Not counted for Transaction / + // WriteRawUnchecked / Declare, which no such test hammers. + private int _writeGateWaiters; + + /// Test seam (#171): see . + internal int WriteGateWaiters => Volatile.Read(ref _writeGateWaiters); + /// /// Internal hook for : raised after a value-mutating write /// ( / , or an Add* call that @@ -394,8 +406,10 @@ public void WriteUnsigned(ushort index, byte subindex, uint value) // The type is resolved under the write gate, so a re-declaration is either fully before // or fully after this write: the value is encoded and range-checked against the // declaration it lands on (Codex on #133). + Interlocked.Increment(ref _writeGateWaiters); lock (_writeGate) { + Interlocked.Decrement(ref _writeGateWaiters); OdDataType type; lock (_sync) { @@ -437,8 +451,10 @@ private OdEntry Add(ushort index, byte subindex, OdDataType type, OdAccess acces bool replaced; // Under the write gate: a write validated against the entry being replaced is stored on // it before the replacement lands, never on the new declaration (Codex on #133). + Interlocked.Increment(ref _writeGateWaiters); lock (_writeGate) { + Interlocked.Decrement(ref _writeGateWaiters); lock (_sync) { var key = Key(index, subindex); From a80e28ef1547e64d283356b9eecc7272018fdae2 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 19:41:27 +0000 Subject: [PATCH 2/6] test(canopen): replace NMT-start, boot-up, SDO-wire, and RPDO-pump sleeps with observables MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Part of the #114/#171 delay-and-sleep audit ("CANopen NMT Start" and "CANopen session, wire, and pump" tables). None of these sleeps waited on a clock the node's actor had armed: they assumed a received frame had been dequeued and ApplyNmtTransition had run, an SDO session had been installed, an upload-init was on the wire, or the RPDO event pump had drained. - ICanOpenNode.State already round-trips through the node's actor, so a bounded poll on it (WaitUntilOperationalAsync) is a real barrier for "the NMT Start this test just sent has been applied" — no product change needed. Used in CanOpenDynamicMappingTests and the four Tpdo_*/Nmt_Broadcast tests in CanOpenNodeIntegrationTests that used to sleep 50-100 ms after sending Start. - A local BootupWatch (the pattern CanOpenCommunicationProfileTests and CanOpenDeviceDescriptionTests already use) replaces the sleeps that waited to consume a node's opening boot-up or a reset's boot-up (Nmt_ResetNode_EmitsBootup, Nmt_ResetCommunication_Emits..., and the heartbeat-consumer arming sleep in Heartbeat_Consumer_FiresTimeout...). - Sdo_ServerSupersede_EmitsWireAbort_ForPriorTransfer, Sdo_ClientResponseWithShortDlc_IsAcceptedAndCompletes, and Sdo_ClientSegmentedUploadResponse_OverMaxTransferBytes_AbortsOutOfMemory now wait for the server's session-install ack / the client's own init frame to be observed on the wire (FrameObserved) instead of guessing when a request "must" have been transmitted. - Sdo_Segmented_Download_Wrong_Toggle_Aborts waits for the segmented download's init-ack (scs 0x60) instead of guessing 100 ms. - Tpdo_Emission_UnderConcurrentOdWrites_NeverTears flushes the consumer's RPDO event pump with one more deterministic TPDO and waits for its own delivery (the pump is a single FIFO reader) instead of a fixed 200 ms, so a late torn payload can no longer slip past the assertions unsampled. Left unconverted, with reasons: - Every "nothing else was transmitted" wall-clock window whose effect the actor round trip cannot observe (Tpdo_ChangeOfState_DoesNotEcho_ On_RpdoUnpack's echo-quiet window, Overlapping_Emcys_..., A_Guarding_ Reply_..., the producer-tick windows) is left as-is per the audit's own guidance: shortening a green "did not happen" window only weakens it. - Nmt_ResetCommunication_EmitsBootup_And_Settles_In_PreOperational's trailing state check needed no barrier at all: PerformNmtReset sets _state = PreOperational before it emits the boot-up frame the test already awaits, on the same actor, so the state is already correct by the time bootups.Second resolves. Mutation-tested: ApplyNmtTransition (state never reaches Operational), the two EmitHeartbeat(0x00) boot-up call sites (construction and reset), AbortSupersededServerSession (stays silent), and the download- segment toggle check each caught their guarded test(s) when broken and were restored; src/ carries no leftover mutation. Ran the full CANopen + ApiApproval filter (424 tests, all green) and the three converted files' tests 10x in a row (116 tests each run, all green). Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../CANopen/CanOpenDynamicMappingTests.cs | 30 +++- .../CANopen/CanOpenNodeIntegrationTests.cs | 166 ++++++++++++++---- 2 files changed, 155 insertions(+), 41 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs index b2bfbb2..94ff70c 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs @@ -51,6 +51,28 @@ private static Task CreatePdoAsync(ICanOpenNode client, byte serverNodeId, ushor return client.SdoDownloadAsync(serverNodeId, commIndex, 0x01, U32Bytes(cobId)).WithTimeoutAsync(ShortTimeout); } + // The barrier for "the NMT Start each node was just sent has been dequeued off its bus and + // applied by ApplyNmtTransition": ICanOpenNode.State round-trips through the node's actor + // (CanOpenNode.State getter posts to the actor and returns what it reads there), so polling + // it cannot observe a state the actor has not actually reached yet. It replaces a fixed + // Task.Delay that only ever guessed how long the dequeue-and-apply hop would take. + private static async Task WaitUntilOperationalAsync(params ICanOpenNode[] nodes) + { + var deadline = DateTime.UtcNow + ShortTimeout; + while (true) + { + var allOperational = true; + foreach (var node in nodes) + { + if (node.State != NmtState.Operational) { allOperational = false; break; } + } + if (allOperational) return; + if (DateTime.UtcNow >= deadline) + throw new TimeoutException($"Node(s) did not reach Operational within {ShortTimeout}."); + await Task.Delay(5); + } + } + private static byte[] MappingEntryBytes(ushort index, byte subindex, byte bitLength) { uint raw = ((uint)index << 16) | ((uint)subindex << 8) | bitLength; @@ -122,7 +144,7 @@ await WriteMappingAsync(consumer, serverNodeId: 0x11, mapIndex: 0x1A00, await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); await producer.TriggerTpdoAsync(1); var payload = await received.Task.WithTimeoutAsync(ShortTimeout); @@ -168,7 +190,7 @@ await WriteMappingAsync(producer, serverNodeId: 0x01, mapIndex: 0x1600, await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); await producer.TriggerTpdoAsync(1); await received.Task.WithTimeoutAsync(ShortTimeout); @@ -371,7 +393,7 @@ public async Task Tpdo_ChangeOfState_Emits_On_ApplicationOdWrite() await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); // Application-originated OD write — no TriggerTpdoAsync. producer.ObjectDictionary.WriteUnsigned(0x2000, 0x00, 0x1234u); @@ -420,7 +442,7 @@ public async Task Tpdo_ChangeOfState_DoesNotEcho_On_RpdoUnpack() await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); await producer.TriggerTpdoAsync(1); await Task.Delay(300); // give any (wrong) echo ample time to appear diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs index 9b26aed..5f4d351 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs @@ -35,6 +35,58 @@ private static ICanBus Open(string session, int channel) => CanBus.Open( $"virtual://{session}/{channel}", cfg => cfg.SetProtocolMode(CanProtocolMode.Can20).Baud(VirtualAdapterFixture.Bitrate)); + // The barrier for "the NMT Start each node was just sent has been dequeued off its bus and + // applied by ApplyNmtTransition": ICanOpenNode.State round-trips through the node's actor + // (CanOpenNode.State getter posts to the actor and returns what it reads there), so polling + // it cannot observe a state the actor has not actually reached yet. It replaces a fixed + // Task.Delay that only ever guessed how long the dequeue-and-apply hop would take. + private static async Task WaitUntilOperationalAsync(params ICanOpenNode[] nodes) + { + var deadline = DateTime.UtcNow + ShortTimeout; + while (true) + { + var allOperational = true; + foreach (var node in nodes) + { + if (node.State != NmtState.Operational) { allOperational = false; break; } + } + if (allOperational) return; + if (DateTime.UtcNow >= deadline) + throw new TimeoutException($"Node(s) did not reach Operational within {ShortTimeout}."); + await Task.Delay(5); + } + } + + /// + /// Counts the boot-up frames (00h on 700h + producer) one node sees from another. + /// A node sends one when it is opened and one on every NMT reset; a test that waits for the + /// reset's boot-up first consumes the opening one, so the two cannot be confused. The observer + /// must be opened before the producer so that the first one is guaranteed to be seen. + /// + private sealed class BootupWatch + { + private readonly TaskCompletionSource _first = new(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly TaskCompletionSource _second = new(TaskCreationOptions.RunContinuationsAsynchronously); + private int _seen; + + public BootupWatch(ICanOpenNode observer, byte producer) + { + observer.HeartbeatReceived += (_, e) => + { + if (e.ProducerNodeId != producer || e.State != NmtState.Initializing) return; + switch (Interlocked.Increment(ref _seen)) + { + case 1: _first.TrySetResult(true); break; + case 2: _second.TrySetResult(true); break; + } + }; + } + + public Task First => _first.Task.WithTimeoutAsync(ShortTimeout); + + public Task Second => _second.Task.WithTimeoutAsync(ShortTimeout); + } + // ----------------------------------------------------------------------------------------- // FR-CO-002 — SDO expedited upload/download over two nodes on the same virtual bus. // ----------------------------------------------------------------------------------------- @@ -361,6 +413,7 @@ public async Task Sdo_ServerSupersede_EmitsWireAbort_ForPriorTransfer() // Observe every SDO frame the slave emits (COB-ID 0x580+0x11 = 0x591) on a third bus // so we can distinguish the slave's own transmits from anything the master sends. var slaveSdoTx = new List(); + var sessionInstalled = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var abortSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); busObserver.FrameObserved += (_, e) => { @@ -369,6 +422,10 @@ public async Task Sdo_ServerSupersede_EmitsWireAbort_ForPriorTransfer() if ((uint)frame.ID != 0x580u + 0x11u) return; var data = frame.Data.ToArray(); lock (slaveSdoTx) slaveSdoTx.Add(data); + // scs=0x60 for (0x2100, 0x00): the server acknowledged the segmented initiate, so + // the session this test needs to be superseded is now genuinely installed. + if (data.Length >= 4 && data[0] == 0x60 && data[1] == 0x00 && data[2] == 0x21 && data[3] == 0x00) + sessionInstalled.TrySetResult(data); // cs=0x80 -> SDO abort. Fire completion for the FIRST abort we see with the // superseded (index, subindex) = (0x2100, 0x00). That is the marker we care // about for this regression; ignore other frames. @@ -390,8 +447,9 @@ public async Task Sdo_ServerSupersede_EmitsWireAbort_ForPriorTransfer() }; busA.Transmit(CanFrame.Classic(0x600 + 0x11, priorInit, isExtendedFrame: false)); - // Give the actor loop a moment to install the segmented session for 0x2100:00. - await Task.Delay(50); + // Wait for the server's own ack rather than a delay: the segmented session for 0x2100 + // is installed exactly when that ack goes out. + await sessionInstalled.Task.WithTimeoutAsync(ShortTimeout); // Now supersede: master runs an expedited download to an unrelated (index, subindex) // on the same server. Per CiA 301 the server must abort the still-open 0x2100 transfer @@ -432,12 +490,21 @@ public async Task Sdo_ClientResponseWithShortDlc_IsAcceptedAndCompletes() using var master = CanOpen.OpenNode(busA, nodeId: 0x01); PeerSdoLaboratory.Bind(master, 0x11); + // Wait for the master's own upload-init request to be on the wire, rather than guessing + // how long that takes: COB-ID 0x600+0x11 is where the master's SDO client transmits. + var initOnWire = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + busB.FrameObserved += (_, e) => + { + var frame = e.CanFrame; + if (frame.IsExtendedFrame || (uint)frame.ID != 0x600u + 0x11u) return; + initOnWire.TrySetResult(true); + }; + // Master initiates an expedited upload from a phantom server 0x11 at (0x2500, 0x00). // We do NOT open a slave; instead we fake the server response on busB. var uploadTask = master.SdoUploadAsync(serverNodeId: 0x11, index: 0x2500, subindex: 0x00); - // Give the master a moment to actually put its init request on the wire. - await Task.Delay(30); + await initOnWire.Task.WithTimeoutAsync(ShortTimeout); // Fake a 5-byte SDO expedited upload response: cs=0x4F selects size-indicated with // n=3 (one valid byte), followed by index (0x2500), subindex (0x00), and one payload @@ -550,6 +617,7 @@ public async Task Sdo_ClientSegmentedUploadResponse_OverMaxTransferBytes_AbortsO // side of the fix (the master must abort the transfer back to the peer, not silently // fail its own task). var abortSeen = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var initOnWire = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); busB.FrameObserved += (_, e) => { var frame = e.CanFrame; @@ -557,12 +625,13 @@ public async Task Sdo_ClientSegmentedUploadResponse_OverMaxTransferBytes_AbortsO if ((uint)frame.ID != 0x600u + 0x11u) return; var data = frame.Data.ToArray(); if (data.Length >= 8 && data[0] == 0x80) abortSeen.TrySetResult(data); + else initOnWire.TrySetResult(true); }; var uploadTask = master.SdoUploadAsync(serverNodeId: 0x11, index: 0x2900, subindex: 0x00); - // Give the master a moment to put its upload-init request on the wire. - await Task.Delay(30); + // Wait for the master's upload-init request to be on the wire, rather than guessing. + await initOnWire.Task.WithTimeoutAsync(ShortTimeout); // Fake a segmented upload-init response: cs=0x41 (size-indicated), then (index, // subindex), then a little-endian 32-bit declared length way above the cap. The client @@ -635,7 +704,7 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); var shortPattern = Enumerable.Repeat((byte)0xAA, 2).ToArray(); // 2 bytes var longPattern = Enumerable.Repeat((byte)0xBB, 8).ToArray(); // 8 bytes @@ -687,11 +756,31 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() catch (Exception) { Interlocked.Increment(ref emitCrashes); } } - // Give the RPDO event pump time to drain before we sample counts. - await Task.Delay(200); cts.Cancel(); await writer; + // Flush the consumer's RPDO event pump before sampling counts: one more deterministic + // TPDO, and wait for ITS delivery. The pump is a single reader that delivers events in + // enqueue order (RunEventPumpAsync), so this one's arrival proves every event the 2000 + // emissions above queued has already been delivered too — not a guess at how long that + // drain takes. + producer.ObjectDictionary.WriteRaw(0x2A00, 0x00, longPattern); + var flushed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + void OnFlush(object? _, RpdoReceivedEventArgs e) + { + if (e.CobId == producerCobId) flushed.TrySetResult(true); + } + consumer.RpdoReceived += OnFlush; + try + { + await producer.TriggerTpdoAsync(1); + await flushed.Task.WithTimeoutAsync(ShortTimeout); + } + finally + { + consumer.RpdoReceived -= OnFlush; + } + observedCount.Should().BeGreaterThan(0, "the consumer must have observed at least one TPDO frame to make the tear check meaningful"); emitCrashes.Should().Be(0, @@ -746,7 +835,7 @@ public async Task Nmt_Broadcast_TransitionsAllNodes() using var slave2 = CanOpen.OpenNode(busC, nodeId: 0x12); await master.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0); // broadcast - await Task.Delay(100); // give the slave loops time to apply + await WaitUntilOperationalAsync(slave1, slave2); slave1.State.Should().Be(NmtState.Operational); slave2.State.Should().Be(NmtState.Operational); } @@ -804,20 +893,14 @@ public async Task Nmt_ResetNode_EmitsBootup() using var busB = Open(session, 1); using var master = CanOpen.OpenNode(busA, nodeId: 0x01); + var bootups = new BootupWatch(master, 0x11); using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); - // Wait until the initial bootup from `slave` is consumed. - await Task.Delay(50); - - var bootup = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - master.HeartbeatReceived += (s, e) => - { - if (e.ProducerNodeId == 0x11 && e.State == NmtState.Initializing) - bootup.TrySetResult(e.State); - }; + // Consume the initial bootup before the reset's own, so the two cannot be confused. + await bootups.First; await master.SendNmtCommandAsync(NmtCommand.ResetNode, targetNodeId: 0x11); - (await bootup.Task.WithTimeoutAsync(ShortTimeout)).Should().Be(NmtState.Initializing); + await bootups.Second; } // FR-CO-007: ResetCommunication (0x82) follows the same re-init path as ResetNode — @@ -830,20 +913,16 @@ public async Task Nmt_ResetCommunication_EmitsBootup_And_Settles_In_PreOperation using var busB = Open(session, 1); using var master = CanOpen.OpenNode(busA, nodeId: 0x01); + var bootups = new BootupWatch(master, 0x11); using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); - await Task.Delay(50); // consume the initial bootup - - var bootup = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - master.HeartbeatReceived += (s, e) => - { - if (e.ProducerNodeId == 0x11 && e.State == NmtState.Initializing) - bootup.TrySetResult(e.State); - }; + await bootups.First; // consume the initial bootup await master.SendNmtCommandAsync(NmtCommand.ResetCommunication, targetNodeId: 0x11); - (await bootup.Task.WithTimeoutAsync(ShortTimeout)).Should().Be(NmtState.Initializing); - await Task.Delay(50); + await bootups.Second; + // PerformNmtReset sets _state = PreOperational before it emits the boot-up frame, on the + // slave's own actor: by the time master's HeartbeatReceived for that boot-up has fired + // (bootups.Second), the slave has already reached PreOperational -- no further wait needed. slave.State.Should().Be(NmtState.PreOperational); } @@ -859,15 +938,27 @@ public async Task Sdo_Segmented_Download_Wrong_Toggle_Aborts() using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); slave.ObjectDictionary.AddDomain(0x2100, 0x00, new byte[20]); + var sdoTxCobId = CanOpenCobId.SdoTx(0x11); + var initAck = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + rawBus.FrameObserved += (_, e) => + { + var frame = e.CanFrame; + if (frame.IsExtendedFrame || (uint)frame.ID != sdoTxCobId) return; + var data = frame.Data; + // scs=0x60 for (0x2100, 0x00): the segmented download initiate was acknowledged, so + // the session the wrong-toggle segment below must land against is genuinely open. + if (data.Length >= 4 && data.Span[0] == 0x60 && data.Span[1] == 0x00 && data.Span[2] == 0x21 && data.Span[3] == 0x00) + initAck.TrySetResult(true); + }; + // Segmented download initiate (cs = 0x21), then the FIRST segment with the wrong // toggle bit (0x10 set instead of 0x00 expected). rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.SdoRx(0x11)), new byte[] { 0x21, 0x00, 0x21, 0x00, 0x14, 0x00, 0x00, 0x00 })); - await Task.Delay(100); // let the init-ack happen + await initAck.Task.WithTimeoutAsync(ShortTimeout); // the init-ack, not a guess at its timing rawBus.Transmit(CanFrame.Classic(unchecked((int)CanOpenCobId.SdoRx(0x11)), new byte[] { 0x10, 1, 2, 3, 4, 5, 6, 7 })); - var sdoTxCobId = CanOpenCobId.SdoTx(0x11); using var cts = new CancellationTokenSource(ShortTimeout); while (true) { @@ -922,10 +1013,11 @@ public async Task Heartbeat_Consumer_FiresTimeoutWhenPeerGoesSilent() using var busB = Open(session, 1); using var master = CanOpen.OpenNode(busA, nodeId: 0x01); + var bootups = new BootupWatch(master, 0x11); using var slave = CanOpen.OpenNode(busB, nodeId: 0x11); - // Wait past the initial bootup so the consumer arms cleanly. - await Task.Delay(100); + // Consume the initial bootup so the consumer arms cleanly. + await bootups.First; var timeout = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); master.HeartbeatTimeout += (s, e) => @@ -1231,7 +1323,7 @@ public async Task Tpdo_EventDriven_Emits_MappedOdValues() // Operational on both ends (RPDO unpack gated in HandleRpdo). await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); await producer.TriggerTpdoAsync(1); var payload = await received.Task.WithTimeoutAsync(ShortTimeout); @@ -1300,7 +1392,7 @@ public async Task Tpdo_DummyMapping_KeepsSubsequentSlotOffsets() await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); await producer.TriggerTpdoAsync(1); var payload = await received.Task.WithTimeoutAsync(ShortTimeout); @@ -1331,7 +1423,7 @@ public async Task Tpdo_SyncTriggered_FiresEverySync() // Bring producer Operational. await consumer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x11); await producer.SendNmtCommandAsync(NmtCommand.Start, targetNodeId: 0x01); - await Task.Delay(50); + await WaitUntilOperationalAsync(producer, consumer); int rpdoCount = 0; var enough = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); From a0bcb53a301fa1fda80d77262d9bbb82c08b34b2 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 19:41:37 +0000 Subject: [PATCH 3/6] test(canopen): replace OD write-gate sleeps with the new write-gate-waiters signal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Part of the #114/#171 delay-and-sleep audit's object-dictionary write-gate sub-table. A_Redeclaration_Waits_For_The_Write_In_Flight_On_The_Entry, A_Typed_Write_Resolves_Its_Type_Under_The_Write_Gate, and A_Direct_Write_Cannot_Land_Inside_A_ConfigureTpdo_Transaction each held the ObjectDictionary's write gate from inside a WriteValidator callback and slept a fixed 100-200 ms, assuming a concurrent writer had reached the same gate by then. They now spin (SpinUntilAtTheWriteGate) on the ObjectDictionary.WriteGateWaiters seam added in the prior commit — real evidence the other writer's own lock (_writeGate) attempt is blocked on this one, not a guess at how long reaching it takes. Left unconverted: A_Save_During_An_Nmt_Reset_Stores_All_Restored_Values_ Not_A_Mix's Thread.Sleep(300) is not the same shape. It runs inside an NMT reset's Transaction(...), which already holds _writeGate for the whole restore; the concurrent save's WriteRaw call is provably blocked on that same lock the instant it is issued, regardless of the sleep's length. The sleep is margin for the "if the save were wrongly not held back it would have completed by now" direction (matching the audit's "green"), not a guess about when the save reaches the gate, so it is left as-is. Mutation-tested: making Add's write-gate lock private (no longer the shared gate) broke A_Redeclaration_Waits_For_The_Write_In_Flight_On_ The_Entry as expected, then reverted. Ran CanOpenCommunicationProfileTests (67 tests) and the full CANopen + ApiApproval filter (424 tests) green, and the three converted files' tests 10x in a row (116 tests/run) green. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../CanOpenCommunicationProfileTests.cs | 27 ++++++++++++++++--- 1 file changed, 23 insertions(+), 4 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenCommunicationProfileTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenCommunicationProfileTests.cs index 723f50c..4769149 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenCommunicationProfileTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenCommunicationProfileTests.cs @@ -91,6 +91,25 @@ private static TaskCompletionSource NewTcs() private static uint HeartbeatEntry(byte nodeId, ushort milliseconds) => ((uint)nodeId << 16) | milliseconds; + /// + /// Spin-waits, synchronously, until (the + /// write-gate test seam, #171) reaches at least — real evidence that + /// another writer's own lock (_writeGate) attempt has reached the gate this call is + /// itself holding, rather than a fixed sleep that only ever guessed how long that takes. Used + /// from inside a WriteValidator callback, which runs synchronously on the holding + /// thread, so the wait is a blocking spin, not an awaited one. + /// + private static void SpinUntilAtTheWriteGate(ObjectDictionary od, int count = 1) + { + var deadline = DateTime.UtcNow + ShortTimeout; + while (od.WriteGateWaiters < count) + { + if (DateTime.UtcNow >= deadline) + throw new TimeoutException($"Fewer than {count} writer(s) reached the write gate within {ShortTimeout}."); + Thread.Sleep(1); + } + } + /// /// Counts the boot-up frames (00h on 700h + producer) one node sees from another. /// A node sends one when it is opened and one on every NMT reset; a test that waits for the @@ -439,7 +458,7 @@ public async Task A_Redeclaration_Waits_For_The_Write_In_Flight_On_The_Entry() if (index == 0x2000 && redeclare is null) { redeclare = Task.Run(() => od.AddU8(0x2000, 0x00, 0x01)); - Thread.Sleep(200); // a re-declaration that did not wait for the gate would land here + SpinUntilAtTheWriteGate(od); // the re-declaration is genuinely blocked on this gate } return inner(index, subindex, value); }; @@ -499,8 +518,8 @@ public async Task A_Typed_Write_Resolves_Its_Type_Under_The_Write_Gate() { gateHeld.Set(); writeStarted.Wait(ShortTimeout); - Thread.Sleep(200); // the typed write has started and waits for the gate - od.AddU32(0x2000, 0x00, 0); // re-declared while the gate is held + SpinUntilAtTheWriteGate(od); // the typed write is genuinely blocked on this gate + od.AddU32(0x2000, 0x00, 0); // re-declared while the gate is held } return inner(index, subindex, value); }; @@ -642,7 +661,7 @@ public async Task A_Direct_Write_Cannot_Land_Inside_A_ConfigureTpdo_Transaction( if (index == 0x1800 && subindex == 0x01 && !sequenceStarted.IsSet) { sequenceStarted.Set(); - Thread.Sleep(100); // the hammers are at the gate before the sequence goes on + SpinUntilAtTheWriteGate(od, count: 4); // all four hammers are genuinely blocked on this gate } return inner(index, subindex, value); }; From 70afa3de2f107dc26b1abf2d2fddc7a4ebdca22b Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 20:51:51 +0000 Subject: [PATCH 4/6] test(canopen): recognise the RPDO flush by its own payload The handler added to wait for the flush TPDO matched any RPDO from the producer. An event still queued from the 2000 emissions before it is raised to every handler subscribed when it runs, so it could complete the flush early and let the counts be sampled while later deliveries, a torn one included, were still pending (Codex on #188). The flush TPDO now carries a payload nothing else emits, and the tear check ignores it. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../CANopen/CanOpenNodeIntegrationTests.cs | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs index 5f4d351..f646d62 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs @@ -717,12 +717,16 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() var expectedShortPayload = new byte[8]; Buffer.BlockCopy(shortPattern, 0, expectedShortPayload, 0, shortPattern.Length); var expectedLongPayload = longPattern; + // Written only once the writer has stopped, for the TPDO that flushes the consumer's + // RPDO event pump below: the one payload nothing else emits, so its delivery is that + // TPDO's own and not a queued earlier one's. + var flushPattern = Enumerable.Repeat((byte)0xCC, 8).ToArray(); int observedCount = 0; int tornCount = 0; consumer.RpdoReceived += (_, e) => { - if (e.CobId != producerCobId) return; + if (e.CobId != producerCobId || e.Payload.SequenceEqual(flushPattern)) return; Interlocked.Increment(ref observedCount); if (!e.Payload.SequenceEqual(expectedShortPayload) && !e.Payload.SequenceEqual(expectedLongPayload)) @@ -763,12 +767,14 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() // TPDO, and wait for ITS delivery. The pump is a single reader that delivers events in // enqueue order (RunEventPumpAsync), so this one's arrival proves every event the 2000 // emissions above queued has already been delivered too — not a guess at how long that - // drain takes. - producer.ObjectDictionary.WriteRaw(0x2A00, 0x00, longPattern); + // drain takes. It is recognised by its payload: an event still queued from the loop + // above is raised to every handler subscribed when it runs, this one included (Codex + // on #188). + producer.ObjectDictionary.WriteRaw(0x2A00, 0x00, flushPattern); var flushed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); void OnFlush(object? _, RpdoReceivedEventArgs e) { - if (e.CobId == producerCobId) flushed.TrySetResult(true); + if (e.CobId == producerCobId && e.Payload.SequenceEqual(flushPattern)) flushed.TrySetResult(true); } consumer.RpdoReceived += OnFlush; try From 00768e7948d3f99098392042c63f763e3f63141b Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 20:54:48 +0000 Subject: [PATCH 5/6] test(canopen): state the operational wait with All CodeQL on #188 flagged the foreach in both copies of WaitUntilOperationalAsync as a missed All(...). Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../TestCases/CANopen/CanOpenDynamicMappingTests.cs | 8 ++------ .../TestCases/CANopen/CanOpenNodeIntegrationTests.cs | 7 +------ 2 files changed, 3 insertions(+), 12 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs index 94ff70c..8d43f92 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenDynamicMappingTests.cs @@ -1,5 +1,6 @@ using System; using System.Collections.Generic; +using System.Linq; using System.Threading; using System.Threading.Tasks; using CanKit.Abstractions.API.Can; @@ -61,12 +62,7 @@ private static async Task WaitUntilOperationalAsync(params ICanOpenNode[] nodes) var deadline = DateTime.UtcNow + ShortTimeout; while (true) { - var allOperational = true; - foreach (var node in nodes) - { - if (node.State != NmtState.Operational) { allOperational = false; break; } - } - if (allOperational) return; + if (nodes.All(node => node.State == NmtState.Operational)) return; if (DateTime.UtcNow >= deadline) throw new TimeoutException($"Node(s) did not reach Operational within {ShortTimeout}."); await Task.Delay(5); diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs index f646d62..22de451 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs @@ -45,12 +45,7 @@ private static async Task WaitUntilOperationalAsync(params ICanOpenNode[] nodes) var deadline = DateTime.UtcNow + ShortTimeout; while (true) { - var allOperational = true; - foreach (var node in nodes) - { - if (node.State != NmtState.Operational) { allOperational = false; break; } - } - if (allOperational) return; + if (nodes.All(node => node.State == NmtState.Operational)) return; if (DateTime.UtcNow >= deadline) throw new TimeoutException($"Node(s) did not reach Operational within {ShortTimeout}."); await Task.Delay(5); From ff50e60be44c7733334015b0732a721be3fedf26 Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 27 Sep 2026 21:02:10 +0000 Subject: [PATCH 6/6] test(canopen): keep the tear test's drain window The RPDO flush added on this branch cannot order anything. EmitTpdo sends each TPDO through its own Task.Run, so the flush TPDO can overtake earlier ones on the wire (Codex on #188). Counting deliveries does not work either: every WriteRaw in the writer loop also emits a change-of-state TPDO, and a probe saw 16,000 to 35,000 of them arrive for the loop's 2000 triggers. With no barrier available, the 200 ms window from main is restored, with the reason recorded next to it. Refs #171 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_013WJ8h1ahw4Nj5dYuEWy34s --- .../CANopen/CanOpenNodeIntegrationTests.cs | 36 ++++--------------- 1 file changed, 7 insertions(+), 29 deletions(-) diff --git a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs index 22de451..795b406 100644 --- a/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs +++ b/tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenNodeIntegrationTests.cs @@ -712,16 +712,12 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() var expectedShortPayload = new byte[8]; Buffer.BlockCopy(shortPattern, 0, expectedShortPayload, 0, shortPattern.Length); var expectedLongPayload = longPattern; - // Written only once the writer has stopped, for the TPDO that flushes the consumer's - // RPDO event pump below: the one payload nothing else emits, so its delivery is that - // TPDO's own and not a queued earlier one's. - var flushPattern = Enumerable.Repeat((byte)0xCC, 8).ToArray(); int observedCount = 0; int tornCount = 0; consumer.RpdoReceived += (_, e) => { - if (e.CobId != producerCobId || e.Payload.SequenceEqual(flushPattern)) return; + if (e.CobId != producerCobId) return; Interlocked.Increment(ref observedCount); if (!e.Payload.SequenceEqual(expectedShortPayload) && !e.Payload.SequenceEqual(expectedLongPayload)) @@ -755,33 +751,15 @@ public async Task Tpdo_Emission_UnderConcurrentOdWrites_NeverTears() catch (Exception) { Interlocked.Increment(ref emitCrashes); } } + // Give the RPDO event pump time to drain before we sample counts. A wall window on + // purpose, not a flush barrier (#171): every WriteRaw above also emits a change-of-state + // TPDO, so how many frames are still in flight is unknown, and EmitTpdo sends each one + // through its own Task.Run, so a later TPDO can overtake an earlier one on the wire + // (Codex on #188). Nothing that is sent can mark "everything before me has arrived". + await Task.Delay(200); cts.Cancel(); await writer; - // Flush the consumer's RPDO event pump before sampling counts: one more deterministic - // TPDO, and wait for ITS delivery. The pump is a single reader that delivers events in - // enqueue order (RunEventPumpAsync), so this one's arrival proves every event the 2000 - // emissions above queued has already been delivered too — not a guess at how long that - // drain takes. It is recognised by its payload: an event still queued from the loop - // above is raised to every handler subscribed when it runs, this one included (Codex - // on #188). - producer.ObjectDictionary.WriteRaw(0x2A00, 0x00, flushPattern); - var flushed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - void OnFlush(object? _, RpdoReceivedEventArgs e) - { - if (e.CobId == producerCobId && e.Payload.SequenceEqual(flushPattern)) flushed.TrySetResult(true); - } - consumer.RpdoReceived += OnFlush; - try - { - await producer.TriggerTpdoAsync(1); - await flushed.Task.WithTimeoutAsync(ShortTimeout); - } - finally - { - consumer.RpdoReceived -= OnFlush; - } - observedCount.Should().BeGreaterThan(0, "the consumer must have observed at least one TPDO frame to make the tear check meaningful"); emitCrashes.Should().Be(0,