908bd7ba4ae3b5973640f24be0be6e87753ecb9a
4 Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
6d7a458c4d |
fix(mqtt): reset Sparkplug ingest state at every session/authoring seam
Task 21 review — two Criticals sharing one root cause, plus four follow-ups.
C1/C2 root cause: `OnReconnectedAsync` opened with `Births.Clear(); _trackers.Clear();`
— "not housekeeping, the correctness step" — and that reset was missing from the
other two 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 Host/Port/ClientId/credentials/
TLS/a timeout gives `!SameSession && SameIngest` — `ReinitializeAsync` tears the
session down and re-establishes on the same instance, and `TeardownAsync` touches
no ingest state (the resilience layer re-running `InitializeAsync` is the second
path). A node that rebound alias 5 during the changeover would publish Pressure's
value under the Temperature RawPath at Good. Hoisted into `ResetSessionState`,
called from `AttachTo` (earliest seam — a persistent session can push between
CONNACK and our SUBSCRIBE), `EstablishAsync` and `OnReconnectedAsync`.
- C2: `Register` pruned trackers but not births. While a node is unauthored its
traffic is dropped at `Dispatch`, so its cached birth FREEZES rather than
refreshing; drop-then-re-add across a week mis-routes at Good, and `Births` grew
monotonically. Now evicted on edge-node departure.
Deliberately NARROWER than "absent from ByScope": a device with no authored tags
under a still-authored node keeps receiving traffic, so its birth is live, not
frozen — evicting it would discard correct state. Pinned both ways.
- I1: a birth whose seq is unreadable produced a permanent debounce-paced rebirth
loop. "A birth arrived" is now tracked separately from "the birth established a
baseline". Writing the test showed the cap belonged wider than the review scoped
it: data-before-birth and unknown-alias loop identically (20 messages → 23 NCMDs
against a cap of 3), so all three missing-metadata reasons share one per-node
consecutive cap, re-armed by any applied birth. A seq gap stays uncapped — it
self-limits.
- I2: the oversize `WarnOnce` key was the raw topic, off a group-wide `#`
subscription. `HandleMessage` now parses and filters to authored nodes BEFORE the
size check (also skipping the protobuf parse for discarded traffic), and `_warned`
carries a hard ceiling so a future bad derivation degrades to silence.
- I3: array metrics slipped the unsupported-type gate — `ToDriverDataType` returns
the ELEMENT type, so a String-typed tag published a packed `byte[]` as base64 at
Good. Gated on `IsSparkplugArray`.
- I4: driver-level Sparkplug coverage for the composition-root lines T22 left
untested — the `OnDataChange` re-raise, `ReadAsync` off the shared cache,
subscribe/unsubscribe dispatch, and both directions of the mode gate.
Minors: M1 the oversize test fed unparseable bytes and passed with the size check
deleted (now a real NBIRTH + an at-the-ceiling control); M2 answered in-place (an
absent seq is refused WITHOUT becoming the baseline, so the NCMD assertion is the
complete check); M3 `HostStateObserved` raise wrapped; M5 the legacy `STATE/{host}`
form is now subscribed, so the parser's tolerance is reachable in production; M8
`ResolveBinding` prefers the NAME over a disagreeing alias; M9 a narrowing float
conversion refuses instead of publishing ±Infinity at Good (a publisher's own
infinity still passes through). M4 was already resolved by T22.
Also fixed a race in the new tests: `publish.Topics.Clear()` after a call that
dispatches rebirths off-thread must drain first — it made the C1 sequence-baseline
pin pass vacuously in the first falsifiability run.
Falsifiability: dropping the reset from AttachTo+EstablishAsync reddens 4 (C1);
dropping the Register eviction reddens 2 (C2); widening the eviction predicate to
the reviewer's literal phrasing reddens the keep-live-births pin (C2b). 545/545 MQTT
unit tests green (was 524); whole-solution `--no-incremental` build clean under TWAE.
Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
|
||
|
|
6917d8805d |
feat(mqtt): Sparkplug UntilStable discovery + rediscover-on-DBIRTH
RediscoverPolicy is now mode-dependent: UntilStable in Sparkplug B (a tag's dataType is optional -- the birth certificate declares it, so the discovered tree fills in asynchronously), Once in Plain (nothing about the authored set arrives on the wire). DiscoverAsync resolves each Sparkplug tag's datatype per pass from the live BirthCache, with the ingest path's own precedence: authored dataType wins, else the birth's, else the record default. An unsupported Sparkplug type (DataSet/Template/PropertySet/Unknown) falls back rather than blanking the tag. The authored tag SET is never conditional on a birth -- an authored tag is part of the declared configuration, not something the plant grants by publishing. SparkplugIngestor.BirthObserved is wired to RaiseRediscoveryNeeded behind a two-part change gate, which is the anti-storm mechanism rather than an optimisation: fire only for a scope an authored tag binds, and only when the birth's metric-name SET (ordered-distinct, ordinal) differs from the last one seen for that scope. Edge nodes re-birth freely -- on their own timer, on every reconnect via the late-join rebirth fan-out, and once per rebirth NCMD the gap policy sends -- and firing on each of those would make a healthy plant a permanent address-space-rebuild loop. No extra debounce: what survives the gate is a real edit to the served tree, and delaying it would need its own trailing-edge flush to avoid dropping the last change. The signature map is NOT cleared on death/reconnect/ingest rebuild (the ingestor drops its whole birth cache on reconnect, but "the driver forgot" is not "the address space changed"); it is pruned when a redeploy stops authoring a scope. ScopeHint is the "Mqtt" discovery folder, deliberately not the Sparkplug scope path the plan named: the discovered tree is flat, so an edge-node/device path would name a subtree that does not exist and a consumer scoping a rebuild on it would rebuild nothing. The scope is carried in Reason instead. Falsifiability: always-fire reddens the identical-rebirth + reordered-rebirth pins; always-authored-datatype reddens the birth-fill pin; dropping the authored-scope filter reddens the unauthored-device pin. 524/524 MQTT unit tests green (was 513). Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW |
||
|
|
a13ae926a8 |
fix(mqtt): gate DisposeAsync against in-flight lifecycle ops (C1) + rebuild-branch tests (I1)
C1 (Critical) — DisposeAsync went straight to TeardownAsync with no lifecycle-gate
wait, reopening the orphaned-connection class. Interleaving: ReinitializeAsync's
rebuild branch holds the gate mid-InitializeCoreAsync, past its own teardown but
before `_connection = connection`; an ungated dispose sees a null connection, does
nothing, and sets `_disposed`; the rebuild then publishes a live connection plus a
host-probe loop that every later dispose short-circuits past. Not proven reachable
today (DriverInstanceActor.PostStop calls the gated ShutdownAsync), but MqttDriver
publicly implements IAsyncDisposable, which invites `await using`.
- DisposeAsync now takes the gate with the same bounded wait + fallback teardown
ShutdownAsync uses, and sets `_disposed` BEFORE the wait.
- InitializeCoreAsync re-checks `_disposed` after the connect and disposes the
connection it just built rather than publishing it — this closes the residual
window on the bounded-timeout fallback path.
- The doc comment no longer claims parity with ShutdownAsync it did not have, and
stops conflating "don't Dispose() the semaphore object" with "don't WaitAsync".
I1 (Important) — the SameSession == false rebuild branch had no test driving it.
Adds three: an endpoint-changing delta rebuilds and Faults on a refused connect;
an ingest-only delta (MaxPayloadBytes) ALSO rebuilds, pinning that IngestIdentity
is nested inside SessionIdentity; and DisposeAsync serializes behind a lifecycle
operation parked at a new internal BeforeConnectHookForTests seam (mirroring
MqttConnection's AfterConnectHookForTests, which exists for the same race class).
Minors: comments recording that IngestIdentity names no Sparkplug field and must
grow one in P2 (tasks 21/22), and why _options/_subscriptions/_authoredRawPaths
get looser memory discipline than _health/_hostState.
Falsifiability: reverting DisposeAsync to the ungated form reddens exactly the new
serialization test ("Shouldly.ShouldAssertException : raced"); forcing SameSession
to always-true reddens exactly the two new rebuild tests. 240/240 MQTT tests pass;
forced rebuild of the driver project is 0 warnings under TreatWarningsAsErrors.
Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
|
||
|
|
420692b6e5 |
feat(mqtt): MqttDriver shell — lifecycle + authored-only discovery (Once)
Composes MqttConnection + MqttSubscriptionManager + LastValueCache into the IDriver / ITagDiscovery / ISubscribable / IReadable / IHostConnectivityProbe / IRediscoverable capability set. - Discovery replays ONLY the authored raw tags (RediscoverPolicy = Once, SupportsOnlineDiscovery = false) — a chatty broker can never auto-provision. - Composition order is register -> AttachTo -> ConnectAsync; the manager's Reconnected handler is passed through unwrapped so its throw-on-total-failure still tears a deaf session down. - ReinitializeAsync applies a tag-only delta in place and never faults on a bad one; a session-changing delta rebuilds. - ReadAsync serves the last-value cache and degrades per reference (BadWaitingForInitialData), never throwing for the batch. - FlushOptionalCachesAsync is a no-op: the last-value cache IS IReadable. - Adds MqttDriverOptions.RawTags (the authored-tag delivery mechanism every other driver's options DTO already has) and promotes MaxPayloadBytes from a manager ctor knob to an operator-facing key. - Converges on MqttDriverProbe.JsonOpts rather than a second, divergent JsonSerializerOptions. Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW |