fix: add graceful MQTT reconnect backoff during broker outages

Dieser Commit ist enthalten in:
Jared Dohrman
2026-04-12 09:01:25 +10:00
Ursprung 4f4964e454
Commit 3bb4f69566
2 geänderte Dateien mit 59 neuen und 12 gelöschten Zeilen
+55 -11
Datei anzeigen
@@ -50,7 +50,8 @@
namespace { namespace {
constexpr unsigned long kWifiRetryMillis = 15000; constexpr unsigned long kWifiRetryMillis = 15000;
constexpr unsigned long kWifiConnectTimeoutMillis = 45000; 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 size_t kBrokerTokenSize = 640;
constexpr time_t kTokenLifetimeSecs = 3600; constexpr time_t kTokenLifetimeSecs = 3600;
constexpr time_t kTokenRefreshSlackSecs = 300; constexpr time_t kTokenRefreshSlackSecs = 300;
@@ -79,6 +80,18 @@ const char* getWifiQualityLabel(int rssi_dbm) {
return "poor"; return "poor";
} }
unsigned long getBrokerRetryDelayMillis(uint8_t failures) {
unsigned long delay_ms = kBrokerRetryBaseMillis;
if (failures > 0) {
uint8_t shifts = min<uint8_t>(failures - 1, 5);
delay_ms <<= shifts;
}
if (delay_ms > kBrokerRetryMaxMillis) {
delay_ms = kBrokerRetryMaxMillis;
}
return delay_ms;
}
char* allocScratchBuffer(size_t size) { char* allocScratchBuffer(size_t size) {
void* ptr = heap_caps_malloc(size, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); void* ptr = heap_caps_malloc(size, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT);
if (ptr == nullptr) { if (ptr == nullptr) {
@@ -384,7 +397,7 @@ bool MQTTUplink::refreshToken(BrokerState& broker) {
return true; return true;
} }
void MQTTUplink::destroyBroker(BrokerState& broker) { void MQTTUplink::destroyBroker(BrokerState& broker, bool reset_retry_state) {
if (broker.client != nullptr) { if (broker.client != nullptr) {
MQTT_LOG("%s destroy broker client", broker.spec->label); MQTT_LOG("%s destroy broker client", broker.spec->label);
esp_mqtt_client_stop(broker.client); esp_mqtt_client_stop(broker.client);
@@ -396,6 +409,12 @@ void MQTTUplink::destroyBroker(BrokerState& broker) {
broker.connected = false; broker.connected = false;
broker.connect_announced = false; broker.connect_announced = false;
broker.token_expires_at = 0; 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) { 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) { switch (event_id) {
case MQTT_EVENT_CONNECTED: case MQTT_EVENT_CONNECTED:
broker->connected = true; broker->connected = true;
broker->reconnect_pending = false;
broker->next_connect_attempt = 0;
broker->reconnect_failures = 0;
MQTT_LOG("%s connected", broker->spec->label); MQTT_LOG("%s connected", broker->spec->label);
break; break;
case MQTT_EVENT_DISCONNECTED: case MQTT_EVENT_DISCONNECTED:
MQTT_LOG("%s disconnected", broker->spec->label); MQTT_LOG("%s disconnected", broker->spec->label);
case MQTT_EVENT_ERROR: case MQTT_EVENT_ERROR:
broker->connected = false; 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_id == MQTT_EVENT_ERROR) {
if (event != nullptr && event->error_handle != nullptr) { 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", 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 error event", broker->spec->label);
} }
} }
MQTT_LOG("%s reconnect in %lu ms (failures=%u)", broker->spec->label,
getBrokerRetryDelayMillis(broker->reconnect_failures),
static_cast<unsigned>(broker->reconnect_failures));
break; break;
case MQTT_EVENT_BEFORE_CONNECT: case MQTT_EVENT_BEFORE_CONNECT:
MQTT_LOG("%s before connect", broker->spec->label); MQTT_LOG("%s before connect", broker->spec->label);
@@ -702,20 +732,29 @@ void MQTTUplink::ensureBroker(BrokerState& broker) {
time_t now = time(nullptr); time_t now = time(nullptr);
if (broker.client != nullptr && broker.token_expires_at > 0 && now + kTokenRefreshSlackSecs >= broker.token_expires_at) { if (broker.client != nullptr && broker.token_expires_at > 0 && now + kTokenRefreshSlackSecs >= broker.token_expires_at) {
destroyBroker(broker); destroyBroker(broker, false);
}
if (broker.client != nullptr) {
return;
} }
unsigned long now_ms = millis(); 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; return;
} }
broker.last_connect_attempt = now_ms; broker.last_connect_attempt = now_ms;
broker.reconnect_pending = false;
if (!refreshToken(broker)) { if (!refreshToken(broker)) {
broker.reconnect_pending = true;
broker.next_connect_attempt = now_ms + kBrokerRetryBaseMillis;
return; return;
} }
@@ -739,7 +778,7 @@ void MQTTUplink::ensureBroker(BrokerState& broker) {
cfg.session.last_will.retain = 1; cfg.session.last_will.retain = 1;
cfg.network.reconnect_timeout_ms = 10000; cfg.network.reconnect_timeout_ms = 10000;
cfg.network.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.size = 768;
cfg.buffer.out_size = 1280; cfg.buffer.out_size = 1280;
#else #else
@@ -753,7 +792,7 @@ void MQTTUplink::ensureBroker(BrokerState& broker) {
cfg.out_buffer_size = 1280; cfg.out_buffer_size = 1280;
cfg.reconnect_timeout_ms = 10000; cfg.reconnect_timeout_ms = 10000;
cfg.network_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.transport = MQTT_TRANSPORT_OVER_WSS;
cfg.cert_pem = mqtt_ca_certs::kCombinedPem; cfg.cert_pem = mqtt_ca_certs::kCombinedPem;
cfg.lwt_topic = broker.status_topic; 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); esp_mqtt_client_register_event(broker.client, MQTT_EVENT_ANY, &MQTTUplink::handleMqttEvent, &broker);
if (esp_mqtt_client_start(broker.client) != ESP_OK) { if (esp_mqtt_client_start(broker.client) != ESP_OK) {
MQTT_LOG("%s mqtt start failed", broker.spec->label); 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 { } else {
MQTT_LOG("%s mqtt start requested", broker.spec->label); 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) { if (broker->client != nullptr) {
return "conn"; return "conn";
} }
if (broker->next_connect_attempt != 0 && broker->next_connect_attempt > millis()) {
return "backoff";
}
return "retry"; return "retry";
}; };
+4 -1
Datei anzeigen
@@ -95,8 +95,11 @@ private:
esp_mqtt_client_handle_t client; esp_mqtt_client_handle_t client;
bool connected; bool connected;
bool connect_announced; bool connect_announced;
bool reconnect_pending;
unsigned long last_connect_attempt; unsigned long last_connect_attempt;
unsigned long next_connect_attempt;
time_t token_expires_at; time_t token_expires_at;
uint8_t reconnect_failures;
char username[70]; char username[70];
char* token; char* token;
char client_id[48]; char client_id[48];
@@ -142,7 +145,7 @@ private:
void refreshBrokerIdentity(BrokerState& broker); void refreshBrokerIdentity(BrokerState& broker);
void refreshBrokerState(BrokerState& broker); void refreshBrokerState(BrokerState& broker);
void ensureBroker(BrokerState& broker); void ensureBroker(BrokerState& broker);
void destroyBroker(BrokerState& broker); void destroyBroker(BrokerState& broker, bool reset_retry_state = true);
bool refreshToken(BrokerState& broker); bool refreshToken(BrokerState& broker);
void publishStatus(bool online); void publishStatus(bool online);
void publishOnlineStatus(BrokerState& broker); void publishOnlineStatus(BrokerState& broker);