test MQTT bridge
This commit is contained in:
@@ -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 <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 <port>`: sets the peer MQTT broker port and restarts the bridge.
|
||||
- `get bridge.secret` / `set bridge.secret <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`
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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 <helpers/mqtt/MQTTUplink.h>
|
||||
#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;
|
||||
|
||||
@@ -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<unsigned char>(host[i]);
|
||||
if (c <= ' ' || c == '/' || c == ':' || c == '\\') {
|
||||
strcpy(reply, "Error: invalid host");
|
||||
return;
|
||||
}
|
||||
cleaned[oi++] = static_cast<char>(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<uint16_t>(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
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
#include <helpers/ClientACL.h>
|
||||
#include <helpers/RegionMap.h>
|
||||
|
||||
#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 {
|
||||
|
||||
@@ -0,0 +1,308 @@
|
||||
#include "MQTTBridge.h"
|
||||
|
||||
#ifdef WITH_MQTT_BRIDGE
|
||||
|
||||
#if defined(ESP_PLATFORM)
|
||||
|
||||
#include <Arduino.h>
|
||||
#include <WiFi.h>
|
||||
#include <ctype.h>
|
||||
#include <string.h>
|
||||
|
||||
#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<uint8_t>(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<unsigned>(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<MQTTBridge *>(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<esp_mqtt_event_handle_t>(event_data);
|
||||
if (event != nullptr && event->data_len > 0) {
|
||||
bridge->handleMqttData(reinterpret_cast<const uint8_t *>(event->data),
|
||||
static_cast<size_t>(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<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
|
||||
if (len > MAX_MQTT_PACKET_SIZE) {
|
||||
BRIDGE_DEBUG_PRINTLN("MQTT RX packet too large, len=%u\n", static_cast<unsigned>(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<unsigned>(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<const char *>(buffer),
|
||||
static_cast<int>(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
|
||||
@@ -0,0 +1,52 @@
|
||||
#pragma once
|
||||
|
||||
#include "helpers/bridges/BridgeBase.h"
|
||||
|
||||
#ifdef WITH_MQTT_BRIDGE
|
||||
|
||||
#if defined(ESP_PLATFORM)
|
||||
#include <mqtt_client.h>
|
||||
#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
|
||||
@@ -1026,6 +1026,15 @@ const char kWebPanelAppHtml[] PROGMEM = R"HTML(
|
||||
<div class="panel-note">A maximum of two MQTT brokers can be enabled at once.</div>
|
||||
<div id="mqttBrokerWarning" class="panel-warning"></div>
|
||||
</div>
|
||||
<div class="field-card" id="bridgePeerConfig" style="display:none">
|
||||
<label class="label" for="bridgePeerEndpoint">Mesh bridge peer MQTT host:port</label>
|
||||
<div class="inline-actions">
|
||||
<input id="bridgePeerEndpoint" placeholder="192.168.1.10:1883" maxlength="80">
|
||||
<button id="refreshBridgePeerEndpointBtn" class="iconbtn" title="Refresh bridge peer MQTT host and port">↻</button>
|
||||
<button id="saveBridgePeerEndpointBtn" class="savebtn">Save</button>
|
||||
</div>
|
||||
<div class="panel-note">Both bridge nodes must use the same peer broker, port, and bridge secret.</div>
|
||||
</div>
|
||||
</div>
|
||||
</section>
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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}
|
||||
+<helpers/bridges/ESPNowBridge.cpp>
|
||||
+<helpers/bridges/MQTTBridge.cpp>
|
||||
|
||||
Reference in New Issue
Block a user