Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 9 additions & 1 deletion src/CanKit.Pro.CANopen/CanOpenNode.Pdo.cs
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,7 @@ private void RebuildRpdo(int n)
var map = (ushort)(Co.RpdoMap + n - 1);
var rp = _rpdos[n] ??= new RpdoRuntime(n);
rp.SyncPending = null;
rp.ShortFrameReported = false;

uint word = _od.TryReadUnsigned(comm, 0x01, out var w) ? w : CanOpenCobId.InvalidBit;
rp.Valid = (word & CanOpenCobId.InvalidBit) == 0;
Expand Down Expand Up @@ -561,14 +562,20 @@ private void HandleRpdo(RpdoRuntime rp, byte[] payload)
// more than mapped → the first bytes up to the mapped length are used.
if (payload.Length < rp.TotalBytes)
{
if (_emcyValid)
// Once per run of short frames: a producer that keeps sending one would otherwise
// put an EMCY on the bus for every frame, and the EMCY traffic is what fills it.
// A frame that fits ends the run, as does reconfiguring the RPDO.
if (_emcyValid && !rp.ShortFrameReported)
{
rp.ShortFrameReported = true;
var errorRegister = (byte)_od.ReadUnsigned(Co.ErrorRegister, 0x00);
_ = EmitEmcy(new EmcyMessage(_nodeId, 0x8210, errorRegister));
}
return;
}

rp.ShortFrameReported = false;

if (CanOpenTransmissionType.IsSynchronous(rp.TransmissionType))
{
rp.SyncPending = payload;
Expand Down Expand Up @@ -826,5 +833,6 @@ public RpdoRuntime(int pdoIndex)
public PdoMappingEntry[] Mapping { get; set; } = Array.Empty<PdoMappingEntry>();
public int TotalBytes { get; set; }
public byte[]? SyncPending { get; set; }
public bool ShortFrameReported { get; set; }
}
}
6 changes: 3 additions & 3 deletions src/CanKit.Pro.CANopen/CanOpenNode.SdoBlock.cs
Original file line number Diff line number Diff line change
Expand Up @@ -502,8 +502,8 @@ private bool HandleBlockUploadSegment(SdoBlockClientSession session, byte[] data
if (session.Payload!.Length - session.Offset < 7)
{
// Grow when the declared size was 0 (unbounded) or when the payload was under-declared.
var grown = new byte[session.Offset + 7];
Buffer.BlockCopy(session.Payload, 0, grown, 0, session.Payload.Length);
var grown = new byte[GrowCapacity(session.Payload.Length, session.Offset + 7, (long)_options.MaxSdoTransferBytes + 7)];
Buffer.BlockCopy(session.Payload, 0, grown, 0, session.Offset);
session.Payload = grown;
}
// Copy the full 7 data bytes; unused bytes in the *last* segment are trimmed off later
Expand Down Expand Up @@ -813,7 +813,7 @@ private void HandleBlockDownloadServerSegment(SdoBlockServerSession session, byt
// Ensure room for 7 bytes; grow if declared size was under-specified or unbounded.
if (session.Offset + 7 > session.Buffer.Length)
{
var grown = new byte[session.Offset + 7];
var grown = new byte[GrowCapacity(session.Buffer.Length, session.Offset + 7, (long)_options.MaxSdoTransferBytes + 7)];
Buffer.BlockCopy(session.Buffer, 0, grown, 0, session.Offset);
session.Buffer = grown;
}
Expand Down
17 changes: 15 additions & 2 deletions src/CanKit.Pro.CANopen/CanOpenNode.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2127,8 +2127,8 @@ private void HandleSdoClientResponse(byte serverNodeId, byte[] data)
}
if (session.Payload!.Length < needed)
{
var grown = new byte[needed];
Buffer.BlockCopy(session.Payload, 0, grown, 0, session.Payload.Length);
var grown = new byte[GrowCapacity(session.Payload.Length, needed, _options.MaxSdoTransferBytes)];
Buffer.BlockCopy(session.Payload, 0, grown, 0, session.Offset);
session.Payload = grown;
}
Buffer.BlockCopy(payload, 0, session.Payload, session.Offset, payload.Length);
Expand All @@ -2154,6 +2154,19 @@ private void HandleSdoClientResponse(byte serverNodeId, byte[] data)
}
}

/// <summary>
/// Next capacity for a receive buffer that has to grow to <paramref name="needed"/> bytes:
/// doubled, so a transfer of N bytes copies O(N) in total instead of O(N²), but never past
/// <paramref name="ceiling"/> (the transfer cap, plus the slack the caller's final segment
/// may overshoot it by) and never below <paramref name="needed"/>.
/// </summary>
internal static int GrowCapacity(int current, int needed, long ceiling)
{
long doubled = Math.Max(8L, (long)current * 2);
var capacity = (int)Math.Min(doubled, Math.Min(ceiling, int.MaxValue));
return Math.Max(capacity, needed);
}

/// <summary>
/// Restarts the client's request timer after a frame that was attributed to
/// <paramref name="session"/> and leaves it open. Not called for frames the session ignores
Expand Down
81 changes: 81 additions & 0 deletions tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenPdoEngineTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1197,6 +1197,87 @@ public async Task Rpdo_Shorter_Than_Its_Mapping_Is_Not_Actuated_And_Raises_Emcy_
od.ReadUnsigned(0x2101, 0x00).Should().Be(0x5678u, "the first data bytes up to the mapped length are used");
}

// FR-CO-024 (C13) — a producer that keeps sending short frames raises EMCY 8210h once, not once
// per frame: the EMCY traffic would otherwise be a multiple of the fault it reports. A frame
// that fits ends the run, so the next short one is reported again. The count is read when the
// third EMCY arrives, after which every EMCY the second short frame could have raised is
// already in the list: the bus delivers them in order.
[Fact]
public async Task Rpdo_Shorter_Than_Its_Mapping_Raises_Emcy_Once_Per_Run_Of_Short_Frames()
{
var session = NewSession();
using var busB = Open(session, 1);
using var busC = Open(session, 2);
using var wire = new Wire(session, 3);
using var device = CanOpen.OpenNode(busB, Device);
using var observer = CanOpen.OpenNode(busC, Observer);
device.ObjectDictionary.AddU16(0x2100, 0x00, 0);
device.ObjectDictionary.AddU16(0x2101, 0x00, 0);
device.ConfigureRpdo(1, new PdoMapping().Add(0x2100, 0x00, 16).Add(0x2101, 0x00, 16));

var emcys = new System.Collections.Concurrent.ConcurrentQueue<EmcyMessage>();
using var arrived = new System.Collections.Concurrent.BlockingCollection<int>();
using var timeout = new CancellationTokenSource(ShortTimeout);
observer.EmcyReceived += (_, e) =>
{
if (e.Message.ProducerNodeId != Device) return;
emcys.Enqueue(e.Message);
arrived.Add(emcys.Count);
};
var received = new TaskCompletionSource<byte[]>(TaskCreationOptions.RunContinuationsAsynchronously);
device.RpdoReceived += (_, e) =>
{
if (e.CobId == Rpdo1) received.TrySetResult(e.Payload);
};
await StartAsync(wire, device);

wire.Transmit(Rpdo1, new byte[] { 0x01 });
arrived.Take(timeout.Token);
wire.Transmit(Rpdo1, new byte[] { 0x02 });
wire.Transmit(Rpdo1, new byte[] { 0x03, 0x04, 0x05, 0x06 });
await received.Task.WithTimeoutAsync(ShortTimeout);
wire.Transmit(Rpdo1, new byte[] { 0x07 });
arrived.Take(timeout.Token);

emcys.Should().HaveCount(2, "the second short frame continues the run the first one reported");
emcys.Should().OnlyContain(m => m.ErrorCode == 0x8210);
}

// FR-CO-024 (C13) — with the EMCY switched off (1014h bit 31) a short frame raises nothing, and
// it does not use up the report either: once the EMCY is back, the next short frame is the
// first of its run. A fitting frame after the first short one tells when the actor is through
// with it, so an EMCY it wrongly raised is on the wire by then.
[Fact]
public async Task Rpdo_Shorter_Than_Its_Mapping_Raises_No_Emcy_While_The_Emcy_Is_Switched_Off()
{
var session = NewSession();
using var busA = Open(session, 0);
using var busB = Open(session, 1);
using var wire = new Wire(session, 3);
using var master = CanOpen.OpenNode(busA, Master);
using var device = CanOpen.OpenNode(busB, Device);
device.ObjectDictionary.AddU16(0x2100, 0x00, 0);
device.ConfigureRpdo(1, new PdoMapping().Add(0x2100, 0x00, 16));
var received = new TaskCompletionSource<byte[]>(TaskCreationOptions.RunContinuationsAsynchronously);
device.RpdoReceived += (_, e) =>
{
if (e.CobId == Rpdo1) received.TrySetResult(e.Payload);
};
await StartAsync(wire, device);

var emcyCobId = CanOpenCobId.Emcy(Device);
await DownloadAsync(master, 0x1014, 0x00, U32Bytes(CanOpenCobId.InvalidBit | emcyCobId));
wire.Transmit(Rpdo1, new byte[] { 0x01 });
wire.Transmit(Rpdo1, new byte[] { 0x02, 0x03 });
await received.Task.WithTimeoutAsync(ShortTimeout);
wire.Count(emcyCobId).Should().Be(0, "the EMCY is switched off");

await DownloadAsync(master, 0x1014, 0x00, U32Bytes(emcyCobId));
wire.Transmit(Rpdo1, new byte[] { 0x04 });
await wire.WaitForCountAsync(emcyCobId, 1);
wire.Payloads(emcyCobId)[0].Take(2).Should().Equal(0x10, 0x82);
}

// =========================================================================================
// Findings of the #133 review.
// =========================================================================================
Expand Down
138 changes: 138 additions & 0 deletions tests/CanKit.Pro.Tests/TestCases/CANopen/CanOpenSdoCorrectnessTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1136,6 +1136,144 @@ public void Sdo_Server_Sizeless_Download_Doubles_Then_Clamps_To_The_Cap()
"the value is the 24 segment bytes; capacity left above Offset is not written");
}

// -----------------------------------------------------------------------------------------
// C11: a receive buffer that grows by exactly the next segment copies everything received so
// far on every segment, O(N^2) for an N-byte transfer. 64 KiB is 9,363 segments; growing by
// one segment each time allocates about 307 MB over the transfer, doubling about 0.3 MB. The
// bound below sits far from both, so it measures the growth policy and not the allocator.
// -----------------------------------------------------------------------------------------
private const int GrowthPayloadBytes = 64 * 1024;
private const long GrowthAllocationBound = 40L * 1024 * 1024;

[Theory]
[InlineData(0, 1, 100, 8)]
[InlineData(8, 9, 100, 16)]
[InlineData(16, 17, 24, 24)]
[InlineData(16, 30, 24, 30)]
[InlineData(int.MaxValue / 2 + 1, 5, int.MaxValue, int.MaxValue)]
[InlineData(int.MaxValue - 100, int.MaxValue - 90, (long)int.MaxValue + 7, int.MaxValue)]
public void GrowCapacity_Doubles_Clamps_To_The_Ceiling_And_Never_Undershoots(
int current, int needed, long ceiling, int expected)
=> CanOpenNode.GrowCapacity(current, needed, ceiling).Should().Be(expected);

#if NET5_0_OR_GREATER
// GC.GetTotalAllocatedBytes does not exist on net48; the growth code is the same on both.
[Fact]
public async Task Sdo_Client_Sizeless_Upload_Grows_Geometrically()
{
var session = NewSession();
using var busA = Open(session, 1);
using var rawBus = Open(session, 2);
using var client = CanOpen.OpenNode(busA, nodeId: 0x01,
new CanOpenNodeOptions().With(maxSdoTransferBytes: 1 << 20));
PeerSdoLaboratory.Bind(client, 0x02);
using var tap = new FrameTap(rawBus, CanOpenCobId.SdoRx(0x02));
var payload = Enumerable.Range(0, GrowthPayloadBytes).Select(i => (byte)(i * 13)).ToArray();

var before = GC.GetTotalAllocatedBytes(precise: true);
var upload = client.SdoUploadAsync(0x02, 0x2100, 0x00);
tap.Next(ShortTimeout)[0].Should().Be(SdoFrames.CcsUploadInit);
// Initiate response without the size indicator: the client cannot size its buffer.
Send(rawBus, CanOpenCobId.SdoTx(0x02), new byte[] { SdoFrames.ScsUploadInitSegmented & 0xFE, 0x00, 0x21, 0x00, 0, 0, 0, 0 });
var toggle = false;
for (var offset = 0; offset < payload.Length; offset += 7)
{
tap.Next(ShortTimeout);
var chunk = payload.AsSpan(offset, Math.Min(7, payload.Length - offset));
Send(rawBus, CanOpenCobId.SdoTx(0x02), SdoFrames.BuildSegment(
SdoFrames.ScsUploadSegmentBase, toggle, offset + 7 >= payload.Length, chunk));
toggle = !toggle;
}
var received = await upload.WithTimeoutAsync(TimeSpan.FromSeconds(30));
var allocated = GC.GetTotalAllocatedBytes(precise: true) - before;

received.Should().Equal(payload);
allocated.Should().BeLessThan(GrowthAllocationBound);
}

[Fact]
public void Sdo_BlockDownload_Server_Sizeless_Grows_Geometrically()
{
var session = NewSession();
using var busB = Open(session, 1);
using var rawBus = Open(session, 2);
using var server = CanOpen.OpenNode(busB, nodeId: 0x02,
new CanOpenNodeOptions().With(maxSdoTransferBytes: 1 << 20));
server.ObjectDictionary.AddDomain(0x2100, 0x00, new byte[4]);
using var tap = new FrameTap(rawBus, CanOpenCobId.SdoTx(0x02));
var payload = Enumerable.Range(0, GrowthPayloadBytes).Select(i => (byte)(i * 13)).ToArray();

var before = GC.GetTotalAllocatedBytes(precise: true);
Send(rawBus, CanOpenCobId.SdoRx(0x02),
SdoBlockFrames.BuildBlockDownloadInit(0x2100, 0x00, clientCrcSupported: false, sizeIndicated: false, totalSize: 0));
var blockSize = tap.Next(ShortTimeout)[4];
var offset = 0;
while (offset < payload.Length)
{
byte seq = 0;
while (seq < blockSize && offset < payload.Length)
{
seq++;
var chunk = payload.AsSpan(offset, Math.Min(7, payload.Length - offset));
offset += chunk.Length;
Send(rawBus, CanOpenCobId.SdoRx(0x02),
SdoBlockFrames.BuildSegment(seq, offset >= payload.Length, chunk));
}
tap.Next(ShortTimeout)[0].Should().Be(SdoBlockFrames.ScsBlockDownloadSubBlockAck);
}
var unused = (byte)((7 - payload.Length % 7) % 7);
Send(rawBus, CanOpenCobId.SdoRx(0x02), SdoBlockFrames.BuildEnd(SdoBlockFrames.CcsBlockDownloadEndBase, unused, 0));
tap.Next(ShortTimeout)[0].Should().Be(SdoBlockFrames.ScsBlockDownloadEndResponse);
var allocated = GC.GetTotalAllocatedBytes(precise: true) - before;

server.ObjectDictionary.ReadRaw(0x2100, 0x00).Should().Equal(payload);
allocated.Should().BeLessThan(GrowthAllocationBound);
}

[Fact]
public async Task Sdo_BlockUpload_Client_Sizeless_Grows_Geometrically()
{
var session = NewSession();
using var busA = Open(session, 1);
using var rawBus = Open(session, 2);
using var client = CanOpen.OpenNode(busA, nodeId: 0x01,
new CanOpenNodeOptions().With(maxSdoTransferBytes: 1 << 20));
PeerSdoLaboratory.Bind(client, 0x02);
using var tap = new FrameTap(rawBus, CanOpenCobId.SdoRx(0x02));
var payload = Enumerable.Range(0, GrowthPayloadBytes).Select(i => (byte)(i * 13)).ToArray();

var before = GC.GetTotalAllocatedBytes(precise: true);
var upload = client.SdoUploadAsync(0x02, 0x2100, 0x00, SdoTransferMode.Block);
var blockSize = tap.Next(ShortTimeout)[4];
Send(rawBus, CanOpenCobId.SdoTx(0x02), SdoBlockFrames.BuildBlockUploadInitResponse(
0x2100, 0x00, serverCrcSupported: false, sizeIndicated: false, totalSize: 0));
tap.Next(ShortTimeout)[0].Should().Be(SdoBlockFrames.CcsBlockUploadStart);
var offset = 0;
while (offset < payload.Length)
{
byte seq = 0;
while (seq < blockSize && offset < payload.Length)
{
seq++;
var chunk = payload.AsSpan(offset, Math.Min(7, payload.Length - offset));
offset += chunk.Length;
Send(rawBus, CanOpenCobId.SdoTx(0x02),
SdoBlockFrames.BuildSegment(seq, offset >= payload.Length, chunk));
}
tap.Next(ShortTimeout)[0].Should().Be(SdoBlockFrames.CcsBlockUploadSubBlockAck);
}
var unused = (byte)((7 - payload.Length % 7) % 7);
Send(rawBus, CanOpenCobId.SdoTx(0x02), SdoBlockFrames.BuildEnd(SdoBlockFrames.ScsBlockUploadEndBase, unused, 0));
tap.Next(ShortTimeout)[0].Should().Be(SdoBlockFrames.CcsBlockUploadEndResponse);
var received = await upload.WithTimeoutAsync(ShortTimeout);
var allocated = GC.GetTotalAllocatedBytes(precise: true) - before;

received.Should().Equal(payload);
allocated.Should().BeLessThan(GrowthAllocationBound);
}

#endif

[Fact]
public void Sdo_Server_Sizeless_Download_Aborts_OutOfMemory_Before_Passing_The_Cap()
{
Expand Down
Loading