From ff2df2b7126ffb6edaf58721d0dca065168cef1d Mon Sep 17 00:00:00 2001 From: agessaman Date: Tue, 14 Apr 2026 18:07:02 -0700 Subject: [PATCH] Refactor MQTTMessageBuilder to accept external JsonDocument parameters for status and packet messages, reducing stack usage and improving memory management. Update MQTTBridge to allocate and reuse DynamicJsonDocument instances for JSON serialization, enhancing performance and efficiency during publish operations. --- src/helpers/MQTTMessageBuilder.cpp | 18 ++++++++---- src/helpers/MQTTMessageBuilder.h | 4 +++ src/helpers/bridges/MQTTBridge.cpp | 45 ++++++++++++++++++------------ src/helpers/bridges/MQTTBridge.h | 5 ++++ 4 files changed, 48 insertions(+), 24 deletions(-) diff --git a/src/helpers/MQTTMessageBuilder.cpp b/src/helpers/MQTTMessageBuilder.cpp index 8597c230..8332108f 100644 --- a/src/helpers/MQTTMessageBuilder.cpp +++ b/src/helpers/MQTTMessageBuilder.cpp @@ -5,6 +5,7 @@ #include "MeshCore.h" int MQTTMessageBuilder::buildStatusMessage( + JsonDocument& doc, const char* origin, const char* origin_id, const char* model, @@ -25,8 +26,9 @@ int MQTTMessageBuilder::buildStatusMessage( int recv_errors, int internal_heap ) { - // Use StaticJsonDocument to avoid heap fragmentation (fixed-size stack allocation) - StaticJsonDocument<768> doc; // Increased size to accommodate stats + // doc is provided by the caller (heap-allocated DynamicJsonDocument in MQTTBridge), + // keeping this 768-byte scratch space off the MQTT task stack. + doc.clear(); JsonObject root = doc.to(); root["status"] = status; @@ -78,6 +80,7 @@ int MQTTMessageBuilder::buildStatusMessage( } int MQTTMessageBuilder::buildPacketMessage( + JsonDocument& doc, const char* origin, const char* origin_id, const char* timestamp, @@ -96,10 +99,9 @@ int MQTTMessageBuilder::buildPacketMessage( char* buffer, size_t buffer_size ) { - // Use StaticJsonDocument with fixed maximum size to avoid heap fragmentation - // Base JSON overhead ~200 bytes, raw hex can be up to 510 chars (255 bytes packet) - // Use maximum size (2048) to handle all packet sizes without heap allocation - StaticJsonDocument<2048> doc; + // doc is provided by the caller (heap-allocated DynamicJsonDocument in MQTTBridge), + // keeping this 2048-byte scratch space off the MQTT task stack. + doc.clear(); JsonObject root = doc.to(); // Format numeric values as strings to avoid String object allocations @@ -165,6 +167,7 @@ int MQTTMessageBuilder::buildRawMessage( } int MQTTMessageBuilder::buildPacketJSON( + JsonDocument& doc, mesh::Packet* packet, bool is_tx, const char* origin, @@ -227,6 +230,7 @@ int MQTTMessageBuilder::buildPacketJSON( } return buildPacketMessage( + doc, origin, origin_id, timestamp, is_tx ? "tx" : "rx", time_str, date_str, @@ -243,6 +247,7 @@ int MQTTMessageBuilder::buildPacketJSON( } int MQTTMessageBuilder::buildPacketJSONFromRaw( + JsonDocument& doc, const uint8_t* raw_data, int raw_len, mesh::Packet* packet, @@ -309,6 +314,7 @@ int MQTTMessageBuilder::buildPacketJSONFromRaw( } return buildPacketMessage( + doc, origin, origin_id, timestamp, is_tx ? "tx" : "rx", time_str, date_str, diff --git a/src/helpers/MQTTMessageBuilder.h b/src/helpers/MQTTMessageBuilder.h index 96efabcd..4930bbaa 100644 --- a/src/helpers/MQTTMessageBuilder.h +++ b/src/helpers/MQTTMessageBuilder.h @@ -42,6 +42,7 @@ public: * @return Length of JSON string, or 0 on error */ static int buildStatusMessage( + JsonDocument& doc, const char* origin, const char* origin_id, const char* model, @@ -86,6 +87,7 @@ public: * @return Length of JSON string, or 0 on error */ static int buildPacketMessage( + JsonDocument& doc, const char* origin, const char* origin_id, const char* timestamp, @@ -137,6 +139,7 @@ public: * @return Length of JSON string, or 0 on error */ static int buildPacketJSON( + JsonDocument& doc, mesh::Packet* packet, bool is_tx, const char* origin, @@ -147,6 +150,7 @@ public: ); static int buildPacketJSONFromRaw( + JsonDocument& doc, const uint8_t* raw_data, int raw_len, mesh::Packet* packet, diff --git a/src/helpers/bridges/MQTTBridge.cpp b/src/helpers/bridges/MQTTBridge.cpp index 4eb0bd45..138c45d0 100644 --- a/src/helpers/bridges/MQTTBridge.cpp +++ b/src/helpers/bridges/MQTTBridge.cpp @@ -305,7 +305,8 @@ MQTTBridge::MQTTBridge(NodePrefs *prefs, mesh::PacketManager *mgr, mesh::RTCCloc #endif _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), - _last_slot_reconnect_ms(0) + _last_slot_reconnect_ms(0), + _packet_json_doc(nullptr), _status_json_doc(nullptr) #ifdef ESP_PLATFORM , _packet_queue_handle(nullptr), _mqtt_task_handle(nullptr), _mqtt_task_stack(nullptr), _packet_queue_storage(nullptr) @@ -368,6 +369,11 @@ MQTTBridge::MQTTBridge(NodePrefs *prefs, mesh::PacketManager *mgr, mesh::RTCCloc // Pre-allocate JSON publish buffer (reused for all publishes to avoid alloc/free churn) _publish_json_buffer = (char*)psram_malloc(PUBLISH_JSON_BUFFER_SIZE); _status_json_buffer = (char*)psram_malloc(STATUS_JSON_BUFFER_SIZE); + + // Allocate JSON document scratch space once; reused via doc.clear() on every publish. + // Keeps the large StaticJsonDocument off the 8 KB MQTT task stack. + _packet_json_doc = new DynamicJsonDocument(PUBLISH_JSON_BUFFER_SIZE); + _status_json_doc = new DynamicJsonDocument(STATUS_JSON_BUFFER_SIZE); } // --------------------------------------------------------------------------- @@ -646,6 +652,10 @@ void MQTTBridge::end() { psram_free(_status_json_buffer); _status_json_buffer = nullptr; + // Free JSON document scratch space + delete _packet_json_doc; _packet_json_doc = nullptr; + delete _status_json_doc; _status_json_doc = nullptr; + _initialized = false; _slots_setup_done = false; // Reset so deferred setup runs again on next begin() MQTT_DEBUG_PRINTLN("MQTT Bridge stopped"); @@ -916,21 +926,12 @@ void MQTTBridge::mqttTaskLoop() { last_slot_status_update = now; } - // Adaptive task delay based on work done - bool has_work = (_queue_count > 0); - if (!has_work && _status_enabled) { - if (_last_status_publish == 0 || - (now - _last_status_publish >= (_status_interval - 10000))) { - has_work = true; - } - } - - // Adaptive delay: shorter when work pending, longer when idle - if (has_work) { - vTaskDelay(pdMS_TO_TICKS(5)); - } else { - vTaskDelay(pdMS_TO_TICKS(50)); - } + // Adaptive delay: 5 ms when packets are queued, 50 ms when idle. + // The previous "status approaching" check (widening to 5 ms for 10 s before each status + // publish) caused 2 000 unnecessary wakeups per interval; the 50 ms idle tick catches + // the status deadline with at most 50 ms of extra latency, which is irrelevant at a + // 5-minute interval. + vTaskDelay(pdMS_TO_TICKS(_queue_count > 0 ? 5 : 50)); } } #endif @@ -1657,6 +1658,7 @@ void MQTTBridge::publishStatusToSlot(int index) { int internal_heap_free = (int)heap_caps_get_free_size(MALLOC_CAP_INTERNAL); int len = MQTTMessageBuilder::buildStatusMessage( + *_status_json_doc, _origin, origin_id, _board_model, _firmware_version, radio_info, client_version, "online", timestamp, json_buffer, STATUS_JSON_BUFFER_SIZE, battery_mv, uptime_secs, errors, _queue_count, noise_floor, @@ -1809,7 +1811,7 @@ bool MQTTBridge::handleWiFiConnection(unsigned long now) { } else if (ps_pref == 2) { ps_mode = WIFI_PS_MAX_MODEM; } else { - ps_mode = WIFI_PS_MIN_MODEM; + ps_mode = WIFI_PS_NONE; // default: no power save; eliminates DTIM wake latency on mains-powered bridges } esp_wifi_set_ps(ps_mode); #ifdef MQTT_WIFI_TX_POWER @@ -2147,6 +2149,7 @@ void MQTTBridge::processPacketQueue() { queued.has_raw_data ? queued.raw_data : nullptr, queued.has_raw_data ? queued.raw_len : 0, queued.snr, queued.rssi); + taskYIELD(); // allow higher-priority tasks to run between packet publishes // Publish raw if enabled if (_raw_enabled) { @@ -2203,6 +2206,7 @@ void MQTTBridge::processPacketQueue() { queued.has_raw_data ? queued.raw_data : nullptr, queued.has_raw_data ? queued.raw_len : 0, queued.snr, queued.rssi); + // No taskYIELD() on non-ESP32 platforms (non-FreeRTOS, cooperative scheduling not needed) if (_raw_enabled) { publishRaw(&queued.packet_copy); @@ -2273,6 +2277,7 @@ bool MQTTBridge::publishStatus() { int internal_heap_free = (int)heap_caps_get_free_size(MALLOC_CAP_INTERNAL); int len = MQTTMessageBuilder::buildStatusMessage( + *_status_json_doc, _origin, origin_id, _board_model, _firmware_version, radio_info, client_version, "online", timestamp, json_buffer, STATUS_JSON_BUFFER_SIZE, battery_mv, uptime_secs, errors, _queue_count, noise_floor, @@ -2354,11 +2359,13 @@ void MQTTBridge::publishPacket(mesh::Packet* packet, bool is_tx, int len; if (raw_data && raw_len > 0) { len = MQTTMessageBuilder::buildPacketJSONFromRaw( + *_packet_json_doc, raw_data, raw_len, packet, is_tx, _origin, origin_id, snr, rssi, _timezone, active_buffer, active_buffer_size ); } else if (_last_raw_data && _last_raw_len > 0 && (millis() - _last_raw_timestamp) < 1000) { len = MQTTMessageBuilder::buildPacketJSONFromRaw( + *_packet_json_doc, _last_raw_data, _last_raw_len, packet, is_tx, _origin, origin_id, _last_snr, _last_rssi, _timezone, active_buffer, active_buffer_size ); @@ -2370,11 +2377,13 @@ void MQTTBridge::publishPacket(mesh::Packet* packet, bool is_tx, uint8_t rlen = packet->writeTo(reconstructed); if (rlen > 0) { len = MQTTMessageBuilder::buildPacketJSONFromRaw( + *_packet_json_doc, reconstructed, rlen, packet, is_tx, _origin, origin_id, snr, rssi, _timezone, active_buffer, active_buffer_size ); } else { len = MQTTMessageBuilder::buildPacketJSON( + *_packet_json_doc, packet, is_tx, _origin, origin_id, _timezone, active_buffer, active_buffer_size ); } @@ -2544,7 +2553,7 @@ void MQTTBridge::storeRawRadioData(const uint8_t* raw_data, int len, float snr, _staged_snr = snr; _staged_rssi = rssi; _staged_raw_valid = true; - MQTT_DEBUG_PRINTLN("Staged raw radio data: %d bytes, SNR=%.1f, RSSI=%.1f", len, snr, rssi); + MQTT_DEBUG_PRINTLN("Stored raw radio data: %d bytes, SNR=%.1f, RSSI=%.1f", len, snr, rssi); } } diff --git a/src/helpers/bridges/MQTTBridge.h b/src/helpers/bridges/MQTTBridge.h index 5ca722a5..eb0fb065 100644 --- a/src/helpers/bridges/MQTTBridge.h +++ b/src/helpers/bridges/MQTTBridge.h @@ -196,6 +196,11 @@ private: static const size_t STATUS_JSON_BUFFER_SIZE = 768; char* _status_json_buffer; + // JSON document scratch space — allocated once on the heap, reused via doc.clear(). + // Keeps StaticJsonDocument off the 8 KB MQTT task stack. + DynamicJsonDocument* _packet_json_doc; // capacity = PUBLISH_JSON_BUFFER_SIZE + DynamicJsonDocument* _status_json_doc; // capacity = STATUS_JSON_BUFFER_SIZE + // Memory pressure monitoring unsigned long _last_memory_check; int _skipped_publishes; // Count of skipped publishes due to memory pressure