From bebde37fb75f7250d99d6bf687c5ab4e84253a4a Mon Sep 17 00:00:00 2001 From: Adam Gessaman Date: Thu, 5 Feb 2026 21:45:23 -0800 Subject: [PATCH] Refactor MQTTBridge memory management with unified critical check - Introduced a new method `runCriticalMemoryCheckAndRecovery()` to streamline memory checks and recovery processes. - Added member variables to track fragmentation pressure and the last critical check run time. - Removed redundant critical memory check logic from `mqttTaskLoop()` and `loop()`, enhancing code clarity and maintainability. - Improved logging and pressure tracking for better diagnostics during memory management. --- src/helpers/bridges/MQTTBridge.cpp | 132 +++++++++++++++-------------- src/helpers/bridges/MQTTBridge.h | 5 ++ 2 files changed, 72 insertions(+), 65 deletions(-) diff --git a/src/helpers/bridges/MQTTBridge.cpp b/src/helpers/bridges/MQTTBridge.cpp index 2a44e22e..184f7a3d 100644 --- a/src/helpers/bridges/MQTTBridge.cpp +++ b/src/helpers/bridges/MQTTBridge.cpp @@ -140,6 +140,7 @@ MQTTBridge::MQTTBridge(NodePrefs *prefs, mesh::PacketManager *mgr, mesh::RTCCloc _analyzer_us_client(nullptr), _analyzer_eu_client(nullptr), _config_valid(false), _cached_has_brokers(false), _cached_has_analyzer_servers(false), _last_memory_check(0), _skipped_publishes(0), _last_fragmentation_recovery(0), + _fragmentation_pressure_since(0), _last_critical_check_run(0), _last_no_broker_log(0), _last_config_warning(0), _dispatcher(nullptr), _radio(nullptr), _board(nullptr), _ms(nullptr), _last_wifi_check(0), _last_wifi_status(WL_DISCONNECTED), _wifi_status_initialized(false), _wifi_disconnected_time(0), _last_wifi_reconnect_attempt(0), _wifi_reconnect_backoff_attempt(0), @@ -733,10 +734,11 @@ void MQTTBridge::mqttTaskLoop() { _last_status_publish = now; _last_status_retry = 0; MQTT_DEBUG_PRINTLN("Status published successfully, next publish in %lu ms", _status_interval); - // If we're in the hole but just proved connectivity, recover sooner than the 15-min critical check + // If we're in the hole but just proved connectivity, recover sooner than the dedicated pressure timer size_t max_alloc = ESP.getMaxAllocHeap(); if (max_alloc < 58000 && (now - _last_fragmentation_recovery) > 300000) { _last_fragmentation_recovery = now; + _fragmentation_pressure_since = 0; // Reset pressure timer so dedicated check doesn't fire again soon MQTT_DEBUG_PRINTLN("Fragmentation recovery after status (max_alloc=%d)", (int)max_alloc); recreateMqttClientsForFragmentationRecovery(); } @@ -748,40 +750,8 @@ void MQTTBridge::mqttTaskLoop() { } } - // Critical memory check (every 15 minutes) - log warnings and run fragmentation recovery on ESP32 - static unsigned long last_critical_check = 0; - if (now - last_critical_check > 900000) { - size_t free_h = ESP.getFreeHeap(); - size_t max_alloc = ESP.getMaxAllocHeap(); - #ifdef MQTT_MEMORY_DEBUG - // #region agent log - unsigned long internal_f = heap_caps_get_free_size(MALLOC_CAP_INTERNAL); - unsigned long spiram_f = 0; - #ifdef BOARD_HAS_PSRAM - spiram_f = heap_caps_get_free_size(MALLOC_CAP_SPIRAM); - #endif - agentLogHeap("MQTTBridge.cpp:716", "critical_memory_check", "H1_H4", free_h, max_alloc, internal_f, spiram_f); - // #endregion - #endif - if (max_alloc < 40000) { - MQTT_DEBUG_PRINTLN("CRITICAL: Low memory! Free: %d, Max: %d", (int)free_h, (int)max_alloc); - } else if (max_alloc < 60000) { - MQTT_DEBUG_PRINTLN("WARNING: Memory pressure. Free: %d, Max: %d", (int)free_h, (int)max_alloc); - } - // Log active MQTT client count (we never have more than 3: main + US analyzer + EU analyzer) - int n_main = (_mqtt_client != nullptr) ? 1 : 0; - int n_us = (_analyzer_us_client != nullptr) ? 1 : 0; - int n_eu = (_analyzer_eu_client != nullptr) ? 1 : 0; - MQTT_DEBUG_PRINTLN("MQTT clients active: %d (main=%d us=%d eu=%d)", n_main + n_us + n_eu, n_main, n_us, n_eu); - // Proactive fragmentation recovery: recreate MQTT clients to recover max_alloc. Throttle once per 5 min. - if (max_alloc < 58000 && (now - _last_fragmentation_recovery) > 300000) { - _last_fragmentation_recovery = now; - MQTT_DEBUG_PRINTLN("Fragmentation recovery: recreating MQTT clients (max_alloc=%d)", (int)max_alloc); - recreateMqttClientsForFragmentationRecovery(); - } - last_critical_check = now; - } - + runCriticalMemoryCheckAndRecovery(); + // Update cached analyzer server status periodically (every 5 seconds) // This ensures cache stays accurate even if callbacks miss updates static unsigned long last_analyzer_status_update = 0; @@ -1073,36 +1043,7 @@ void MQTTBridge::loop() { #endif #ifdef ESP_PLATFORM - // Critical memory check (every 15 minutes) - log warnings and optionally recover from fragmentation - static unsigned long last_critical_check = 0; - unsigned long now_crit = millis(); - if (now_crit - last_critical_check > 900000) { - size_t free_h = ESP.getFreeHeap(); - size_t max_alloc = ESP.getMaxAllocHeap(); - #ifdef MQTT_MEMORY_DEBUG - // #region agent log - unsigned long internal_f = heap_caps_get_free_size(MALLOC_CAP_INTERNAL); - unsigned long spiram_f = 0; - #ifdef BOARD_HAS_PSRAM - spiram_f = heap_caps_get_free_size(MALLOC_CAP_SPIRAM); - #endif - agentLogHeap("MQTTBridge.cpp:1015", "critical_memory_check", "H1_H4", free_h, max_alloc, internal_f, spiram_f); - // #endregion - #endif - if (max_alloc < 40000) { - MQTT_DEBUG_PRINTLN("CRITICAL: Low memory! Free: %d, Max: %d", (int)free_h, (int)max_alloc); - } else if (max_alloc < 60000) { - MQTT_DEBUG_PRINTLN("WARNING: Memory pressure. Free: %d, Max: %d", (int)free_h, (int)max_alloc); - } - // Proactive fragmentation recovery: if max_alloc collapsed (e.g. after poor WiFi then reconnect), - // recreate MQTT clients so they allocate fresh blocks. Throttle to once per 5 minutes. - if (max_alloc < 58000 && (now_crit - _last_fragmentation_recovery) > 300000) { - _last_fragmentation_recovery = now_crit; - MQTT_DEBUG_PRINTLN("Fragmentation recovery: recreating MQTT clients (max_alloc=%d)", (int)max_alloc); - recreateMqttClientsForFragmentationRecovery(); - } - last_critical_check = now_crit; - } + runCriticalMemoryCheckAndRecovery(); #endif } @@ -1190,6 +1131,67 @@ void MQTTBridge::ensureMainMqttClient() { MQTT_DEBUG_PRINTLN("Main MQTT client recreated (fresh buffers)"); } +#ifdef ESP_PLATFORM +void MQTTBridge::runCriticalMemoryCheckAndRecovery() { + const unsigned long CRITICAL_CHECK_INTERVAL_MS = 60000; // Sample heap at most every 60s + const unsigned long PRESSURE_WINDOW_MS = 180000; // Recover if under pressure for 3 min + const unsigned long RECOVERY_THROTTLE_MS = 300000; // 5 min between recovery runs + const unsigned long CRITICAL_LOG_INTERVAL_MS = 900000; // Log CRITICAL/WARNING/client count at most every 15 min + const size_t PRESSURE_THRESHOLD = 58000; // max_alloc below this = under pressure + + unsigned long now = millis(); + if (now - _last_critical_check_run < CRITICAL_CHECK_INTERVAL_MS) { + return; + } + _last_critical_check_run = now; + + size_t free_h = ESP.getFreeHeap(); + size_t max_alloc = ESP.getMaxAllocHeap(); + #ifdef MQTT_MEMORY_DEBUG + unsigned long internal_f = heap_caps_get_free_size(MALLOC_CAP_INTERNAL); + unsigned long spiram_f = 0; + #ifdef BOARD_HAS_PSRAM + spiram_f = heap_caps_get_free_size(MALLOC_CAP_SPIRAM); + #endif + agentLogHeap("MQTTBridge.cpp:runCriticalMemoryCheckAndRecovery", "critical_memory_check", "H1_H4", free_h, max_alloc, internal_f, spiram_f); + #endif + + // Pressure timer: track how long max_alloc has been below threshold + if (max_alloc >= PRESSURE_THRESHOLD) { + _fragmentation_pressure_since = 0; + } else { + if (_fragmentation_pressure_since == 0) { + _fragmentation_pressure_since = now; + } + } + + // Rate-limited diagnostic logging (every 15 min) + static unsigned long last_critical_log = 0; + if (now - last_critical_log >= CRITICAL_LOG_INTERVAL_MS) { + last_critical_log = now; + if (max_alloc < 40000) { + MQTT_DEBUG_PRINTLN("CRITICAL: Low memory! Free: %d, Max: %d", (int)free_h, (int)max_alloc); + } else if (max_alloc < 60000) { + MQTT_DEBUG_PRINTLN("WARNING: Memory pressure. Free: %d, Max: %d", (int)free_h, (int)max_alloc); + } + int n_main = (_mqtt_client != nullptr) ? 1 : 0; + int n_us = (_analyzer_us_client != nullptr) ? 1 : 0; + int n_eu = (_analyzer_eu_client != nullptr) ? 1 : 0; + MQTT_DEBUG_PRINTLN("MQTT clients active: %d (main=%d us=%d eu=%d)", n_main + n_us + n_eu, n_main, n_us, n_eu); + } + + // Dedicated recovery timer: recover if under pressure for PRESSURE_WINDOW_MS and throttle passed + if (_fragmentation_pressure_since != 0 && + (now - _fragmentation_pressure_since) >= PRESSURE_WINDOW_MS && + (now - _last_fragmentation_recovery) >= RECOVERY_THROTTLE_MS) { + _last_fragmentation_recovery = now; + _fragmentation_pressure_since = 0; + MQTT_DEBUG_PRINTLN("Fragmentation recovery: recreating MQTT clients (max_alloc=%d, pressure 3 min)", (int)max_alloc); + recreateMqttClientsForFragmentationRecovery(); + } +} +#endif + void MQTTBridge::recreateMqttClientsForFragmentationRecovery() { // Disconnect, delete, and recreate all MQTT clients so they allocate fresh buffers. // This can recover max_alloc when the internal heap was fragmented (e.g. after poor diff --git a/src/helpers/bridges/MQTTBridge.h b/src/helpers/bridges/MQTTBridge.h index 2657fd3b..0afc2217 100644 --- a/src/helpers/bridges/MQTTBridge.h +++ b/src/helpers/bridges/MQTTBridge.h @@ -158,6 +158,8 @@ private: unsigned long _last_memory_check; int _skipped_publishes; // Count of skipped publishes due to memory pressure unsigned long _last_fragmentation_recovery; // Throttle: 5 min between recovery runs (task + loop) + unsigned long _fragmentation_pressure_since; // 0 = not under pressure; else first time max_alloc < threshold + unsigned long _last_critical_check_run; // Throttle: run unified check at most every 60s unsigned long _last_token_renewal_attempt_us; unsigned long _last_token_renewal_attempt_eu; unsigned long _last_reconnect_attempt_us; @@ -213,6 +215,9 @@ private: // Internal methods void ensureMainMqttClient(); // Create main MQTT client if _config_valid and _mqtt_client is null (e.g. after reinit) + #ifdef ESP_PLATFORM + void runCriticalMemoryCheckAndRecovery(); // Unified heap check, pressure timer, optional recovery + #endif void recreateMqttClientsForFragmentationRecovery(); // Disconnect, delete, recreate all MQTT clients to recover max_alloc void connectToBrokers(); void processPacketQueue();