Skip to content
5 changes: 3 additions & 2 deletions src/helpers/MQTTMessageBuilder.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -57,13 +57,14 @@ int MQTTMessageBuilder::buildStatusMessage(
int internal_heap,
int packets_sent,
int packets_received,
const char* repeat
const char* repeat,
const MQTTConnHealth& conn_health
) {
return MQTTPayloadBuilder::buildStatusMessage(
doc, origin, origin_id, model, firmware_version, radio, client_version,
status, timestamp, buffer, buffer_size, battery_mv, uptime_secs, errors,
queue_len, noise_floor, tx_air_secs, rx_air_secs, recv_errors, internal_heap,
packets_sent, packets_received, repeat);
packets_sent, packets_received, repeat, conn_health);
}

int MQTTMessageBuilder::buildPacketMessage(
Expand Down
4 changes: 3 additions & 1 deletion src/helpers/MQTTMessageBuilder.h
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ class MQTTMessageBuilder {
* @param recv_errors Radio receive/CRC errors (optional, -1 to omit)
* @param internal_heap Internal heap free bytes (optional, -1 to omit)
* @param repeat Repeat/forwarding status ("on" or "off"); nullptr omits the field
* @param conn_health MQTT slot connection health; defaults omit every field
* @return Length of JSON string, or 0 on error
*/
static int buildStatusMessage(
Expand All @@ -93,7 +94,8 @@ class MQTTMessageBuilder {
int internal_heap = -1,
int packets_sent = -1,
int packets_received = -1,
const char* repeat = nullptr
const char* repeat = nullptr,
const MQTTConnHealth& conn_health = MQTTConnHealth()
);

/**
Expand Down
25 changes: 23 additions & 2 deletions src/helpers/MQTTPayloadBuilder.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ int MQTTPayloadBuilder::buildStatusMessage(
int internal_heap,
int packets_sent,
int packets_received,
const char* repeat
const char* repeat,
const MQTTConnHealth& conn_health
) {
doc.clear();
JsonObject root = doc.to<JsonObject>();
Expand All @@ -62,9 +63,13 @@ int MQTTPayloadBuilder::buildStatusMessage(
root["repeat"] = repeat;
}

const bool has_conn_health = conn_health.slots_up >= 0 || conn_health.slots_total >= 0 ||
conn_health.worst_outage_secs >= 0 || conn_health.heap_largest >= 0 ||
conn_health.connect_failures > 0 || conn_health.slots_breaker > 0;

if (battery_mv >= 0 || uptime_secs >= 0 || errors >= 0 || queue_len >= 0 ||
noise_floor > -999 || tx_air_secs >= 0 || rx_air_secs >= 0 || recv_errors >= 0 ||
internal_heap >= 0 || packets_sent >= 0 || packets_received >= 0) {
internal_heap >= 0 || packets_sent >= 0 || packets_received >= 0 || has_conn_health) {
JsonObject stats = root["stats"].to<JsonObject>();

if (battery_mv >= 0) stats["battery_mv"] = battery_mv;
Expand All @@ -78,6 +83,22 @@ int MQTTPayloadBuilder::buildStatusMessage(
if (rx_air_secs >= 0) stats["rx_air_secs"] = rx_air_secs;
if (recv_errors >= 0) stats["recv_errors"] = recv_errors;
if (internal_heap >= 0) stats["internal_heap"] = internal_heap;
if (conn_health.heap_largest >= 0) stats["heap_largest"] = conn_health.heap_largest;
if (conn_health.slots_up >= 0) stats["mqtt_slots_up"] = conn_health.slots_up;
if (conn_health.slots_total >= 0) stats["mqtt_slots_total"] = conn_health.slots_total;
// Emitted only during an outage of known duration, so the steady-state payload
// does not grow. Presence implies a slot is down, but absence does NOT imply
// health: a slot has no outage start time until its first attempt ends.
// `mqtt_slots_up < mqtt_slots_total` is the signal to alert on; this is detail.
if (conn_health.worst_outage_secs >= 0) {
stats["mqtt_outage_secs"] = conn_health.worst_outage_secs;
}
if (conn_health.connect_failures > 0) {
stats["mqtt_connect_failures"] = conn_health.connect_failures;
}
if (conn_health.slots_breaker > 0) {
stats["mqtt_slots_breaker"] = conn_health.slots_breaker;
}
}

return serializeComplete(root, buffer, buffer_size);
Expand Down
32 changes: 31 additions & 1 deletion src/helpers/MQTTPayloadBuilder.h
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,35 @@
#include <stddef.h>
#include <stdint.h>

// Connection health of the MQTT slots themselves, reported inside a status
// publication. A slot that cannot reach its broker cannot report its own outage,
// so this travels on every *healthy* slot's status message and is how a fleet
// learns that a sibling slot is down. Publish-success/error counters cannot show
// this: a slot that never connects never attempts a publish, so its error count
// stays at zero for the entire outage.
//
// `heap_largest` is the leading indicator. Observed on non-PSRAM hardware: a slot
// needs one contiguous 16,384-byte block for its TLS record buffers, and measured
// immediately before an attempt this value predicted the outcome every time. Free
// heap alone does not — an 85 KB-free board sat unable to connect for hours.
// A low value with all slots up is normal, not a fault: established sessions
// already hold their buffers. It means only that the next slot to ask will fail.
// Sampled at status-publish time, not at the attempt, so it is a trend, not a
// per-attempt reading.
struct MQTTConnHealth {
int slots_up = -1; // -1 = not supplied, field omitted
int slots_total = -1; // enabled, fully configured slots, connected or not
int worst_outage_secs = -1; // longest timed outage; -1 = none timed (see below)
int heap_largest = -1; // largest free internal block, bytes
// Attempts that ended without reaching onConnect, summed across all slots since
// boot. A slot stuck retrying moves this while its publish counters stay frozen,
// which is what separates "actively failing" from "idle and quiet".
int connect_failures = -1;
// Slots parked by the circuit breaker. Materially worse than a retrying slot:
// the breaker only probes every 30 minutes, so recovery is far slower.
int slots_breaker = -1;
};

// Mesh-independent JSON serialization core for MQTT publication payloads.
// MQTTMessageBuilder keeps the firmware-facing API and delegates these three
// deterministic contracts here so they can be exercised by native tests.
Expand Down Expand Up @@ -37,7 +66,8 @@ class MQTTPayloadBuilder {
int internal_heap = -1,
int packets_sent = -1,
int packets_received = -1,
const char* repeat = nullptr
const char* repeat = nullptr,
const MQTTConnHealth& conn_health = MQTTConnHealth()
);

static int buildPacketMessage(
Expand Down
102 changes: 97 additions & 5 deletions src/helpers/bridges/MQTTBridge.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -373,12 +373,26 @@ void MQTTBridge::formatMqttStatsReply(char* buf, size_t bufsize) {
if (b->_filtered_packets > 0) {
replyAppendf(buf, bufsize, &pos, " filt=%lu", b->_filtered_packets);
}
// down=<n>: ready slots that are not connected. Without this the reply cannot
// show an outage at all — sN= counts publishes, and a slot that never connects
// never publishes, so its ok freezes and its err stays 0 for the whole outage.
// Omitted while zero, like filt=, so a healthy node's reply keeps its length.
int down = 0;
for (int i = 0; i < RUNTIME_MQTT_SLOTS; i++) {
if (b->_slots[i].enabled && !b->_slots[i].connected && b->isSlotReady(i)) down++;
}
if (down > 0) {
replyAppendf(buf, bufsize, &pos, " down=%d", down);
}
replyAppendf(buf, bufsize, &pos, " |");
for (int i = 0; i < RUNTIME_MQTT_SLOTS; i++) {
if (!b->_slots[i].enabled || !b->_slots[i].client) continue;
replyAppendf(buf, bufsize, &pos, " s%d=%lu/%lu", i + 1,
// Trailing '!' marks a disconnected slot. Appended after the counts so the
// established "s(\d)=(\d+)/(\d+)" parsers keep matching unchanged.
replyAppendf(buf, bufsize, &pos, " s%d=%lu/%lu%s", i + 1,
b->_slots[i].client->getPublishOk(),
b->_slots[i].client->getPublishErr());
b->_slots[i].client->getPublishErr(),
b->_slots[i].connected ? "" : "!");
}
}

Expand Down Expand Up @@ -621,6 +635,14 @@ void MQTTBridge::formatSlotDiagReply(char* buf, size_t bufsize, int slot_index)
} else if (!slot.connected) {
replyAppendf(buf, bufsize, &pos, ", no error info");
}
// cf: attempts that never reached onConnect, plus starts that failed outright. A
// slot retrying into a heap that cannot satisfy the TLS allocation moves only this;
// the publish counters cannot move at all while the slot is down.
// After the error tail so a clamped reply loses this count, not the error detail.
const uint32_t cf = slot.connect_failures + slot.start_failures;
if (cf > 0) {
replyAppendf(buf, bufsize, &pos, ", cf:%lu", (unsigned long)cf);
}

// Appended last so it never displaces connection diagnostics. replyAppendf
// clamps rather than overflows, but a clipped type list is worse than no
Expand Down Expand Up @@ -752,6 +774,7 @@ MQTTBridge::MQTTBridge(NodePrefs *prefs, MQTTPrefs *obs, mesh::PacketManager *mg
_slot_reconfigure_pending[i] = false;
_slot_force_jwt_mint[i] = false;
_status_publish_pending[i] = false;
_slot_attempt_pending[i] = false;
}

// Reset CLI-requested forced NTP sync handshake (bridge object is reused across restarts)
Expand Down Expand Up @@ -1707,9 +1730,11 @@ bool MQTTBridge::ensureSlotClient(int index) {
MQTT_DEBUG_PRINTLN("MQTT%d ignoring late CONNECTED (state=%s, gen=%lu, enabled=%d)",
index + 1, clientStateName(st),
(unsigned long)_slots[index].generation, (int)_slots[index].enabled);
_slot_attempt_pending[index] = false; // still the end of that attempt
return;
}
MQTT_DEBUG_PRINTLN("MQTT%d connected", index + 1);
_slot_attempt_pending[index] = false;
_slots[index].client_state = ClientState::Connected;
_slots[index].connected = true;
_slot_force_jwt_mint[index] = false;
Expand Down Expand Up @@ -1743,6 +1768,13 @@ bool MQTTBridge::ensureSlotClient(int index) {
});
slot.client->onDisconnect([this, index](bool sessionPresent) {
MQTT_DEBUG_PRINTLN("MQTT%d disconnected", index + 1);
// Ended an attempt that never reached onConnect. A deliberate stop clears the flag
// first; a bounced live session is still marked connected until this handler runs,
// so its late DISCONNECTED after a timed-out softDisconnect() is not counted either.
if (_slot_attempt_pending[index] && !_slots[index].connected) {
_slot_attempt_pending[index] = false;
_slots[index].connect_failures++;
}
// Only a live client's disconnect is news. One arriving for a client we
// already stopped (or quarantined) must not resurrect its state.
if (clientStateIsLive(_slots[index].client_state)) {
Expand Down Expand Up @@ -2386,6 +2418,7 @@ void MQTTBridge::closeLiveClientForLinkTransition(int index) {
// the rest of the boot.
void MQTTBridge::stopSlotClient(int index) {
MQTTSlot& slot = _slots[index];
_slot_attempt_pending[index] = false;
const esp_err_t r = slot.client->disconnect();

// What the SDK's results actually mean here (mqtt_client.h documents
Expand Down Expand Up @@ -2480,7 +2513,10 @@ esp_err_t MQTTBridge::reconnectSlotClient(int index) {
}

esp_err_t r;
if (!slot.client->isStarted()) {
const bool was_pending = _slot_attempt_pending[index];
const bool starting = !slot.client->isStarted();
_slot_attempt_pending[index] = true;
if (starting) {
MQTT_DEBUG_PRINTLN("MQTT%d start (client was stopped)", index + 1);
r = slot.client->connect();
} else {
Expand All @@ -2492,6 +2528,13 @@ esp_err_t MQTTBridge::reconnectSlotClient(int index) {
// The attempt carries whatever credential is configured right now, so this
// is the point at which a freshly minted token becomes the one in use.
slot.applied_token_expires_at = slot.token_expires_at;
} else {
// A failed start leaves nothing in flight; a rejected reconnect leaves the prior one.
_slot_attempt_pending[index] = starting ? false : was_pending;
// A failed start (buffer/task allocation, config) means no attempt could begin.
// reconnect() on a started client returns ESP_FAIL while one is already in
// progress, which is not a failure (seen on hardware after a reconfigure).
if (starting) slot.start_failures++;
}
return r;
}
Expand Down Expand Up @@ -3135,7 +3178,8 @@ void MQTTBridge::publishStatusToSlot(int index) {
battery_mv, uptime_secs, errors, _queue_count, noise_floor,
tx_air_secs, rx_air_secs, recv_errors, internal_heap_free,
packets_sent, packets_received,
_prefs->disable_fwd ? "off" : "on"
_prefs->disable_fwd ? "off" : "on",
collectConnHealth()
);

if (len > 0) {
Expand All @@ -3151,6 +3195,53 @@ void MQTTBridge::publishStatusToSlot(int index) {
}
}

MQTTConnHealth MQTTBridge::collectConnHealth() const {
MQTTConnHealth health;
health.slots_up = 0;
health.slots_total = 0;

const unsigned long now = millis();
unsigned long worst_outage_ms = 0;
bool any_outage_timed = false;
uint32_t connect_failures = 0;
int breaker_slots = 0;

for (int i = 0; i < RUNTIME_MQTT_SLOTS; i++) {
const MQTTSlot& slot = _slots[i];
// Summed before the skip so disabling a slot never makes the total go backwards.
connect_failures += slot.connect_failures + slot.start_failures;
// Unconfigured slots and slots waiting on a token/IATA/credential are not outages.
if (!slot.enabled || !isSlotReady(i)) continue;
health.slots_total++;
if (slot.circuit_breaker_tripped) breaker_slots++;
if (slot.connected) {
health.slots_up++;
continue;
}
// The outage clock starts at a slot's first DISCONNECTED, so a slot that has not yet
// finished its first attempt has no start time. It still counts against slots_up,
// which is why that pair — not this duration — is the signal to alert on.
if (slot.current_outage_started_ms != 0) {
const unsigned long outage_ms = now - slot.current_outage_started_ms; // wrap-safe under 49.7 days
if (!any_outage_timed || outage_ms > worst_outage_ms) {
worst_outage_ms = outage_ms;
any_outage_timed = true;
}
}
}

if (any_outage_timed) {
health.worst_outage_secs = static_cast<int>(worst_outage_ms / 1000UL);
}
// Clamp: the payload field is int, and these are lifetime counters.
health.connect_failures = connect_failures > static_cast<uint32_t>(INT32_MAX)
? INT32_MAX
: static_cast<int>(connect_failures);
health.slots_breaker = breaker_slots;
health.heap_largest = static_cast<int>(heap_caps_get_largest_free_block(MALLOC_CAP_INTERNAL));
return health;
}

void MQTTBridge::updateCachedConnectionStatus() {
bool any_connected = false;
for (int i = 0; i < RUNTIME_MQTT_SLOTS; i++) {
Expand Down Expand Up @@ -3993,7 +4084,8 @@ bool MQTTBridge::publishStatus() {
battery_mv, uptime_secs, errors, _queue_count, noise_floor,
tx_air_secs, rx_air_secs, recv_errors, internal_heap_free,
packets_sent, packets_received,
_prefs->disable_fwd ? "off" : "on"
_prefs->disable_fwd ? "off" : "on",
collectConnHealth()
);

if (len > 0) {
Expand Down
31 changes: 28 additions & 3 deletions src/helpers/bridges/MQTTBridge.h
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
#include <Timezone.h>
#include "helpers/JWTHelper.h"
#include "helpers/MQTTPacketFilter.h"
#include "helpers/MQTTPayloadBuilder.h" // MQTTConnHealth
#include "helpers/MQTTPresets.h"
#include "helpers/MQTTPresetPolicy.h"
#include "helpers/MQTTLifecycle.h"
Expand Down Expand Up @@ -182,6 +183,14 @@ class MQTTBridge : public BridgeBase {
uint8_t last_connack_code;
unsigned long last_error_time; // millis() of last error
uint32_t disconnect_count; // Number of disconnect callbacks since boot
// Attempts that ended without ever reaching onConnect. Distinct from
// disconnect_count (which also counts dropped live sessions) and from the
// publish error counter, which cannot move at all while a slot is down.
uint32_t connect_failures;
// connect() calls on a stopped client that returned an error, so no attempt
// started and no event will arrive. Written only by the bridge task; reported
// summed with connect_failures.
uint32_t start_failures;
unsigned long first_disconnect_time; // millis() of first disconnect after boot

// Current-outage timer (used by AlertReporter to fire faults after a sustained
Expand Down Expand Up @@ -281,6 +290,11 @@ class MQTTBridge : public BridgeBase {
// _slot_reconfigure_pending; a single-byte volatile store/load is atomic.
volatile bool _status_publish_pending[RUNTIME_MQTT_SLOTS];

// Set by the bridge task before connect()/reconnect(), so a failure the event task
// delivers before that call returns is still seen; cleared by the attempt's outcome
// (including an ignored late CONNECTED) or a deliberate stop.
volatile bool _slot_attempt_pending[RUNTIME_MQTT_SLOTS];

// CLI-requested forced NTP sync, marshalled onto the MQTT task (Core 0).
// All NTP I/O must run on Core 0; the CLI thread
// (Core 1) sets _ntp_force_requested and blocks in requestForcedNtpSync()
Expand Down Expand Up @@ -365,9 +379,16 @@ class MQTTBridge : public BridgeBase {
// internal heap. On non-PSRAM: inline in the class object so the allocation doesn't
// interleave with large TLS buffers at startup.
static const size_t PUBLISH_JSON_BUFFER_SIZE = 2048;
// Status keeps its own smaller ceiling: raising it would change which oversized
// status documents get published instead of dropped.
static const size_t STATUS_JSON_BUFFER_SIZE = 768;
// Status keeps its own smaller ceiling: raising it changes which oversized status
// documents get published instead of dropped. 768 left only 6 bytes spare at the real
// field maxima (_origin[32] of quotes, which isValidName permits and JSON doubles, plus
// _device_id[65], _board_model[64], _firmware_version[64] and a 63-byte client version):
// measured 762/768 before the connection-health fields existed. Overflow is silent —
// serializeComplete() returns 0 and the bridge publishes only len>0, so status stops
// entirely. Raised to 1024: no extra RAM while status serializes into the shared scratch
// buffer, but +256 B of bridge-task stack in the fallback used if the PSRAM allocation
// failed. See the worst-case test in test_mqtt_payload_builder.
static const size_t STATUS_JSON_BUFFER_SIZE = 1024;
static_assert(STATUS_JSON_BUFFER_SIZE <= PUBLISH_JSON_BUFFER_SIZE,
"status payloads serialize into the shared publish buffer");
#if defined(BOARD_HAS_PSRAM)
Expand Down Expand Up @@ -542,6 +563,10 @@ class MQTTBridge : public BridgeBase {
bool publishToSlot(int index, const char* topic, const char* payload, size_t payload_len, bool retained = false, uint8_t qos = 0);
bool publishToAllSlots(const char* topic, const char* payload, size_t payload_len, bool retained = false, uint8_t qos = 0);
void publishStatusToSlot(int index);
// Snapshot of slot connection health for the status payload. A slot that cannot
// reach its broker cannot report its own outage, so every healthy slot carries
// this on its status publication.
MQTTConnHealth collectConnHealth() const;
void updateCachedConnectionStatus();

void processPacketQueue();
Expand Down
Loading
Loading