From 90d2c9302fd2268c710dee2cc87d27c66fc07a3d Mon Sep 17 00:00:00 2001 From: prabathbr Date: Wed, 24 Jun 2026 21:42:22 +1000 Subject: [PATCH] test MQTT bridge --- eastmesh-docs/custom-cli.md | 28 ++- eastmesh-docs/web-panel.md | 2 + examples/simple_repeater/MyMesh.cpp | 8 + examples/simple_repeater/MyMesh.h | 7 + src/helpers/CommonCLI.cpp | 61 ++++- src/helpers/CommonCLI.h | 6 +- src/helpers/bridges/MQTTBridge.cpp | 308 ++++++++++++++++++++++++++ src/helpers/bridges/MQTTBridge.h | 52 +++++ src/helpers/web/WebPanelServer.cpp | 53 ++++- variants/eastmesh_mqtt/platformio.ini | 4 +- 10 files changed, 521 insertions(+), 8 deletions(-) create mode 100644 src/helpers/bridges/MQTTBridge.cpp create mode 100644 src/helpers/bridges/MQTTBridge.h diff --git a/eastmesh-docs/custom-cli.md b/eastmesh-docs/custom-cli.md index 04c230cf..4d06fab4 100644 --- a/eastmesh-docs/custom-cli.md +++ b/eastmesh-docs/custom-cli.md @@ -126,7 +126,7 @@ Default servers are `au.pool.ntp.org`, `time.google.com`, and `time.cloudflare.c ### ESP-NOW Bridge Settings For Observer ESP-NOW Builds -These commands are available on `*_repeater_observer_espnow` firmware targets. +These commands are available on `*_repeater_observer_espnow` firmware targets that use the ESP-NOW bridge transport. Bridge commands are for local ESP-NOW bridge use between nearby repeaters, such as linking repeaters on `Australia (Narrow)` and `Australia (Mid)`. They are not MQTT-over-WAN, VPN, or internet bridge controls. @@ -156,6 +156,32 @@ Example: OK ``` +### MQTT Bridge Settings For Observer MQTT Bridge Builds + +These commands are available on observer builds that compile with the MQTT mesh bridge transport (for example `Xiao_S3_WIO_repeater_observer_espnow` when built with `WITH_MQTT_BRIDGE`). + +The MQTT bridge forwards raw mesh packets over a shared topic at a peer MQTT broker. It is separate from MQTT uplink publishing to EastMesh or MeshMapper brokers. + +- `get bridge.type`: returns `mqtt` on MQTT bridge builds. +- `get bridge.peer.host`: shows the configured peer MQTT broker host. +- `set bridge.peer.host `: sets the peer MQTT broker host and restarts the bridge. +- `get bridge.peer.port`: shows the configured peer MQTT broker port (defaults to `1883` when unset). +- `set bridge.peer.port `: sets the peer MQTT broker port and restarts the bridge. +- `get bridge.secret` / `set bridge.secret `: shared XOR key used by all bridge nodes on the same bridge network. + +Both bridge nodes must use the same peer broker address, port, and `bridge.secret`. Mesh packets are published and subscribed on topic `meshcore/bridge/packets`. + +Example: + +```text +> set bridge.peer.host 192.168.1.10 +OK +> set bridge.peer.port 1883 +OK +> set bridge.secret my-shared-secret +OK +``` + ### Web Panel Controls - `get web` diff --git a/eastmesh-docs/web-panel.md b/eastmesh-docs/web-panel.md index dce014b8..399cf135 100644 --- a/eastmesh-docs/web-panel.md +++ b/eastmesh-docs/web-panel.md @@ -255,6 +255,7 @@ This section includes: - `mqtt.email`: owner contact email. - MQTT brokers: **Primary MQTT** and **Secondary MQTT** dropdowns, each selecting one of `eastmesh-au`, `meshmapper`, `Custom`, the retired `letsmesh-eu`/`letsmesh-us`, or `None`. The two slots enforce the two-broker maximum, and a broker chosen in one slot is disabled in the other. - custom MQTT `host:port`, TCP/WSS transport, username, and password fields, shown when `Custom` is selected in either slot. +- mesh bridge peer MQTT `host:port`, shown on MQTT bridge builds (`get bridge.type` returns `mqtt`). This is separate from MQTT uplink brokers and points at the shared peer broker used for bidirectional mesh packet bridging. `UNSET - To be configured` is the default for new observer installs until a real saved value exists. @@ -269,6 +270,7 @@ Notes: - turning off a connected MQTT server publishes retained offline status before the client disconnects - changing `mqtt.iata` away from a configured value publishes retained offline status to the old status topic, restarts connected broker clients, and reconnects under the new topic path - at most two MQTT brokers can be enabled at once +- on MQTT bridge builds, both bridge nodes must use the same peer broker host, port, and `bridge.secret` ## `/stats` Overview diff --git a/examples/simple_repeater/MyMesh.cpp b/examples/simple_repeater/MyMesh.cpp index 9996b884..1dd57add 100644 --- a/examples/simple_repeater/MyMesh.cpp +++ b/examples/simple_repeater/MyMesh.cpp @@ -533,6 +533,8 @@ uint8_t MyMesh::handleAnonClockReq(const mesh::Identity& sender, uint32_t sender reply_data[8] |= 0x01; // is bridge, type UART #elif WITH_ESPNOW_BRIDGE reply_data[8] |= 0x03; // is bridge, type ESP-NOW +#elif WITH_MQTT_BRIDGE + reply_data[8] |= 0x04; // is bridge, type MQTT #endif if (_prefs.disable_fwd) { // is this repeater currently disabled reply_data[8] |= 0x80; // is disabled @@ -1198,6 +1200,9 @@ MyMesh::MyMesh(mesh::MainBoard &board, mesh::Radio &radio, mesh::MillisecondCloc #if defined(WITH_ESPNOW_BRIDGE) , bridge(&_prefs, _mgr, &rtc) #endif +#if defined(WITH_MQTT_BRIDGE) + , bridge(&_prefs, _mgr, &rtc) +#endif #if defined(WITH_MQTT_UPLINK) , mqtt(rtc, self_id) #endif @@ -1254,6 +1259,9 @@ MyMesh::MyMesh(mesh::MainBoard &board, mesh::Radio &radio, mesh::MillisecondCloc _prefs.bridge_channel = 1; // channel 1 StrHelper::strncpy(_prefs.bridge_secret, "LVSITANOS", sizeof(_prefs.bridge_secret)); +#if defined(WITH_MQTT_BRIDGE) + _prefs.bridge_peer_port = 1883; +#endif // GPS defaults _prefs.gps_enabled = 0; diff --git a/examples/simple_repeater/MyMesh.h b/examples/simple_repeater/MyMesh.h index 192b937b..b8bd5657 100644 --- a/examples/simple_repeater/MyMesh.h +++ b/examples/simple_repeater/MyMesh.h @@ -23,6 +23,11 @@ #define WITH_BRIDGE #endif +#ifdef WITH_MQTT_BRIDGE +#include "helpers/bridges/MQTTBridge.h" +#define WITH_BRIDGE +#endif + #ifdef WITH_MQTT_UPLINK #include #endif @@ -170,6 +175,8 @@ class MyMesh : public mesh::Mesh, public CommonCLICallbacks, public WebPanelComm RS232Bridge bridge; #elif defined(WITH_ESPNOW_BRIDGE) ESPNowBridge bridge; +#elif defined(WITH_MQTT_BRIDGE) + MQTTBridge bridge; #endif #ifdef WITH_MQTT_UPLINK MQTTUplink mqtt; diff --git a/src/helpers/CommonCLI.cpp b/src/helpers/CommonCLI.cpp index e5ed752b..7afb705d 100644 --- a/src/helpers/CommonCLI.cpp +++ b/src/helpers/CommonCLI.cpp @@ -103,7 +103,13 @@ void CommonCLI::loadPrefsInt(FILESYSTEM* fs, const char* filename) { if (file.available() >= (int)sizeof(_prefs->flood_max_advert)) { file.read((uint8_t *)&_prefs->flood_max_advert, sizeof(_prefs->flood_max_advert)); // 295 } - // next: 296 + if (file.available() >= (int)sizeof(_prefs->bridge_peer_host)) { + file.read((uint8_t *)&_prefs->bridge_peer_host, sizeof(_prefs->bridge_peer_host)); // 296 + } + if (file.available() >= (int)sizeof(_prefs->bridge_peer_port)) { + file.read((uint8_t *)&_prefs->bridge_peer_port, sizeof(_prefs->bridge_peer_port)); // 360 + } + // next: 362 // sanitise bad pref values _prefs->rx_delay_base = constrain(_prefs->rx_delay_base, 0, 20.0f); @@ -125,6 +131,9 @@ void CommonCLI::loadPrefsInt(FILESYSTEM* fs, const char* filename) { _prefs->bridge_pkt_src = constrain(_prefs->bridge_pkt_src, 0, 1); _prefs->bridge_baud = constrain(_prefs->bridge_baud, 9600, BRIDGE_MAX_BAUD); _prefs->bridge_channel = constrain(_prefs->bridge_channel, 0, 14); + if (_prefs->bridge_peer_port > 65535) { + _prefs->bridge_peer_port = 1883; + } _prefs->powersaving_enabled = constrain(_prefs->powersaving_enabled, 0, 1); @@ -200,7 +209,9 @@ void CommonCLI::savePrefs(FILESYSTEM* fs) { file.write((uint8_t *)&_prefs->fan_mode, sizeof(_prefs->fan_mode)); // 292 file.write((uint8_t *)&_prefs->fan_timeout_secs, sizeof(_prefs->fan_timeout_secs)); // 293 file.write((uint8_t *)&_prefs->flood_max_advert, sizeof(_prefs->flood_max_advert)); // 295 - // next: 296 + file.write((uint8_t *)&_prefs->bridge_peer_host, sizeof(_prefs->bridge_peer_host)); // 296 + file.write((uint8_t *)&_prefs->bridge_peer_port, sizeof(_prefs->bridge_peer_port)); // 360 + // next: 362 file.close(); } @@ -770,6 +781,42 @@ void CommonCLI::handleSetCmd(uint32_t sender_timestamp, char* command, char* rep _callbacks->restartBridge(); savePrefs(); strcpy(reply, "OK"); +#endif +#ifdef WITH_MQTT_BRIDGE + } else if (memcmp(config, "bridge.peer.host ", 17) == 0) { + const char* host = &config[17]; + size_t oi = 0; + char cleaned[sizeof(_prefs->bridge_peer_host)]; + memset(cleaned, 0, sizeof(cleaned)); + for (size_t i = 0; host[i] != 0 && oi + 1 < sizeof(cleaned); ++i) { + unsigned char c = static_cast(host[i]); + if (c <= ' ' || c == '/' || c == ':' || c == '\\') { + strcpy(reply, "Error: invalid host"); + return; + } + cleaned[oi++] = static_cast(c); + } + cleaned[oi] = 0; + StrHelper::strncpy(_prefs->bridge_peer_host, cleaned, sizeof(_prefs->bridge_peer_host)); + _callbacks->restartBridge(); + savePrefs(); + strcpy(reply, "OK"); + } else if (memcmp(config, "bridge.peer.port ", 17) == 0) { + char* end = nullptr; + unsigned long parsed = strtoul(&config[17], &end, 10); + if (end == &config[17] || *end != 0 || parsed == 0 || parsed > 65535UL) { + strcpy(reply, "Error: port must be 1-65535"); + return; + } + _prefs->bridge_peer_port = static_cast(parsed); + _callbacks->restartBridge(); + savePrefs(); + strcpy(reply, "OK"); + } else if (memcmp(config, "bridge.secret ", 14) == 0) { + StrHelper::strncpy(_prefs->bridge_secret, &config[14], sizeof(_prefs->bridge_secret)); + _callbacks->restartBridge(); + savePrefs(); + strcpy(reply, "OK"); #endif } else if (memcmp(config, "adc.multiplier ", 15) == 0) { _prefs->adc_multiplier = atof(&config[15]); @@ -884,6 +931,8 @@ void CommonCLI::handleGetCmd(uint32_t sender_timestamp, char* command, char* rep "rs232" #elif WITH_ESPNOW_BRIDGE "espnow" +#elif WITH_MQTT_BRIDGE + "mqtt" #else "none" #endif @@ -905,6 +954,14 @@ void CommonCLI::handleGetCmd(uint32_t sender_timestamp, char* command, char* rep sprintf(reply, "> %d", (uint32_t)_prefs->bridge_channel); } else if (memcmp(config, "bridge.secret", 13) == 0) { sprintf(reply, "> %s", _prefs->bridge_secret); +#endif +#ifdef WITH_MQTT_BRIDGE + } else if (memcmp(config, "bridge.peer.host", 16) == 0) { + sprintf(reply, "> %s", _prefs->bridge_peer_host); + } else if (memcmp(config, "bridge.peer.port", 16) == 0) { + sprintf(reply, "> %u", _prefs->bridge_peer_port != 0 ? _prefs->bridge_peer_port : 1883); + } else if (memcmp(config, "bridge.secret", 13) == 0) { + sprintf(reply, "> %s", _prefs->bridge_secret); #endif } else if (memcmp(config, "bootloader.ver", 14) == 0) { #ifdef NRF52_PLATFORM diff --git a/src/helpers/CommonCLI.h b/src/helpers/CommonCLI.h index 7310ac78..a7673766 100644 --- a/src/helpers/CommonCLI.h +++ b/src/helpers/CommonCLI.h @@ -6,7 +6,7 @@ #include #include -#if defined(WITH_RS232_BRIDGE) || defined(WITH_ESPNOW_BRIDGE) +#if defined(WITH_RS232_BRIDGE) || defined(WITH_ESPNOW_BRIDGE) || defined(WITH_MQTT_BRIDGE) #define WITH_BRIDGE #endif @@ -50,7 +50,7 @@ struct NodePrefs { // persisted to file uint8_t bridge_pkt_src; // 0 = logTx, 1 = logRx (default logTx) uint32_t bridge_baud; // 9600, 19200, 38400, 57600, 115200 (default 115200) uint8_t bridge_channel; // 1-14 (ESP-NOW only) - char bridge_secret[16]; // for XOR encryption of bridge packets (ESP-NOW only) + char bridge_secret[16]; // for XOR encryption of bridge packets (ESP-NOW / MQTT) // Power setting uint8_t powersaving_enabled; // boolean // Gps settings @@ -65,6 +65,8 @@ struct NodePrefs { // persisted to file uint8_t loop_detect; uint8_t fan_mode; uint16_t fan_timeout_secs; + char bridge_peer_host[64]; // peer MQTT broker host (MQTT bridge only) + uint16_t bridge_peer_port; // peer MQTT broker port (MQTT bridge only, default 1883) }; class CommonCLICallbacks { diff --git a/src/helpers/bridges/MQTTBridge.cpp b/src/helpers/bridges/MQTTBridge.cpp new file mode 100644 index 00000000..362fca4f --- /dev/null +++ b/src/helpers/bridges/MQTTBridge.cpp @@ -0,0 +1,308 @@ +#include "MQTTBridge.h" + +#ifdef WITH_MQTT_BRIDGE + +#if defined(ESP_PLATFORM) + +#include +#include +#include +#include + +#ifndef BRIDGE_DEBUG + #define BRIDGE_DEBUG 0 +#endif + +#if BRIDGE_DEBUG + #define BRIDGE_DEBUG_PRINTLN(...) Serial.printf(__VA_ARGS__) +#else + #define BRIDGE_DEBUG_PRINTLN(...) do { } while (0) +#endif + +namespace { +constexpr unsigned long kConnectRetryBaseMillis = 10000; +constexpr unsigned long kConnectRetryMaxMillis = 120000; + +unsigned long connectRetryDelayMillis(uint8_t failures) { + unsigned long delay_ms = kConnectRetryBaseMillis; + if (failures > 0) { + uint8_t shifts = min(failures - 1, 3); + delay_ms <<= shifts; + } + if (delay_ms > kConnectRetryMaxMillis) { + delay_ms = kConnectRetryMaxMillis; + } + return delay_ms; +} + +bool peerHostConfigured(const NodePrefs *prefs) { + return prefs != nullptr && prefs->bridge_peer_host[0] != 0; +} + +uint16_t peerPort(const NodePrefs *prefs) { + return prefs->bridge_peer_port != 0 ? prefs->bridge_peer_port : 1883; +} +} // namespace + +MQTTBridge *MQTTBridge::_instance = nullptr; +const char *MQTTBridge::kBridgeTopic = "meshcore/bridge/packets"; + +MQTTBridge::MQTTBridge(NodePrefs *prefs, mesh::PacketManager *mgr, mesh::RTCClock *rtc) + : BridgeBase(prefs, mgr, rtc), _client(nullptr), _connected(false), _started(false), + _next_connect_attempt(0), _reconnect_failures(0) { + _instance = this; + _client_id[0] = 0; +} + +void MQTTBridge::xorCrypt(uint8_t *data, size_t len) { + size_t keyLen = strlen(_prefs->bridge_secret); + if (keyLen == 0) { + return; + } + for (size_t i = 0; i < len; i++) { + data[i] ^= _prefs->bridge_secret[i % keyLen]; + } +} + +void MQTTBridge::destroyClient() { + if (_client != nullptr) { + esp_mqtt_client_stop(_client); + esp_mqtt_client_destroy(_client); + _client = nullptr; + } + _connected = false; + _started = false; +} + +void MQTTBridge::begin() { + BRIDGE_DEBUG_PRINTLN("MQTT bridge initializing\n"); + if (!peerHostConfigured(_prefs)) { + BRIDGE_DEBUG_PRINTLN("MQTT bridge peer host not configured\n"); + return; + } + + uint8_t mac[6]; + WiFi.macAddress(mac); + snprintf(_client_id, sizeof(_client_id), "mc-br-%02x%02x%02x", mac[3], mac[4], mac[5]); + + _initialized = true; + _next_connect_attempt = 0; + _reconnect_failures = 0; +} + +void MQTTBridge::end() { + BRIDGE_DEBUG_PRINTLN("MQTT bridge stopping\n"); + destroyClient(); + _initialized = false; +} + +void MQTTBridge::onMqttConnected() { + _connected = true; + _reconnect_failures = 0; + _next_connect_attempt = 0; + if (_client != nullptr) { + esp_mqtt_client_subscribe(_client, kBridgeTopic, 0); + } + BRIDGE_DEBUG_PRINTLN("MQTT bridge connected to %s:%u\n", _prefs->bridge_peer_host, + static_cast(peerPort(_prefs))); +} + +void MQTTBridge::onMqttDisconnected() { + _connected = false; + _started = false; + destroyClient(); + _next_connect_attempt = millis() + connectRetryDelayMillis(_reconnect_failures); + if (_reconnect_failures < 255) { + ++_reconnect_failures; + } +} + +void MQTTBridge::mqttEventHandler(void *handler_args, esp_event_base_t, int32_t event_id, void *event_data) { + auto *bridge = static_cast(handler_args); + if (bridge == nullptr) { + return; + } + + switch (event_id) { + case MQTT_EVENT_CONNECTED: + bridge->onMqttConnected(); + break; + case MQTT_EVENT_DISCONNECTED: + bridge->onMqttDisconnected(); + break; + case MQTT_EVENT_DATA: { + auto *event = static_cast(event_data); + if (event != nullptr && event->data_len > 0) { + bridge->handleMqttData(reinterpret_cast(event->data), + static_cast(event->data_len)); + } + break; + } + default: + break; + } +} + +bool MQTTBridge::ensureClient() { + if (_client != nullptr) { + return true; + } + if (!peerHostConfigured(_prefs)) { + return false; + } + if (WiFi.status() != WL_CONNECTED) { + return false; + } + + esp_mqtt_client_config_t cfg = {}; +#if ESP_IDF_VERSION_MAJOR >= 5 + cfg.broker.address.hostname = _prefs->bridge_peer_host; + cfg.broker.address.port = peerPort(_prefs); + cfg.broker.address.transport = MQTT_TRANSPORT_OVER_TCP; + cfg.credentials.client_id = _client_id; + cfg.session.keepalive = 30; + cfg.network.reconnect_timeout_ms = 10000; + cfg.network.timeout_ms = 10000; + cfg.network.disable_auto_reconnect = true; + cfg.buffer.size = MAX_MQTT_PACKET_SIZE; + cfg.buffer.out_size = MAX_MQTT_PACKET_SIZE; +#else + cfg.host = _prefs->bridge_peer_host; + cfg.port = peerPort(_prefs); + cfg.transport = MQTT_TRANSPORT_OVER_TCP; + cfg.client_id = _client_id; + cfg.keepalive = 30; + cfg.buffer_size = MAX_MQTT_PACKET_SIZE; + cfg.out_buffer_size = MAX_MQTT_PACKET_SIZE; + cfg.reconnect_timeout_ms = 10000; + cfg.network_timeout_ms = 10000; + cfg.disable_auto_reconnect = true; +#endif + + _client = esp_mqtt_client_init(&cfg); + if (_client == nullptr) { + BRIDGE_DEBUG_PRINTLN("MQTT bridge client init failed\n"); + return false; + } + + esp_mqtt_client_register_event(_client, MQTT_EVENT_ANY, &MQTTBridge::mqttEventHandler, this); + if (esp_mqtt_client_start(_client) != ESP_OK) { + BRIDGE_DEBUG_PRINTLN("MQTT bridge client start failed\n"); + destroyClient(); + return false; + } + + _started = true; + return true; +} + +void MQTTBridge::loop() { + if (!_initialized) { + return; + } + + if (WiFi.status() != WL_CONNECTED) { + if (_client != nullptr) { + destroyClient(); + } + return; + } + + if (_client == nullptr) { + if (_next_connect_attempt != 0 && millis() < _next_connect_attempt) { + return; + } + ensureClient(); + } +} + +void MQTTBridge::handleMqttData(const uint8_t *data, size_t len) { + if (len < (BRIDGE_MAGIC_SIZE + BRIDGE_CHECKSUM_SIZE)) { + BRIDGE_DEBUG_PRINTLN("MQTT RX packet too small, len=%u\n", static_cast(len)); + return; + } + + if (len > MAX_MQTT_PACKET_SIZE) { + BRIDGE_DEBUG_PRINTLN("MQTT RX packet too large, len=%u\n", static_cast(len)); + return; + } + + uint16_t received_magic = (data[0] << 8) | data[1]; + if (received_magic != BRIDGE_PACKET_MAGIC) { + BRIDGE_DEBUG_PRINTLN("MQTT RX invalid magic 0x%04X\n", received_magic); + return; + } + + uint8_t decrypted[MAX_MQTT_PACKET_SIZE]; + const size_t encryptedDataLen = len - BRIDGE_MAGIC_SIZE; + memcpy(decrypted, data + BRIDGE_MAGIC_SIZE, encryptedDataLen); + + xorCrypt(decrypted, encryptedDataLen); + + uint16_t received_checksum = (decrypted[0] << 8) | decrypted[1]; + const size_t payloadLen = encryptedDataLen - BRIDGE_CHECKSUM_SIZE; + + if (!validateChecksum(decrypted + BRIDGE_CHECKSUM_SIZE, payloadLen, received_checksum)) { + BRIDGE_DEBUG_PRINTLN("MQTT RX checksum mismatch, rcv=0x%04X\n", received_checksum); + return; + } + + mesh::Packet *pkt = _mgr->allocNew(); + if (!pkt) { + return; + } + + if (pkt->readFrom(decrypted + BRIDGE_CHECKSUM_SIZE, payloadLen)) { + onPacketReceived(pkt); + } else { + _mgr->free(pkt); + } +} + +void MQTTBridge::sendPacket(mesh::Packet *packet) { + if (!_initialized || !_connected || _client == nullptr || packet == nullptr) { + return; + } + + if (_seen_packets.hasSeen(packet)) { + return; + } + + uint8_t sizingBuffer[MAX_PAYLOAD_SIZE]; + uint16_t meshPacketLen = packet->writeTo(sizingBuffer); + if (meshPacketLen > MAX_PAYLOAD_SIZE) { + BRIDGE_DEBUG_PRINTLN("MQTT TX packet too large (payload=%u, max=%u)\n", meshPacketLen, + static_cast(MAX_PAYLOAD_SIZE)); + return; + } + + uint8_t buffer[MAX_MQTT_PACKET_SIZE]; + buffer[0] = (BRIDGE_PACKET_MAGIC >> 8) & 0xFF; + buffer[1] = BRIDGE_PACKET_MAGIC & 0xFF; + + const size_t packetOffset = BRIDGE_MAGIC_SIZE + BRIDGE_CHECKSUM_SIZE; + memcpy(buffer + packetOffset, sizingBuffer, meshPacketLen); + + uint16_t checksum = fletcher16(buffer + packetOffset, meshPacketLen); + buffer[2] = (checksum >> 8) & 0xFF; + buffer[3] = checksum & 0xFF; + + xorCrypt(buffer + BRIDGE_MAGIC_SIZE, meshPacketLen + BRIDGE_CHECKSUM_SIZE); + + const size_t totalPacketSize = BRIDGE_MAGIC_SIZE + BRIDGE_CHECKSUM_SIZE + meshPacketLen; + int msg_id = esp_mqtt_client_publish(_client, kBridgeTopic, reinterpret_cast(buffer), + static_cast(totalPacketSize), 0, 0); + if (msg_id >= 0) { + BRIDGE_DEBUG_PRINTLN("MQTT TX, len=%u\n", meshPacketLen); + } else { + BRIDGE_DEBUG_PRINTLN("MQTT TX FAILED\n"); + } +} + +void MQTTBridge::onPacketReceived(mesh::Packet *packet) { + handleReceivedPacket(packet); +} + +#endif // ESP_PLATFORM + +#endif // WITH_MQTT_BRIDGE diff --git a/src/helpers/bridges/MQTTBridge.h b/src/helpers/bridges/MQTTBridge.h new file mode 100644 index 00000000..bfa9880a --- /dev/null +++ b/src/helpers/bridges/MQTTBridge.h @@ -0,0 +1,52 @@ +#pragma once + +#include "helpers/bridges/BridgeBase.h" + +#ifdef WITH_MQTT_BRIDGE + +#if defined(ESP_PLATFORM) +#include +#endif + +/** + * @brief Bridge implementation using MQTT for bidirectional mesh packet transport + * + * Publishes and subscribes on a shared topic at a peer MQTT broker. Uses the same + * binary framing and XOR encryption as other bridge types for network isolation. + */ +class MQTTBridge : public BridgeBase { +private: + static MQTTBridge *_instance; + + static void mqttEventHandler(void *handler_args, esp_event_base_t base, int32_t event_id, + void *event_data); + + static const char *kBridgeTopic; + static const size_t MAX_MQTT_PACKET_SIZE = 512; + static const size_t MAX_PAYLOAD_SIZE = MAX_MQTT_PACKET_SIZE - (BRIDGE_MAGIC_SIZE + BRIDGE_CHECKSUM_SIZE); + + esp_mqtt_client_handle_t _client; + bool _connected; + bool _started; + unsigned long _next_connect_attempt; + uint8_t _reconnect_failures; + char _client_id[24]; + + void xorCrypt(uint8_t *data, size_t len); + void destroyClient(); + bool ensureClient(); + void handleMqttData(const uint8_t *data, size_t len); + void onMqttConnected(); + void onMqttDisconnected(); + +public: + MQTTBridge(NodePrefs *prefs, mesh::PacketManager *mgr, mesh::RTCClock *rtc); + + void begin() override; + void end() override; + void loop() override; + void onPacketReceived(mesh::Packet *packet) override; + void sendPacket(mesh::Packet *packet) override; +}; + +#endif diff --git a/src/helpers/web/WebPanelServer.cpp b/src/helpers/web/WebPanelServer.cpp index c58c7cbc..4c463283 100644 --- a/src/helpers/web/WebPanelServer.cpp +++ b/src/helpers/web/WebPanelServer.cpp @@ -1026,6 +1026,15 @@ const char kWebPanelAppHtml[] PROGMEM = R"HTML(
A maximum of two MQTT brokers can be enabled at once.
+ @@ -2833,6 +2842,39 @@ const char kWebPanelAppHtml[] PROGMEM = R"HTML( if (!portResult.ok) return; input.value = `${parsed.host}:${parsed.port}`; } + async function loadBridgePeerEndpoint(options = {}) { + const config = document.getElementById("bridgePeerConfig"); + if (!config) return; + const typeResult = await runCommand("get bridge.type", options); + if (!typeResult.ok || parseReplyValue(typeResult.text) !== "mqtt") { + config.style.display = "none"; + return; + } + config.style.display = ""; + const hostResult = await runCommand("get bridge.peer.host", options); + if (!hostResult.ok) return; + const portResult = await runCommand("get bridge.peer.port", options); + if (!portResult.ok) return; + const host = parseReplyValue(hostResult.text); + const port = parseReplyValue(portResult.text); + const input = document.getElementById("bridgePeerEndpoint"); + if (!input) return; + input.value = host ? `${host}:${port}` : ""; + } + async function saveBridgePeerEndpoint() { + const input = document.getElementById("bridgePeerEndpoint"); + if (!input) return; + const parsed = parseCustomEndpoint(input.value); + if (!parsed) { + statusEl.textContent = "Use host:port, for example 192.168.1.10:1883"; + return; + } + const hostResult = await runCommand("set bridge.peer.host " + parsed.host); + if (!hostResult.ok) return; + const portResult = await runCommand("set bridge.peer.port " + parsed.port); + if (!portResult.ok) return; + input.value = `${parsed.host}:${parsed.port}`; + } async function loadRadioConfig(options = {}) { const result = await runCommand("get radio", options); if (!result.ok) { @@ -3015,6 +3057,14 @@ const char kWebPanelAppHtml[] PROGMEM = R"HTML( }); document.getElementById("refreshCustomEndpointBtn").onclick = () => loadCustomEndpoint(); document.getElementById("saveCustomEndpointBtn").onclick = () => saveCustomEndpoint(); + const refreshBridgePeerEndpointBtn = document.getElementById("refreshBridgePeerEndpointBtn"); + if (refreshBridgePeerEndpointBtn) { + refreshBridgePeerEndpointBtn.onclick = () => loadBridgePeerEndpoint(); + } + const saveBridgePeerEndpointBtn = document.getElementById("saveBridgePeerEndpointBtn"); + if (saveBridgePeerEndpointBtn) { + saveBridgePeerEndpointBtn.onclick = () => saveBridgePeerEndpoint(); + } const customTransportSlider = document.getElementById("mqttCustomTransport"); if (customTransportSlider) { customTransportSlider.addEventListener("input", () => { @@ -3197,7 +3247,8 @@ const char kWebPanelAppHtml[] PROGMEM = R"HTML( () => loadBrokerState("get mqtt.custom", "mqttCustom", quiet), () => loadCustomEndpoint(quiet), () => loadCustomTransport(quiet), - () => loadField("get mqtt.custom.username", "mqttCustomUsername", null, quiet) + () => loadField("get mqtt.custom.username", "mqttCustomUsername", null, quiet), + () => loadBridgePeerEndpoint(quiet) ]); if (!isCurrentPageLoad(generation)) return; statusEl.textContent = "Ready"; diff --git a/variants/eastmesh_mqtt/platformio.ini b/variants/eastmesh_mqtt/platformio.ini index e64f4aac..694fd1eb 100644 --- a/variants/eastmesh_mqtt/platformio.ini +++ b/variants/eastmesh_mqtt/platformio.ini @@ -1156,6 +1156,6 @@ build_src_filter = ${env:WHY2025_badge_repeater_observer.build_src_filter} extends = env:Xiao_S3_WIO_repeater_observer build_flags = ${env:Xiao_S3_WIO_repeater_observer.build_flags} - -D WITH_ESPNOW_BRIDGE=1 + -D WITH_MQTT_BRIDGE=1 build_src_filter = ${env:Xiao_S3_WIO_repeater_observer.build_src_filter} - + + +