diff --git a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/Sparkplug/SparkplugIngestor.cs b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/Sparkplug/SparkplugIngestor.cs
index 0133813d..cb14d71f 100644
--- a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/Sparkplug/SparkplugIngestor.cs
+++ b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/Sparkplug/SparkplugIngestor.cs
@@ -108,6 +108,25 @@ public sealed class SparkplugIngestor
/// The subscribed reference is not an authored Sparkplug tag.
private const uint StatusBadNodeIdUnknown = 0x80340000u;
+ ///
+ /// Hard ceiling on . Every key fed to is derived
+ /// from authored configuration, so the set is bounded by construction — this is the belt to that
+ /// braces. A key set that grows off wire input is a memory leak on the dispatcher path,
+ /// and getting one of the derivations subtly wrong later must degrade to "stops warning" rather
+ /// than to unbounded growth.
+ ///
+ private const int MaxWarnKeys = 1024;
+
+ ///
+ /// How many consecutive metadata rebirths one edge node may be asked for before the
+ /// driver gives up and says so. Bounds the rebirth loop the per-node debounce cannot break: the
+ /// debounce caps the RATE, not the LIFETIME, so a node whose birth this driver can never use —
+ /// one over , or one that simply never comes — would otherwise be
+ /// commanded every debounce period for the life of the process, answering each time with the
+ /// same unusable payload.
+ ///
+ private const int MaxConsecutiveMetadataRebirths = 3;
+
private readonly string _driverId;
private readonly ILogger? _logger;
private readonly TimeSpan _rebirthDebounce;
@@ -125,6 +144,20 @@ public sealed class SparkplugIngestor
/// Last rebirth-NCMD instant per edge node, in UTC ticks — the debounce floor.
private readonly ConcurrentDictionary _lastRebirthTicks = new();
+ ///
+ /// Edge nodes whose NBIRTH has been seen in the current session, whether or not its
+ /// seq was readable. Distinct from ,
+ /// which additionally requires a usable seq — see .
+ ///
+ private readonly ConcurrentDictionary _nodeBirthSeen = new();
+
+ ///
+ /// Consecutive metadata rebirths requested from an edge node since its last applied
+ /// birth. Capped by — see
+ /// .
+ ///
+ private readonly ConcurrentDictionary _metadataRebirths = new();
+
/// Keys whose anomaly has already been logged loudly once, cleared by .
private readonly ConcurrentDictionary _warned = new(StringComparer.Ordinal);
@@ -249,6 +282,20 @@ public sealed class SparkplugIngestor
}
}
+ ///
+ /// The pre-3.0 primary-host STATE filter (STATE/{hostId}, no namespace prefix), or
+ /// when no host id is configured. Subscribed alongside
+ /// so the legacy form this driver already parses can actually arrive.
+ ///
+ public string? LegacyStateFilter
+ {
+ get
+ {
+ var hostId = _options.Sparkplug?.HostId;
+ return string.IsNullOrWhiteSpace(hostId) ? null : $"STATE/{hostId}";
+ }
+ }
+
///
/// Wires this ingestor onto a live connection: inbound messages route to
/// , the connection becomes both transports, and every reconnect runs
@@ -259,6 +306,11 @@ public sealed class SparkplugIngestor
{
ArgumentNullException.ThrowIfNull(connection);
+ // Binding to a session means the previous one is gone. Reset FIRST, before delivery is wired:
+ // with a persistent session a broker may push queued messages the moment CONNACK lands, ahead
+ // of our own SUBSCRIBE, and those must not resolve against the old session's aliases.
+ ResetSessionState("attach to a new broker session");
+
_subscribeTransport = connection;
_publishTransport = connection;
connection.MessageReceived += HandleMessage;
@@ -311,6 +363,32 @@ public sealed class SparkplugIngestor
{
_trackers.TryRemove(key, out _);
_lastRebirthTicks.TryRemove(key, out _);
+ _nodeBirthSeen.TryRemove(key, out _);
+ _metadataRebirths.TryRemove(key, out _);
+ }
+ }
+
+ // …and their BIRTHS, which is the half that actually mis-routes. While an edge node is
+ // unauthored its traffic is dropped at Dispatch, so its cached alias table stops refreshing —
+ // it FREEZES rather than going stale-and-obvious. Deploy A authors Plant1/EdgeA/Filler1,
+ // deploy B drops it, deploy C re-adds it a week later having missed every rebirth in between:
+ // Births.Find still answers with the week-old table and the next DDATA publishes the wrong
+ // metric's value under the right RawPath at Good quality. It is also why Births would otherwise
+ // grow monotonically across redeploys.
+ //
+ // The predicate is edge-node departure ONLY, deliberately narrower than "absent from ByScope".
+ // A device with no authored tags but whose EDGE NODE is still authored keeps receiving traffic
+ // (the Dispatch filter is node-level, and OnDeviceBirth runs ahead of the authored-scope
+ // filter), so its birth is live, not frozen — evicting it would throw away correct state and
+ // cost a needless rebirth round-trip the moment a later deploy authored a tag on it.
+ foreach (var scope in Births.Scopes)
+ {
+ var owner = new SparkplugNodeKey(scope.GroupId, scope.EdgeNodeId);
+ if (!table.EdgeNodes.Contains(owner))
+ {
+ // Forgets the node scope and every device under it in one call; idempotent, so a node
+ // contributing several scopes to this loop costs one eviction and some no-ops.
+ Births.ForgetNode(owner.GroupId, owner.EdgeNodeId);
}
}
@@ -359,6 +437,22 @@ public sealed class SparkplugIngestor
/// No group id is configured, or every filter was refused.
public async Task EstablishAsync(CancellationToken cancellationToken)
{
+ // The SAME reset OnReconnectedAsync performs, and for the same reason — this ingestor OUTLIVES
+ // a session rebuild. `MqttDriver.AdoptOptions` only replaces it when `!SameIngest`, but
+ // `SessionIdentity` is a strict SUPERSET of `IngestIdentity`: changing Host, Port, ClientId,
+ // credentials, TLS or a timeout alone yields `SameSession == false` AND `SameIngest == true`,
+ // so `ReinitializeAsync` tears the session down and re-establishes on this very instance.
+ // `TeardownAsync` touches no ingest state. The driver-host resilience layer re-running
+ // `InitializeAsync` after a fault is the second path onto this line.
+ //
+ // Without the reset, an edge node that restarted during the changeover and rebound alias 5 from
+ // Temperature to Pressure would have its first DDATA on the new session resolved against the
+ // OLD table — publishing Pressure's value under the Temperature RawPath at Good quality. The
+ // late-join rebirth below does not close that: it is dispatched off-thread (DATA can arrive
+ // first), passes the debounce, and is a no-op under `requestRebirthOnGap: false` or against a
+ // node that ignores the NCMD.
+ ResetSessionState("establish a broker session");
+
await SubscribeGroupAsync(cancellationToken).ConfigureAwait(false);
RequestLateJoinRebirths("initial connect");
}
@@ -457,13 +551,59 @@ public sealed class SparkplugIngestor
/// Every filter was refused — the caller tears the session down.
public async Task OnReconnectedAsync(CancellationToken cancellationToken)
{
- Births.Clear();
- _trackers.Clear();
+ ResetSessionState("reconnect");
await SubscribeGroupAsync(cancellationToken).ConfigureAwait(false);
RequestLateJoinRebirths("reconnect");
}
+ ///
+ /// Drops every scrap of per-session ingest state: births, sequence baselines, the birth-seen
+ /// set and the late-join counters. Called from every seam at which the session underneath
+ /// this ingestor changes — , and
+ /// .
+ ///
+ ///
+ ///
+ /// This is the correctness step, not housekeeping — and the reason it lives in one
+ /// method called from three places rather than inline in one. Every piece of state cleared
+ /// here is a claim about a session that no longer exists: an alias table binds names to
+ /// numbers the previous edge-node session chose, and a sequence baseline is only
+ /// evidence relative to the stream it was taken from. Carried across, an alias resolves to
+ /// whichever metric used to own it (good quality, plausible value, wrong tag) and a fresh
+ /// stream can look contiguous purely by coincidence.
+ ///
+ ///
+ /// Unconditionally safe. On a freshly-constructed ingestor every collection is
+ /// already empty, so calling it three times during one connect costs three no-ops.
+ ///
+ ///
+ /// Values are deliberately not staled here. A session change is a driver-side event,
+ /// not a statement about the plant; last observed values remain the best available answer
+ /// until a birth lands, and the connection's own health surface already reports the outage.
+ /// Only a death — the source saying it is offline — stales a tag.
+ ///
+ ///
+ /// Why the reset is happening, for the diagnostic log.
+ private void ResetSessionState(string reason)
+ {
+ var births = Births.Count;
+ Births.Clear();
+ _trackers.Clear();
+ _nodeBirthSeen.Clear();
+ _metadataRebirths.Clear();
+
+ if (births > 0)
+ {
+ _logger?.LogInformation(
+ "MQTT driver '{DriverId}': dropped {Count} Sparkplug birth scope(s) and every sequence "
+ + "baseline — {Reason}. Aliases will be rebuilt from the next birth certificate.",
+ _driverId,
+ births,
+ reason);
+ }
+ }
+
///
/// Routes one inbound MQTT message: parses the topic, decodes the Sparkplug payload and hands it
/// to . Runs on MQTTnet's dispatcher thread and never throws.
@@ -489,26 +629,68 @@ public sealed class SparkplugIngestor
return;
}
+ if (parsed.Type == SparkplugMessageType.STATE)
+ {
+ // A STATE body is a handful of bytes; an oversized one is a broken publisher, not a
+ // dispatcher-thread hazard, and it is logged at Debug rather than through WarnOnce
+ // because its key would be a host id this driver did not choose.
+ if (payload.Length > MaxPayloadBytes)
+ {
+ _logger?.LogDebug(
+ "MQTT driver '{DriverId}': STATE on '{Topic}' exceeded the {Max}-byte ceiling; dropped.",
+ _driverId,
+ topic,
+ MaxPayloadBytes);
+ return;
+ }
+
+ HandleState(parsed, payload);
+ return;
+ }
+
+ // The authored-node filter runs HERE as well as inside Dispatch, and the order matters for
+ // two reasons. It keeps the oversize warning's key bounded — keyed on the raw topic it grew
+ // once per distinct topic off a group-wide '#' subscription, i.e. off wire input a plant (or
+ // a hostile publisher) controls. And it drops an unauthored node's body before the protobuf
+ // parse rather than after, which on a busy group is most of the decode work this driver
+ // would otherwise do for traffic it discards.
+ if (parsed.GroupId is not { } groupId || parsed.EdgeNodeId is not { } edgeNodeId)
+ {
+ return;
+ }
+
+ var node = new SparkplugNodeKey(groupId, edgeNodeId);
+ if (!_authored.EdgeNodes.Contains(node))
+ {
+ _logger?.LogDebug(
+ "MQTT driver '{DriverId}': {Type} for unauthored edge node '{Node}'; ignored.",
+ _driverId,
+ parsed.Type,
+ node);
+ return;
+ }
+
if (payload.Length > MaxPayloadBytes)
{
// Refused BEFORE any decode: the whole point of the bound is that one publisher cannot
// impose an unbounded protobuf parse on the shared dispatcher thread. Nothing mutates.
+ //
+ // An oversized NBIRTH is the worst case and is called out by name: the driver cannot
+ // learn this node's aliases at all, so every subsequent DATA message is unroutable. The
+ // late-join cap in CheckSequence is what stops that becoming a permanent NCMD loop.
WarnOnce(
- $"oversize:{topic}",
- "MQTT driver '{DriverId}': '{Topic}' carried {Bytes} bytes, over the {Max}-byte ceiling; dropped.",
+ $"oversize:{node}",
+ "MQTT driver '{DriverId}': {Type} from '{Node}' carried {Bytes} bytes, over the "
+ + "{Max}-byte ceiling; dropped. Raise maxPayloadBytes if this node's birth is "
+ + "legitimately this large.",
_driverId,
- topic,
+ parsed.Type,
+ node,
payload.Length,
MaxPayloadBytes);
return;
}
- if (parsed.Type == SparkplugMessageType.STATE)
- {
- HandleState(parsed, payload);
- return;
- }
-
Dispatch(parsed, SparkplugCodec.Decode(payload));
}
catch (Exception ex)
@@ -645,6 +827,12 @@ public sealed class SparkplugIngestor
// AcceptNodeBirth, never Accept: an NBIRTH restarts the sequence by definition, so running it
// through the gap detector would flag a gap on the very message that resynchronized everything
// and answer a birth with a demand for another one.
+ // Recorded BEFORE the seq is judged, and independently of the verdict: "a birth arrived" and
+ // "the birth established a sequence baseline" are different facts, and conflating them is the
+ // rebirth loop CheckSequence documents.
+ _nodeBirthSeen[node] = 0;
+ _metadataRebirths.TryRemove(node, out _);
+
var bdSeq = SequenceTracker.TryReadBdSeq(payload, out var token) ? token : (ulong?)null;
if (!TrackerFor(node).AcceptNodeBirth(payload.Seq, bdSeq))
{
@@ -682,6 +870,10 @@ public sealed class SparkplugIngestor
// A DBIRTH is a sequenced member of the edge node's stream — only the NBIRTH restarts it.
CheckSequence(node, payload, "DBIRTH");
+ // Metadata arrived, so the give-up counter re-arms: a device that re-births after a bad patch
+ // must not stay permanently un-askable because of requests made before it recovered.
+ _metadataRebirths.TryRemove(node, out _);
+
var result = Births.ApplyDeviceBirth(node.GroupId, node.EdgeNodeId, deviceId, payload.Metrics);
var scope = new SparkplugScope(node.GroupId, node.EdgeNodeId, deviceId);
LogBirthAnomaly(scope.ToString(), result);
@@ -725,7 +917,7 @@ public sealed class SparkplugIngestor
+ "requesting a rebirth.",
_driverId,
scope);
- RequestRebirth(node, "data before birth");
+ RequestMetadataRebirth(node, "data before birth");
return;
}
@@ -754,7 +946,7 @@ public sealed class SparkplugIngestor
+ "declare; requesting a rebirth.",
_driverId,
scope);
- RequestRebirth(node, "unknown alias");
+ RequestMetadataRebirth(node, "unknown alias");
}
}
@@ -765,12 +957,19 @@ public sealed class SparkplugIngestor
///
private static SparkplugMetricBinding? ResolveBinding(AliasTable table, SparkplugMetric metric)
{
- if (metric.Alias is { } alias && table.Resolve(alias) is { } byAlias)
+ // NAME FIRST, deliberately. A DATA metric that carries a name has stated its identity in the
+ // one vocabulary that survives a rebirth; the alias is a per-birth compression detail. When the
+ // two disagree — a publisher that recycled an alias mid-stream, or a birth this driver applied
+ // out of order — preferring the alias would resolve to whichever metric currently occupies that
+ // slot, which is precisely the mis-route this type's binding rule exists to prevent. The common
+ // case costs nothing: a real post-birth DATA metric carries no name at all, so this is a null
+ // check before the alias lookup.
+ if (metric.Name is { } name && table.ResolveByName(name) is { } byName)
{
- return byAlias;
+ return byName;
}
- return metric.Name is { } name ? table.ResolveByName(name) : null;
+ return metric.Alias is { } alias ? table.Resolve(alias) : null;
}
// ---- deaths ----
@@ -878,7 +1077,15 @@ public sealed class SparkplugIngestor
online ? "ONLINE" : "OFFLINE");
}
- HostStateObserved?.Invoke(this, new SparkplugHostStateEventArgs(state));
+ try
+ {
+ HostStateObserved?.Invoke(this, new SparkplugHostStateEventArgs(state));
+ }
+ catch (Exception ex)
+ {
+ // Same containment as every other raise here — this runs on MQTTnet's dispatcher thread.
+ _logger?.LogError(ex, "MQTT driver '{DriverId}': a host-state subscriber threw.", _driverId);
+ }
}
/// Reads the online flag out of a STATE body — v3.0 JSON first, then the legacy text form.
@@ -976,6 +1183,25 @@ public sealed class SparkplugIngestor
return; // A published metric nobody authored a tag for. Not an error — that is most of them.
}
+ // Arrays slip past the unsupported-type gate below, because ToDriverDataType maps every
+ // *Array variant to its ELEMENT type and is therefore never null for one. Design §3.5 defers
+ // array support, and the failure is not benign: the value rides in bytes_value as a byte[] at
+ // ValueKind.Scalar, so a numeric target throws its way to BadTypeMismatch (safe, wrong reason)
+ // while a String/Reference target reaches Convert.ToBase64String and publishes a packed binary
+ // buffer as a plausible string at GOOD quality. Gate on the array bit itself.
+ if (binding.DataType.IsSparkplugArray())
+ {
+ WarnOnce(
+ $"array-type:{key}",
+ "MQTT driver '{DriverId}': metric '{Metric}' on '{Scope}' is a Sparkplug array "
+ + "({DataType}); array metrics are not supported in v1 and its tag(s) are not being fed.",
+ _driverId,
+ binding.Name,
+ scope,
+ binding.DataType);
+ return;
+ }
+
var birthType = binding.DataType.ToDriverDataType();
if (birthType is null)
{
@@ -1187,17 +1413,76 @@ public sealed class SparkplugIngestor
return;
}
- if (!tracker.IsBirthSynchronized)
+ if (tracker.IsBirthSynchronized)
+ {
+ return;
+ }
+
+ // "Birth-synchronized" in SequenceTracker's sense means a birth established a usable SEQUENCE
+ // BASELINE — it is false both when no NBIRTH has arrived (a genuine late join, which a rebirth
+ // fixes) and when one arrived carrying a seq this driver could not read (which a rebirth does
+ // NOT fix: the node answers with the same unusable birth). Reacting to the two identically is
+ // what turns the second into a permanent NCMD loop at the debounce cadence, defeating the guard
+ // OnNodeBirth deliberately applies one method over. Whether a birth ARRIVED is this type's own
+ // fact to keep, so it keeps it.
+ if (_nodeBirthSeen.ContainsKey(node))
+ {
+ return;
+ }
+
+ WarnOnce(
+ $"latejoin:{node}",
+ "MQTT driver '{DriverId}': {Kind} from '{Node}' arrived with no NBIRTH seen since connect; "
+ + "requesting a rebirth.",
+ _driverId,
+ kind,
+ node);
+
+ // The residual loop the birth-seen flag cannot see: an NBIRTH that never reaches Dispatch at
+ // all, because it exceeded MaxPayloadBytes or because the node never sends one. Nothing about
+ // asking again changes that, so it goes through the capped path.
+ RequestMetadataRebirth(node, "late join (no NBIRTH observed)");
+ }
+
+ ///
+ /// Requests a rebirth for a missing-metadata reason — late join, data before birth, or an
+ /// alias the current birth does not declare — under a per-node cap on consecutive attempts.
+ ///
+ ///
+ ///
+ /// Why these three and not a sequence gap. A gap is self-limiting: the tracker
+ /// resynchronizes onto the observed value, so the very next message is contiguous and asks
+ /// for nothing. The three reasons routed here are not — every one of them means "this
+ /// driver holds no usable catalog for this node", and if the node cannot supply one (its
+ /// birth is over the payload ceiling, or it ignores NCMD entirely) then every
+ /// subsequent message re-raises the same condition. The debounce caps the rate at one
+ /// request per period; only this caps the lifetime.
+ ///
+ ///
+ /// Consecutive, not cumulative. Any applied birth clears the counter, so a node that
+ /// recovers is not permanently penalised for having once been unreachable.
+ ///
+ ///
+ /// The edge node to command.
+ /// Why, for the log.
+ private void RequestMetadataRebirth(SparkplugNodeKey node, string reason)
+ {
+ var attempts = _metadataRebirths.AddOrUpdate(node, 1, static (_, n) => n + 1);
+ if (attempts > MaxConsecutiveMetadataRebirths)
{
WarnOnce(
- $"latejoin:{node}",
- "MQTT driver '{DriverId}': {Kind} from '{Node}' arrived with no NBIRTH seen since connect; "
- + "requesting a rebirth.",
+ $"metadata-exhausted:{node}",
+ "MQTT driver '{DriverId}': '{Node}' has not delivered a usable birth certificate after "
+ + "{Attempts} rebirth request(s); giving up until one arrives on its own. Its tags will "
+ + "not update. Check the node's birth payload size against maxPayloadBytes, its seq "
+ + "field, and that it honours NCMD Node Control/Rebirth.",
_driverId,
- kind,
- node);
- RequestRebirth(node, "late join (no NBIRTH observed)");
+ node,
+ MaxConsecutiveMetadataRebirths);
+ return;
}
+
+ RequestRebirth(node, reason);
}
///
@@ -1359,6 +1644,15 @@ public sealed class SparkplugIngestor
if (StateFilter is { } stateFilter)
{
filters.Add(new MqttTopicSubscription(stateFilter, Qos: 1, SeedRetained: true));
+
+ // The pre-3.0 form, which HandleState already parses. Subscribing only the v3.0 topic made
+ // that tolerance unreachable in production — the legacy branch could be exercised by a test
+ // calling HandleMessage directly and by nothing else. Design §3.1 is explicit: emit v3.0,
+ // tolerate legacy on receive, and tolerating requires asking for it.
+ filters.Add(new MqttTopicSubscription(
+ LegacyStateFilter!,
+ Qos: 1,
+ SeedRetained: true));
}
var transport = _subscribeTransport;
@@ -1427,7 +1721,10 @@ public sealed class SparkplugIngestor
///
private void WarnOnce(string key, string message, params object?[] args)
{
- if (_warned.TryAdd(key, 0))
+ // The ceiling is unreachable by design — see MaxWarnKeys. It exists so that a derivation that
+ // turns out NOT to be bounded (as the oversize path's raw-topic key was) degrades to silence
+ // instead of to unbounded growth on the dispatcher path.
+ if (_warned.Count < MaxWarnKeys && _warned.TryAdd(key, 0))
{
#pragma warning disable CA2254 // Template is a compile-time constant at every call site.
_logger?.LogWarning(message, args);
@@ -1643,7 +1940,18 @@ internal static class SparkplugValueCoercion
&& Assign(pf, out value);
}
- value = Convert.ToSingle(raw, CultureInfo.InvariantCulture);
+ var asSingle = Convert.ToSingle(raw, CultureInfo.InvariantCulture);
+
+ // Convert.ToSingle SATURATES rather than throwing: a Double of 1e300 onto a Float32
+ // tag becomes +Infinity, which would publish at GOOD quality as a number no gauge
+ // can render and no comparison behaves sensibly against. A source that was itself
+ // infinite is passed through — that is the publisher's own value, not our overflow.
+ if (!float.IsFinite(asSingle) && IsFiniteSource(raw))
+ {
+ return false;
+ }
+
+ value = asSingle;
return true;
case DriverDataType.Float64:
@@ -1653,7 +1961,13 @@ internal static class SparkplugValueCoercion
&& Assign(pd, out value);
}
- value = Convert.ToDouble(raw, CultureInfo.InvariantCulture);
+ var asDouble = Convert.ToDouble(raw, CultureInfo.InvariantCulture);
+ if (!double.IsFinite(asDouble) && IsFiniteSource(raw))
+ {
+ return false;
+ }
+
+ value = asDouble;
return true;
case DriverDataType.DateTime:
@@ -1679,6 +1993,17 @@ internal static class SparkplugValueCoercion
return true;
}
+ ///
+ /// Whether the wire value was itself a finite number — i.e. whether a non-finite result is
+ /// our narrowing overflow rather than the publisher's own Infinity/NaN.
+ ///
+ private static bool IsFiniteSource(object raw) => raw switch
+ {
+ float f => float.IsFinite(f),
+ double d => double.IsFinite(d),
+ _ => true,
+ };
+
private static bool TryCoerceBoolean(object raw, out object? value)
{
switch (raw)
diff --git a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttDriverDiscoveryTests.cs b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttDriverDiscoveryTests.cs
index 60a10cbd..7a5f28df 100644
--- a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttDriverDiscoveryTests.cs
+++ b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttDriverDiscoveryTests.cs
@@ -584,6 +584,124 @@ public sealed class MqttDriverDiscoveryTests
b.Variables[0].Info.DriverDataType.ShouldBe(DriverDataType.String);
}
+ ///
+ /// The composition-root line. Task 21 wires the ingestor's OnDataChange to the
+ /// driver's own in CreateSparkplugIngestor; every other Sparkplug test in the repo drives
+ /// directly and would pass with that lambda deleted. This repo
+ /// has twice shipped a component correct in isolation and inert in production for exactly that
+ /// reason (DeferredAddressSpaceSink, GatewayTagProvisioner) — so the re-raise is
+ /// asserted through 's own event, under the RawPath.
+ ///
+ [Fact]
+ public async Task SparkplugMode_DriverRepublishesIngestDataChanges_UnderTheRawPath()
+ {
+ var driver = SparkplugDriver(SpTag(SpTempPath, SpBlob(SpDevice, "Temperature", "Float32")));
+ await driver.SubscribeAsync([SpTempPath], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ object? sender = null;
+ string? firedRef = null;
+ object? firedValue = null;
+ ((ISubscribable)driver).OnDataChange += (s, e) =>
+ {
+ sender = s;
+ firedRef = e.FullReference;
+ firedValue = e.Snapshot.Value;
+ };
+
+ FeedNodeBirth(driver);
+ FeedDeviceBirth(driver, seq: 1, ("Temperature", 5UL, TahuDataType.Float));
+ FeedDeviceData(driver, seq: 2, (5UL, 21.5f));
+
+ firedRef.ShouldBe(SpTempPath);
+ firedValue.ShouldBe(21.5f);
+
+ // The driver re-raises with ITSELF as sender, not the ingestor — DriverHostActor attributes
+ // publishes to the driver instance.
+ sender.ShouldBeSameAs(driver);
+ }
+
+ ///
+ /// ReadAsync serves from the driver-owned LastValueCache, which the Sparkplug
+ /// ingestor is handed at construction. If the two ever stopped sharing one instance, every
+ /// Sparkplug read would answer BadWaitingForInitialData forever while subscriptions
+ /// looked perfectly healthy.
+ ///
+ [Fact]
+ public async Task SparkplugMode_ReadAsyncServesTheIngestedValue()
+ {
+ var driver = SparkplugDriver(SpTag(SpTempPath, SpBlob(SpDevice, "Temperature", "Float32")));
+
+ FeedNodeBirth(driver);
+ FeedDeviceBirth(driver, seq: 1, ("Temperature", 5UL, TahuDataType.Float));
+ FeedDeviceData(driver, seq: 2, (5UL, 33.5f));
+
+ var results = await driver.ReadAsync([SpTempPath], TestContext.Current.CancellationToken);
+
+ results[0].Value.ShouldBe(33.5f);
+ results[0].StatusCode.ShouldBe(0x00000000u);
+ }
+
+ ///
+ /// Subscribe/unsubscribe dispatch by mode. After an unsubscribe the reference is no longer
+ /// attributed to a handle, so ingest stops raising for it — while the cache keeps answering, so
+ /// a read still works. Both halves matter: a driver that kept publishing after unsubscribe
+ /// would leak notifications into a torn-down monitored item.
+ ///
+ [Fact]
+ public async Task SparkplugMode_UnsubscribeStopsNotifications_ButNotReads()
+ {
+ var driver = SparkplugDriver(SpTag(SpTempPath, SpBlob(SpDevice, "Temperature", "Float32")));
+ var handle = await driver.SubscribeAsync([SpTempPath], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ FeedNodeBirth(driver);
+ FeedDeviceBirth(driver, seq: 1, ("Temperature", 5UL, TahuDataType.Float));
+
+ var fired = 0;
+ ((ISubscribable)driver).OnDataChange += (_, _) => Interlocked.Increment(ref fired);
+
+ await driver.UnsubscribeAsync(handle, TestContext.Current.CancellationToken);
+ FeedDeviceData(driver, seq: 2, (5UL, 7.5f));
+
+ fired.ShouldBe(0);
+ (await driver.ReadAsync([SpTempPath], TestContext.Current.CancellationToken))[0].Value.ShouldBe(7.5f);
+ }
+
+ ///
+ /// The mode gate from the other side — the reviewer's reading of the branching, turned into a
+ /// test. A Plain-mode driver builds NO Sparkplug ingestor at all, so no Sparkplug traffic can
+ /// reach it and a Sparkplug-shaped blob (which carries no topic) resolves to nothing.
+ ///
+ [Fact]
+ public void PlainMode_BuildsNoSparkplugIngestor()
+ {
+ var driver = new MqttDriver(
+ new MqttDriverOptions
+ {
+ Mode = MqttMode.Plain,
+ RawTags = [SpTag(SpTempPath, SpBlob(SpDevice, "Temperature", "Float32"))],
+ },
+ "d",
+ null);
+
+ driver.Sparkplug.ShouldBeNull();
+ driver.Subscriptions.TryResolve(SpTempPath, out _).ShouldBeFalse();
+ }
+
+ ///
+ /// …and the converse: a Sparkplug-mode driver registers into the Sparkplug path, not the plain
+ /// manager. Registering into both would double-warn on every deploy and, worse, let a plain
+ /// topic route a Sparkplug tag.
+ ///
+ [Fact]
+ public void SparkplugMode_RegistersIntoTheSparkplugPathOnly()
+ {
+ var driver = SparkplugDriver(SpTag(SpTempPath, SpBlob(SpDevice, "Temperature", "Float32")));
+
+ driver.Sparkplug!.TryResolve(SpTempPath, out var def).ShouldBeTrue();
+ def.MetricName.ShouldBe("Temperature");
+ driver.Subscriptions.TryResolve(SpTempPath, out _).ShouldBeFalse();
+ }
+
// ---- Sparkplug fixtures ----
private const string SpGroup = "Plant1";
@@ -657,6 +775,22 @@ public sealed class MqttDriverDiscoveryTests
SparkplugCodec.Decode(payload.ToByteArray()));
}
+ ///
+ /// Feeds a real DDATA — alias-only metrics, exactly as a Sparkplug node publishes after a birth.
+ ///
+ private static void FeedDeviceData(MqttDriver driver, ulong seq, params (ulong Alias, float Value)[] metrics)
+ {
+ var payload = new Payload { Seq = seq, Timestamp = 1721822405000UL };
+ foreach (var (alias, value) in metrics)
+ {
+ payload.Metrics.Add(new Payload.Types.Metric { Alias = alias, FloatValue = value });
+ }
+
+ driver.Sparkplug!.Dispatch(
+ new SparkplugTopic(SparkplugMessageType.DDATA, SpGroup, SpNode, SpDevice, null),
+ SparkplugCodec.Decode(payload.ToByteArray()));
+ }
+
///
/// A birth metric declaring name/alias/datatype and no value — a birth's job here is to
/// declare the catalog, and these tests assert on the catalog, never on a published value.
diff --git a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/SparkplugIngestorTests.cs b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/SparkplugIngestorTests.cs
index 4dde20b9..1bd3b830 100644
--- a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/SparkplugIngestorTests.cs
+++ b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/SparkplugIngestorTests.cs
@@ -482,6 +482,15 @@ public sealed class SparkplugIngestorTests
ing.Dispatch(Topic(SparkplugMessageType.NDEATH), Ndeath(bdSeq: 1));
await ing.DrainPendingRebirthsAsync();
+ // This single assertion IS the complete check, and the review question "does it also verify the
+ // baseline?" has an answer rather than a gap: SequenceTracker.Accept refuses an absent seq
+ // WITHOUT adopting it as the baseline, so routing the death through it has exactly one
+ // observable consequence — a reported gap, and therefore this NCMD. Falsifiability run (b)
+ // confirmed it: inserting CheckSequence into OnNodeDeath reddens this test and only this test.
+ //
+ // A follow-up DDATA would NOT strengthen it. An accepted NDEATH evicts the node's births, so
+ // the next data message is legitimately data-before-birth and asks for a rebirth on its own
+ // merits — asserting silence there would pin the wrong behaviour.
publish.Topics.ShouldBeEmpty();
}
@@ -723,21 +732,42 @@ public sealed class SparkplugIngestorTests
Should.NotThrow(() => FeedBirths(ing));
}
- /// An oversize body is refused BEFORE any protobuf parse, and mutates nothing.
+ ///
+ /// An oversize body is refused BEFORE any protobuf parse, and mutates nothing.
+ ///
+ ///
+ /// The payload is a genuine, decodable NBIRTH — the control below proves the very same
+ /// bytes birth the scope once the ceiling is raised. An earlier version of this test fed
+ /// new byte[64], which fails to parse anyway, so it held with the size check deleted.
+ ///
[Fact]
public async Task OversizePayload_IsDroppedBeforeDecoding()
{
- var ing = NewIngestor(maxPayloadBytes: 16);
+ var wire = NbirthBytes(seq: 0, bdSeq: 1);
+ var ing = NewIngestor(maxPayloadBytes: wire.Length - 1);
await ing.SubscribeAsync([TempPath], TimeSpan.Zero, TestContext.Current.CancellationToken);
- var bytes = Nbirth(seq: 0, bdSeq: 1);
- bytes.Metrics.Count.ShouldBeGreaterThan(0);
-
- ing.HandleMessage("spBv1.0/Plant1/NBIRTH/EdgeA", new byte[64], retained: false);
+ ing.HandleMessage("spBv1.0/Plant1/NBIRTH/EdgeA", wire, retained: false);
ing.Births.Count.ShouldBe(0);
}
+ ///
+ /// Falsifiability control for the test above: the identical bytes, one byte of headroom, birth
+ /// the scope. Without this the oversize assertion is satisfied by any reason the payload failed.
+ ///
+ [Fact]
+ public async Task PayloadExactlyAtTheCeiling_IsAccepted()
+ {
+ var wire = NbirthBytes(seq: 0, bdSeq: 1);
+ var ing = NewIngestor(maxPayloadBytes: wire.Length);
+ await ing.SubscribeAsync([TempPath], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ ing.HandleMessage("spBv1.0/Plant1/NBIRTH/EdgeA", wire, retained: false);
+
+ ing.Births.Count.ShouldBe(1);
+ }
+
///
/// The driver subscribes spBv1.0/{group}/#, so a large plant delivers traffic for every
/// node in the group. Only authored edge nodes are processed — that is what bounds the tracker
@@ -877,7 +907,8 @@ public sealed class SparkplugIngestorTests
await ing.EstablishAsync(TestContext.Current.CancellationToken);
await ing.DrainPendingRebirthsAsync();
- subscribe.Filters.Select(f => f.Topic).ShouldBe(["spBv1.0/Plant1/#", "spBv1.0/STATE/otopcua-host-1"]);
+ subscribe.Filters.Select(f => f.Topic).ShouldBe(
+ ["spBv1.0/Plant1/#", "spBv1.0/STATE/otopcua-host-1", "STATE/otopcua-host-1"]);
publish.Topics.ShouldContain("spBv1.0/Plant1/NCMD/EdgeA");
}
@@ -968,6 +999,7 @@ public sealed class SparkplugIngestorTests
FeedBirths(ing);
await ing.OnReconnectedAsync(TestContext.Current.CancellationToken);
+ await ing.DrainPendingRebirthsAsync();
publish.Topics.Clear();
var fired = false;
@@ -1038,6 +1070,378 @@ public sealed class SparkplugIngestorTests
observed[1].MetricNames.ShouldContain("Temperature");
}
+ // ---------------------------------------------------------------------------------
+ // Session-state reset — the seams where the view underneath a surviving cache changes
+ // ---------------------------------------------------------------------------------
+
+ ///
+ /// C1. The ingestor OUTLIVES a session rebuild: SessionIdentity is a strict
+ /// superset of IngestIdentity, so changing only the broker Host / Port / ClientId /
+ /// credentials / TLS / a timeout re-establishes on this very instance. If the birth cache
+ /// survived that, an edge node which restarted during the changeover and rebound alias 5 from
+ /// Temperature to Pressure would have its first DDATA on the NEW session resolved against the
+ /// OLD table — Pressure's value under the Temperature RawPath, at Good quality.
+ ///
+ [Fact]
+ public async Task EstablishAsync_DropsBirthsFromThePreviousSession()
+ {
+ var subscribe = new FakeSubscribeTransport();
+ var ing = await NewSubscribedIngestorAsync(new RecordingPublishTransport(), subscribe);
+ FeedBirths(ing);
+ ing.Births.Find(new SparkplugScope(Group, Node, Device)).ShouldNotBeNull();
+
+ await ing.EstablishAsync(TestContext.Current.CancellationToken);
+
+ ing.Births.Find(new SparkplugScope(Group, Node, Device)).ShouldBeNull();
+ }
+
+ ///
+ /// The consequence, asserted on the wire rather than on the cache: after a session rebuild the
+ /// old alias must not resolve, so the first DDATA publishes nothing and asks for a rebirth.
+ ///
+ [Fact]
+ public async Task AfterEstablish_OldAliasNoLongerResolves()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish, new FakeSubscribeTransport());
+ FeedBirths(ing);
+
+ await ing.EstablishAsync(TestContext.Current.CancellationToken);
+
+ // Drain BEFORE clearing: Establish's own late-join rebirth is dispatched off-thread, so a bare
+ // Clear() races it and the assertions below could be satisfied by that NCMD rather than by the
+ // one the reset is supposed to cause.
+ await ing.DrainPendingRebirthsAsync();
+ publish.Topics.Clear();
+
+ var fired = false;
+ ing.OnDataChange += (_, _) => fired = true;
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: 2, (TempAlias, TahuDataType.Float, 1f)));
+ await ing.DrainPendingRebirthsAsync();
+
+ fired.ShouldBeFalse();
+ publish.Topics.ShouldContain("spBv1.0/Plant1/NCMD/EdgeA");
+ }
+
+ ///
+ /// The sequence half of the same reset. A surviving baseline lets a fresh stream look contiguous
+ /// by coincidence — and, worse, leaves IsBirthSynchronized true, so the late-join branch
+ /// never fires and nothing asks for the metadata the driver no longer has.
+ ///
+ [Fact]
+ public async Task EstablishAsync_DropsTheSequenceBaseline()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish, new FakeSubscribeTransport());
+ FeedBirths(ing);
+
+ await ing.EstablishAsync(TestContext.Current.CancellationToken);
+
+ // Drain BEFORE clearing: Establish's own late-join rebirth is dispatched off-thread, so a bare
+ // Clear() races it and the assertions below could be satisfied by that NCMD rather than by the
+ // one the reset is supposed to cause.
+ await ing.DrainPendingRebirthsAsync();
+ publish.Topics.Clear();
+
+ // Contiguous against the OLD baseline (the DBIRTH was seq 1) — must still be treated as a late
+ // join, because nothing about the new session's stream has been observed.
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: 2, (TempAlias, TahuDataType.Float, 1f)));
+ await ing.DrainPendingRebirthsAsync();
+
+ publish.Topics.ShouldContain("spBv1.0/Plant1/NCMD/EdgeA");
+ }
+
+ ///
+ /// AttachTo is the earliest seam at which the session changes, and a persistent-session
+ /// broker can push queued messages between CONNACK and our own SUBSCRIBE — so the reset happens
+ /// there too, not only at EstablishAsync.
+ ///
+ [Fact]
+ public async Task AttachTo_DropsBirthsFromThePreviousSession()
+ {
+ var ing = await NewSubscribedIngestorAsync();
+ FeedBirths(ing);
+
+ await using var connection = new MqttConnection(
+ new MqttDriverOptions { Host = "127.0.0.1", Port = 1, UseTls = false },
+ "mqtt-spb-1");
+ ing.AttachTo(connection);
+
+ ing.Births.Count.ShouldBe(0);
+ }
+
+ ///
+ /// C2. While an edge node is unauthored its traffic is dropped at Dispatch, so its
+ /// cached birth FREEZES rather than refreshing. Deploy A authors it, deploy B drops it, deploy C
+ /// re-adds it a week later having missed every rebirth in between — a surviving cache answers
+ /// with the week-old alias table and the next DDATA mis-routes at Good quality.
+ ///
+ [Fact]
+ public async Task Register_DroppingAnEdgeNode_EvictsItsBirths()
+ {
+ var ing = await NewSubscribedIngestorAsync();
+ FeedBirths(ing);
+ ing.Births.Count.ShouldBeGreaterThan(0);
+
+ ing.Register([]); // deploy B: the node is no longer authored
+ ing.Register(DefaultTags()); // deploy C: it comes back
+
+ ing.Births.Count.ShouldBe(0);
+ }
+
+ /// The C2 consequence on the wire: the re-added node's old alias must not resolve.
+ [Fact]
+ public async Task AfterEdgeNodeLeftAndReturned_OldAliasNoLongerResolves()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish);
+ FeedBirths(ing);
+
+ ing.Register([]);
+ ing.Register(DefaultTags());
+ await ing.SubscribeAsync(
+ [.. ing.AuthoredRawPaths],
+ TimeSpan.Zero,
+ TestContext.Current.CancellationToken);
+ publish.Topics.Clear();
+
+ var fired = false;
+ ing.OnDataChange += (_, _) => fired = true;
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: 2, (TempAlias, TahuDataType.Float, 1f)));
+ await ing.DrainPendingRebirthsAsync();
+
+ fired.ShouldBeFalse();
+ publish.Topics.ShouldContain("spBv1.0/Plant1/NCMD/EdgeA");
+ }
+
+ ///
+ /// The narrower half of the C2 rule, asserted so a future "tidy-up" cannot widen it. A device
+ /// with no authored tags whose EDGE NODE is still authored keeps receiving traffic (the Dispatch
+ /// filter is node-level), so its birth is live, not frozen — evicting it would throw away
+ /// correct state and cost a rebirth round-trip the moment a later deploy authored a tag on it.
+ ///
+ [Fact]
+ public async Task Register_KeepsBirthsForStillAuthoredEdgeNodes()
+ {
+ var ing = await NewSubscribedIngestorAsync();
+ FeedBirths(ing);
+
+ // Drops the Filler1 tags but keeps the node-level Uptime tag, so EdgeA stays authored.
+ ing.Register([Tag(NodeUptimePath, Blob(Group, Node, null, "Uptime", "Int64"))]);
+
+ ing.Births.Find(new SparkplugScope(Group, Node, Device)).ShouldNotBeNull();
+ ing.Births.Find(new SparkplugScope(Group, Node, null)).ShouldNotBeNull();
+ }
+
+ // ---------------------------------------------------------------------------------
+ // Rebirth-loop bounds
+ // ---------------------------------------------------------------------------------
+
+ ///
+ /// An NBIRTH whose seq cannot be read leaves the tracker un-synchronised forever — asking
+ /// for another birth cannot fix it, because the node answers with the same unusable payload. The
+ /// driver must record that a birth ARRIVED and stop asking.
+ ///
+ [Fact]
+ public async Task NbirthWithUnreadableSeq_DoesNotLoopRequestingRebirths()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish);
+
+ // seq is present but out of the 0-255 Sparkplug range, so AcceptNodeBirth refuses it.
+ var payload = new Payload { Seq = 9999UL, Timestamp = 1UL };
+ payload.Metrics.Add(BirthMetric("Uptime", 7UL, TahuDataType.Int64, 0L));
+ ing.Dispatch(Topic(SparkplugMessageType.NBIRTH), Decode(payload));
+ publish.Topics.Clear();
+
+ for (var i = 0; i < 5; i++)
+ {
+ ing.Dispatch(Topic(SparkplugMessageType.DBIRTH, Device), Dbirth(seq: (ulong)i));
+ }
+
+ await ing.DrainPendingRebirthsAsync();
+
+ publish.Topics.ShouldBeEmpty();
+ }
+
+ ///
+ /// The residual loop the birth-seen flag cannot see: an NBIRTH that never reaches the state
+ /// machine — because it exceeded the payload ceiling, or because the node simply never sends
+ /// one. Nothing about asking again changes that.
+ ///
+ ///
+ /// Each of these messages raises two uncapped conditions in the pre-fix code — late join
+ /// AND data-before-birth — which is what made the first version of this test observe 23 NCMDs
+ /// against a cap of 3. The debounce caps the RATE, not the LIFETIME; only the per-node
+ /// consecutive cap ends the loop, and it has to cover every missing-metadata reason to do it.
+ ///
+ [Fact]
+ public async Task NbirthNeverArriving_CapsConsecutiveMetadataRebirths()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish);
+
+ for (var i = 0; i < 20; i++)
+ {
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: (ulong)i, (TempAlias, TahuDataType.Float, 1f)));
+ }
+
+ await ing.DrainPendingRebirthsAsync();
+
+ // Bounded, and bounded LOW — the cap, not "eventually stops".
+ publish.Topics.Count.ShouldBeLessThanOrEqualTo(3);
+ publish.Topics.ShouldNotBeEmpty();
+ }
+
+ /// A birth that DOES arrive re-arms the cap — the count is consecutive, not cumulative.
+ [Fact]
+ public async Task ArrivingNbirth_ResetsTheMetadataRebirthCap()
+ {
+ var publish = new RecordingPublishTransport();
+ var ing = await NewSubscribedIngestorAsync(publish);
+
+ for (var i = 0; i < 20; i++)
+ {
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: (ulong)i, (TempAlias, TahuDataType.Float, 1f)));
+ }
+
+ await ing.OnReconnectedAsync(TestContext.Current.CancellationToken);
+ await ing.DrainPendingRebirthsAsync();
+ publish.Topics.Clear();
+
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Ddata(seq: 0, (TempAlias, TahuDataType.Float, 1f)));
+ await ing.DrainPendingRebirthsAsync();
+
+ publish.Topics.ShouldNotBeEmpty();
+ }
+
+ // ---------------------------------------------------------------------------------
+ // Value-level gates
+ // ---------------------------------------------------------------------------------
+
+ ///
+ /// Array metrics slip past the unsupported-type gate, because ToDriverDataType maps every
+ /// *Array variant to its ELEMENT type and is therefore never null for one. The value
+ /// rides as a packed byte[], so a String-typed tag would publish base64 at GOOD quality.
+ ///
+ [Fact]
+ public async Task ArrayMetric_IsSkipped_NotPublishedAsBase64()
+ {
+ const string ArrTag = "Plant/Mqtt/spb/TempArray";
+ var ing = NewIngestor(tags: [Tag(ArrTag, Blob(Group, Node, Device, "Temperature", "String"))]);
+ await ing.SubscribeAsync([ArrTag], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ var seen = new List();
+ ing.OnDataChange += (_, e) => seen.Add(e.Snapshot);
+
+ ing.Dispatch(Topic(SparkplugMessageType.NBIRTH), Nbirth(seq: 0, bdSeq: 1));
+
+ var birth = new Payload { Seq = 1UL, Timestamp = 1UL };
+ birth.Metrics.Add(new Payload.Types.Metric
+ {
+ Name = "Temperature",
+ Alias = TempAlias,
+ Datatype = (uint)TahuDataType.FloatArray,
+ BytesValue = ByteString.CopyFrom([0x00, 0x00, 0x80, 0x3F]),
+ });
+ ing.Dispatch(Topic(SparkplugMessageType.DBIRTH, Device), Decode(birth));
+
+ seen.ShouldBeEmpty();
+ ing.Values.Read(ArrTag).StatusCode.ShouldNotBe(Good);
+ }
+
+ ///
+ /// Convert.ToSingle saturates rather than throwing, so a Double of 1e300 onto a Float32
+ /// tag would publish +Infinity at GOOD quality — a number no gauge renders and no comparison
+ /// handles. A narrowing overflow is a refusal.
+ ///
+ [Fact]
+ public async Task NarrowingOverflowToFloat32_IsRefused_NotPublishedAsInfinity()
+ {
+ const string F32 = "Plant/Mqtt/spb/Narrowed";
+ var ing = NewIngestor(tags: [Tag(F32, Blob(Group, Node, Device, "Big", "Float32"))]);
+ await ing.SubscribeAsync([F32], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ DataValueSnapshot? snapshot = null;
+ ing.OnDataChange += (_, e) => snapshot = e.Snapshot;
+
+ ing.Dispatch(Topic(SparkplugMessageType.NBIRTH), Nbirth(seq: 0, bdSeq: 1));
+ ing.Dispatch(
+ Topic(SparkplugMessageType.DBIRTH, Device),
+ DbirthCustom(seq: 1, ("Big", 21UL, TahuDataType.Double, 1e300d)));
+
+ snapshot!.StatusCode.ShouldBe(BadTypeMismatch);
+ snapshot.Value.ShouldBeNull();
+ }
+
+ ///
+ /// A publisher's OWN infinity is its value, not our overflow — it passes through, so the guard
+ /// above cannot be satisfied by blanket-refusing non-finite numbers.
+ ///
+ [Fact]
+ public async Task PublishersOwnInfinity_IsPassedThrough()
+ {
+ const string F32 = "Plant/Mqtt/spb/Inf";
+ var ing = NewIngestor(tags: [Tag(F32, Blob(Group, Node, Device, "Big", "Float32"))]);
+ await ing.SubscribeAsync([F32], TimeSpan.Zero, TestContext.Current.CancellationToken);
+
+ DataValueSnapshot? snapshot = null;
+ ing.OnDataChange += (_, e) => snapshot = e.Snapshot;
+
+ ing.Dispatch(Topic(SparkplugMessageType.NBIRTH), Nbirth(seq: 0, bdSeq: 1));
+ ing.Dispatch(
+ Topic(SparkplugMessageType.DBIRTH, Device),
+ DbirthCustom(seq: 1, ("Big", 21UL, TahuDataType.Float, float.PositiveInfinity)));
+
+ snapshot!.StatusCode.ShouldBe(Good);
+ snapshot.Value.ShouldBe(float.PositiveInfinity);
+ }
+
+ ///
+ /// The metric NAME is the binding key, so when a DATA metric carries both a name and an alias
+ /// that disagree, the name wins. Preferring the alias resolves to whichever metric currently
+ /// occupies that slot — the mis-route the binding rule exists to prevent.
+ ///
+ [Fact]
+ public async Task DataMetricCarryingNameAndDisagreeingAlias_BindsByName()
+ {
+ var ing = await NewSubscribedIngestorAsync();
+ FeedBirths(ing);
+
+ var seen = new Dictionary(StringComparer.Ordinal);
+ ing.OnDataChange += (_, e) => seen[e.FullReference] = e.Snapshot.Value;
+
+ // Alias 6 is Pressure; the name says Temperature. The name is the identity that survives a
+ // rebirth, so this must land on the Temperature RawPath.
+ var payload = new Payload { Seq = 2UL, Timestamp = 1UL };
+ payload.Metrics.Add(new Payload.Types.Metric
+ {
+ Name = "Temperature",
+ Alias = PressAlias,
+ FloatValue = 3.5f,
+ });
+ ing.Dispatch(Topic(SparkplugMessageType.DDATA, Device), Decode(payload));
+
+ seen.ShouldContainKey(TempPath);
+ seen[TempPath].ShouldBe(3.5f);
+ seen.ShouldNotContainKey(PressPath);
+ }
+
+ ///
+ /// The legacy pre-3.0 STATE form is parsed on receive — so it must also be SUBSCRIBED, or that
+ /// tolerance is reachable only from a test calling HandleMessage directly.
+ ///
+ [Fact]
+ public async Task EstablishAsync_SubscribesTheLegacyStateFilterToo()
+ {
+ var subscribe = new FakeSubscribeTransport();
+ var ing = NewIngestor(subscribe: subscribe, hostId: "otopcua-host-1");
+
+ await ing.EstablishAsync(TestContext.Current.CancellationToken);
+
+ subscribe.Filters.Select(f => f.Topic).ShouldContain("STATE/otopcua-host-1");
+ }
+
// =================================================================================
// Fixtures
// =================================================================================
@@ -1132,6 +1536,20 @@ public sealed class SparkplugIngestorTests
private static SparkplugPayload Decode(Payload payload) => SparkplugCodec.Decode(payload.ToByteArray());
+ /// The standard NBIRTH as WIRE BYTES — for the payload-ceiling tests, which need a real one.
+ private static byte[] NbirthBytes(ulong seq, ulong bdSeq)
+ {
+ var payload = new Payload { Seq = seq, Timestamp = 1721822400000UL };
+ payload.Metrics.Add(new Payload.Types.Metric
+ {
+ Name = "bdSeq",
+ Datatype = (uint)TahuDataType.Uint64,
+ LongValue = bdSeq,
+ });
+ payload.Metrics.Add(BirthMetric("Uptime", 7UL, TahuDataType.Int64, 0L));
+ return payload.ToByteArray();
+ }
+
/// An NBIRTH: the bdSeq session token plus the node's own metric catalog.
private static SparkplugPayload Nbirth(ulong seq, ulong bdSeq)
{