MQTTUplink.cpp 35 KB

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