Task #143 — Modbus block-read coalescing (with max-gap knob)
Adds a coalescing read planner that merges nearby tags into single FC03/FC04 PDUs, opt-in via ModbusDriverOptions.MaxReadGap. Default 0 = no coalescing (every tag gets its own PDU — preserves pre-#143 wire output). Worked example with MaxReadGap=10: T1 @ HR 100 (Int16, 1 reg) T2 @ HR 102 (Int16, 1 reg, gap 1 → joins block) T3 @ HR 110 (Float32, 2 regs, gap 7 → joins block) T4 @ HR 200 (Int16, 1 reg, gap 89 → splits, separate read) → 2 PDUs total: FC03 start=100 quantity=12 + FC03 start=200 quantity=1. Planner: - Eligible tags: known + register region (HR/IR) + scalar + not String / BitInRegister / array + not CoalesceProhibited. - Groups by (UnitId, Region) — never coalesces across slaves or regions. - Sorts by start address; merges when (next.start - last.end - 1) ≤ MaxReadGap AND the resulting span ≤ MaxRegistersPerRead. Otherwise opens a new block. - Single-tag blocks are deferred to the per-tag path so WriteOnChange cache semantics stay correct without duplication. - Per-block failure marks every member tag Bad and degrades health — same semantics the per-tag path has, but at the block granularity. Per-tag escape hatch ModbusTagDefinition.CoalesceProhibited (bool, default false) — when true, the tag is read in isolation regardless of MaxReadGap. For PLCs with protected register holes between adjacent tags. Tests (7 new ModbusCoalescingTests): - MaxReadGap=0 keeps the per-tag behavior (2 reads for 2 tags). - MaxReadGap=2 merges 3 tags within 5 registers into 1 read of qty=5. - MaxReadGap=10 splits T1+T2 from T3 when the gap exceeds the threshold. - CoalesceProhibited tag reads alone even when neighbours are eligible. - Coalescing never crosses UnitId boundaries (multi-slave gateway safety). - MaxRegistersPerRead caps a would-be block; planner falls back to separate reads when the merged span would exceed the cap. - Per-tag values surface independently after coalescing (slice-math sanity). Existing 220 unit tests still green; total 224 pass with the new file (tests are additive, no regressions). Follow-up: auto-split-on-protected-hole isn't shipped — a coalesced read that hits an Illegal Data Address right now marks every member Bad until the operator sets CoalesceProhibited on the offending tag. Tracked implicitly by #138's e2e drill against a pymodbus profile with a protected hole mid-block.
This commit is contained in:
@@ -196,8 +196,17 @@ public sealed class ModbusDriver
|
||||
var transport = RequireTransport();
|
||||
var now = DateTime.UtcNow;
|
||||
var results = new DataValueSnapshot[fullReferences.Count];
|
||||
|
||||
// #143 block-read coalescing: when MaxReadGap is non-zero, route eligible tags through
|
||||
// the coalescing planner first. Tags it can't coalesce (arrays, coils, prohibited,
|
||||
// unknown) fall through to the per-tag loop below with results[i] still default.
|
||||
var coalesced = _options.MaxReadGap > 0
|
||||
? await ReadCoalescedAsync(transport, fullReferences, results, now, cancellationToken).ConfigureAwait(false)
|
||||
: new HashSet<int>();
|
||||
|
||||
for (var i = 0; i < fullReferences.Count; i++)
|
||||
{
|
||||
if (coalesced.Contains(i)) continue;
|
||||
if (!_tagsByName.TryGetValue(fullReferences[i], out var tag))
|
||||
{
|
||||
results[i] = new DataValueSnapshot(null, StatusBadNodeIdUnknown, null, now);
|
||||
@@ -377,6 +386,130 @@ public sealed class ModbusDriver
|
||||
/// <summary>Resolve the UnitId for a tag — per-tag override (#142) or driver-level fallback.</summary>
|
||||
private byte ResolveUnitId(ModbusTagDefinition tag) => tag.UnitId ?? _options.UnitId;
|
||||
|
||||
/// <summary>
|
||||
/// #143 block-read coalescing planner. Groups eligible tags by (UnitId, Region), sorts
|
||||
/// by start address, and merges adjacent / near-adjacent (gap ≤ MaxReadGap) into single
|
||||
/// FC03/FC04 reads. Per-block: emit one Modbus PDU, slice the response back into per-tag
|
||||
/// values, populate <paramref name="results"/> and the WriteOnChangeOnly cache. Returns
|
||||
/// the set of <paramref name="fullReferences"/> indices the planner handled — the
|
||||
/// caller falls back to the per-tag path for the rest (arrays, coils, prohibited, unknown).
|
||||
/// </summary>
|
||||
private async Task<HashSet<int>> ReadCoalescedAsync(
|
||||
IModbusTransport transport,
|
||||
IReadOnlyList<string> fullReferences,
|
||||
DataValueSnapshot[] results,
|
||||
DateTime timestamp,
|
||||
CancellationToken ct)
|
||||
{
|
||||
// Eligible: known tag, register region (HoldingRegisters or InputRegisters), scalar
|
||||
// (no array, no string), not CoalesceProhibited, not BitInRegister.
|
||||
var eligible = new List<(int Index, string Ref, ModbusTagDefinition Tag)>();
|
||||
for (var i = 0; i < fullReferences.Count; i++)
|
||||
{
|
||||
if (!_tagsByName.TryGetValue(fullReferences[i], out var tag)) continue;
|
||||
if (tag.CoalesceProhibited) continue;
|
||||
if (tag.ArrayCount.HasValue) continue;
|
||||
if (tag.Region is not (ModbusRegion.HoldingRegisters or ModbusRegion.InputRegisters)) continue;
|
||||
if (tag.DataType is ModbusDataType.String or ModbusDataType.BitInRegister) continue;
|
||||
eligible.Add((i, fullReferences[i], tag));
|
||||
}
|
||||
if (eligible.Count == 0) return new HashSet<int>();
|
||||
|
||||
var handled = new HashSet<int>();
|
||||
|
||||
// Group by (UnitId, Region) — coalescing across slaves or regions is unsafe.
|
||||
foreach (var group in eligible.GroupBy(e => (Unit: ResolveUnitId(e.Tag), e.Tag.Region)))
|
||||
{
|
||||
var fc = group.Key.Region == ModbusRegion.HoldingRegisters ? (byte)0x03 : (byte)0x04;
|
||||
var cap = _options.MaxRegistersPerRead == 0 ? (ushort)125 : _options.MaxRegistersPerRead;
|
||||
var sorted = group.OrderBy(e => e.Tag.Address).ToList();
|
||||
|
||||
// Build merged blocks. A "block" is (start, lastEnd, members[]) where lastEnd is
|
||||
// the inclusive end address of the last tag's register span. A new tag joins the
|
||||
// block if its start ≤ lastEnd + 1 + MaxReadGap AND the resulting span ≤ cap.
|
||||
var blocks = new List<(ushort Start, ushort End, List<(int Index, ModbusTagDefinition Tag)> Members)>();
|
||||
foreach (var (idx, _, tag) in sorted)
|
||||
{
|
||||
var tagStart = tag.Address;
|
||||
var tagEnd = (ushort)(tag.Address + RegisterCount(tag) - 1);
|
||||
if (blocks.Count > 0)
|
||||
{
|
||||
var last = blocks[^1];
|
||||
var gap = tagStart - last.End - 1;
|
||||
var newEnd = Math.Max(tagEnd, last.End);
|
||||
var newSpan = newEnd - last.Start + 1;
|
||||
if (gap <= _options.MaxReadGap && newSpan <= cap)
|
||||
{
|
||||
last.Members.Add((idx, tag));
|
||||
blocks[^1] = (last.Start, (ushort)newEnd, last.Members);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
blocks.Add((tagStart, tagEnd, new List<(int, ModbusTagDefinition)> { (idx, tag) }));
|
||||
}
|
||||
|
||||
// Issue one PDU per block. On block-level failure mark every member Bad — caller's
|
||||
// per-tag fallback won't re-try since handled-set already includes them; auto-split-
|
||||
// on-failure is a follow-up.
|
||||
foreach (var block in blocks)
|
||||
{
|
||||
if (block.Members.Count == 1)
|
||||
{
|
||||
// Lone tag — let the per-tag path handle it for symmetry with WriteOnChange
|
||||
// cache invalidation. Costs nothing because one tag = one PDU either way.
|
||||
continue;
|
||||
}
|
||||
|
||||
var qty = (ushort)(block.End - block.Start + 1);
|
||||
try
|
||||
{
|
||||
var data = await ReadRegisterBlockAsync(transport, group.Key.Unit, fc, block.Start, qty, ct).ConfigureAwait(false);
|
||||
foreach (var (idx, tag) in block.Members)
|
||||
{
|
||||
var sliceOffsetBytes = (tag.Address - block.Start) * 2;
|
||||
var sliceLenBytes = RegisterCount(tag) * 2;
|
||||
var value = DecodeRegister(data.AsSpan(sliceOffsetBytes, sliceLenBytes), tag);
|
||||
results[idx] = new DataValueSnapshot(value, 0u, timestamp, timestamp);
|
||||
handled.Add(idx);
|
||||
InvalidateWriteCacheIfDiverged(fullReferences[idx], value);
|
||||
}
|
||||
_health = new DriverHealth(DriverState.Healthy, timestamp, null);
|
||||
}
|
||||
catch (ModbusException mex)
|
||||
{
|
||||
var status = MapModbusExceptionToStatus(mex.ExceptionCode);
|
||||
foreach (var (idx, _) in block.Members)
|
||||
{
|
||||
results[idx] = new DataValueSnapshot(null, status, null, timestamp);
|
||||
handled.Add(idx);
|
||||
}
|
||||
_health = new DriverHealth(DriverState.Degraded, _health.LastSuccessfulRead, mex.Message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
foreach (var (idx, _) in block.Members)
|
||||
{
|
||||
results[idx] = new DataValueSnapshot(null, StatusBadCommunicationError, null, timestamp);
|
||||
handled.Add(idx);
|
||||
}
|
||||
_health = new DriverHealth(DriverState.Degraded, _health.LastSuccessfulRead, ex.Message);
|
||||
}
|
||||
}
|
||||
}
|
||||
return handled;
|
||||
}
|
||||
|
||||
private void InvalidateWriteCacheIfDiverged(string fullRef, object value)
|
||||
{
|
||||
if (!_options.WriteOnChangeOnly) return;
|
||||
lock (_lastWrittenLock)
|
||||
{
|
||||
if (_lastWrittenByRef.TryGetValue(fullRef, out var prev) && !Equals(prev, value))
|
||||
_lastWrittenByRef.Remove(fullRef);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private async Task<byte[]> ReadRegisterBlockAsync(
|
||||
IModbusTransport transport, byte unitId, byte fc, ushort address, ushort quantity, CancellationToken ct)
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user