From a13ae926a8bde324874861472171d71c5ce86267 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Fri, 24 Jul 2026 17:31:04 -0400 Subject: [PATCH] fix(mqtt): gate DisposeAsync against in-flight lifecycle ops (C1) + rebuild-branch tests (I1) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../MqttDriver.cs | 102 ++++++++++++++++- .../MqttDriverDiscoveryTests.cs | 107 ++++++++++++++++++ 2 files changed, 203 insertions(+), 6 deletions(-) diff --git a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttDriver.cs b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttDriver.cs index b7db7439..da500db5 100644 --- a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttDriver.cs +++ b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttDriver.cs @@ -72,6 +72,14 @@ public sealed class MqttDriver /// Guards / / against each other. private readonly SemaphoreSlim _lifecycleGate = new(1, 1); + // _options / _subscriptions / _authoredRawPaths are written only under _lifecycleGate (or in the + // ctor) and read freely elsewhere, WITHOUT the Volatile discipline _health / _hostState get. + // That is deliberate, not an oversight: all three are reference/interface fields, so a reader + // sees either the whole old object or the whole new one — never a torn one. The looser fields + // are the ones a reader may legitimately observe one publish stale (a read served from the + // previous authored table is a correct read of a value that was correct a microsecond ago); + // _health and _hostState feed ServiceLevel and the status dashboard, where a stale per-core copy + // is an operator-visible lie about whether the driver is up. private MqttDriverOptions _options; private MqttSubscriptionManager _subscriptions; @@ -133,6 +141,16 @@ public sealed class MqttDriver /// internal MqttSubscriptionManager Subscriptions => _subscriptions; + /// + /// Test seam: awaited inside while the lifecycle gate is + /// held, immediately before the broker connect. Lets a test park a lifecycle operation + /// mid-flight and prove serializes behind it rather than racing + /// past a half-built session. Mirrors 's own + /// AfterConnectHookForTests, which exists for the same class of race. Always + /// in production. + /// + internal Func? BeforeConnectHookForTests { get; set; } + // ---- IDriver: lifecycle ---- /// @@ -437,9 +455,37 @@ public sealed class MqttDriver public void Dispose() => DisposeAsync().AsTask().GetAwaiter().GetResult(); /// - /// Performs the same teardown as so a caller that uses - /// await using without an explicit shutdown does not leak a live broker session. + /// Performs the same teardown as , under the same lifecycle + /// gate, so a caller that uses await using without an explicit shutdown does not + /// leak a live broker session. /// + /// + /// + /// The gate wait is the whole point, not ceremony. An ungated dispose that ran + /// while 's rebuild branch held the gate would observe + /// _connection == null — the rebuild has torn the old session down but not yet + /// assigned the new one — do nothing, and set _disposed. The rebuild would then + /// complete, assign a live, connected connection and start a host-probe loop that + /// no later dispose can ever reach, because the disposed guard short-circuits them all. + /// That is the orphaned-connection class 's own remarks + /// describe having already fixed once. Waiting for the gate makes this dispose tear down + /// whatever the in-flight operation ends up assigning. + /// + /// + /// _disposed is set before the wait, so an in-flight + /// sees it at its post-connect re-check and refuses to + /// publish the connection at all — which is what closes the residual window on the + /// bounded-timeout fallback path below. + /// + /// + /// Separately (and for an unrelated reason): the gate object itself is deliberately never + /// Dispose()d — same call as 's semaphores. A caller + /// that disposes before its last would otherwise get an + /// naming , which says + /// nothing useful; the gate holds no timer and no surviving registration, so leaving it + /// undisposed costs nothing. + /// + /// /// A task that represents the asynchronous dispose. public async ValueTask DisposeAsync() { @@ -448,10 +494,27 @@ public sealed class MqttDriver return; } - // The lifecycle gate is deliberately NOT disposed — same call as MqttConnection's semaphores. - // A caller that disposes before its last ShutdownAsync would otherwise get an - // ObjectDisposedException naming SemaphoreSlim, which says nothing useful; the gate holds no - // timer and no surviving registration, so leaving it undisposed costs nothing. + if (await _lifecycleGate.WaitAsync(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)) + .ConfigureAwait(false)) + { + try + { + await TeardownAsync().ConfigureAwait(false); + } + finally + { + _lifecycleGate.Release(); + } + + return; + } + + // Bounded even here: a wedged initialize must not turn dispose into a hang. The post-connect + // disposed re-check in InitializeCoreAsync is what stops that operation from publishing a + // connection behind this teardown's back. + _logger?.LogWarning( + "MQTT driver '{DriverId}': lifecycle gate still held at dispose; tearing the session down anyway.", + _driverInstanceId); await TeardownAsync().ConfigureAwait(false); } @@ -480,7 +543,24 @@ public sealed class MqttDriver // throw-on-total-failure is what tears a deaf session down and retries it. _subscriptions.AttachTo(connection); + if (BeforeConnectHookForTests is { } hook) + { + await hook(cancellationToken).ConfigureAwait(false); + } + await connection.ConnectAsync(cancellationToken).ConfigureAwait(false); + + // Re-check disposal AFTER the connect, before publishing the connection. DisposeAsync + // normally waits for this gate, but its wait is bounded — on the wedged-gate fallback + // path it can run concurrently with this method, and assigning a live connection behind + // a completed teardown is exactly the orphaned-session bug. Losing the race means + // disposing what we just built, never leaking it. + if (Volatile.Read(ref _disposed) != 0) + { + await connection.DisposeAsync().ConfigureAwait(false); + throw new ObjectDisposedException(nameof(MqttDriver)); + } + _connection = connection; WriteHealth(new DriverHealth(DriverState.Healthy, connection.LastMessageUtc, null)); @@ -597,6 +677,16 @@ public sealed class MqttDriver => SessionIdentity(a).Equals(SessionIdentity(b)); /// Whether two option sets produce an identical . + /// + /// ⚠️ P2 (Sparkplug, tasks 21/22) must extend . It names + /// only the settings the P1 manager captures at construction — Mode, + /// Plain.DefaultQos, MaxPayloadBytes — and deliberately no + /// field, because nothing reads one yet. The moment a + /// Sparkplug handler captures its own options (group id, host id, rebirth policy, birth + /// window) at construction, a delta changing one of them would be judged "same ingest", + /// applied in place, and silently never reach the component that needed rebuilding — the + /// driver would keep subscribing to the OLD group id and report perfect health. + /// private static bool SameIngest(MqttDriverOptions a, MqttDriverOptions b) => IngestIdentity(a).Equals(IngestIdentity(b)); 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 70a92223..4fbd1c96 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 @@ -28,6 +28,33 @@ public sealed class MqttDriverDiscoveryTests private static MqttDriver PlainDriver(params RawTagEntry[] tags) => new(new MqttDriverOptions { Mode = MqttMode.Plain, RawTags = tags }, "d", null); + /// + /// Options aimed at a definitely-closed port so a connect is refused immediately rather + /// than hanging — the tests that must prove a rebuild actually dials need the dial to fail + /// fast, not to burn a connect deadline. + /// + private static MqttDriverOptions ClosedPortOptions(params RawTagEntry[] tags) => new() + { + Mode = MqttMode.Plain, + Host = "127.0.0.1", + Port = 1, + UseTls = false, + ConnectTimeoutSeconds = 2, + RawTags = tags, + }; + + private static MqttDriver ClosedPortDriver(params RawTagEntry[] tags) + => new(ClosedPortOptions(tags), "d", null); + + /// + /// Serializes a full options record as the reinitialize blob. Round-tripping the whole record + /// (rather than hand-writing a partial JSON object) guarantees that only the field the caller + /// changed differs — a hand-written delta silently omitting a session field would take the + /// rebuild branch for the wrong reason and the test would pass vacuously. + /// + private static string DeltaJson(MqttDriverOptions options) + => System.Text.Json.JsonSerializer.Serialize(options, MqttJson.Options); + /// /// Plain mode replays ONLY the authored tag set: one variable per authored raw tag, a /// single discovery pass (), and no online @@ -192,6 +219,86 @@ public sealed class MqttDriverDiscoveryTests .ShouldBe(["Plant/Mqtt/dev1/Flow", "Plant/Mqtt/dev1/Pressure"]); } + /// + /// A delta that changes the broker endpoint takes the rebuild branch: it must actually + /// dial the new endpoint (proven here by the refused connect), and a rebuild whose connect + /// fails must land rather than letting the failure escape + /// with the health surface still claiming Healthy. + /// + [Fact] + public async Task ReinitializeAsync_SessionChangingDelta_Rebuilds_AndFaultsOnUnreachableBroker() + { + var driver = ClosedPortDriver(Tag("Plant/Mqtt/dev1/Temp", "f/t")); + + // Only Host differs from the ctor options — everything else round-trips identically. + var delta = DeltaJson(ClosedPortOptions() with { Host = "127.0.0.2" }); + + await Should.ThrowAsync(() => driver.ReinitializeAsync(delta, CancellationToken.None)); + + driver.GetHealth().State.ShouldBe(DriverState.Faulted); + } + + /// + /// Pins that IngestIdentity is nested inside SessionIdentity: a delta touching + /// ONLY an ingest setting still rebuilds. Without the nesting it would be judged "same + /// session", applied in place, and the subscription manager that actually reads + /// MaxPayloadBytes would never be rebuilt — a config change that silently does nothing. + /// + [Fact] + public async Task ReinitializeAsync_IngestOnlyDelta_AlsoRebuilds() + { + var driver = ClosedPortDriver(Tag("Plant/Mqtt/dev1/Temp", "f/t")); + + var delta = DeltaJson(ClosedPortOptions() with { MaxPayloadBytes = 4096 }); + + // Reaching the (refused) connect at all is the proof the rebuild branch was taken; the + // tag-only test above is the control that shows an in-place delta never dials. + await Should.ThrowAsync(() => driver.ReinitializeAsync(delta, CancellationToken.None)); + } + + /// + /// The falsifiability pin for the orphaned-connection class: + /// must serialize behind an in-flight lifecycle operation, not race past it. + /// + /// + /// An ungated dispose running while a rebuild holds the gate observes _connection == null + /// (the rebuild has torn the old session down but not yet assigned the new one), does nothing, + /// and sets the disposed flag — after which the rebuild assigns a live connection plus a + /// host-probe loop that no later dispose can reach. Here the driver is parked inside + /// InitializeCoreAsync at the pre-connect hook; the assertion is that dispose does not + /// complete while it is parked, and does complete once it is released. + /// + [Fact] + public async Task DisposeAsync_SerializesBehindAnInFlightLifecycleOperation() + { + var driver = ClosedPortDriver(Tag("Plant/Mqtt/dev1/Temp", "f/t")); + + var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var release = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + driver.BeforeConnectHookForTests = async _ => + { + entered.TrySetResult(); + await release.Task; + }; + + var initialize = Task.Run(() => driver.InitializeAsync("", CancellationToken.None)); + await entered.Task.WaitAsync(TimeSpan.FromSeconds(5)); + + var dispose = driver.DisposeAsync().AsTask(); + + // The gate is held by the parked initialize. Dispose must still be waiting. (The gate wait is + // bounded by ConnectTimeoutSeconds = 2 s, so this 300 ms window is comfortably inside it — + // an ungated dispose returns essentially instantly.) + var raced = await Task.WhenAny(dispose, Task.Delay(TimeSpan.FromMilliseconds(300))); + raced.ShouldNotBe(dispose, "DisposeAsync returned while a lifecycle operation held the gate"); + + release.SetResult(); + await Should.ThrowAsync(() => initialize); // connect refused on the closed port + + await dispose.WaitAsync(TimeSpan.FromSeconds(10)); + dispose.IsCompletedSuccessfully.ShouldBeTrue(); + } + /// Identity is fixed at construction and must match the persisted DriverInstance.DriverType. [Fact] public void Identity_IsMqtt()