From 3bb4f695660e35665468f9f53328eb8384f29b2c Mon Sep 17 00:00:00 2001 From: Jared Dohrman Date: Sun, 12 Apr 2026 09:01:25 +1000 Subject: [PATCH] fix: add graceful MQTT reconnect backoff during broker outages --- src/helpers/mqtt/MQTTUplink.cpp | 66 +++++++++++++++++++++++++++------ src/helpers/mqtt/MQTTUplink.h | 5 ++- 2 files changed, 59 insertions(+), 12 deletions(-) diff --git a/src/helpers/mqtt/MQTTUplink.cpp b/src/helpers/mqtt/MQTTUplink.cpp index fe745afd..49c32fad 100644 --- a/src/helpers/mqtt/MQTTUplink.cpp +++ b/src/helpers/mqtt/MQTTUplink.cpp @@ -50,7 +50,8 @@ namespace { constexpr unsigned long kWifiRetryMillis = 15000; constexpr unsigned long kWifiConnectTimeoutMillis = 45000; -constexpr unsigned long kBrokerRetryMillis = 10000; +constexpr unsigned long kBrokerRetryBaseMillis = 10000; +constexpr unsigned long kBrokerRetryMaxMillis = 300000; constexpr size_t kBrokerTokenSize = 640; constexpr time_t kTokenLifetimeSecs = 3600; constexpr time_t kTokenRefreshSlackSecs = 300; @@ -79,6 +80,18 @@ const char* getWifiQualityLabel(int rssi_dbm) { return "poor"; } +unsigned long getBrokerRetryDelayMillis(uint8_t failures) { + unsigned long delay_ms = kBrokerRetryBaseMillis; + if (failures > 0) { + uint8_t shifts = min(failures - 1, 5); + delay_ms <<= shifts; + } + if (delay_ms > kBrokerRetryMaxMillis) { + delay_ms = kBrokerRetryMaxMillis; + } + return delay_ms; +} + char* allocScratchBuffer(size_t size) { void* ptr = heap_caps_malloc(size, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); if (ptr == nullptr) { @@ -384,7 +397,7 @@ bool MQTTUplink::refreshToken(BrokerState& broker) { return true; } -void MQTTUplink::destroyBroker(BrokerState& broker) { +void MQTTUplink::destroyBroker(BrokerState& broker, bool reset_retry_state) { if (broker.client != nullptr) { MQTT_LOG("%s destroy broker client", broker.spec->label); esp_mqtt_client_stop(broker.client); @@ -396,6 +409,12 @@ void MQTTUplink::destroyBroker(BrokerState& broker) { broker.connected = false; broker.connect_announced = false; broker.token_expires_at = 0; + if (reset_retry_state) { + broker.reconnect_pending = false; + broker.next_connect_attempt = 0; + broker.reconnect_failures = 0; + broker.last_connect_attempt = 0; + } } void MQTTUplink::queuePublish(BrokerState& broker, const char* topic, const char* payload, bool retain) { @@ -573,12 +592,20 @@ void MQTTUplink::handleMqttEvent(void* handler_args, esp_event_base_t, int32_t e switch (event_id) { case MQTT_EVENT_CONNECTED: broker->connected = true; + broker->reconnect_pending = false; + broker->next_connect_attempt = 0; + broker->reconnect_failures = 0; MQTT_LOG("%s connected", broker->spec->label); break; case MQTT_EVENT_DISCONNECTED: MQTT_LOG("%s disconnected", broker->spec->label); case MQTT_EVENT_ERROR: broker->connected = false; + if (broker->reconnect_failures < 10) { + broker->reconnect_failures++; + } + broker->reconnect_pending = true; + broker->next_connect_attempt = millis() + getBrokerRetryDelayMillis(broker->reconnect_failures); if (event_id == MQTT_EVENT_ERROR) { if (event != nullptr && event->error_handle != nullptr) { MQTT_LOG("%s error type=%d tls_esp=0x%x tls_stack=0x%x cert_flags=0x%x sock_errno=%d conn_refused=%d", @@ -589,6 +616,9 @@ void MQTTUplink::handleMqttEvent(void* handler_args, esp_event_base_t, int32_t e MQTT_LOG("%s error event", broker->spec->label); } } + MQTT_LOG("%s reconnect in %lu ms (failures=%u)", broker->spec->label, + getBrokerRetryDelayMillis(broker->reconnect_failures), + static_cast(broker->reconnect_failures)); break; case MQTT_EVENT_BEFORE_CONNECT: MQTT_LOG("%s before connect", broker->spec->label); @@ -702,20 +732,29 @@ void MQTTUplink::ensureBroker(BrokerState& broker) { time_t now = time(nullptr); if (broker.client != nullptr && broker.token_expires_at > 0 && now + kTokenRefreshSlackSecs >= broker.token_expires_at) { - destroyBroker(broker); - } - - if (broker.client != nullptr) { - return; + destroyBroker(broker, false); } unsigned long now_ms = millis(); - if (now_ms - broker.last_connect_attempt < kBrokerRetryMillis) { + if (broker.client != nullptr) { + if (broker.connected) { + return; + } + if (!broker.reconnect_pending || now_ms < broker.next_connect_attempt) { + return; + } + destroyBroker(broker, false); + } + + if (broker.next_connect_attempt != 0 && now_ms < broker.next_connect_attempt) { return; } broker.last_connect_attempt = now_ms; + broker.reconnect_pending = false; if (!refreshToken(broker)) { + broker.reconnect_pending = true; + broker.next_connect_attempt = now_ms + kBrokerRetryBaseMillis; return; } @@ -739,7 +778,7 @@ void MQTTUplink::ensureBroker(BrokerState& broker) { cfg.session.last_will.retain = 1; cfg.network.reconnect_timeout_ms = 10000; cfg.network.timeout_ms = 10000; - cfg.network.disable_auto_reconnect = false; + cfg.network.disable_auto_reconnect = true; cfg.buffer.size = 768; cfg.buffer.out_size = 1280; #else @@ -753,7 +792,7 @@ void MQTTUplink::ensureBroker(BrokerState& broker) { cfg.out_buffer_size = 1280; cfg.reconnect_timeout_ms = 10000; cfg.network_timeout_ms = 10000; - cfg.disable_auto_reconnect = false; + cfg.disable_auto_reconnect = true; cfg.transport = MQTT_TRANSPORT_OVER_WSS; cfg.cert_pem = mqtt_ca_certs::kCombinedPem; cfg.lwt_topic = broker.status_topic; @@ -772,7 +811,9 @@ void MQTTUplink::ensureBroker(BrokerState& broker) { esp_mqtt_client_register_event(broker.client, MQTT_EVENT_ANY, &MQTTUplink::handleMqttEvent, &broker); if (esp_mqtt_client_start(broker.client) != ESP_OK) { MQTT_LOG("%s mqtt start failed", broker.spec->label); - destroyBroker(broker); + broker.reconnect_pending = true; + broker.next_connect_attempt = now_ms + kBrokerRetryBaseMillis; + destroyBroker(broker, false); } else { MQTT_LOG("%s mqtt start requested", broker.spec->label); } @@ -918,6 +959,9 @@ void MQTTUplink::formatStatusReply(char* reply, size_t reply_size) const { if (broker->client != nullptr) { return "conn"; } + if (broker->next_connect_attempt != 0 && broker->next_connect_attempt > millis()) { + return "backoff"; + } return "retry"; }; diff --git a/src/helpers/mqtt/MQTTUplink.h b/src/helpers/mqtt/MQTTUplink.h index fea5bbf0..6b89c998 100644 --- a/src/helpers/mqtt/MQTTUplink.h +++ b/src/helpers/mqtt/MQTTUplink.h @@ -95,8 +95,11 @@ private: esp_mqtt_client_handle_t client; bool connected; bool connect_announced; + bool reconnect_pending; unsigned long last_connect_attempt; + unsigned long next_connect_attempt; time_t token_expires_at; + uint8_t reconnect_failures; char username[70]; char* token; char client_id[48]; @@ -142,7 +145,7 @@ private: void refreshBrokerIdentity(BrokerState& broker); void refreshBrokerState(BrokerState& broker); void ensureBroker(BrokerState& broker); - void destroyBroker(BrokerState& broker); + void destroyBroker(BrokerState& broker, bool reset_retry_state = true); bool refreshToken(BrokerState& broker); void publishStatus(bool online); void publishOnlineStatus(BrokerState& broker);