From 768fd87774ab5b7866d1696eeb75172262482233 Mon Sep 17 00:00:00 2001 From: Joseph Doherty Date: Fri, 24 Jul 2026 14:58:37 -0400 Subject: [PATCH] feat(mqtt): hand-rolled reconnect loop with bounded backoff + resubscribe Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW --- .../MqttConnection.cs | 675 +++++++++++++++--- .../MqttConnectionTests.cs | 513 +++++++++++++ 2 files changed, 1101 insertions(+), 87 deletions(-) diff --git a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttConnection.cs b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttConnection.cs index 1a6e9072..d01d2444 100644 --- a/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttConnection.cs +++ b/src/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt/MqttConnection.cs @@ -6,17 +6,43 @@ using MQTTnet; namespace ZB.MOM.WW.OtOpcUa.Driver.Mqtt; +/// Lifecycle state of an . +public enum MqttConnectionState +{ + /// Constructed, or the first connect has not yet succeeded. No supervisor is running. + Disconnected = 0, + + /// An MQTT session is established. + Connected = 1, + + /// + /// The session dropped and the supervisor is retrying under backoff. A broker that is merely + /// down stays here indefinitely — being unreachable is never a fault. + /// + Reconnecting = 2, + + /// + /// Unrecoverable configuration: retrying cannot help, so the supervisor stops. Only + /// failures raised while assembling the client options (before any I/O) land here. + /// + Faulted = 3, + + /// has run. Terminal. + Disposed = 4, +} + /// /// Owns the driver's single live MQTTnet-5 client: assembles the broker connection options /// from (endpoint, protocol version, session, credentials, -/// TLS + CA pin) and connects under a bounded deadline. +/// TLS + CA pin), connects under a bounded deadline, and keeps the session alive with a +/// hand-rolled reconnect supervisor. /// /// /// /// Connection-free construction. The constructor stores configuration only — it -/// opens no socket and creates no client. The universal-browser CanBrowse pattern -/// constructs a throwaway instance purely to ask what driver type it is, so a -/// constructor that dialled a broker would stall the AdminUI. +/// opens no socket, creates no client and starts no supervisor. The universal-browser +/// CanBrowse pattern constructs a throwaway instance purely to ask what driver type +/// it is, so a constructor that dialled a broker would stall the AdminUI. /// /// /// Bounded connect. links the caller's token with a @@ -24,6 +50,26 @@ namespace ZB.MOM.WW.OtOpcUa.Driver.Mqtt; /// fed to MQTTnet's own . A broker that accepts the /// TCP connection and then never answers CONNACK — the frozen-peer shape that wedged the /// S7 read path (arch-review R2-01) — fails at the deadline instead of hanging forever. +/// Every wait in this type is bounded the same way, teardown included. +/// +/// +/// Hand-rolled reconnect. MQTTnet v5 dropped v4's ManagedMqttClient, so there +/// is no library-managed reconnect: this type supplies it. The supervisor starts on the +/// first successful connect (a failed first attempt is the caller's to retry — the +/// driver-host resilience layer re-runs InitializeAsync — and must not leave a +/// background task hammering an endpoint the caller has given up on). MQTTnet's +/// DisconnectedAsync callback only signals a semaphore and returns, so the library's +/// own dispatcher thread is never blocked by a backoff sleep. +/// +/// +/// is load-bearing. MQTT subscriptions do not survive a +/// clean session, and a reconnect that completes without re-subscribing produces a +/// healthy-looking connection that receives nothing, forever, with no error, no exception +/// and no bad status. So the callback fires on every successful reconnect — including +/// persistent-session reconnects, where re-subscribing is merely redundant. Re-subscribing +/// an already-subscribed topic is harmless; missing one is silent death. If the callback +/// fails, the freshly-established session is torn down and retried rather than left +/// connected-but-deaf. /// /// /// Credentials never leak. Nothing here logs or formats @@ -31,29 +77,49 @@ namespace ZB.MOM.WW.OtOpcUa.Driver.Mqtt; /// ToString(). /// /// -/// NOT thread-safe — single-caller by contract. and -/// take no lock and must never run concurrently. Calling them -/// concurrently leaks a live broker connection: can -/// observe _disposed == false, set it, exchange a still-null _client for -/// null, and return having disposed nothing — after which the in-flight -/// creates a client, assigns it and connects successfully. -/// That client holds an open socket that no later can ever -/// reach, because _disposed is already true. Serialising the lifecycle is -/// the job of the reconnect supervisor added in Task 4; it is deliberately not solved -/// here. +/// Lifecycle is serialised. Connect (caller-initiated or supervisor-initiated) and +/// dispose run under a single gate, and a connect that completes re-checks the disposed flag +/// while still holding that gate. This closes the connection leak the Task-3 review +/// traced: dispose could observe a still-null client, dispose nothing and return, after +/// which the in-flight connect assigned a client and connected successfully — leaving a live +/// socket no later dispose could ever reach. Now the losing side of that race disposes the +/// client it created and reports . /// -/// Reconnect/backoff (Task 4) and subscription (Task 6) are deliberately not implemented here. +/// Subscription (Task 6) and Sparkplug rebirth-on-reconnect (Task 21, which hangs off +/// ) are deliberately not implemented here. /// public sealed class MqttConnection : IAsyncDisposable { private readonly string _driverId; + + /// + /// Serialises connect (both callers of it) against dispose. Never disposed: a late MQTTnet + /// callback can still reach this instance after teardown, and a disposed + /// would throw inside the library's own dispatcher. + /// + private readonly SemaphoreSlim _lifecycleGate = new(1, 1); + + /// Cancelled by ; aborts the supervisor and any in-flight connect. + private readonly CancellationTokenSource _lifetimeCts = new(); + private readonly ILogger? _logger; private readonly MqttDriverOptions _options; - private IMqttClient? _client; - private bool _disposed; + /// + /// Wakes the supervisor. Unbounded max count: spurious releases (a failed connect attempt + /// also raises DisconnectedAsync) are drained harmlessly by the still-connected guard, + /// whereas a bounded semaphore would throw on the extra release. Never disposed, for the same + /// reason as . + /// + private readonly SemaphoreSlim _reconnectWake = new(0); - /// Stores configuration only — no network, no client instantiation. + private IMqttClient? _client; + private int _disposed; + private long _lastMessageTicksUtc; + private Task? _reconnectWorker; + private int _state = (int)MqttConnectionState.Disconnected; + + /// Stores configuration only — no network, no client instantiation, no supervisor. /// Broker connection settings. /// Driver instance id, used only for log/diagnostic correlation. /// Optional logger; never receives credentials. @@ -67,25 +133,520 @@ public sealed class MqttConnection : IAsyncDisposable _logger = logger; } + /// + /// Raised after every successful reconnect (never after the initial + /// , which the caller already sequences its own subscribe behind). + /// Subscribers must re-establish their MQTT subscriptions here — see the type remarks for + /// why skipping it is silent death. Every subscriber is invoked even if an earlier one + /// throws; a throwing subscriber causes the session to be torn down and retried. + /// + public event Func? Reconnected; + /// Whether the underlying client currently holds an established MQTT session. public bool IsConnected => _client?.IsConnected ?? false; + /// + /// UTC timestamp of the most recent inbound application message, or null if none has + /// arrived on this instance. Connection health only — nothing here routes or parses the + /// message; that is the subscription manager's job (Task 6). + /// + /// + /// The plan calls this member LastMessageAgeUtc, which conflates a timestamp with an + /// age; it is split here into the instant () and the elapsed span + /// (). + /// + public DateTime? LastMessageUtc + { + get + { + var ticks = Interlocked.Read(ref _lastMessageTicksUtc); + return ticks == 0 ? null : new DateTime(ticks, DateTimeKind.Utc); + } + } + + /// + /// How long ago the most recent inbound application message arrived, or null if none + /// has. A connection that is with a large age is + /// the "reconnected but never re-subscribed" shape. + /// + public TimeSpan? LastMessageAge => LastMessageUtc is { } at ? DateTime.UtcNow - at : null; + + /// Current lifecycle state. See . + public MqttConnectionState State => (MqttConnectionState)Volatile.Read(ref _state); + + /// + /// Test seam: awaited inside the gated connect immediately after the broker session is + /// established and before the disposed re-check, so a test can reproduce the exact + /// connect/dispose interleaving that used to leak a live connection. Always null in + /// production. + /// + internal Func? AfterConnectHookForTests { get; set; } + + private bool Disposed => Volatile.Read(ref _disposed) != 0; + /// public async ValueTask DisposeAsync() { - if (_disposed) + if (Interlocked.Exchange(ref _disposed, 1) != 0) { return; } - _disposed = true; + SetState(MqttConnectionState.Disposed); - var client = Interlocked.Exchange(ref _client, null); + // Stop the supervisor and abort any in-flight connect FIRST, so the gate below is not held + // for a full connect deadline by work that is already pointless. + try + { + await _lifetimeCts.CancelAsync().ConfigureAwait(false); + } + catch (ObjectDisposedException) + { + // Nothing to cancel. + } + + var worker = Volatile.Read(ref _reconnectWorker); + if (worker is not null) + { + try + { + await worker.WaitAsync(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)).ConfigureAwait(false); + } + catch (Exception ex) + { + _logger?.LogDebug(ex, "MQTT driver '{DriverId}': reconnect supervisor did not stop cleanly.", _driverId); + } + } + + // Bounded even here: a wedged connect must not turn dispose into a hang. If the gate cannot + // be taken we still claim the client — double-disposing an MqttClient is harmless, leaking a + // live one is not. + var gated = await _lifecycleGate + .WaitAsync(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)) + .ConfigureAwait(false); + if (!gated) + { + _logger?.LogWarning( + "MQTT driver '{DriverId}': lifecycle gate still held at dispose; tearing the client down anyway.", + _driverId); + } + + try + { + var client = Interlocked.Exchange(ref _client, null); + if (client is not null) + { + await CloseAsync(client).ConfigureAwait(false); + } + } + finally + { + if (gated) + { + _lifecycleGate.Release(); + } + } + } + + /// + /// The backoff a supervisor waits before reconnect attempt : + /// minSeconds × 2^attempt, saturated at . Pure — no + /// clock, no state, no jitter (deliberately: the schedule is asserted by exact equality, and + /// with one connection per driver instance there is no thundering herd to spread out). + /// + /// + /// Computed in rather than by shifting, because min << attempt + /// wraps — 1 << 32 is 1 — which would silently turn a long outage back + /// into a once-a-second hammering of a dead broker. A maxSeconds below + /// is misconfiguration; the floor wins, because the floor + /// is what stops the loop going hot. + /// + /// Zero-based attempt index; negatives are treated as the first attempt. + /// Delay before the first attempt, and the floor for every later one. + /// Cap on the exponential growth. + public static TimeSpan NextBackoff(int attempt, int minSeconds, int maxSeconds) + { + var min = Math.Max(1d, minSeconds); + var max = Math.Max(min, maxSeconds); + + if (attempt <= 0) + { + return TimeSpan.FromSeconds(min); + } + + // Math.Pow saturates to +Infinity instead of wrapping, so a huge attempt count clamps to max. + var seconds = min * Math.Pow(2d, attempt); + return TimeSpan.FromSeconds(seconds >= max || double.IsNaN(seconds) ? max : seconds); + } + + /// + /// Connects to the broker, failing at + /// rather than waiting indefinitely on an unresponsive peer. On success the reconnect + /// supervisor takes over keeping the session alive. + /// + /// + /// May be retried on the same instance after a failed attempt — the underlying + /// is created once and reused across attempts. Safe to call + /// concurrently with : the two are serialised, and a connect that + /// loses the race disposes whatever it built and throws + /// . + /// + /// Caller cancellation; linked with the connect deadline. + /// The connect deadline elapsed. + /// was cancelled. + /// + /// The connection was disposed, either before the attempt started or while it was in flight. + /// + public Task ConnectAsync(CancellationToken cancellationToken) + { + ObjectDisposedException.ThrowIf(Disposed, this); + + return ConnectCoreAsync(startSupervisor: true, cancellationToken); + } + + /// + /// Invokes every subscriber. Walks the invocation list explicitly: + /// awaiting a multicast directly yields only the last + /// subscriber's task and abandons the earlier ones, so a single throwing subscriber would + /// silently skip every subscriber behind it — and a skipped re-subscribe is a topic that goes + /// dark. Failures are collected and rethrown after every subscriber has had its turn. + /// + internal async Task FireReconnectedAsync() + { + var subscribers = Reconnected; + if (subscribers is null) + { + return; + } + + List? failures = null; + foreach (var subscriber in subscribers.GetInvocationList().Cast>()) + { + try + { + await subscriber().ConfigureAwait(false); + } + catch (Exception ex) + { + (failures ??= []).Add(ex); + } + } + + if (failures is not null) + { + throw new AggregateException( + $"MQTT driver '{_driverId}': {failures.Count} re-subscribe callback(s) failed after reconnect.", + failures); + } + } + + private async Task ConnectCoreAsync(bool startSupervisor, CancellationToken cancellationToken) + { + // The options are rebuilt per attempt so a rotated CA file is picked up on the next connect + // rather than being pinned for the life of the process. This is also the only step that can + // fail unrecoverably — it is pure config assembly, so retrying it can never help. + MqttClientOptions clientOptions; + try + { + clientOptions = BuildClientOptions(_options, clientIdSuffix: null, _logger); + } + catch (Exception ex) + { + SetState(MqttConnectionState.Faulted); + _logger?.LogError( + ex, + "MQTT driver '{DriverId}': broker configuration is unusable; no amount of retrying will help.", + _driverId); + throw; + } + + using var deadline = new CancellationTokenSource(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)); + using var linked = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken, + deadline.Token, + _lifetimeCts.Token); + + _logger?.LogDebug( + "MQTT driver '{DriverId}': connecting to {Host}:{Port} (tls={UseTls}, timeout={TimeoutSeconds}s).", + _driverId, + _options.Host, + _options.Port, + _options.UseTls, + _options.ConnectTimeoutSeconds); + + try + { + await _lifecycleGate.WaitAsync(linked.Token).ConfigureAwait(false); + } + catch (Exception ex) + { + throw Classify(ex, cancellationToken, deadline); + } + + try + { + // Dispose won the race outright: it has already flagged the instance and disposed + // whatever existed. Creating a client now would be creating one nobody can reach. + ObjectDisposedException.ThrowIf(Disposed, this); + + var client = _client; + if (client is null) + { + client = new MqttClientFactory().CreateMqttClient(); + AttachHandlers(client); + _client = client; + } + + await client.ConnectAsync(clientOptions, linked.Token).ConfigureAwait(false); + + if (AfterConnectHookForTests is { } hook) + { + await hook().ConfigureAwait(false); + } + + // The race the Task-3 review traced: dispose flagged the instance while this connect was + // in flight, then blocked on the very gate this call holds — so it disposed nothing and + // the session below would have been unreachable forever. Hand it back here instead. + if (Disposed) + { + Interlocked.CompareExchange(ref _client, null, client); + await CloseAsync(client).ConfigureAwait(false); + throw new ObjectDisposedException( + nameof(MqttConnection), + $"MQTT driver '{_driverId}': the connection was disposed while a connect to " + + $"{_options.Host}:{_options.Port} was in flight; the established session was torn down."); + } + + SetState(MqttConnectionState.Connected); + + if (startSupervisor) + { + _reconnectWorker ??= Task.Run(() => ReconnectSupervisorAsync(_lifetimeCts.Token)); + } + } + catch (ObjectDisposedException) + { + throw; + } + catch (Exception ex) + { + throw Classify(ex, cancellationToken, deadline); + } + finally + { + _lifecycleGate.Release(); + } + } + + /// + /// Maps a connect failure onto this type's documented contract. MQTTnet wraps a cancelled + /// connect in MqttConnectingFailedException rather than letting the + /// surface, so the legs are told apart by which + /// token fired, not by exception type. Precedence: teardown, then caller intent, then the + /// deadline. + /// + private Exception Classify(Exception ex, CancellationToken cancellationToken, CancellationTokenSource deadline) + { + if (Disposed) + { + return new ObjectDisposedException( + nameof(MqttConnection), + new InvalidOperationException( + $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} was aborted because the " + + "connection was disposed while the attempt was in flight.", + ex)); + } + + if (cancellationToken.IsCancellationRequested) + { + return new OperationCanceledException( + $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} was cancelled.", + ex, + cancellationToken); + } + + if (deadline.IsCancellationRequested) + { + return new TimeoutException( + $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} did not complete within " + + $"{_options.ConnectTimeoutSeconds}s.", + ex); + } + + return ex; + } + + private void AttachHandlers(IMqttClient client) + { + client.DisconnectedAsync += OnDisconnectedAsync; + client.ApplicationMessageReceivedAsync += OnApplicationMessageReceivedAsync; + } + + /// + /// Runs on MQTTnet's own dispatcher, so it does exactly two non-blocking things and returns. + /// Sleeping the backoff here would stall the client's internal pump. + /// + private Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs args) + { + if (Disposed || State is MqttConnectionState.Faulted or MqttConnectionState.Disposed) + { + return Task.CompletedTask; + } + + if (State == MqttConnectionState.Connected) + { + SetState(MqttConnectionState.Reconnecting); + } + + // Released unconditionally rather than only when the client had been connected: a missed + // wake is a permanently dark connection, whereas a spurious one costs a loop iteration. + _reconnectWake.Release(); + return Task.CompletedTask; + } + + private Task OnApplicationMessageReceivedAsync(MqttApplicationMessageReceivedEventArgs args) + { + Interlocked.Exchange(ref _lastMessageTicksUtc, DateTime.UtcNow.Ticks); + return Task.CompletedTask; + } + + /// + /// The reconnect loop MQTTnet v5 no longer provides. Started on the first successful connect + /// and stopped by cancelling the lifetime token; it owns every + /// reconnect attempt so the library's dispatcher never waits on one. + /// + private async Task ReconnectSupervisorAsync(CancellationToken cancellationToken) + { + try + { + while (!cancellationToken.IsCancellationRequested) + { + await _reconnectWake.WaitAsync(cancellationToken).ConfigureAwait(false); + + if (Disposed || State == MqttConnectionState.Faulted) + { + return; + } + + if (IsConnected) + { + continue; // Spurious wake (e.g. a failed attempt's own disconnect callback). + } + + SetState(MqttConnectionState.Reconnecting); + + for (var attempt = 0; !cancellationToken.IsCancellationRequested && !Disposed; attempt++) + { + await Task + .Delay( + NextBackoff( + attempt, + _options.ReconnectMinBackoffSeconds, + _options.ReconnectMaxBackoffSeconds), + cancellationToken) + .ConfigureAwait(false); + + try + { + await ConnectCoreAsync(startSupervisor: false, cancellationToken).ConfigureAwait(false); + } + catch (ObjectDisposedException) + { + return; + } + catch (Exception ex) + { + if (State == MqttConnectionState.Faulted) + { + _logger?.LogError( + ex, + "MQTT driver '{DriverId}': reconnect abandoned — the broker configuration is unusable.", + _driverId); + return; + } + + _logger?.LogDebug( + ex, + "MQTT driver '{DriverId}': reconnect attempt {Attempt} to {Host}:{Port} failed.", + _driverId, + attempt + 1, + _options.Host, + _options.Port); + continue; + } + + try + { + // ALWAYS, on every reconnect — including persistent sessions, where this is + // merely redundant. Subscriptions do not survive a clean session, and a + // reconnect that skips this leaves a healthy-looking, permanently deaf client. + await FireReconnectedAsync().ConfigureAwait(false); + } + catch (Exception ex) + { + _logger?.LogError( + ex, + "MQTT driver '{DriverId}': re-subscribe after reconnect failed; tearing the session down " + + "and retrying rather than serving a connection that receives nothing.", + _driverId); + SetState(MqttConnectionState.Reconnecting); + await ForceDisconnectAsync().ConfigureAwait(false); + continue; + } + + _logger?.LogInformation( + "MQTT driver '{DriverId}': reconnected to {Host}:{Port} after {Attempts} attempt(s); " + + "subscriptions re-established.", + _driverId, + _options.Host, + _options.Port, + attempt + 1); + break; + } + } + } + catch (OperationCanceledException) + { + // Dispose cancelled the lifetime token — the expected way this loop ends. + } + catch (Exception ex) + { + // Reaching here means the connection is permanently dark, which is exactly the failure + // this supervisor exists to prevent — it must never be swallowed quietly. + _logger?.LogError( + ex, + "MQTT driver '{DriverId}': reconnect supervisor stopped unexpectedly; the broker session will not " + + "recover without a driver restart.", + _driverId); + } + } + + /// + /// Drops a session that is connected but unusable, so the next supervisor pass starts from a + /// clean CONNECT. Best-effort and bounded — failure here just means the next attempt sees a + /// client that is already down. + /// + private async Task ForceDisconnectAsync() + { + var client = _client; if (client is null) { return; } + try + { + using var deadline = new CancellationTokenSource(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)); + await client.DisconnectAsync(new MqttClientDisconnectOptions(), deadline.Token).ConfigureAwait(false); + } + catch (Exception ex) + { + _logger?.LogDebug(ex, "MQTT driver '{DriverId}': forced disconnect failed.", _driverId); + } + } + + private async Task CloseAsync(IMqttClient client) + { try { if (client.IsConnected) @@ -105,75 +666,15 @@ public sealed class MqttConnection : IAsyncDisposable } } - /// - /// Connects to the broker, failing at - /// rather than waiting indefinitely on an unresponsive peer. - /// - /// - /// May be retried on the same instance after a failed attempt — the underlying - /// is created once and reused across attempts, which is what the - /// Task-4 reconnect loop depends on. Not safe to call concurrently with itself or with - /// ; see the type-level remarks. - /// - /// Caller cancellation; linked with the connect deadline. - /// The connect deadline elapsed. - /// was cancelled. - /// - /// The connection was disposed, either before the attempt started or while it was in flight. - /// - public async Task ConnectAsync(CancellationToken cancellationToken) + /// is terminal — nothing may move off it. + private void SetState(MqttConnectionState state) { - ObjectDisposedException.ThrowIf(_disposed, this); - - // The options are rebuilt per attempt so a rotated CA file is picked up on the next - // connect rather than being pinned for the life of the process. - var client = _client ??= new MqttClientFactory().CreateMqttClient(); - var clientOptions = BuildClientOptions(_options, clientIdSuffix: null, _logger); - - using var deadline = new CancellationTokenSource(TimeSpan.FromSeconds(_options.ConnectTimeoutSeconds)); - using var linked = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, deadline.Token); - - _logger?.LogDebug( - "MQTT driver '{DriverId}': connecting to {Host}:{Port} (tls={UseTls}, timeout={TimeoutSeconds}s).", - _driverId, - _options.Host, - _options.Port, - _options.UseTls, - _options.ConnectTimeoutSeconds); - - try + if (Volatile.Read(ref _state) == (int)MqttConnectionState.Disposed) { - await client.ConnectAsync(clientOptions, linked.Token).ConfigureAwait(false); - } - // A concurrent DisposeAsync disposes the client out from under the awaiting connect, and - // whatever MQTTnet throws for that matches neither token filter. Classify it rather than - // letting an undocumented exception escape the contract below. - catch (Exception ex) when (_disposed) - { - throw new ObjectDisposedException( - nameof(MqttConnection), - new InvalidOperationException( - $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} was aborted because the " - + "connection was disposed while the attempt was in flight.", - ex)); - } - // MQTTnet wraps a cancelled connect in MqttConnectingFailedException rather than letting - // the OperationCanceledException surface, so the two legs are told apart by which token - // fired, not by the exception type. Caller intent wins over the deadline. - catch (Exception ex) when (cancellationToken.IsCancellationRequested) - { - throw new OperationCanceledException( - $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} was cancelled.", - ex, - cancellationToken); - } - catch (Exception ex) when (deadline.IsCancellationRequested) - { - throw new TimeoutException( - $"MQTT driver '{_driverId}': connect to {_options.Host}:{_options.Port} did not complete within " - + $"{_options.ConnectTimeoutSeconds}s.", - ex); + return; } + + Volatile.Write(ref _state, (int)state); } /// diff --git a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttConnectionTests.cs b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttConnectionTests.cs index c95cb714..b4b5a3b9 100644 --- a/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttConnectionTests.cs +++ b/tests/Drivers/ZB.MOM.WW.OtOpcUa.Driver.Mqtt.Tests/MqttConnectionTests.cs @@ -376,10 +376,366 @@ public sealed class MqttConnectionTests conn.IsConnected.ShouldBeFalse(); } + // --------------------------------------------------------------------------------- + // Task 4 — NextBackoff (pure): bounded, monotone, cannot hot-loop, cannot overflow + // --------------------------------------------------------------------------------- + + [Theory] + [InlineData(0, 1)] + [InlineData(1, 2)] + [InlineData(2, 4)] + [InlineData(3, 8)] + [InlineData(4, 16)] + [InlineData(5, 30)] // 32 → clamped + [InlineData(10, 30)] + public void NextBackoff_IsExponentialClampedToMax(int attempt, int expectedSeconds) + => MqttConnection.NextBackoff(attempt, minSeconds: 1, maxSeconds: 30) + .ShouldBe(TimeSpan.FromSeconds(expectedSeconds)); + + /// + /// A shift-based implementation (min << attempt) wraps: 1 << 32 is + /// 1, not 4294967296, so a broker that has been down for half an hour would suddenly + /// be hammered once a second. The calculator must saturate, not wrap. + /// + [Fact] + public void NextBackoff_HugeAttemptCount_SaturatesAtMax_NeverWrapsBackToMin() + { + foreach (var attempt in new[] { 31, 32, 33, 40, 62, 63, 64, 1_000, int.MaxValue }) + { + MqttConnection.NextBackoff(attempt, minSeconds: 1, maxSeconds: 30) + .ShouldBe(TimeSpan.FromSeconds(30), $"attempt {attempt}"); + } + } + + [Fact] + public void NextBackoff_IsMonotoneNonDecreasing_AndAlwaysWithinMinMax() + { + var previous = TimeSpan.Zero; + for (var attempt = 0; attempt <= 64; attempt++) + { + var backoff = MqttConnection.NextBackoff(attempt, minSeconds: 2, maxSeconds: 30); + + backoff.ShouldBeGreaterThanOrEqualTo(previous, $"attempt {attempt} went backwards"); + backoff.ShouldBeGreaterThanOrEqualTo(TimeSpan.FromSeconds(2), $"attempt {attempt} below min"); + backoff.ShouldBeLessThanOrEqualTo(TimeSpan.FromSeconds(30), $"attempt {attempt} above max"); + previous = backoff; + } + } + + /// + /// A misconfigured max < min must not collapse the floor — the floor is the only + /// thing standing between a dead broker and a hot reconnect loop. + /// + [Fact] + public void NextBackoff_MaxBelowMin_HonoursMin_SoTheLoopCannotGoHot() + => MqttConnection.NextBackoff(0, minSeconds: 10, maxSeconds: 5).ShouldBe(TimeSpan.FromSeconds(10)); + + [Fact] + public void NextBackoff_NonPositiveAttempt_IsTheFirstAttemptDelay() + { + MqttConnection.NextBackoff(0, minSeconds: 3, maxSeconds: 30).ShouldBe(TimeSpan.FromSeconds(3)); + MqttConnection.NextBackoff(-5, minSeconds: 3, maxSeconds: 30).ShouldBe(TimeSpan.FromSeconds(3)); + } + + // --------------------------------------------------------------------------------- + // Task 4 — the Reconnected callback (the load-bearing re-subscribe hook) + // --------------------------------------------------------------------------------- + + [Fact] + public async Task FireReconnected_InvokesEverySubscribedCallback() + { + await using var conn = new MqttConnection(LoopbackOpts(), driverId: "t", logger: null); + var first = 0; + var second = 0; + conn.Reconnected += () => { Interlocked.Increment(ref first); return Task.CompletedTask; }; + conn.Reconnected += () => { Interlocked.Increment(ref second); return Task.CompletedTask; }; + + await conn.FireReconnectedAsync(); + + first.ShouldBe(1); + second.ShouldBe(1); + } + + /// + /// Awaiting a multicast Func<Task> directly returns only the LAST handler's task + /// and abandons the earlier ones — so one throwing subscriber would silently skip every + /// subscriber behind it. The fan-out must walk the invocation list and still surface the + /// failure, because a re-subscribe that failed is a connection that has gone dark. + /// + [Fact] + public async Task FireReconnected_OneHandlerThrows_StillInvokesTheRest_AndSurfacesTheFailure() + { + await using var conn = new MqttConnection(LoopbackOpts(), driverId: "t", logger: null); + var behindTheThrower = 0; + conn.Reconnected += () => throw new InvalidOperationException("SUBACK refused"); + conn.Reconnected += () => { Interlocked.Increment(ref behindTheThrower); return Task.CompletedTask; }; + + await Should.ThrowAsync(async () => await conn.FireReconnectedAsync()); + + behindTheThrower.ShouldBe(1); + } + + [Fact] + public async Task FireReconnected_NoSubscribers_IsANoOp() + { + await using var conn = new MqttConnection(LoopbackOpts(), driverId: "t", logger: null); + + await Should.NotThrowAsync(async () => await conn.FireReconnectedAsync()); + } + + // --------------------------------------------------------------------------------- + // Task 4 — State + LastMessage + // --------------------------------------------------------------------------------- + + [Fact] + public void State_AfterCtor_IsDisconnected() + => new MqttConnection(LoopbackOpts(), driverId: "t", logger: null) + .State.ShouldBe(MqttConnectionState.Disconnected); + + [Fact] + public async Task State_AfterDispose_IsDisposed() + { + var conn = new MqttConnection(LoopbackOpts(), driverId: "t", logger: null); + + await conn.DisposeAsync(); + + conn.State.ShouldBe(MqttConnectionState.Disposed); + } + + /// + /// Faulted is reserved for config that cannot succeed no matter how long we retry. A broker + /// that is merely down is , indefinitely. + /// + [Fact] + public async Task State_UnrecoverableProtocolConfig_IsFaulted() + { + var opts = new MqttDriverOptions + { + Host = "127.0.0.1", + Port = 1, + UseTls = false, + ConnectTimeoutSeconds = 1, + ProtocolVersion = (MqttProtocolVersion)99, + }; + await using var conn = new MqttConnection(opts, driverId: "t", logger: null); + + await Should.ThrowAsync(async () => await conn.ConnectAsync(CancellationToken.None)); + + conn.State.ShouldBe(MqttConnectionState.Faulted); + } + + /// A broker that is simply unreachable is retryable, never Faulted. + [Fact] + public async Task State_UnreachableBroker_IsNotFaulted() + { + using var blackhole = new BlackholeBroker(); + var opts = new MqttDriverOptions + { + Host = "127.0.0.1", Port = blackhole.Port, UseTls = false, ConnectTimeoutSeconds = 1, + }; + await using var conn = new MqttConnection(opts, driverId: "t", logger: null); + + await Should.ThrowAsync(async () => await conn.ConnectAsync(CancellationToken.None)); + + conn.State.ShouldNotBe(MqttConnectionState.Faulted); + } + + [Fact] + public void LastMessage_BeforeAnyTraffic_IsNull() + { + var conn = new MqttConnection(LoopbackOpts(), driverId: "t", logger: null); + + conn.LastMessageUtc.ShouldBeNull(); + conn.LastMessageAge.ShouldBeNull(); + } + + /// + /// A failed FIRST connect is the caller's problem (driver-host resilience retries + /// InitializeAsync) — it must not silently leave a background reconnect worker + /// hammering an endpoint the caller has already given up on. + /// + [Fact] + public async Task ConnectAsync_InitialAttemptFails_LeavesNoBackgroundReconnectStorm() + { + using var blackhole = new BlackholeBroker(); + var opts = new MqttDriverOptions + { + Host = "127.0.0.1", + Port = blackhole.Port, + UseTls = false, + ConnectTimeoutSeconds = 1, + ReconnectMinBackoffSeconds = 1, + ReconnectMaxBackoffSeconds = 1, + }; + await using var conn = new MqttConnection(opts, driverId: "t", logger: null); + + await Should.ThrowAsync(async () => await conn.ConnectAsync(CancellationToken.None)); + var acceptedAfterTheAttempt = blackhole.AcceptedCount; + + await Task.Delay(TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + + blackhole.AcceptedCount.ShouldBe(acceptedAfterTheAttempt); + } + + // --------------------------------------------------------------------------------- + // Task 4 — the reconnect supervisor, against a real (minimal) broker + // --------------------------------------------------------------------------------- + + /// + /// THE load-bearing invariant. MQTT subscriptions do not survive a clean session, so a + /// reconnect that completes without firing Reconnected leaves a healthy-looking + /// connection that receives nothing, forever, with no error to show for it. Driven through + /// the real socket → real MQTTnet DisconnectedAsync → real re-CONNECT path, not a + /// test seam. + /// + [Fact] + public async Task BrokerDropsTheConnection_ReconnectsAndFiresReconnected() + { + using var broker = new MiniBroker(); + await using var conn = new MqttConnection(BrokerOpts(broker.Port), driverId: "t", logger: null); + var resubscribes = 0; + conn.Reconnected += () => { Interlocked.Increment(ref resubscribes); return Task.CompletedTask; }; + + await conn.ConnectAsync(CancellationToken.None); + + conn.State.ShouldBe(MqttConnectionState.Connected); + conn.IsConnected.ShouldBeTrue(); + Volatile.Read(ref resubscribes).ShouldBe(0, "the initial connect is not a reconnect"); + + broker.DropAll(); + + (await WaitUntilAsync(() => Volatile.Read(ref resubscribes) >= 1, TimeSpan.FromSeconds(30))) + .ShouldBeTrue("the reconnect never fired Reconnected — subscriptions would be silently gone"); + (await WaitUntilAsync(() => conn.IsConnected, TimeSpan.FromSeconds(10))).ShouldBeTrue(); + conn.State.ShouldBe(MqttConnectionState.Connected); + broker.AcceptedCount.ShouldBeGreaterThanOrEqualTo(2); + } + + /// + /// A broker that stays down keeps the supervisor in Reconnecting indefinitely (never + /// Faulted) — and the backoff keeps the retry rate low enough that a handful of + /// attempts, not hundreds, land in the observation window. + /// + [Fact] + public async Task BrokerStaysDown_RetriesUnderBackoff_WithoutHotLooping() + { + using var broker = new MiniBroker(); + var opts = BrokerOpts(broker.Port) with + { + ConnectTimeoutSeconds = 1, ReconnectMinBackoffSeconds = 1, ReconnectMaxBackoffSeconds = 2, + }; + await using var conn = new MqttConnection(opts, driverId: "t", logger: null); + + await conn.ConnectAsync(CancellationToken.None); + var acceptedWhileHealthy = broker.AcceptedCount; + + broker.Blackhole = true; // accepts TCP, never answers CONNACK again + broker.DropAll(); + + (await WaitUntilAsync(() => conn.State == MqttConnectionState.Reconnecting, TimeSpan.FromSeconds(10))) + .ShouldBeTrue(); + await Task.Delay(TimeSpan.FromSeconds(8), TestContext.Current.CancellationToken); + + conn.State.ShouldBe(MqttConnectionState.Reconnecting, "a down broker is retryable, not Faulted"); + var retries = broker.AcceptedCount - acceptedWhileHealthy; + retries.ShouldBeGreaterThanOrEqualTo(1, "the supervisor gave up on a down broker"); + retries.ShouldBeLessThan(15, $"backoff is not bounding the retry rate ({retries} attempts in ~8s)"); + } + + /// + /// Dispose during a reconnect backoff must stop the supervisor dead — no leaked task, no + /// later attempt resurrecting a connection the caller has torn down. + /// + [Fact] + public async Task Dispose_DuringReconnectBackoff_StopsTheSupervisor_NoFurtherAttempts() + { + using var broker = new MiniBroker(); + var opts = BrokerOpts(broker.Port) with + { + ConnectTimeoutSeconds = 1, ReconnectMinBackoffSeconds = 1, ReconnectMaxBackoffSeconds = 2, + }; + var conn = new MqttConnection(opts, driverId: "t", logger: null); + + await conn.ConnectAsync(CancellationToken.None); + broker.Blackhole = true; + broker.DropAll(); + (await WaitUntilAsync(() => conn.State == MqttConnectionState.Reconnecting, TimeSpan.FromSeconds(10))) + .ShouldBeTrue(); + + await conn.DisposeAsync(); + var acceptedAtDispose = broker.AcceptedCount; + + conn.State.ShouldBe(MqttConnectionState.Disposed); + conn.IsConnected.ShouldBeFalse(); + await Task.Delay(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + broker.AcceptedCount.ShouldBe(acceptedAtDispose, "the supervisor outlived DisposeAsync"); + } + + /// + /// The T3 connect/dispose leak, reproduced at the exact interleaving the reviewer traced: + /// the client has already CONNECTED when dispose flags the instance and blocks on the + /// lifecycle gate this connect is holding — so dispose disposes nothing. Without the + /// post-connect re-check the client stays live with an open socket that no later + /// DisposeAsync can ever reach; the broker sees the leak as a connection that never + /// closes. + /// + [Fact] + public async Task DisposeRacingAnInFlightConnect_LeavesNoLiveBrokerConnection() + { + using var broker = new MiniBroker(); + var conn = new MqttConnection(BrokerOpts(broker.Port), driverId: "t", logger: null); + conn.AfterConnectHookForTests = async () => + { + _ = Task.Run(async () => await conn.DisposeAsync()); + (await WaitUntilAsync(() => conn.State == MqttConnectionState.Disposed, TimeSpan.FromSeconds(5))) + .ShouldBeTrue("dispose never flagged the instance — the race was not reproduced"); + }; + + await Should.ThrowAsync(async () => await conn.ConnectAsync(CancellationToken.None)); + + (await WaitUntilAsync(() => broker.LiveConnections == 0, TimeSpan.FromSeconds(15))) + .ShouldBeTrue($"a live broker connection leaked ({broker.LiveConnections} still open)"); + conn.IsConnected.ShouldBeFalse(); + } + // --------------------------------------------------------------------------------- // helpers // --------------------------------------------------------------------------------- + private static MqttDriverOptions LoopbackOpts() + => new() { Host = "127.0.0.1", Port = 1, UseTls = false, ConnectTimeoutSeconds = 1 }; + + /// Options aimed at : plaintext v3.1.1, keep-alive parked out of the way. + private static MqttDriverOptions BrokerOpts(int port) + => new() + { + Host = "127.0.0.1", + Port = port, + UseTls = false, + ProtocolVersion = MqttProtocolVersion.V311, + CleanSession = true, + KeepAliveSeconds = 300, + ConnectTimeoutSeconds = 5, + ReconnectMinBackoffSeconds = 1, + ReconnectMaxBackoffSeconds = 2, + }; + + private static async Task WaitUntilAsync(Func condition, TimeSpan timeout) + { + var deadline = Stopwatch.StartNew(); + while (deadline.Elapsed < timeout) + { + if (condition()) + { + return true; + } + + await Task.Delay(25); + } + + return condition(); + } + private static MqttClientTlsOptions BuildTls(MqttDriverOptions opts) => MqttConnection.BuildClientOptions(opts, clientIdSuffix: null) .ChannelOptions.ShouldBeOfType().TlsOptions; @@ -475,6 +831,7 @@ public sealed class MqttConnectionTests private readonly List _accepted = []; private readonly CancellationTokenSource _cts = new(); private readonly TcpListener _listener; + private int _acceptedCount; public BlackholeBroker() { @@ -486,6 +843,9 @@ public sealed class MqttConnectionTests public int Port { get; } + /// Total TCP connections accepted — how a caller detects a background retry storm. + public int AcceptedCount => Volatile.Read(ref _acceptedCount); + public void Dispose() { _cts.Cancel(); @@ -510,6 +870,7 @@ public sealed class MqttConnectionTests while (!_cts.IsCancellationRequested) { var client = await _listener.AcceptTcpClientAsync(_cts.Token); + Interlocked.Increment(ref _acceptedCount); lock (_accepted) { _accepted.Add(client); @@ -522,4 +883,156 @@ public sealed class MqttConnectionTests } } } + + /// + /// The smallest broker that can complete an MQTT 3.1.1 handshake: accept TCP, answer CONNECT + /// with a success CONNACK, answer PINGREQ with PINGRESP, and nothing else. That is enough for + /// MQTTnet to report a real connection — which is what makes a real drop + /// () drive a real DisconnectedAsync and a real reconnect, rather + /// than a test seam standing in for one. + /// + private sealed class MiniBroker : IDisposable + { + private static readonly byte[] ConnAckAccepted = [0x20, 0x02, 0x00, 0x00]; + private static readonly byte[] PingResp = [0xD0, 0x00]; + + private readonly CancellationTokenSource _cts = new(); + private readonly HashSet _live = []; + private readonly TcpListener _listener; + private int _acceptedCount; + + public MiniBroker() + { + _listener = new TcpListener(IPAddress.Loopback, 0); + _listener.Start(); + Port = ((IPEndPoint)_listener.LocalEndpoint).Port; + _ = Task.Run(AcceptLoopAsync); + } + + public int Port { get; } + + /// When set, TCP is still accepted but CONNECT is never answered (frozen peer). + public bool Blackhole { get; set; } + + /// Total TCP connections accepted since construction — counts reconnect attempts. + public int AcceptedCount => Volatile.Read(ref _acceptedCount); + + /// + /// Sockets the broker still has open. Drops to zero only once the peer actually closes, + /// so a leaked-but-live client keeps this above zero. + /// + public int LiveConnections + { + get + { + lock (_live) + { + return _live.Count; + } + } + } + + /// Kills every established session — the "broker went away" event. + public void DropAll() + { + lock (_live) + { + foreach (var client in _live) + { + try + { + client.Close(); + } + catch (Exception) + { + // Already gone. + } + } + + _live.Clear(); + } + } + + public void Dispose() + { + _cts.Cancel(); + _listener.Stop(); + DropAll(); + _cts.Dispose(); + } + + private async Task AcceptLoopAsync() + { + try + { + while (!_cts.IsCancellationRequested) + { + var client = await _listener.AcceptTcpClientAsync(_cts.Token); + Interlocked.Increment(ref _acceptedCount); + lock (_live) + { + _live.Add(client); + } + + _ = Task.Run(() => ServeAsync(client)); + } + } + catch (Exception) + { + // Listener stopped / token cancelled — the fixture is going away. + } + } + + private async Task ServeAsync(TcpClient client) + { + var buffer = new byte[1024]; + try + { + var stream = client.GetStream(); + while (!_cts.IsCancellationRequested) + { + var read = await stream.ReadAsync(buffer, _cts.Token); + if (read == 0) + { + break; // peer closed — this is how a disposed client stops counting as live + } + + if (Blackhole) + { + continue; + } + + var packetType = buffer[0] & 0xF0; + if (packetType == 0x10) + { + await stream.WriteAsync(ConnAckAccepted, _cts.Token); + } + else if (packetType == 0xC0) + { + await stream.WriteAsync(PingResp, _cts.Token); + } + } + } + catch (Exception) + { + // Socket torn down by either side. + } + finally + { + lock (_live) + { + _live.Remove(client); + } + + try + { + client.Dispose(); + } + catch (Exception) + { + // Already gone. + } + } + } + } }