MQTTUplink.cpp 35 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030
  1. #include "MQTTUplink.h"
  2. #include "MQTTCaCerts.h"
  3. #ifdef WITH_MQTT_UPLINK
  4. #if defined(ESP_PLATFORM)
  5. #include <helpers/ESP32Board.h>
  6. #include <WiFi.h>
  7. #include <WiFiClientSecure.h>
  8. #include <esp_idf_version.h>
  9. #include <esp_heap_caps.h>
  10. #include <esp_system.h>
  11. #include <helpers/TxtDataHelpers.h>
  12. #include <ctype.h>
  13. #include <string.h>
  14. #include <time.h>
  15. #ifndef FIRMWARE_VERSION
  16. #define FIRMWARE_VERSION "v1.14.1"
  17. #endif
  18. #ifndef FIRMWARE_BUILD_DATE
  19. #define FIRMWARE_BUILD_DATE "20 Mar 2026"
  20. #endif
  21. #ifndef CLIENT_VERSION
  22. #define CLIENT_VERSION "eastmesh-repeater-mqtt"
  23. #endif
  24. #ifndef MQTT_DEBUG
  25. #define MQTT_DEBUG 0
  26. #endif
  27. #if MQTT_DEBUG
  28. #define LOG_CAT(tag, fmt, ...) Serial.printf("[" tag "] " fmt "\n", ##__VA_ARGS__)
  29. #define MQTT_LOG(fmt, ...) LOG_CAT("MQTT", fmt, ##__VA_ARGS__)
  30. #else
  31. #define LOG_CAT(...) do { } while (0)
  32. #define MQTT_LOG(...) do { } while (0)
  33. #endif
  34. namespace {
  35. constexpr unsigned long kBrokerRetryBaseMillis = 10000;
  36. constexpr unsigned long kBrokerRetryMaxMillis = 300000;
  37. constexpr size_t kBrokerTokenSize = 640;
  38. constexpr time_t kTokenLifetimeSecs = 3600;
  39. constexpr time_t kTokenRefreshSlackSecs = 300;
  40. constexpr time_t kMinSaneEpoch = 1735689600; // 2025-01-01T00:00:00Z
  41. unsigned long getBrokerRetryDelayMillis(uint8_t failures) {
  42. unsigned long delay_ms = kBrokerRetryBaseMillis;
  43. if (failures > 0) {
  44. uint8_t shifts = min<uint8_t>(failures - 1, 5);
  45. delay_ms <<= shifts;
  46. }
  47. if (delay_ms > kBrokerRetryMaxMillis) {
  48. delay_ms = kBrokerRetryMaxMillis;
  49. }
  50. return delay_ms;
  51. }
  52. char* allocScratchBuffer(size_t size) {
  53. void* ptr = heap_caps_malloc(size, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT);
  54. if (ptr == nullptr) {
  55. ptr = heap_caps_malloc(size, MALLOC_CAP_8BIT);
  56. }
  57. return static_cast<char*>(ptr);
  58. }
  59. void freeScratchBuffer(void* ptr) {
  60. if (ptr != nullptr) {
  61. heap_caps_free(ptr);
  62. }
  63. }
  64. #if MQTT_DEBUG
  65. void logMqttMemorySnapshot(const char* phase, const char* broker_label = nullptr) {
  66. MQTT_LOG("mem phase=%s broker=%s uptime_ms=%lu heap_free=%u heap_min=%u heap_max=%u psram_free=%u psram_min=%u "
  67. "psram_max=%u",
  68. phase != nullptr ? phase : "-",
  69. broker_label != nullptr ? broker_label : "-",
  70. millis(),
  71. ESP.getFreeHeap(),
  72. ESP.getMinFreeHeap(),
  73. ESP.getMaxAllocHeap(),
  74. ESP.getFreePsram(),
  75. ESP.getMinFreePsram(),
  76. ESP.getMaxAllocPsram());
  77. }
  78. #else
  79. void logMqttMemorySnapshot(const char*, const char* = nullptr) {
  80. }
  81. #endif
  82. }
  83. const MQTTUplink::BrokerSpec MQTTUplink::kBrokerSpecs[3] = {
  84. {"eastmesh-au", "eastmesh-au", "mqtt2.eastmesh.au", "wss://mqtt2.eastmesh.au:443/mqtt", kEastmeshBit},
  85. {"letsmesh-eu", "letsmesh-eu", "mqtt-eu-v1.letsmesh.net", "wss://mqtt-eu-v1.letsmesh.net:443/mqtt",
  86. kLetsmeshEuBit},
  87. {"letsmesh-us", "letsmesh-us", "mqtt-us-v1.letsmesh.net", "wss://mqtt-us-v1.letsmesh.net:443/mqtt",
  88. kLetsmeshUsBit},
  89. };
  90. MQTTUplink::MQTTUplink(mesh::RTCClock& rtc, mesh::LocalIdentity& identity)
  91. : _fs(nullptr), _rtc(&rtc), _identity(&identity), _running(false), _last_status_publish(0), _last_status{},
  92. _node_name(nullptr), _network(nullptr)
  93. {
  94. memset(_device_id, 0, sizeof(_device_id));
  95. MQTTPrefsStore::setDefaults(_prefs);
  96. for (size_t i = 0; i < 3; ++i) {
  97. memset(&_brokers[i], 0, sizeof(_brokers[i]));
  98. _brokers[i].spec = &kBrokerSpecs[i];
  99. }
  100. MQTT_LOG("uplink init");
  101. }
  102. const char* MQTTUplink::getClientVersion() const {
  103. return CLIENT_VERSION;
  104. }
  105. bool MQTTUplink::savePrefs() {
  106. return MQTTPrefsStore::save(_fs, _prefs);
  107. }
  108. bool MQTTUplink::hasEnabledBroker() const {
  109. return (_prefs.enabled_mask & 0x07) != 0;
  110. }
  111. bool MQTTUplink::isUnsetIataValue(const char* iata) {
  112. return iata == nullptr || iata[0] == 0 || strcmp(iata, MQTT_UNSET_IATA) == 0;
  113. }
  114. uint8_t MQTTUplink::normalizeEnabledMask(uint8_t mask) {
  115. uint8_t normalized = 0;
  116. uint8_t count = 0;
  117. for (const BrokerSpec& spec : kBrokerSpecs) {
  118. if ((mask & spec.bit) != 0) {
  119. if (count >= kMaxEnabledBrokers) {
  120. break;
  121. }
  122. normalized |= spec.bit;
  123. ++count;
  124. }
  125. }
  126. return normalized;
  127. }
  128. void MQTTUplink::formatTopic(char* dst, size_t dst_size, const char* leaf) const {
  129. if (dst == nullptr || dst_size == 0) {
  130. return;
  131. }
  132. snprintf(dst, dst_size, "meshcore/%s/%s/%s", _prefs.iata, _device_id, leaf);
  133. }
  134. bool MQTTUplink::isActive() const {
  135. return _running && hasEnabledBroker();
  136. }
  137. bool MQTTUplink::sendStatusNow() {
  138. if (!_running || _network == nullptr || !_network->hasTimeSync() || !_network->isWifiConnected()) {
  139. return false;
  140. }
  141. bool any_connected = false;
  142. for (const BrokerState& broker : _brokers) {
  143. if (broker.spec != nullptr && broker.connected && broker.client != nullptr) {
  144. any_connected = true;
  145. break;
  146. }
  147. }
  148. if (!any_connected) {
  149. return false;
  150. }
  151. publishStatus(true);
  152. _last_status_publish = millis();
  153. return true;
  154. }
  155. void MQTTUplink::makeSafeToken(const char* input, char* output, size_t output_size) {
  156. if (output_size == 0) {
  157. return;
  158. }
  159. size_t oi = 0;
  160. for (size_t i = 0; input != nullptr && input[i] != 0 && oi + 1 < output_size; ++i) {
  161. char c = input[i];
  162. if ((c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z') || (c >= '0' && c <= '9') || c == '-' || c == '_') {
  163. output[oi++] = c;
  164. } else {
  165. output[oi++] = '_';
  166. }
  167. }
  168. output[oi] = 0;
  169. }
  170. void MQTTUplink::bytesToHexUpper(const uint8_t* src, size_t len, char* dst, size_t dst_size) {
  171. if (dst_size == 0) {
  172. return;
  173. }
  174. size_t di = 0;
  175. for (size_t i = 0; i < len && di + 2 < dst_size; ++i) {
  176. snprintf(&dst[di], dst_size - di, "%02X", src[i]);
  177. di += 2;
  178. }
  179. dst[min(di, dst_size - 1)] = 0;
  180. }
  181. void MQTTUplink::formatIsoTimestamp(time_t ts, char* dst, size_t dst_size) {
  182. if (dst_size == 0) {
  183. return;
  184. }
  185. struct tm tm_local;
  186. localtime_r(&ts, &tm_local);
  187. strftime(dst, dst_size, "%Y-%m-%dT%H:%M:%S", &tm_local);
  188. size_t len = strlen(dst);
  189. if (len + 8 < dst_size) {
  190. memcpy(&dst[len], ".000000", 8);
  191. }
  192. }
  193. void MQTTUplink::escapeJsonString(const char* input, char* output, size_t output_size) {
  194. if (output_size == 0) {
  195. return;
  196. }
  197. size_t oi = 0;
  198. for (size_t i = 0; input != nullptr && input[i] != 0 && oi + 1 < output_size; ++i) {
  199. char c = input[i];
  200. const char* escape = nullptr;
  201. switch (c) {
  202. case '\"':
  203. escape = "\\\"";
  204. break;
  205. case '\\':
  206. escape = "\\\\";
  207. break;
  208. case '\n':
  209. escape = "\\n";
  210. break;
  211. case '\r':
  212. escape = "\\r";
  213. break;
  214. case '\t':
  215. escape = "\\t";
  216. break;
  217. default:
  218. break;
  219. }
  220. if (escape != nullptr) {
  221. for (size_t j = 0; escape[j] != 0 && oi + 1 < output_size; ++j) {
  222. output[oi++] = escape[j];
  223. }
  224. continue;
  225. }
  226. unsigned char uc = static_cast<unsigned char>(c);
  227. if (uc < 0x20) {
  228. output[oi++] = '_';
  229. continue;
  230. }
  231. output[oi++] = c;
  232. }
  233. output[oi] = 0;
  234. }
  235. void MQTTUplink::refreshIdentityStrings() {
  236. bytesToHexUpper(_identity->pub_key, PUB_KEY_SIZE, _device_id, sizeof(_device_id));
  237. for (BrokerState& broker : _brokers) {
  238. refreshBrokerIdentity(broker);
  239. }
  240. }
  241. void MQTTUplink::refreshBrokerIdentity(BrokerState& broker) {
  242. if (broker.spec == nullptr) {
  243. return;
  244. }
  245. snprintf(broker.username, sizeof(broker.username), "v1_%s", _device_id);
  246. snprintf(broker.client_id, sizeof(broker.client_id), "mqtt_%s-%.6s", broker.spec->key, _device_id);
  247. formatTopic(broker.status_topic, sizeof(broker.status_topic), "status");
  248. }
  249. void MQTTUplink::refreshBrokerState(BrokerState& broker) {
  250. char safe_name[40];
  251. makeSafeToken(board.getManufacturerName(), safe_name, sizeof(safe_name));
  252. char origin[80];
  253. const char* node_name = (_node_name != nullptr && _node_name[0] != 0) ? _node_name : _device_id;
  254. escapeJsonString(node_name, origin, sizeof(origin));
  255. char client_version[96];
  256. snprintf(client_version, sizeof(client_version), "%s", CLIENT_VERSION);
  257. char radio_info[48];
  258. snprintf(radio_info, sizeof(radio_info), "%.6f,%.1f,%u,%u", static_cast<double>(_last_status.radio_freq),
  259. static_cast<double>(_last_status.radio_bw), _last_status.radio_sf, _last_status.radio_cr);
  260. char ts[32];
  261. formatIsoTimestamp(time(nullptr), ts, sizeof(ts));
  262. snprintf(broker.offline_payload, sizeof(broker.offline_payload),
  263. "{\"status\":\"offline\",\"timestamp\":\"%s\",\"origin\":\"%s\",\"origin_id\":\"%s\",\"model\":\"%s\","
  264. "\"firmware_version\":\"%s\",\"radio\":\"%s\",\"client_version\":\"%s\"}",
  265. ts, origin, _device_id, safe_name, FIRMWARE_VERSION, radio_info, client_version);
  266. }
  267. bool MQTTUplink::refreshToken(BrokerState& broker) {
  268. time_t now = time(nullptr);
  269. if (now < kMinSaneEpoch) {
  270. MQTT_LOG("%s token skipped: clock not ready (%lu)", broker.spec->label, static_cast<unsigned long>(now));
  271. return false;
  272. }
  273. if (broker.token == nullptr) {
  274. broker.token = allocScratchBuffer(kBrokerTokenSize);
  275. if (broker.token == nullptr) {
  276. MQTT_LOG("%s token alloc failed", broker.spec->label);
  277. return false;
  278. }
  279. }
  280. time_t expires_at = now + kTokenLifetimeSecs;
  281. const char* owner = _prefs.owner_public_key[0] ? _prefs.owner_public_key : nullptr;
  282. const char* email = _prefs.owner_email[0] ? _prefs.owner_email : nullptr;
  283. if (!JWTHelper::createAuthToken(*_identity, broker.spec->host, now, expires_at, broker.token, kBrokerTokenSize,
  284. owner, email)) {
  285. freeScratchBuffer(broker.token);
  286. broker.token = nullptr;
  287. MQTT_LOG("%s token creation failed", broker.spec->label);
  288. return false;
  289. }
  290. broker.token_expires_at = expires_at;
  291. MQTT_LOG("%s token ready exp=%lu owner=%s email=%s", broker.spec->label,
  292. static_cast<unsigned long>(expires_at), owner != nullptr ? "yes" : "no", email != nullptr ? "yes" : "no");
  293. return true;
  294. }
  295. void MQTTUplink::destroyBroker(BrokerState& broker, bool reset_retry_state) {
  296. bool had_runtime_state = broker.client != nullptr || broker.token != nullptr || broker.connected ||
  297. broker.connect_announced || broker.reconnect_pending || broker.next_connect_attempt != 0 ||
  298. broker.last_connect_attempt != 0 || broker.reconnect_failures != 0 ||
  299. broker.token_expires_at != 0;
  300. if (!had_runtime_state) {
  301. return;
  302. }
  303. logMqttMemorySnapshot("destroy-pre", broker.spec != nullptr ? broker.spec->label : nullptr);
  304. if (broker.client != nullptr) {
  305. MQTT_LOG("%s destroy broker client", broker.spec->label);
  306. esp_mqtt_client_stop(broker.client);
  307. esp_mqtt_client_destroy(broker.client);
  308. broker.client = nullptr;
  309. }
  310. freeScratchBuffer(broker.token);
  311. broker.token = nullptr;
  312. broker.connected = false;
  313. broker.connect_announced = false;
  314. broker.connected_since_ms = 0;
  315. broker.token_expires_at = 0;
  316. if (reset_retry_state) {
  317. broker.reconnect_pending = false;
  318. broker.next_connect_attempt = 0;
  319. broker.reconnect_failures = 0;
  320. broker.last_connect_attempt = 0;
  321. }
  322. logMqttMemorySnapshot("destroy-post", broker.spec != nullptr ? broker.spec->label : nullptr);
  323. }
  324. void MQTTUplink::queuePublish(BrokerState& broker, const char* topic, const char* payload, bool retain) {
  325. if (broker.client == nullptr || !broker.connected) {
  326. return;
  327. }
  328. MQTT_LOG("%s publish topic=%s retain=%d bytes=%u", broker.spec->label, topic, retain ? 1 : 0,
  329. static_cast<unsigned>(strlen(payload)));
  330. int enqueue_rc = esp_mqtt_client_enqueue(broker.client, topic, payload, 0, 1, retain, true);
  331. MQTT_LOG("%s enqueue topic=%s rc=%d connected=%d", broker.spec->label, topic, enqueue_rc, broker.connected ? 1 : 0);
  332. }
  333. int MQTTUplink::buildStatusJson(char* buffer, size_t buffer_size, bool online) const {
  334. char ts[32];
  335. formatIsoTimestamp(time(nullptr), ts, sizeof(ts));
  336. char model[48];
  337. makeSafeToken(board.getManufacturerName(), model, sizeof(model));
  338. char origin[80];
  339. const char* node_name = (_node_name != nullptr && _node_name[0] != 0) ? _node_name : _device_id;
  340. escapeJsonString(node_name, origin, sizeof(origin));
  341. char client_version[96];
  342. snprintf(client_version, sizeof(client_version), "%s", CLIENT_VERSION);
  343. char radio_info[48];
  344. snprintf(radio_info, sizeof(radio_info), "%.6f,%.1f,%u,%u", static_cast<double>(_last_status.radio_freq),
  345. static_cast<double>(_last_status.radio_bw), _last_status.radio_sf, _last_status.radio_cr);
  346. return snprintf(buffer, buffer_size,
  347. "{\"status\":\"%s\",\"timestamp\":\"%s\",\"origin\":\"%s\",\"origin_id\":\"%s\","
  348. "\"model\":\"%s\",\"firmware_version\":\"%s\",\"radio\":\"%s\",\"client_version\":\"%s\","
  349. "\"stats\":{\"battery_mv\":%d,\"uptime_secs\":%lu,\"errors\":%u,\"queue_len\":%u,"
  350. "\"noise_floor\":%d,\"tx_air_secs\":%lu,\"rx_air_secs\":%lu,\"recv_errors\":%lu}}",
  351. online ? "online" : "offline", ts, origin, _device_id, model, FIRMWARE_VERSION, radio_info, client_version,
  352. _last_status.battery_mv, static_cast<unsigned long>(_last_status.uptime_secs), _last_status.error_flags,
  353. _last_status.queue_len, _last_status.noise_floor, static_cast<unsigned long>(_last_status.tx_air_secs),
  354. static_cast<unsigned long>(_last_status.rx_air_secs),
  355. static_cast<unsigned long>(_last_status.recv_errors));
  356. }
  357. int MQTTUplink::buildPacketJson(char* buffer, size_t buffer_size, const mesh::Packet& packet, bool is_tx, int rssi,
  358. float snr, int score, int duration) const {
  359. uint8_t raw[256];
  360. int raw_len = packet.writeTo(raw);
  361. char* raw_hex = allocScratchBuffer(520);
  362. if (raw_hex == nullptr) {
  363. return -1;
  364. }
  365. bytesToHexUpper(raw, raw_len, raw_hex, 520);
  366. uint8_t packet_hash[MAX_HASH_SIZE];
  367. packet.calculatePacketHash(packet_hash);
  368. char hash_hex[(MAX_HASH_SIZE * 2) + 1];
  369. bytesToHexUpper(packet_hash, MAX_HASH_SIZE, hash_hex, sizeof(hash_hex));
  370. time_t now = time(nullptr);
  371. char ts[32];
  372. formatIsoTimestamp(now, ts, sizeof(ts));
  373. struct tm tm_utc;
  374. gmtime_r(&now, &tm_utc);
  375. char time_only[16];
  376. char date_only[16];
  377. strftime(time_only, sizeof(time_only), "%H:%M:%S", &tm_utc);
  378. strftime(date_only, sizeof(date_only), "%d/%m/%Y", &tm_utc);
  379. char origin[80];
  380. const char* node_name = (_node_name != nullptr && _node_name[0] != 0) ? _node_name : _device_id;
  381. escapeJsonString(node_name, origin, sizeof(origin));
  382. if (packet.isRouteDirect() && packet.path_len > 0) {
  383. char path_info[128];
  384. snprintf(path_info, sizeof(path_info), "path_%dx%d_%db", (int)packet.getPathHashCount(),
  385. (int)packet.getPathHashSize(), (int)packet.getPathByteLen());
  386. int len;
  387. if (score >= 0) {
  388. len = snprintf(buffer, buffer_size,
  389. "{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\",\"type\":\"PACKET\","
  390. "\"direction\":\"%s\",\"time\":\"%s\",\"date\":\"%s\",\"len\":\"%d\",\"packet_type\":\"%u\","
  391. "\"route\":\"D\",\"payload_len\":\"%u\",\"raw\":\"%s\",\"SNR\":\"%.1f\",\"RSSI\":\"%d\","
  392. "\"score\":\"%d\",\"duration\":\"%d\",\"hash\":\"%s\",\"path\":\"%s\"}",
  393. origin, _device_id, ts, is_tx ? "tx" : "rx", time_only, date_only, raw_len,
  394. packet.getPayloadType(), packet.payload_len, raw_hex, snr, rssi, score, duration, hash_hex,
  395. path_info);
  396. } else {
  397. len = snprintf(buffer, buffer_size,
  398. "{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\",\"type\":\"PACKET\","
  399. "\"direction\":\"%s\",\"time\":\"%s\",\"date\":\"%s\",\"len\":\"%d\",\"packet_type\":\"%u\","
  400. "\"route\":\"D\",\"payload_len\":\"%u\",\"raw\":\"%s\",\"SNR\":\"%.1f\",\"RSSI\":\"%d\","
  401. "\"hash\":\"%s\",\"path\":\"%s\"}",
  402. origin, _device_id, ts, is_tx ? "tx" : "rx", time_only, date_only, raw_len,
  403. packet.getPayloadType(), packet.payload_len, raw_hex, snr, rssi, hash_hex, path_info);
  404. }
  405. freeScratchBuffer(raw_hex);
  406. return len;
  407. }
  408. int len;
  409. if (score >= 0) {
  410. len = snprintf(buffer, buffer_size,
  411. "{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\",\"type\":\"PACKET\","
  412. "\"direction\":\"%s\",\"time\":\"%s\",\"date\":\"%s\",\"len\":\"%d\",\"packet_type\":\"%u\","
  413. "\"route\":\"F\",\"payload_len\":\"%u\",\"raw\":\"%s\",\"SNR\":\"%.1f\",\"RSSI\":\"%d\","
  414. "\"score\":\"%d\",\"duration\":\"%d\",\"hash\":\"%s\"}",
  415. origin, _device_id, ts, is_tx ? "tx" : "rx", time_only, date_only, raw_len,
  416. packet.getPayloadType(), packet.payload_len, raw_hex, snr, rssi, score, duration, hash_hex);
  417. } else {
  418. len = snprintf(buffer, buffer_size,
  419. "{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\",\"type\":\"PACKET\","
  420. "\"direction\":\"%s\",\"time\":\"%s\",\"date\":\"%s\",\"len\":\"%d\",\"packet_type\":\"%u\","
  421. "\"route\":\"F\",\"payload_len\":\"%u\",\"raw\":\"%s\",\"SNR\":\"%.1f\",\"RSSI\":\"%d\","
  422. "\"hash\":\"%s\"}",
  423. origin, _device_id, ts, is_tx ? "tx" : "rx", time_only, date_only, raw_len,
  424. packet.getPayloadType(), packet.payload_len, raw_hex, snr, rssi, hash_hex);
  425. }
  426. freeScratchBuffer(raw_hex);
  427. return len;
  428. }
  429. int MQTTUplink::buildRawJson(char* buffer, size_t buffer_size, const mesh::Packet& packet, bool is_tx, int rssi,
  430. float snr) const {
  431. (void)is_tx;
  432. (void)rssi;
  433. (void)snr;
  434. uint8_t raw[256];
  435. int raw_len = packet.writeTo(raw);
  436. char* raw_hex = allocScratchBuffer(520);
  437. if (raw_hex == nullptr) {
  438. return -1;
  439. }
  440. bytesToHexUpper(raw, raw_len, raw_hex, 520);
  441. char ts[32];
  442. formatIsoTimestamp(time(nullptr), ts, sizeof(ts));
  443. char origin[80];
  444. const char* node_name = (_node_name != nullptr && _node_name[0] != 0) ? _node_name : _device_id;
  445. escapeJsonString(node_name, origin, sizeof(origin));
  446. int len = snprintf(buffer, buffer_size,
  447. "{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\",\"type\":\"RAW\",\"data\":\"%s\"}",
  448. origin, _device_id, ts, raw_hex);
  449. freeScratchBuffer(raw_hex);
  450. return len;
  451. }
  452. void MQTTUplink::publishOnlineStatus(BrokerState& broker) {
  453. char* payload = allocScratchBuffer(768);
  454. if (payload == nullptr) {
  455. return;
  456. }
  457. int len = buildStatusJson(payload, 768, true);
  458. if (len > 0 && static_cast<size_t>(len) < 768) {
  459. queuePublish(broker, broker.status_topic, payload, true);
  460. }
  461. freeScratchBuffer(payload);
  462. }
  463. void MQTTUplink::publishStatus(bool online) {
  464. logMqttMemorySnapshot(online ? "status-pre" : "status-offline-pre");
  465. char* payload = allocScratchBuffer(768);
  466. if (payload == nullptr) {
  467. return;
  468. }
  469. int len = buildStatusJson(payload, 768, online);
  470. if (len <= 0 || static_cast<size_t>(len) >= 768) {
  471. freeScratchBuffer(payload);
  472. return;
  473. }
  474. for (BrokerState& broker : _brokers) {
  475. if (broker.spec != nullptr) {
  476. queuePublish(broker, broker.status_topic, payload, true);
  477. }
  478. }
  479. freeScratchBuffer(payload);
  480. logMqttMemorySnapshot(online ? "status-post" : "status-offline-post");
  481. }
  482. void MQTTUplink::handleMqttEvent(void* handler_args, esp_event_base_t, int32_t event_id, void* event_data) {
  483. auto* broker = static_cast<BrokerState*>(handler_args);
  484. if (broker == nullptr) {
  485. return;
  486. }
  487. auto* event = static_cast<esp_mqtt_event_handle_t>(event_data);
  488. unsigned long now_ms = millis();
  489. unsigned long connected_for_ms = broker->connected_since_ms != 0 ? (now_ms - broker->connected_since_ms) : 0;
  490. switch (event_id) {
  491. case MQTT_EVENT_CONNECTED:
  492. broker->connected = true;
  493. broker->reconnect_pending = false;
  494. broker->next_connect_attempt = 0;
  495. broker->reconnect_failures = 0;
  496. broker->connected_since_ms = now_ms;
  497. MQTT_LOG("%s connected", broker->spec->label);
  498. logMqttMemorySnapshot("connected", broker->spec->label);
  499. break;
  500. case MQTT_EVENT_DISCONNECTED:
  501. broker->connected = false;
  502. broker->connected_since_ms = 0;
  503. MQTT_LOG("%s disconnected wifi_status=%d rssi=%d connected_for_ms=%lu", broker->spec->label,
  504. static_cast<int>(WiFi.status()), WiFi.RSSI(), connected_for_ms);
  505. logMqttMemorySnapshot("disconnected", broker->spec->label);
  506. if (!broker->reconnect_pending) {
  507. if (broker->reconnect_failures < 10) {
  508. broker->reconnect_failures++;
  509. }
  510. broker->reconnect_pending = true;
  511. broker->next_connect_attempt = now_ms + getBrokerRetryDelayMillis(broker->reconnect_failures);
  512. MQTT_LOG("%s reconnect in %lu ms (failures=%u)", broker->spec->label,
  513. getBrokerRetryDelayMillis(broker->reconnect_failures),
  514. static_cast<unsigned>(broker->reconnect_failures));
  515. }
  516. break;
  517. case MQTT_EVENT_ERROR:
  518. broker->connected = false;
  519. broker->connected_since_ms = 0;
  520. if (event != nullptr && event->error_handle != nullptr) {
  521. MQTT_LOG("%s error type=%d tls_esp=0x%x tls_stack=0x%x cert_flags=0x%x sock_errno=%d conn_refused=%d "
  522. "connected_for_ms=%lu",
  523. broker->spec->label, event->error_handle->error_type, event->error_handle->esp_tls_last_esp_err,
  524. event->error_handle->esp_tls_stack_err, event->error_handle->esp_tls_cert_verify_flags,
  525. event->error_handle->esp_transport_sock_errno, event->error_handle->connect_return_code,
  526. connected_for_ms);
  527. } else {
  528. MQTT_LOG("%s error event connected_for_ms=%lu", broker->spec->label, connected_for_ms);
  529. }
  530. MQTT_LOG("%s wifi_status=%d rssi=%d", broker->spec->label, static_cast<int>(WiFi.status()), WiFi.RSSI());
  531. logMqttMemorySnapshot("error", broker->spec->label);
  532. if (!broker->reconnect_pending) {
  533. if (broker->reconnect_failures < 10) {
  534. broker->reconnect_failures++;
  535. }
  536. broker->reconnect_pending = true;
  537. broker->next_connect_attempt = now_ms + getBrokerRetryDelayMillis(broker->reconnect_failures);
  538. MQTT_LOG("%s reconnect in %lu ms (failures=%u)", broker->spec->label,
  539. getBrokerRetryDelayMillis(broker->reconnect_failures),
  540. static_cast<unsigned>(broker->reconnect_failures));
  541. }
  542. break;
  543. case MQTT_EVENT_BEFORE_CONNECT:
  544. MQTT_LOG("%s before connect", broker->spec->label);
  545. logMqttMemorySnapshot("before-connect", broker->spec->label);
  546. break;
  547. default:
  548. break;
  549. }
  550. }
  551. void MQTTUplink::ensureBroker(BrokerState& broker, bool allow_new_connect) {
  552. if (broker.spec == nullptr) {
  553. return;
  554. }
  555. bool enabled = (_prefs.enabled_mask & broker.spec->bit) != 0;
  556. bool iata_configured = !isUnsetIataValue(_prefs.iata);
  557. if (!enabled || !iata_configured) {
  558. if (broker.client != nullptr || broker.token != nullptr || broker.connected || broker.connect_announced ||
  559. broker.reconnect_pending || broker.next_connect_attempt != 0 || broker.last_connect_attempt != 0 ||
  560. broker.reconnect_failures != 0 || broker.token_expires_at != 0) {
  561. destroyBroker(broker);
  562. }
  563. return;
  564. }
  565. if (_network == nullptr || !_network->hasTimeSync() || !_network->isWifiConnected()) {
  566. return;
  567. }
  568. time_t now = time(nullptr);
  569. if (broker.client != nullptr && broker.token_expires_at > 0 && now + kTokenRefreshSlackSecs >= broker.token_expires_at) {
  570. destroyBroker(broker, false);
  571. }
  572. unsigned long now_ms = millis();
  573. if (broker.client != nullptr) {
  574. if (broker.connected) {
  575. return;
  576. }
  577. if (!broker.reconnect_pending || now_ms < broker.next_connect_attempt) {
  578. return;
  579. }
  580. destroyBroker(broker, false);
  581. }
  582. if (broker.next_connect_attempt != 0 && now_ms < broker.next_connect_attempt) {
  583. return;
  584. }
  585. if (!allow_new_connect) {
  586. return;
  587. }
  588. broker.last_connect_attempt = now_ms;
  589. broker.reconnect_pending = false;
  590. if (!refreshToken(broker)) {
  591. broker.reconnect_pending = true;
  592. broker.next_connect_attempt = now_ms + kBrokerRetryBaseMillis;
  593. return;
  594. }
  595. refreshBrokerState(broker);
  596. MQTT_LOG("%s mqtt init host=%s port=%d path=%s client_id=%s", broker.spec->label, broker.spec->host, 443, "/mqtt",
  597. broker.client_id);
  598. logMqttMemorySnapshot("init-pre", broker.spec->label);
  599. esp_mqtt_client_config_t cfg = {};
  600. #if ESP_IDF_VERSION_MAJOR >= 5
  601. cfg.broker.address.hostname = broker.spec->host;
  602. cfg.broker.address.port = 443;
  603. cfg.broker.address.transport = MQTT_TRANSPORT_OVER_WSS;
  604. cfg.broker.address.path = "/mqtt";
  605. cfg.broker.verification.certificate = mqtt_ca_certs::kCombinedPem;
  606. cfg.credentials.username = broker.username;
  607. cfg.credentials.client_id = broker.client_id;
  608. cfg.credentials.authentication.password = broker.token;
  609. cfg.session.keepalive = 30;
  610. cfg.session.last_will.topic = broker.status_topic;
  611. cfg.session.last_will.msg = broker.offline_payload;
  612. cfg.session.last_will.qos = 1;
  613. cfg.session.last_will.retain = 1;
  614. cfg.network.reconnect_timeout_ms = 10000;
  615. cfg.network.timeout_ms = 10000;
  616. cfg.network.disable_auto_reconnect = true;
  617. cfg.buffer.size = 768;
  618. cfg.buffer.out_size = 1280;
  619. #else
  620. cfg.host = broker.spec->host;
  621. cfg.port = 443;
  622. cfg.username = broker.username;
  623. cfg.password = broker.token;
  624. cfg.client_id = broker.client_id;
  625. cfg.keepalive = 30;
  626. cfg.buffer_size = 768;
  627. cfg.out_buffer_size = 1280;
  628. cfg.reconnect_timeout_ms = 10000;
  629. cfg.network_timeout_ms = 10000;
  630. cfg.disable_auto_reconnect = true;
  631. cfg.transport = MQTT_TRANSPORT_OVER_WSS;
  632. cfg.cert_pem = mqtt_ca_certs::kCombinedPem;
  633. cfg.lwt_topic = broker.status_topic;
  634. cfg.lwt_msg = broker.offline_payload;
  635. cfg.lwt_qos = 1;
  636. cfg.lwt_retain = 1;
  637. cfg.path = "/mqtt";
  638. #endif
  639. broker.client = esp_mqtt_client_init(&cfg);
  640. if (broker.client == nullptr) {
  641. MQTT_LOG("%s mqtt init failed", broker.spec->label);
  642. logMqttMemorySnapshot("init-failed", broker.spec->label);
  643. return;
  644. }
  645. logMqttMemorySnapshot("init-post", broker.spec->label);
  646. esp_mqtt_client_register_event(broker.client, MQTT_EVENT_ANY, &MQTTUplink::handleMqttEvent, &broker);
  647. if (esp_mqtt_client_start(broker.client) != ESP_OK) {
  648. MQTT_LOG("%s mqtt start failed", broker.spec->label);
  649. logMqttMemorySnapshot("start-failed", broker.spec->label);
  650. broker.reconnect_pending = true;
  651. broker.next_connect_attempt = now_ms + kBrokerRetryBaseMillis;
  652. destroyBroker(broker, false);
  653. } else {
  654. MQTT_LOG("%s mqtt start requested", broker.spec->label);
  655. logMqttMemorySnapshot("start-requested", broker.spec->label);
  656. }
  657. }
  658. void MQTTUplink::begin(FILESYSTEM* fs) {
  659. _fs = fs;
  660. MQTTPrefsStore::load(_fs, _prefs);
  661. uint8_t normalized_mask = normalizeEnabledMask(_prefs.enabled_mask & 0x07);
  662. if (normalized_mask != _prefs.enabled_mask) {
  663. _prefs.enabled_mask = normalized_mask;
  664. savePrefs();
  665. }
  666. refreshIdentityStrings();
  667. _running = true;
  668. _last_status_publish = millis();
  669. MQTT_LOG("begin iata=%s enabled_mask=0x%02X", _prefs.iata, _prefs.enabled_mask);
  670. }
  671. void MQTTUplink::end() {
  672. MQTT_LOG("end");
  673. publishStatus(false);
  674. for (BrokerState& broker : _brokers) {
  675. destroyBroker(broker);
  676. }
  677. _running = false;
  678. }
  679. void MQTTUplink::loop(const MQTTStatusSnapshot& snapshot) {
  680. if (!_running) {
  681. return;
  682. }
  683. _last_status = snapshot;
  684. BrokerState* active_connecting_broker = nullptr;
  685. for (BrokerState& broker : _brokers) {
  686. if (broker.client != nullptr && !broker.connected && !broker.reconnect_pending) {
  687. active_connecting_broker = &broker;
  688. break;
  689. }
  690. }
  691. bool connect_started = false;
  692. for (BrokerState& broker : _brokers) {
  693. bool allow_new_connect = active_connecting_broker == nullptr && !connect_started;
  694. ensureBroker(broker, allow_new_connect);
  695. if (active_connecting_broker == nullptr && broker.client != nullptr && !broker.connected && !broker.reconnect_pending) {
  696. connect_started = true;
  697. }
  698. if (broker.connected && !broker.connect_announced) {
  699. publishOnlineStatus(broker);
  700. broker.connect_announced = true;
  701. _last_status_publish = millis();
  702. } else if (!broker.connected) {
  703. broker.connect_announced = false;
  704. }
  705. }
  706. if (_prefs.status_enabled && hasEnabledBroker() && _network != nullptr && _network->hasTimeSync() &&
  707. millis() - _last_status_publish >= _prefs.status_interval_ms) {
  708. publishStatus(true);
  709. _last_status_publish = millis();
  710. }
  711. }
  712. void MQTTUplink::publishPacket(const mesh::Packet& packet, bool is_tx, int rssi, float snr, int score, int duration) {
  713. if (!_running || _network == nullptr || !_network->hasTimeSync() || !_network->isWifiConnected() || !hasEnabledBroker() ||
  714. !_prefs.packets_enabled) {
  715. return;
  716. }
  717. if (is_tx && !_prefs.tx_enabled) {
  718. return;
  719. }
  720. MQTT_LOG("packet dir=%s type=%u payload_len=%u rssi=%d snr=%.1f score=%d duration=%d",
  721. is_tx ? "tx" : "rx", packet.getPayloadType(), packet.payload_len, rssi, snr, score, duration);
  722. char* payload = allocScratchBuffer(1280);
  723. if (payload == nullptr) {
  724. return;
  725. }
  726. int len = buildPacketJson(payload, 1280, packet, is_tx, rssi, snr, score, duration);
  727. if (len <= 0 || static_cast<size_t>(len) >= 1280) {
  728. freeScratchBuffer(payload);
  729. return;
  730. }
  731. char topic[128];
  732. formatTopic(topic, sizeof(topic), "packets");
  733. for (BrokerState& broker : _brokers) {
  734. if (broker.spec != nullptr) {
  735. queuePublish(broker, topic, payload, false);
  736. }
  737. }
  738. freeScratchBuffer(payload);
  739. if (!_prefs.raw_enabled) {
  740. return;
  741. }
  742. char* raw_payload = allocScratchBuffer(896);
  743. if (raw_payload == nullptr) {
  744. return;
  745. }
  746. len = buildRawJson(raw_payload, 896, packet, is_tx, rssi, snr);
  747. if (len <= 0 || static_cast<size_t>(len) >= 896) {
  748. freeScratchBuffer(raw_payload);
  749. return;
  750. }
  751. formatTopic(topic, sizeof(topic), "raw");
  752. for (BrokerState& broker : _brokers) {
  753. if (broker.spec != nullptr) {
  754. queuePublish(broker, topic, raw_payload, false);
  755. }
  756. }
  757. freeScratchBuffer(raw_payload);
  758. }
  759. void MQTTUplink::formatStatusReply(char* reply, size_t reply_size) const {
  760. auto broker_state = [this](uint8_t bit) -> const char* {
  761. if ((_prefs.enabled_mask & bit) == 0) {
  762. return "off";
  763. }
  764. if (isUnsetIataValue(_prefs.iata)) {
  765. return "invalid iata";
  766. }
  767. const BrokerState* broker = nullptr;
  768. for (const BrokerState& candidate : _brokers) {
  769. if (candidate.spec != nullptr && candidate.spec->bit == bit) {
  770. broker = &candidate;
  771. break;
  772. }
  773. }
  774. if (broker == nullptr) {
  775. return "retry";
  776. }
  777. if (broker->connected) {
  778. return "up";
  779. }
  780. if (_network == nullptr || !_network->isWifiConnected() || !_network->hasTimeSync()) {
  781. return "wait";
  782. }
  783. if (broker->client != nullptr) {
  784. return "conn";
  785. }
  786. if (broker->next_connect_attempt != 0 && broker->next_connect_attempt > millis()) {
  787. return "backoff";
  788. }
  789. return "retry";
  790. };
  791. snprintf(reply, reply_size, "> wifi:%s ntp:%s iata:%s eastmesh-au:%s letsmesh-eu:%s letsmesh-us:%s status:%s tx:%s",
  792. (_network != nullptr && _network->isWifiConnected()) ? "up" : "down",
  793. (_network != nullptr && _network->hasTimeSync()) ? "up" : "wait",
  794. _prefs.iata,
  795. broker_state(kEastmeshBit), broker_state(kLetsmeshEuBit), broker_state(kLetsmeshUsBit),
  796. _prefs.status_enabled ? "on" : "off", _prefs.tx_enabled ? "on" : "off");
  797. }
  798. bool MQTTUplink::setEndpointEnabled(uint8_t bit, bool enabled) {
  799. uint8_t next_mask = _prefs.enabled_mask & 0x07;
  800. if (enabled) {
  801. next_mask = normalizeEnabledMask(next_mask | bit);
  802. if ((next_mask & bit) == 0) {
  803. return false;
  804. }
  805. } else {
  806. next_mask &= ~bit;
  807. }
  808. _prefs.enabled_mask = next_mask;
  809. savePrefs();
  810. return true;
  811. }
  812. bool MQTTUplink::isEndpointEnabled(uint8_t bit) const {
  813. return (_prefs.enabled_mask & bit) != 0;
  814. }
  815. bool MQTTUplink::setPacketsEnabled(bool enabled) {
  816. _prefs.packets_enabled = enabled ? 1 : 0;
  817. return savePrefs();
  818. }
  819. bool MQTTUplink::setRawEnabled(bool enabled) {
  820. _prefs.raw_enabled = enabled ? 1 : 0;
  821. return savePrefs();
  822. }
  823. bool MQTTUplink::setStatusEnabled(bool enabled) {
  824. _prefs.status_enabled = enabled ? 1 : 0;
  825. return savePrefs();
  826. }
  827. bool MQTTUplink::setTxEnabled(bool enabled) {
  828. _prefs.tx_enabled = enabled ? 1 : 0;
  829. return savePrefs();
  830. }
  831. bool MQTTUplink::setIata(const char* iata) {
  832. if (iata == nullptr || *iata == 0) {
  833. return false;
  834. }
  835. char cleaned[sizeof(_prefs.iata)];
  836. memset(cleaned, 0, sizeof(cleaned));
  837. makeSafeToken(iata, cleaned, sizeof(cleaned));
  838. for (size_t i = 0; cleaned[i] != 0; ++i) {
  839. cleaned[i] = toupper(static_cast<unsigned char>(cleaned[i]));
  840. }
  841. if (strcmp(cleaned, MQTT_UNSET_IATA) == 0) {
  842. StrHelper::strncpy(_prefs.iata, MQTT_UNSET_IATA, sizeof(_prefs.iata));
  843. refreshIdentityStrings();
  844. return savePrefs();
  845. }
  846. StrHelper::strncpy(_prefs.iata, cleaned, sizeof(_prefs.iata));
  847. refreshIdentityStrings();
  848. return savePrefs();
  849. }
  850. bool MQTTUplink::setOwnerPublicKey(const char* owner_public_key) {
  851. if (owner_public_key == nullptr) {
  852. return false;
  853. }
  854. if (owner_public_key[0] == 0) {
  855. _prefs.owner_public_key[0] = 0;
  856. return savePrefs();
  857. }
  858. if (strlen(owner_public_key) != 64) {
  859. return false;
  860. }
  861. for (size_t i = 0; i < 64; ++i) {
  862. if (!isxdigit(static_cast<unsigned char>(owner_public_key[i]))) {
  863. return false;
  864. }
  865. _prefs.owner_public_key[i] = toupper(static_cast<unsigned char>(owner_public_key[i]));
  866. }
  867. _prefs.owner_public_key[64] = 0;
  868. return savePrefs();
  869. }
  870. bool MQTTUplink::setOwnerEmail(const char* owner_email) {
  871. if (owner_email == nullptr) {
  872. return false;
  873. }
  874. StrHelper::strncpy(_prefs.owner_email, owner_email, sizeof(_prefs.owner_email));
  875. return savePrefs();
  876. }
  877. bool MQTTUplink::isAnyBrokerConnected() const {
  878. for (const BrokerState& broker : _brokers) {
  879. if (broker.spec != nullptr && broker.connected) {
  880. return true;
  881. }
  882. }
  883. return false;
  884. }
  885. const char* MQTTUplink::getAggregateBrokerState() const {
  886. uint8_t enabled_count = 0;
  887. uint8_t connected_count = 0;
  888. for (const BrokerState& broker : _brokers) {
  889. if (broker.spec == nullptr || (broker.spec->bit & _prefs.enabled_mask) == 0) {
  890. continue;
  891. }
  892. enabled_count++;
  893. if (broker.connected) {
  894. connected_count++;
  895. }
  896. }
  897. if (enabled_count == 0 || connected_count == 0) {
  898. return "down";
  899. }
  900. if (connected_count < enabled_count) {
  901. return "degraded";
  902. }
  903. return "up";
  904. }
  905. #else
  906. MQTTUplink::MQTTUplink(mesh::RTCClock&, mesh::LocalIdentity&)
  907. : _fs(nullptr), _rtc(nullptr), _identity(nullptr), _running(false), _last_status_publish(0), _last_status{},
  908. _node_name(nullptr), _network(nullptr) {
  909. MQTTPrefsStore::setDefaults(_prefs);
  910. }
  911. bool MQTTUplink::savePrefs() { return false; }
  912. void MQTTUplink::begin(FILESYSTEM*) {}
  913. void MQTTUplink::end() {}
  914. void MQTTUplink::loop(const MQTTStatusSnapshot&) {}
  915. void MQTTUplink::publishPacket(const mesh::Packet&, bool, int, float, int, int) {}
  916. void MQTTUplink::formatStatusReply(char* reply, size_t reply_size) const { snprintf(reply, reply_size, "> unsupported"); }
  917. bool MQTTUplink::setEndpointEnabled(uint8_t, bool) { return false; }
  918. bool MQTTUplink::isEndpointEnabled(uint8_t) const { return false; }
  919. bool MQTTUplink::setPacketsEnabled(bool) { return false; }
  920. bool MQTTUplink::setRawEnabled(bool) { return false; }
  921. bool MQTTUplink::setStatusEnabled(bool) { return false; }
  922. bool MQTTUplink::setTxEnabled(bool) { return false; }
  923. bool MQTTUplink::setIata(const char*) { return false; }
  924. bool MQTTUplink::isActive() const { return false; }
  925. bool MQTTUplink::setOwnerPublicKey(const char*) { return false; }
  926. bool MQTTUplink::setOwnerEmail(const char*) { return false; }
  927. bool MQTTUplink::sendStatusNow() { return false; }
  928. bool MQTTUplink::isAnyBrokerConnected() const { return false; }
  929. const char* MQTTUplink::getAggregateBrokerState() const { return "down"; }
  930. #endif
  931. #endif