mirror of
https://github.com/ALLFATHER-BV/wadamesh.git
synced 2026-09-28 15:38:09 +00:00
feat: add MQTT Observer profiles
This commit is contained in:
@@ -0,0 +1,143 @@
|
||||
# MQTT Observer testing on ThinkNode M9
|
||||
|
||||
The Observer feature has 39 selectable profiles. Common starting points are:
|
||||
|
||||
- **Custom broker** publishes standard MeshCore Observer v1 JSON over plain MQTT.
|
||||
- **NebraskaMesh** connects directly to `wss://mqtt.nebraskamesh.net:443/mqtt`
|
||||
using the M9 identity's Ed25519-signed JWT.
|
||||
- **MichMesh** connects to the Let's Mesh Analyzer US endpoint at
|
||||
`wss://mqtt-us-v1.letsmesh.net:443/mqtt` with the same identity authentication.
|
||||
|
||||
The profile list also includes Analyzer US/EU, NZ Analyzer, MeshMapper,
|
||||
MeshRank, WAEV, Meshomatic, CascadiaMesh, TennMesh, NashMesh, CTMesh, ChiMesh,
|
||||
Meshat.se, East Idaho Mesh, ColoradoMesh, Dutch MeshCore 1/2, MeshCore Canada
|
||||
1/2, MeshCore Finland, OkiMesh 1/2, INWMesh, BostonMesh, RFLab, IP Network UK,
|
||||
FLMesh, CoreComms, MeshTexas, Mesh Chaun14, WCMesh, Atviras Tinklas, GoMesh,
|
||||
IdahoMesh, NTXMesh and BSMesh.
|
||||
|
||||
Most built-in profiles need only a location/IATA code. Exceptions:
|
||||
|
||||
- **MeshRank** needs its account token in **Profile token** and does not use IATA.
|
||||
- **INWMesh** needs the assigned broker username and password.
|
||||
- **Mesh Chaun14** uses the device public key as username and needs its assigned password.
|
||||
- TennMesh, NashMesh and CTMesh carry their published community credentials automatically.
|
||||
|
||||
Observer reporting is receive-only. It does not transmit packets over LoRa or
|
||||
change normal WADAMESH messaging behavior.
|
||||
|
||||
Community profiles use public internet endpoints. The M9 can be on any Wi-Fi or
|
||||
phone hotspot that permits outbound DNS, HTTPS/WSS on TCP port 443, and time
|
||||
sync; it does not need to share a LAN with the Observer broker or analyzer.
|
||||
|
||||
## Open the MQTT settings
|
||||
|
||||
MQTT is hidden by default because it is experimental:
|
||||
|
||||
1. Open the app drawer, then **Store**.
|
||||
2. Select **Built-in**.
|
||||
3. Turn on **MQTT bridge**.
|
||||
4. Open **Settings > MQTT bridge**.
|
||||
|
||||
Accept the privacy warning and enable the MQTT bridge. Leave **Publish direct
|
||||
messages** off unless decoded-message forwarding is also under test.
|
||||
|
||||
## Test with a custom broker
|
||||
|
||||
Use a broker reachable from the M9's Wi-Fi. On another computer, install the
|
||||
test subscriber once:
|
||||
|
||||
```bash
|
||||
python3 -m pip install paho-mqtt
|
||||
```
|
||||
|
||||
Start it before saving the M9 settings:
|
||||
|
||||
```bash
|
||||
python3 scripts/mqtt/observer_verify.py --host 192.168.1.10 --iata TEST --count 3
|
||||
```
|
||||
|
||||
Add `--username NAME --password VALUE` if the broker requires credentials.
|
||||
|
||||
On the M9 MQTT screen set:
|
||||
|
||||
1. **Broker host / IP** to the broker address and **Port** to `1883`.
|
||||
2. **Observer profile** to **Custom broker**.
|
||||
3. Turn on **Enable Observer reporting**.
|
||||
4. Set **Observer name** to `M9_TEST`.
|
||||
5. Set **Location / IATA code** to `TEST`.
|
||||
6. Leave **Observer topic template** as `meshcore/{iata}/{device}`.
|
||||
7. Select **Save**.
|
||||
|
||||
Generate traffic with another MeshCore node: send an advert, channel message,
|
||||
or direct message over RF. The verifier passes after receiving retained online
|
||||
status and three valid packet records. Each packet is checked for a full public
|
||||
key, UTC timestamp, raw frame, packet type, route, RSSI, SNR, packet hash, and
|
||||
path formatting.
|
||||
|
||||
Open the M9 **Terminal** and run:
|
||||
|
||||
```text
|
||||
mqtt status
|
||||
```
|
||||
|
||||
Expected: `mqtt: connected`, profile `custom`, Observer `on`, sent increasing,
|
||||
queued returning to zero, and dropped remaining zero.
|
||||
|
||||
## Test persistence and queue limits
|
||||
|
||||
1. Reboot the M9 and confirm the verifier receives a new retained online status.
|
||||
2. Generate another RF packet and confirm the sent counter increases.
|
||||
3. Stop the broker or disconnect M9 Wi-Fi.
|
||||
4. Generate more than eight RF packets.
|
||||
5. Run `mqtt status`: queued must never exceed 8 and dropped should increase.
|
||||
6. Restore the broker/Wi-Fi. The queue should drain to zero and new packets
|
||||
should continue publishing without freezing navigation.
|
||||
7. Turn off **Enable Observer reporting**, save, and confirm no further packet
|
||||
messages are published.
|
||||
|
||||
## Test NebraskaMesh
|
||||
|
||||
Obtain a registered IATA-style location code from the Nebraska Mesh team first.
|
||||
Then configure the M9:
|
||||
|
||||
1. Select **NebraskaMesh** as the Observer profile.
|
||||
2. Turn on **Enable Observer reporting**.
|
||||
3. Set the assigned **Location / IATA code** and an identifiable Observer name.
|
||||
4. Leave the custom broker, credentials, topic template, and message encryption
|
||||
fields unchanged; this profile uses its built-in WSS endpoint and standard topic.
|
||||
5. Save and run `mqtt status` in Terminal.
|
||||
|
||||
Expected progression: `waiting for Wi-Fi` or `waiting for valid UTC time`, then
|
||||
`connecting`, then `connected`. Generate RF traffic and verify the sent counter
|
||||
increases with no drops. Confirm the observer appears in the Nebraska Mesh
|
||||
Analyzer under the assigned location code.
|
||||
|
||||
If it reaches `retrying`, record the complete `mqtt status` output. Error `4` or
|
||||
`5` indicates broker authentication/authorization; a persistent TCP failure
|
||||
indicates DNS, firewall, Wi-Fi, or TLS reachability rather than packet parsing.
|
||||
The status output also includes `tls`, `stack`, and `socket` error numbers.
|
||||
|
||||
## Test MichMesh
|
||||
|
||||
Choose the Michigan location code that covers the observer site. MichMesh
|
||||
currently documents these active zones:
|
||||
|
||||
- `DET` - greater Detroit
|
||||
- `FNT` - Flint / Genesee County
|
||||
- `GDW` - Gladwin and northern Michigan
|
||||
- `AZO` - Kalamazoo
|
||||
- `GRR` - Kent County / Grand Rapids
|
||||
- `MBS` - Midland / Bay City / Saginaw
|
||||
|
||||
On the M9 MQTT screen:
|
||||
|
||||
1. Select **MichMesh** as the Observer profile.
|
||||
2. Turn on **Enable Observer reporting**.
|
||||
3. Enter the appropriate zone in **Location / IATA code**.
|
||||
4. Enter an identifiable Observer name and save.
|
||||
5. Run `mqtt status` in Terminal.
|
||||
|
||||
Expected broker: `mqtt-us-v1.letsmesh.net:443`, profile `MichMesh`, followed by
|
||||
`connected`. Generate RF traffic and confirm the sent count rises while queued
|
||||
returns to zero and dropped remains zero. The observer should then appear in the
|
||||
[Let's Mesh Analyzer](https://analyzer.letsmesh.net/) under the selected zone.
|
||||
@@ -468,11 +468,12 @@
|
||||
<div class="scard">
|
||||
<div class="thumb"><img src="docs-img/set_mqtt_bridge.png" alt="MQTT bridge settings"></div>
|
||||
<h3>MQTT bridge <span class="tag">experimental</span></h3>
|
||||
<div class="lede">Forward received messages to your own MQTT broker.</div>
|
||||
<div class="lede">Forward decoded messages or MeshCore Observer packets over MQTT.</div>
|
||||
<ul>
|
||||
<li>Publishes the <b>text, sender and time</b> of messages your node hears to an MQTT broker</li>
|
||||
<li>Broker <b>host, port and topic</b>, plus an <b>encryption key</b></li>
|
||||
<li>Stays off until you <b>accept the privacy note</b> — anyone who can read the broker sees the messages, so use one you control, never a public broker</li>
|
||||
<li>Optional <b>MeshCore Observer v1</b> reporting includes raw packets, RSSI, SNR, routes, paths and packet hashes</li>
|
||||
<li>Choose from <b>39 profiles</b>: configurable MQTT, Analyzer US/EU, MeshMapper, NebraskaMesh, MichMesh and regional community brokers</li>
|
||||
<li>Stays off until you <b>accept the privacy note</b>; decoded direct-message forwarding remains separately opt-in</li>
|
||||
</ul>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
Executable
+196
@@ -0,0 +1,196 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Subscribe to a MeshCore Observer topic and validate WADAMESH payloads."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
from datetime import datetime
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
from threading import Event
|
||||
|
||||
|
||||
HEX_RE = re.compile(r"^[0-9A-Fa-f]+$")
|
||||
|
||||
|
||||
def require(condition: bool, message: str) -> None:
|
||||
if not condition:
|
||||
raise ValueError(message)
|
||||
|
||||
|
||||
def validate_common(payload: dict) -> None:
|
||||
require(isinstance(payload.get("origin"), str) and payload["origin"], "origin is missing")
|
||||
origin_id = payload.get("origin_id")
|
||||
require(isinstance(origin_id, str) and len(origin_id) == 64 and HEX_RE.fullmatch(origin_id),
|
||||
"origin_id must be a 64-character hexadecimal public key")
|
||||
timestamp = payload.get("timestamp")
|
||||
require(isinstance(timestamp, str), "timestamp is missing")
|
||||
datetime.fromisoformat(timestamp.replace("Z", "+00:00"))
|
||||
|
||||
|
||||
def validate_status(payload: dict) -> None:
|
||||
validate_common(payload)
|
||||
require(payload.get("status") in {"online", "offline"}, "invalid status value")
|
||||
require(isinstance(payload.get("firmware_version"), str) and payload["firmware_version"],
|
||||
"firmware_version is missing")
|
||||
require(payload.get("client_version") == "wadamesh-observer/1",
|
||||
"unexpected client_version")
|
||||
radio = payload.get("radio")
|
||||
require(isinstance(radio, str) and len(radio.split(",")) == 4,
|
||||
"radio must contain frequency, bandwidth, spreading factor and coding rate")
|
||||
frequency, bandwidth, spreading_factor, coding_rate = radio.split(",")
|
||||
float(frequency)
|
||||
float(bandwidth)
|
||||
int(spreading_factor)
|
||||
int(coding_rate)
|
||||
|
||||
|
||||
def validate_packet(payload: dict) -> None:
|
||||
validate_common(payload)
|
||||
require(payload.get("type") == "PACKET", "type must be PACKET")
|
||||
require(payload.get("direction") == "rx", "direction must be rx")
|
||||
require(payload.get("route") in {"F", "D", "T", "U"}, "invalid route")
|
||||
|
||||
for field in ("len", "packet_type", "payload_len", "SNR", "RSSI"):
|
||||
require(isinstance(payload.get(field), str), f"{field} must be a JSON string")
|
||||
|
||||
raw = payload.get("raw")
|
||||
require(isinstance(raw, str) and len(raw) >= 4 and len(raw) % 2 == 0 and HEX_RE.fullmatch(raw),
|
||||
"raw must be an even-length hexadecimal packet")
|
||||
wire_len = int(payload["len"])
|
||||
require(len(raw) == wire_len * 2, "raw length does not match len")
|
||||
raw_bytes = bytes.fromhex(raw)
|
||||
header = raw_bytes[0]
|
||||
route_type = header & 0x03
|
||||
packet_type = (header >> 2) & 0x0F
|
||||
require(int(payload["packet_type"]) == packet_type,
|
||||
"packet_type does not match the raw header")
|
||||
require(payload["route"] == ("D" if route_type in {2, 3} else "F"),
|
||||
"route does not match the raw header")
|
||||
|
||||
path_offset = 5 if route_type in {0, 3} else 1
|
||||
require(path_offset < len(raw_bytes), "raw packet has no path-length byte")
|
||||
path_descriptor = raw_bytes[path_offset]
|
||||
path_count = path_descriptor & 0x3F
|
||||
path_hash_size = ((path_descriptor >> 6) & 0x03) + 1
|
||||
require(path_hash_size <= 3, "reserved path hash size")
|
||||
path_start = path_offset + 1
|
||||
payload_start = path_start + path_count * path_hash_size
|
||||
require(payload_start < len(raw_bytes), "raw packet has no payload")
|
||||
packet_payload = raw_bytes[payload_start:]
|
||||
require(int(payload["payload_len"]) == len(packet_payload),
|
||||
"payload_len does not match the raw packet")
|
||||
float(payload["SNR"])
|
||||
int(payload["RSSI"])
|
||||
|
||||
packet_hash = payload.get("hash")
|
||||
require(isinstance(packet_hash, str) and len(packet_hash) == 16 and HEX_RE.fullmatch(packet_hash),
|
||||
"hash must be the 8-byte MeshCore packet hash in hexadecimal")
|
||||
hash_input = bytes([packet_type])
|
||||
if packet_type == 9:
|
||||
hash_input += bytes([path_descriptor, 0])
|
||||
expected_hash = hashlib.sha256(hash_input + packet_payload).hexdigest()[:16]
|
||||
require(packet_hash.lower() == expected_hash, "hash does not match the MeshCore packet hash")
|
||||
if "path" in payload:
|
||||
require(isinstance(payload["path"], list), "path must be an array")
|
||||
expected_path = [raw_bytes[path_start + index * path_hash_size:
|
||||
path_start + (index + 1) * path_hash_size].hex()
|
||||
for index in range(path_count)]
|
||||
require(payload["path"] == expected_path, "path does not match the raw packet")
|
||||
elif route_type in {2, 3} and path_count:
|
||||
raise ValueError("direct packet path is missing")
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--host", required=True, help="MQTT broker hostname or IP")
|
||||
parser.add_argument("--port", type=int, default=1883)
|
||||
parser.add_argument("--username")
|
||||
parser.add_argument("--password")
|
||||
parser.add_argument("--iata", default="TEST", help="Location code configured on the device")
|
||||
parser.add_argument("--topic", help="Override the subscription topic")
|
||||
parser.add_argument("--count", type=int, default=3, help="Packet payloads required for PASS")
|
||||
parser.add_argument("--timeout", type=int, default=120)
|
||||
parser.add_argument("--tls", action="store_true")
|
||||
parser.add_argument("--websockets", action="store_true")
|
||||
parser.add_argument("--ws-path", default="/mqtt")
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main() -> int:
|
||||
args = parse_args()
|
||||
try:
|
||||
import paho.mqtt.client as mqtt
|
||||
except ImportError:
|
||||
print("Install the verifier dependency: python3 -m pip install paho-mqtt", file=sys.stderr)
|
||||
return 2
|
||||
|
||||
topic = args.topic or f"meshcore/{args.iata.upper()}/+/+"
|
||||
complete = Event()
|
||||
result = {"online": False, "packets": 0, "errors": []}
|
||||
transport = "websockets" if args.websockets else "tcp"
|
||||
try:
|
||||
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, transport=transport)
|
||||
except (AttributeError, TypeError):
|
||||
client = mqtt.Client(transport=transport)
|
||||
|
||||
if args.username:
|
||||
client.username_pw_set(args.username, args.password)
|
||||
if args.tls:
|
||||
client.tls_set()
|
||||
if args.websockets:
|
||||
client.ws_set_options(path=args.ws_path)
|
||||
|
||||
def on_connect(client, _userdata, _flags, reason_code, _properties=None):
|
||||
code = getattr(reason_code, "value", reason_code)
|
||||
if code != 0:
|
||||
result["errors"].append(f"broker rejected connection: {reason_code}")
|
||||
complete.set()
|
||||
return
|
||||
client.subscribe(topic)
|
||||
print(f"Subscribed to {topic}")
|
||||
|
||||
def on_message(_client, _userdata, message):
|
||||
try:
|
||||
payload = json.loads(message.payload)
|
||||
if message.topic.endswith("/status"):
|
||||
validate_status(payload)
|
||||
if payload["status"] == "online":
|
||||
result["online"] = True
|
||||
print(f"STATUS {payload['status']} {payload['origin']} {payload['origin_id'][:12]}")
|
||||
elif message.topic.endswith("/packets"):
|
||||
validate_packet(payload)
|
||||
result["packets"] += 1
|
||||
print(f"PACKET {result['packets']}/{args.count} type={payload['packet_type']} "
|
||||
f"route={payload['route']} RSSI={payload['RSSI']} SNR={payload['SNR']}")
|
||||
if result["online"] and result["packets"] >= args.count:
|
||||
complete.set()
|
||||
except (ValueError, KeyError, TypeError, json.JSONDecodeError) as exc:
|
||||
result["errors"].append(f"{message.topic}: {exc}")
|
||||
complete.set()
|
||||
|
||||
client.on_connect = on_connect
|
||||
client.on_message = on_message
|
||||
client.connect(args.host, args.port, 60)
|
||||
client.loop_start()
|
||||
complete.wait(args.timeout)
|
||||
client.loop_stop()
|
||||
client.disconnect()
|
||||
|
||||
if result["errors"]:
|
||||
print(f"FAIL: {result['errors'][0]}", file=sys.stderr)
|
||||
return 1
|
||||
if not result["online"]:
|
||||
print("FAIL: no retained online status received", file=sys.stderr)
|
||||
return 1
|
||||
if result["packets"] < args.count:
|
||||
print(f"FAIL: received {result['packets']} of {args.count} required packets", file=sys.stderr)
|
||||
return 1
|
||||
print(f"PASS: valid online status and {result['packets']} Observer packets")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
+30
-8
@@ -339,17 +339,23 @@ static void mqttFormatStatus(char* out, size_t cap) {
|
||||
return;
|
||||
}
|
||||
if (status.phase == MqttBridgePhase::Misconfigured) {
|
||||
snprintf(out, cap, "mqtt: misconfigured (broker host is empty)");
|
||||
if (status.observerProfile == MqttObserverProfile::Custom)
|
||||
snprintf(out, cap, "mqtt: misconfigured (broker host is empty)");
|
||||
else
|
||||
snprintf(out, cap, "mqtt: misconfigured (profile needs IATA, token, or credentials)");
|
||||
return;
|
||||
}
|
||||
|
||||
const uint32_t now = millis();
|
||||
char state[64];
|
||||
char timing[96] = {};
|
||||
char timing[160] = {};
|
||||
switch (status.phase) {
|
||||
case MqttBridgePhase::WaitingForWifi:
|
||||
snprintf(state, sizeof(state), "mqtt: waiting for Wi-Fi");
|
||||
break;
|
||||
case MqttBridgePhase::WaitingForTime:
|
||||
snprintf(state, sizeof(state), "mqtt: waiting for valid UTC time");
|
||||
break;
|
||||
case MqttBridgePhase::Connecting: {
|
||||
char age[32];
|
||||
mqttFormatAge(age, sizeof(age), status.lastAttemptMs, now);
|
||||
@@ -371,8 +377,16 @@ static void mqttFormatStatus(char* out, size_t cap) {
|
||||
if (status.lastResultMs) {
|
||||
char age[32];
|
||||
mqttFormatAge(age, sizeof(age), status.lastResultMs, now);
|
||||
snprintf(timing, sizeof(timing), "last error: %d (%s), %s",
|
||||
(int)status.lastResult, mqttResultText(status.lastResult), age);
|
||||
if (status.lastTlsError || status.lastTlsStackError || status.lastSocketError) {
|
||||
snprintf(timing, sizeof(timing),
|
||||
"last error: %ld (%s), %s\ntls: %ld stack: %ld socket: %ld",
|
||||
(long)status.lastResult, mqttResultText(status.lastResult), age,
|
||||
(long)status.lastTlsError, (long)status.lastTlsStackError,
|
||||
(long)status.lastSocketError);
|
||||
} else {
|
||||
snprintf(timing, sizeof(timing), "last error: %ld (%s), %s",
|
||||
(long)status.lastResult, mqttResultText(status.lastResult), age);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
@@ -380,10 +394,15 @@ static void mqttFormatStatus(char* out, size_t cap) {
|
||||
snprintf(state, sizeof(state), "mqtt: unavailable on this build");
|
||||
break;
|
||||
}
|
||||
snprintf(out, cap, "%s\nbroker: %s:%u\npublish: channels=%s dm=%s encrypted=%s%s%s",
|
||||
snprintf(out, cap,
|
||||
"%s\nbroker: %s:%u\nprofile: %s\npublish: channels=%s dm=%s observer=%s encrypted=%s"
|
||||
"\nobserver: %lu sent, %u queued, %lu dropped%s%s",
|
||||
state, status.host, (unsigned)status.port,
|
||||
MqttBridge::observerProfileLabel(status.observerProfile),
|
||||
status.publishChannel ? "on" : "off", status.publishDm ? "on" : "off",
|
||||
status.encrypted ? "yes" : "no", timing[0] ? "\n" : "", timing);
|
||||
status.publishObserver ? "on" : "off", status.encrypted ? "yes" : "no",
|
||||
(unsigned long)status.observerPublished, (unsigned)status.observerQueued,
|
||||
(unsigned long)status.observerDropped, timing[0] ? "\n" : "", timing);
|
||||
}
|
||||
#endif
|
||||
|
||||
@@ -827,7 +846,7 @@ bool MyMesh::handleMeshcomodCommand(const char* text, int text_len) {
|
||||
while (*q == ' ' || *q == '\t') q++;
|
||||
if (isCmd(q, "status")) {
|
||||
#if defined(ESP32) && defined(MULTI_TRANSPORT_COMPANION)
|
||||
char reply[320];
|
||||
char reply[512];
|
||||
mqttFormatStatus(reply, sizeof(reply));
|
||||
pushMeshcomodReply(reply);
|
||||
#else
|
||||
@@ -2319,6 +2338,9 @@ void MyMesh::companionRetryService() {
|
||||
|
||||
void MyMesh::logRxRaw(float snr, float rssi, const uint8_t raw[], int len) {
|
||||
companionRetryObserveRaw(raw, len);
|
||||
#if defined(ESP32) && defined(MULTI_TRANSPORT_COMPANION)
|
||||
mqtt_bridge.observeRx(snr, rssi, raw, len);
|
||||
#endif
|
||||
|
||||
const int8_t snr_q4 = (int8_t)(snr * 4.0f);
|
||||
const uint32_t now_ms = millis();
|
||||
@@ -5456,7 +5478,7 @@ void MyMesh::checkCLIRescueCmd() {
|
||||
strcmp(cli_command, "get mqtt.status") == 0 ||
|
||||
strcmp(cli_command, "mqtt status") == 0) {
|
||||
#if defined(ESP32) && defined(MULTI_TRANSPORT_COMPANION)
|
||||
char reply[320];
|
||||
char reply[512];
|
||||
mqttFormatStatus(reply, sizeof(reply));
|
||||
Serial.printf(" > %s\n", reply);
|
||||
#else
|
||||
|
||||
@@ -5,19 +5,212 @@
|
||||
// makes MqttBridge a no-op stub on the Tanmatsu, and the real bridge below on all other boards.
|
||||
MqttBridge mqtt_bridge;
|
||||
|
||||
namespace {
|
||||
enum class ObserverAuth : uint8_t { None, UserPass, Jwt };
|
||||
enum class ObserverTopicStyle : uint8_t { MeshCore, MeshRank };
|
||||
|
||||
struct ObserverPreset {
|
||||
const char* label;
|
||||
const char* name;
|
||||
const char* uri;
|
||||
const char* host;
|
||||
uint16_t port;
|
||||
const char* audience;
|
||||
ObserverAuth auth;
|
||||
ObserverTopicStyle topicStyle;
|
||||
uint32_t tokenLifetime;
|
||||
uint16_t keepalive;
|
||||
bool allowRetain;
|
||||
bool requiresIata;
|
||||
const char* username;
|
||||
const char* password;
|
||||
};
|
||||
|
||||
constexpr const char* OBSERVER_PUBKEY_USERNAME = "{pubkey}";
|
||||
constexpr ObserverPreset OBSERVER_PRESETS[] = {
|
||||
{"Custom broker", "custom", nullptr, nullptr, 0, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 60, true, false, nullptr, nullptr},
|
||||
{"NebraskaMesh", "nebraskamesh", "wss://mqtt.nebraskamesh.net:443/mqtt", "mqtt.nebraskamesh.net", 443, "mqtt.nebraskamesh.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MichMesh (Analyzer US)", "michmesh", "wss://mqtt-us-v1.letsmesh.net:443/mqtt", "mqtt-us-v1.letsmesh.net", 443, "mqtt-us-v1.letsmesh.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Analyzer US", "analyzer-us", "wss://mqtt-us-v1.letsmesh.net:443/mqtt", "mqtt-us-v1.letsmesh.net", 443, "mqtt-us-v1.letsmesh.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Analyzer EU", "analyzer-eu", "wss://mqtt-eu-v1.letsmesh.net:443/mqtt", "mqtt-eu-v1.letsmesh.net", 443, "mqtt-eu-v1.letsmesh.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"NZ Analyzer", "nz-analyzer", "wss://meshcore-mqtt-1.baird.io:443", "meshcore-mqtt-1.baird.io", 443, "meshcore-mqtt-1.baird.io", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshMapper", "meshmapper", "wss://mqtt.meshmapper.net:443/mqtt", "mqtt.meshmapper.net", 443, "mqtt.meshmapper.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshRank", "meshrank", "mqtts://meshrank.net:8883", "meshrank.net", 8883, nullptr, ObserverAuth::None, ObserverTopicStyle::MeshRank, 0, 60, false, false, nullptr, nullptr},
|
||||
{"WAEV", "waev", "wss://mqtt.waev.app:443/mqtt", "mqtt.waev.app", 443, "mqtt.waev.app", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 3300, 55, false, true, nullptr, nullptr},
|
||||
{"Meshomatic", "meshomatic", "wss://us-east.meshomatic.net:443/mqtt", "us-east.meshomatic.net", 443, "us-east.meshomatic.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"CascadiaMesh", "cascadiamesh", "wss://mqtt-v1.cascadiamesh.org:443/mqtt", "mqtt-v1.cascadiamesh.org", 443, "mqtt-v1.cascadiamesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"TennMesh", "tennmesh", "mqtt://mqtt.tennmesh.com:1883", "mqtt.tennmesh.com", 1883, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 55, true, true, "mqttfeed", "tc2live"},
|
||||
{"NashMesh", "nashmesh", "mqtt://mqtt.nashme.sh:1883", "mqtt.nashme.sh", 1883, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 55, true, true, "meshdev", "large4cats"},
|
||||
{"CTMesh", "ctmesh", "mqtt://mqtt.ctmesh.org:1883", "mqtt.ctmesh.org", 1883, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 60, true, true, "meshdev", "large4cats"},
|
||||
{"ChiMesh", "chimesh", "wss://mqtt.chimesh.org:443", "mqtt.chimesh.org", 443, "mqtt.chimesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Meshat.se", "meshat.se", "wss://meshcore-mqtt.meshat.se:443", "meshcore-mqtt.meshat.se", 443, "meshcore-mqtt.meshat.se", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"East Idaho Mesh", "eastidahomesh", "mqtt://live.eastidahomesh.com:1883", "live.eastidahomesh.com", 1883, nullptr, ObserverAuth::None, ObserverTopicStyle::MeshCore, 0, 55, true, true, nullptr, nullptr},
|
||||
{"ColoradoMesh", "coloradomesh", "wss://mqtt.meshcore.coloradomesh.org:443", "mqtt.meshcore.coloradomesh.org", 443, "mqtt.meshcore.coloradomesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Dutch MeshCore 1", "dutchmeshcore-1", "wss://collector1.dutchmeshcore.nl:443/mqtt", "collector1.dutchmeshcore.nl", 443, "collector1.dutchmeshcore.nl", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Dutch MeshCore 2", "dutchmeshcore-2", "wss://collector2.dutchmeshcore.nl:443/mqtt", "collector2.dutchmeshcore.nl", 443, "collector2.dutchmeshcore.nl", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshCore Canada 1", "meshcore-ca-1", "wss://mqtt1.meshcore.ca:443/mqtt", "mqtt1.meshcore.ca", 443, "mqtt1.meshcore.ca", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshCore Canada 2", "meshcore-ca-2", "wss://mqtt2.meshcore.ca:443/mqtt", "mqtt2.meshcore.ca", 443, "mqtt2.meshcore.ca", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshCore Finland", "meshcore-fi", "wss://mc-mqtt.meshcore.fi:443/", "mc-mqtt.meshcore.fi", 443, "mc-mqtt.meshcore.fi", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"OkiMesh 1", "okimesh-1", "wss://mqtt1.okimesh.org:9002/mqtt", "mqtt1.okimesh.org", 9002, "mqtt1.okimesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"OkiMesh 2", "okimesh-2", "wss://mqtt2.okimesh.org:9002/mqtt", "mqtt2.okimesh.org", 9002, "mqtt2.okimesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"INWMesh", "inwmesh", "mqtts://scope.inwmesh.org:8883", "scope.inwmesh.org", 8883, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 55, true, true, nullptr, nullptr},
|
||||
{"BostonMesh", "bostonmesh", "wss://mqttmc01.bostonme.sh:443/mqtt", "mqttmc01.bostonme.sh", 443, "mqttmc01.bostonme.sh", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"RFLab", "rflab", "wss://mqtt.rflab.io:443", "mqtt.rflab.io", 443, "mqtt.rflab.io", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"IP Network UK", "ipnt.uk", "wss://mqtt.ipnt.uk:443", "mqtt.ipnt.uk", 443, "mqtt.ipnt.uk", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"FLMesh", "flmesh", "wss://mcmqtt.jntconnections.com:443", "mcmqtt.jntconnections.com", 443, "mcmqtt.jntconnections.com", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"CoreComms", "corecomms", "wss://mqtt.corecomms.net:443/mqtt", "mqtt.corecomms.net", 443, "mqtt.corecomms.net", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"MeshTexas", "meshtexas", "wss://mqtt.meshtexas.org:443/mqtt", "mqtt.meshtexas.org", 443, "mqtt.meshtexas.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Mesh Chaun14", "mesh-chaun14", "mqtt://mqtt.mesh.chaun14.fr:1884", "mqtt.mesh.chaun14.fr", 1884, nullptr, ObserverAuth::UserPass, ObserverTopicStyle::MeshCore, 0, 60, true, true, OBSERVER_PUBKEY_USERNAME, nullptr},
|
||||
{"WCMesh", "wcmesh", "wss://mqtt.wcmesh.com:443", "mqtt.wcmesh.com", 443, "mqtt.wcmesh.com", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"Atviras Tinklas", "atvirastinklas", "wss://mqtt-mc.atvirastinklas.lt:443", "mqtt-mc.atvirastinklas.lt", 443, "mqtt-mc.atvirastinklas.lt", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"GoMesh", "gomesh", "wss://mqtt.gomesh.dev:443", "mqtt.gomesh.dev", 443, "mqtt.gomesh.dev", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"IdahoMesh", "idahomesh", "wss://mqtt.idahomesh.org:443/mqtt", "mqtt.idahomesh.org", 443, "mqtt.idahomesh.org", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"NTXMesh", "ntxmesh", "wss://ntxmesh.dhovin.me:8883", "ntxmesh.dhovin.me", 8883, "ntxmesh.dhovin.me", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
{"BSMesh", "bsmesh", "wss://mqtt.bsmesh.de:8885", "mqtt.bsmesh.de", 8885, "mqtt.bsmesh.de", ObserverAuth::Jwt, ObserverTopicStyle::MeshCore, 86400, 55, true, true, nullptr, nullptr},
|
||||
};
|
||||
|
||||
static_assert(sizeof(OBSERVER_PRESETS) / sizeof(OBSERVER_PRESETS[0]) ==
|
||||
(size_t)MqttObserverProfile::Count,
|
||||
"Observer profile enum and preset table must stay aligned");
|
||||
|
||||
const ObserverPreset& observerPreset(MqttObserverProfile profile) {
|
||||
const uint8_t index = (uint8_t)profile;
|
||||
return OBSERVER_PRESETS[index < (uint8_t)MqttObserverProfile::Count ? index : 0];
|
||||
}
|
||||
|
||||
bool isBuiltInObserverProfile(MqttObserverProfile profile) {
|
||||
return profile != MqttObserverProfile::Custom &&
|
||||
(uint8_t)profile < (uint8_t)MqttObserverProfile::Count;
|
||||
}
|
||||
} // namespace
|
||||
|
||||
uint8_t MqttBridge::observerProfileCount() {
|
||||
return (uint8_t)MqttObserverProfile::Count;
|
||||
}
|
||||
|
||||
const char* MqttBridge::observerProfileOptions() {
|
||||
static char options[768] = {};
|
||||
if (!options[0]) {
|
||||
size_t pos = 0;
|
||||
for (uint8_t i = 0; i < observerProfileCount(); ++i) {
|
||||
const char* label = OBSERVER_PRESETS[i].label;
|
||||
const int written = snprintf(options + pos, sizeof(options) - pos,
|
||||
"%s%s", i ? "\n" : "", label);
|
||||
if (written <= 0 || (size_t)written >= sizeof(options) - pos) break;
|
||||
pos += (size_t)written;
|
||||
}
|
||||
}
|
||||
return options;
|
||||
}
|
||||
|
||||
const char* MqttBridge::observerProfileLabel(MqttObserverProfile profile) {
|
||||
return observerPreset(profile).label;
|
||||
}
|
||||
|
||||
MqttObserverProfile MqttBridge::observerProfileFromIndex(uint16_t index) {
|
||||
return index < observerProfileCount() ? (MqttObserverProfile)index
|
||||
: MqttObserverProfile::Custom;
|
||||
}
|
||||
|
||||
uint16_t MqttBridge::observerProfileIndex(MqttObserverProfile profile) {
|
||||
const uint8_t index = (uint8_t)profile;
|
||||
return index < observerProfileCount() ? index : 0;
|
||||
}
|
||||
|
||||
bool MqttBridge::observerProfileNeedsIata(MqttObserverProfile profile) {
|
||||
return observerPreset(profile).requiresIata;
|
||||
}
|
||||
|
||||
bool MqttBridge::observerProfileNeedsToken(MqttObserverProfile profile) {
|
||||
return observerPreset(profile).topicStyle == ObserverTopicStyle::MeshRank;
|
||||
}
|
||||
|
||||
bool MqttBridge::observerProfileNeedsUsername(MqttObserverProfile profile) {
|
||||
const ObserverPreset& preset = observerPreset(profile);
|
||||
return preset.auth == ObserverAuth::UserPass && preset.username == nullptr;
|
||||
}
|
||||
|
||||
bool MqttBridge::observerProfileNeedsPassword(MqttObserverProfile profile) {
|
||||
const ObserverPreset& preset = observerPreset(profile);
|
||||
return preset.auth == ObserverAuth::UserPass && preset.password == nullptr;
|
||||
}
|
||||
|
||||
#if !defined(HAS_TANMATSU) && !defined(HAS_TDISPLAY_P4) // ---- real PubSubClient/WiFiClient implementation; NOT built on the P4 boards ----
|
||||
#include "SdNvsPrefs.h" // NOT raw NVS Preferences: the touch firmware abandoned NVS (tiny, shared,
|
||||
// doesn't survive Launcher) for a file-backed store. MQTT was the last setting
|
||||
// still on NVS, so its writes silently failed / didn't persist (GH #128).
|
||||
#include <Identity.h>
|
||||
#include <Packet.h>
|
||||
#include <WiFi.h>
|
||||
#include <ctype.h>
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <time.h>
|
||||
#include <lwip/sockets.h>
|
||||
#include "esp_random.h"
|
||||
#include "mqtt_client.h"
|
||||
#include "mbedtls/gcm.h"
|
||||
#include "mbedtls/md.h"
|
||||
#include "mbedtls/base64.h"
|
||||
|
||||
extern "C" esp_err_t esp_crt_bundle_attach(void* conf);
|
||||
|
||||
namespace {
|
||||
constexpr const char* OBSERVER_DEFAULT_TOPIC = "meshcore/{iata}/{device}";
|
||||
constexpr uint32_t OBSERVER_STATUS_INTERVAL_MS = 300000;
|
||||
constexpr uint32_t OBSERVER_TOKEN_REFRESH_SECS = 300;
|
||||
|
||||
void appendTopicText(char* out, size_t outCap, size_t& pos, const char* text) {
|
||||
if (!out || outCap == 0 || !text) return;
|
||||
while (*text && pos + 1 < outCap) out[pos++] = *text++;
|
||||
out[pos] = '\0';
|
||||
}
|
||||
|
||||
void formatObserverTimestamp(uint32_t epoch, char* out, size_t outCap,
|
||||
char* timeOut = nullptr, size_t timeCap = 0,
|
||||
char* dateOut = nullptr, size_t dateCap = 0) {
|
||||
time_t value = epoch ? (time_t)epoch : time(nullptr);
|
||||
struct tm utc{};
|
||||
if (!gmtime_r(&value, &utc)) {
|
||||
snprintf(out, outCap, "1970-01-01T00:00:00.000000+00:00");
|
||||
if (timeOut && timeCap) snprintf(timeOut, timeCap, "00:00:00");
|
||||
if (dateOut && dateCap) snprintf(dateOut, dateCap, "01/01/1970");
|
||||
return;
|
||||
}
|
||||
char base[24];
|
||||
strftime(base, sizeof(base), "%Y-%m-%dT%H:%M:%S", &utc);
|
||||
snprintf(out, outCap, "%s.000000+00:00", base);
|
||||
if (timeOut && timeCap) strftime(timeOut, timeCap, "%H:%M:%S", &utc);
|
||||
if (dateOut && dateCap) strftime(dateOut, dateCap, "%d/%m/%Y", &utc);
|
||||
}
|
||||
|
||||
void bytesToHex(const uint8_t* data, size_t len, char* out, size_t outCap, bool uppercase) {
|
||||
const char* digits = uppercase ? "0123456789ABCDEF" : "0123456789abcdef";
|
||||
if (!out || outCap == 0) return;
|
||||
out[0] = '\0';
|
||||
if (!data || outCap < len * 2 + 1) return;
|
||||
for (size_t i = 0; i < len; ++i) {
|
||||
out[i * 2] = digits[data[i] >> 4];
|
||||
out[i * 2 + 1] = digits[data[i] & 0x0F];
|
||||
}
|
||||
out[len * 2] = '\0';
|
||||
}
|
||||
|
||||
size_t base64UrlEncode(const uint8_t* input, size_t inputLen, char* output, size_t outputCap) {
|
||||
if (!input || !output || outputCap == 0) return 0;
|
||||
size_t written = 0;
|
||||
if (mbedtls_base64_encode((unsigned char*)output, outputCap - 1, &written,
|
||||
input, inputLen) != 0) return 0;
|
||||
for (size_t i = 0; i < written; ++i) {
|
||||
if (output[i] == '+') output[i] = '-';
|
||||
else if (output[i] == '/') output[i] = '_';
|
||||
}
|
||||
while (written > 0 && output[written - 1] == '=') --written;
|
||||
output[written] = '\0';
|
||||
return written;
|
||||
}
|
||||
} // namespace
|
||||
|
||||
// Only hand bytes to lwIP when the socket can accept them right now (zero-timeout
|
||||
// select). Otherwise WiFiClient::write() blocks in 1 s select() retries against a
|
||||
// broker that stopped ACKing — on the loop thread that is a visible UI freeze per
|
||||
@@ -47,11 +240,22 @@ void MqttBridge::syncStatusConfig(MqttBridgePhase phase) {
|
||||
next.phase = phase;
|
||||
next.available = true;
|
||||
next.requestedEnabled = _requestedEnabled;
|
||||
next.publishDm = _pubDm;
|
||||
next.publishChannel = _pubChannel;
|
||||
next.encrypted = _encOn;
|
||||
strncpy(next.host, _host, sizeof(next.host) - 1);
|
||||
next.port = _port;
|
||||
next.publishDm = _plainEnabled && _pubDm;
|
||||
next.publishChannel = _plainEnabled && _pubChannel;
|
||||
next.publishObserver = _pubObserver;
|
||||
next.observerProfile = _observerProfile;
|
||||
next.encrypted = _plainEnabled && _encOn;
|
||||
if (isBuiltInObserverProfile(_observerProfile)) {
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
strncpy(next.host, preset.host, sizeof(next.host) - 1);
|
||||
next.port = preset.port;
|
||||
} else {
|
||||
strncpy(next.host, _host, sizeof(next.host) - 1);
|
||||
next.port = _port;
|
||||
}
|
||||
next.observerPublished = _observerPublished;
|
||||
next.observerDropped = _observerDropped;
|
||||
next.observerQueued = _observerQueueCount;
|
||||
portENTER_CRITICAL(&_statusMux);
|
||||
_status = next;
|
||||
portEXIT_CRITICAL(&_statusMux);
|
||||
@@ -76,7 +280,11 @@ void MqttBridge::setStatusPhase(MqttBridgePhase phase, int result, uint32_t now,
|
||||
// USB-CDC companion stream (see TouchPrefsStore prefsGetStr note).
|
||||
void MqttBridge::loadConfig() {
|
||||
_requestedEnabled = false; _enabled = false; _pubDm = false; _pubChannel = true;
|
||||
_host[0] = _user[0] = _pwd[0] = _psk[0] = '\0';
|
||||
_pubObserver = false; _plainEnabled = false;
|
||||
_observerProfile = MqttObserverProfile::Custom;
|
||||
_host[0] = _user[0] = _pwd[0] = _psk[0] = _observerOrigin[0] = _observerIata[0] =
|
||||
_observerToken[0] = '\0';
|
||||
snprintf(_observerTopic, sizeof(_observerTopic), "%s", OBSERVER_DEFAULT_TOPIC);
|
||||
_port = 1883;
|
||||
|
||||
SdNvsPrefs p;
|
||||
@@ -85,16 +293,39 @@ void MqttBridge::loadConfig() {
|
||||
_port = (uint16_t)p.getUInt("port", 1883);
|
||||
_pubDm = p.getBool("dm", false);
|
||||
_pubChannel = p.getBool("ch", true);
|
||||
_pubObserver = p.getBool("obs", false);
|
||||
const uint32_t profile = p.getUInt("obs_profile", 0);
|
||||
if (profile < (uint32_t)MqttObserverProfile::Count)
|
||||
_observerProfile = (MqttObserverProfile)profile;
|
||||
if (p.isKey("host")) p.getString("host", _host, sizeof(_host));
|
||||
if (p.isKey("user")) p.getString("user", _user, sizeof(_user));
|
||||
if (p.isKey("pwd")) p.getString("pwd", _pwd, sizeof(_pwd));
|
||||
if (p.isKey("psk")) p.getString("psk", _psk, sizeof(_psk));
|
||||
if (p.isKey("obs_origin")) p.getString("obs_origin", _observerOrigin, sizeof(_observerOrigin));
|
||||
if (p.isKey("obs_iata")) p.getString("obs_iata", _observerIata, sizeof(_observerIata));
|
||||
if (p.isKey("obs_topic")) p.getString("obs_topic", _observerTopic, sizeof(_observerTopic));
|
||||
if (p.isKey("obs_token")) p.getString("obs_token", _observerToken, sizeof(_observerToken));
|
||||
p.end();
|
||||
}
|
||||
if (_observerOrigin[0] == '\0') snprintf(_observerOrigin, sizeof(_observerOrigin), "%s", _nodeName);
|
||||
if (_observerTopic[0] == '\0') snprintf(_observerTopic, sizeof(_observerTopic), "%s", OBSERVER_DEFAULT_TOPIC);
|
||||
deriveKey();
|
||||
_enabled = _requestedEnabled && _host[0] != '\0';
|
||||
_plainEnabled = _requestedEnabled && _observerProfile == MqttObserverProfile::Custom && _host[0] != '\0';
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
const bool hasLocation = !preset.requiresIata || _observerIata[0] != '\0';
|
||||
bool hasCredentials = true;
|
||||
if (preset.topicStyle == ObserverTopicStyle::MeshRank) {
|
||||
hasCredentials = _observerToken[0] != '\0';
|
||||
} else if (preset.auth == ObserverAuth::UserPass) {
|
||||
const bool hasUsername = preset.username != nullptr || _user[0] != '\0';
|
||||
const bool hasPassword = preset.password != nullptr || _pwd[0] != '\0';
|
||||
hasCredentials = hasUsername && hasPassword;
|
||||
}
|
||||
const bool cloudEnabled = _requestedEnabled && _pubObserver && hasLocation && hasCredentials &&
|
||||
isBuiltInObserverProfile(_observerProfile);
|
||||
_enabled = _plainEnabled || cloudEnabled;
|
||||
MqttBridgePhase phase = MqttBridgePhase::Disabled;
|
||||
if (_requestedEnabled && _host[0] == '\0') phase = MqttBridgePhase::Misconfigured;
|
||||
if (_requestedEnabled && !_enabled) phase = MqttBridgePhase::Misconfigured;
|
||||
else if (_enabled) phase = WiFi.status() == WL_CONNECTED
|
||||
? MqttBridgePhase::RetryWait
|
||||
: MqttBridgePhase::WaitingForWifi;
|
||||
@@ -113,42 +344,77 @@ void MqttBridge::deriveKey() {
|
||||
}
|
||||
}
|
||||
|
||||
void MqttBridge::begin(const char* nodeHex) {
|
||||
strncpy(_nodeHex, nodeHex, sizeof(_nodeHex) - 1);
|
||||
_nodeHex[sizeof(_nodeHex) - 1] = '\0';
|
||||
void MqttBridge::begin(const mesh::LocalIdentity* identity, const char* nodeName,
|
||||
float frequencyMhz, float bandwidthKhz, uint8_t spreadingFactor,
|
||||
uint8_t codingRate, const char* firmwareVersion) {
|
||||
_identity = identity;
|
||||
if (_identity) {
|
||||
bytesToHex(_identity->pub_key, 6, _nodeHex, sizeof(_nodeHex), false);
|
||||
bytesToHex(_identity->pub_key, 32, _nodeId, sizeof(_nodeId), true);
|
||||
}
|
||||
snprintf(_nodeName, sizeof(_nodeName), "%s", nodeName && nodeName[0] ? nodeName : "WADAMESH");
|
||||
snprintf(_observerRadio, sizeof(_observerRadio), "%.6f,%.1f,%u,%u",
|
||||
(double)frequencyMhz, (double)bandwidthKhz,
|
||||
(unsigned)spreadingFactor, (unsigned)codingRate);
|
||||
snprintf(_firmwareVersion, sizeof(_firmwareVersion), "%s",
|
||||
firmwareVersion && firmwareVersion[0] ? firmwareVersion : "wadamesh");
|
||||
|
||||
loadConfig();
|
||||
if (!_enabled) return;
|
||||
|
||||
_mqtt.setServer(_host, _port);
|
||||
_mqtt.setKeepAlive(60);
|
||||
_mqtt.setSocketTimeout(2); // bound connect/read — a dead broker must not stall loop()
|
||||
_mqtt.setBufferSize(_encOn ? 1024 : 512); // sealed base64 needs the bigger buffer
|
||||
Serial.printf("[MQTT] configured -> %s:%u dm=%d ch=%d enc=%d\n",
|
||||
_host, _port, (int)_pubDm, (int)_pubChannel, (int)_encOn);
|
||||
if (_plainEnabled) {
|
||||
_mqtt.setServer(_host, _port);
|
||||
_mqtt.setKeepAlive(60);
|
||||
_mqtt.setSocketTimeout(2); // bound connect/read — a dead broker must not stall loop()
|
||||
_mqtt.setBufferSize(_pubObserver ? 1536 : (_encOn ? 1024 : 512));
|
||||
}
|
||||
Serial.printf("[MQTT] configured -> profile=%s host=%s:%u dm=%d ch=%d obs=%d enc=%d\n",
|
||||
observerPreset(_observerProfile).name,
|
||||
isBuiltInObserverProfile(_observerProfile) ? observerPreset(_observerProfile).host : _host,
|
||||
isBuiltInObserverProfile(_observerProfile) ?
|
||||
(unsigned)observerPreset(_observerProfile).port : (unsigned)_port,
|
||||
(int)_pubDm, (int)_pubChannel, (int)_pubObserver, (int)_encOn);
|
||||
}
|
||||
|
||||
bool MqttBridge::reconnect() {
|
||||
if (!_enabled) return false;
|
||||
if (!_plainEnabled) return false;
|
||||
if (WiFi.status() != WL_CONNECTED) {
|
||||
setStatusPhase(MqttBridgePhase::WaitingForWifi, MQTT_DISCONNECTED,
|
||||
millis(), false);
|
||||
return false;
|
||||
}
|
||||
|
||||
char clientId[32], lwtTopic[80];
|
||||
char clientId[80], lwtTopic[192], lwtPayload[512];
|
||||
snprintf(clientId, sizeof(clientId), "wadamesh-%s", _nodeHex);
|
||||
snprintf(lwtTopic, sizeof(lwtTopic), "wadamesh/%s/status", _nodeHex);
|
||||
if (_pubObserver) {
|
||||
observerTopic(lwtTopic, sizeof(lwtTopic), "status");
|
||||
char timestamp[40], safeOrigin[72];
|
||||
formatObserverTimestamp((uint32_t)time(nullptr), timestamp, sizeof(timestamp));
|
||||
escapeJson(_observerOrigin, safeOrigin, sizeof(safeOrigin));
|
||||
snprintf(lwtPayload, sizeof(lwtPayload),
|
||||
"{\"status\":\"offline\",\"timestamp\":\"%s\",\"origin\":\"%s\","
|
||||
"\"origin_id\":\"%s\",\"firmware_version\":\"%s\","
|
||||
"\"radio\":\"%s\",\"client_version\":\"wadamesh-observer/1\"}",
|
||||
timestamp, safeOrigin, _nodeId, _firmwareVersion, _observerRadio);
|
||||
} else {
|
||||
snprintf(lwtTopic, sizeof(lwtTopic), "wadamesh/%s/status", _nodeHex);
|
||||
snprintf(lwtPayload, sizeof(lwtPayload), "offline");
|
||||
}
|
||||
|
||||
bool ok = _user[0]
|
||||
? _mqtt.connect(clientId, _user, _pwd, lwtTopic, 0, true, "offline")
|
||||
: _mqtt.connect(clientId, nullptr, nullptr, lwtTopic, 0, true, "offline");
|
||||
? _mqtt.connect(clientId, _user, _pwd, lwtTopic, 0, true, lwtPayload)
|
||||
: _mqtt.connect(clientId, nullptr, nullptr, lwtTopic, 0, true, lwtPayload);
|
||||
|
||||
const uint32_t now = millis();
|
||||
const int result = _mqtt.state();
|
||||
if (ok) {
|
||||
setStatusPhase(MqttBridgePhase::Connected, result, now, true);
|
||||
_mqtt.publish(lwtTopic, "online", true);
|
||||
if (_pubObserver) publishObserverStatus(true);
|
||||
else {
|
||||
char legacyStatusTopic[80];
|
||||
snprintf(legacyStatusTopic, sizeof(legacyStatusTopic), "wadamesh/%s/status", _nodeHex);
|
||||
_mqtt.publish(legacyStatusTopic, "online", true);
|
||||
}
|
||||
Serial.printf("[MQTT] connected as %s\n", clientId);
|
||||
} else {
|
||||
MqttBridgeStatus current;
|
||||
@@ -173,9 +439,208 @@ void MqttBridge::reconnectTask(void* arg) {
|
||||
vTaskDelete(nullptr);
|
||||
}
|
||||
|
||||
bool MqttBridge::createObserverJwt(char* token, size_t tokenCap) const {
|
||||
if (!_identity || !token || tokenCap == 0) return false;
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
if (preset.auth != ObserverAuth::Jwt || !preset.audience) return false;
|
||||
const uint32_t now = (uint32_t)time(nullptr);
|
||||
if (now < 1700000000UL) return false;
|
||||
|
||||
static const char headerJson[] = "{\"alg\":\"Ed25519\",\"typ\":\"JWT\"}";
|
||||
char payloadJson[256];
|
||||
const int payloadLen = snprintf(payloadJson, sizeof(payloadJson),
|
||||
"{\"publicKey\":\"%s\",\"aud\":\"%s\",\"iat\":%lu,\"exp\":%lu,"
|
||||
"\"client\":\"wadamesh-observer/1\"}",
|
||||
_nodeId, preset.audience, (unsigned long)now,
|
||||
(unsigned long)(now + preset.tokenLifetime));
|
||||
if (payloadLen <= 0 || (size_t)payloadLen >= sizeof(payloadJson)) return false;
|
||||
|
||||
char header[64], payload[384], signingInput[512];
|
||||
const size_t headerLen = base64UrlEncode((const uint8_t*)headerJson,
|
||||
strlen(headerJson), header, sizeof(header));
|
||||
const size_t encodedPayloadLen = base64UrlEncode((const uint8_t*)payloadJson,
|
||||
(size_t)payloadLen, payload, sizeof(payload));
|
||||
if (headerLen == 0 || encodedPayloadLen == 0) return false;
|
||||
const int signingLen = snprintf(signingInput, sizeof(signingInput), "%s.%s", header, payload);
|
||||
if (signingLen <= 0 || (size_t)signingLen >= sizeof(signingInput)) return false;
|
||||
|
||||
uint8_t signature[SIGNATURE_SIZE];
|
||||
_identity->sign(signature, (const uint8_t*)signingInput, signingLen);
|
||||
char signatureHex[SIGNATURE_SIZE * 2 + 1];
|
||||
bytesToHex(signature, sizeof(signature), signatureHex, sizeof(signatureHex), true);
|
||||
const int tokenLen = snprintf(token, tokenCap, "%s.%s", signingInput, signatureHex);
|
||||
return tokenLen > 0 && (size_t)tokenLen < tokenCap;
|
||||
}
|
||||
|
||||
void MqttBridge::observerCloudEvent(void* arg, const char* base, int32_t eventId, void* eventData) {
|
||||
(void)base;
|
||||
MqttBridge* self = static_cast<MqttBridge*>(arg);
|
||||
if (!self) return;
|
||||
esp_mqtt_event_handle_t event = static_cast<esp_mqtt_event_handle_t>(eventData);
|
||||
const uint32_t now = millis();
|
||||
switch ((esp_mqtt_event_id_t)eventId) {
|
||||
case MQTT_EVENT_CONNECTED:
|
||||
self->_observerCloudConnected = true;
|
||||
self->_observerCloudRetryAtMs = 0;
|
||||
self->setStatusPhase(MqttBridgePhase::Connected, 0, now, true);
|
||||
portENTER_CRITICAL(&self->_statusMux);
|
||||
self->_status.lastTlsError = 0;
|
||||
self->_status.lastTlsStackError = 0;
|
||||
self->_status.lastSocketError = 0;
|
||||
portEXIT_CRITICAL(&self->_statusMux);
|
||||
self->publishObserverStatus(true);
|
||||
break;
|
||||
case MQTT_EVENT_DISCONNECTED:
|
||||
self->_observerCloudConnected = false;
|
||||
self->setStatusPhase(MqttBridgePhase::RetryWait, -1, now, false,
|
||||
now + RECONNECT_INTERVAL_MS);
|
||||
break;
|
||||
case MQTT_EVENT_ERROR: {
|
||||
self->_observerCloudConnected = false;
|
||||
int result = -2;
|
||||
int32_t tlsError = 0, tlsStackError = 0, socketError = 0;
|
||||
if (event && event->error_handle &&
|
||||
event->error_handle->error_type == MQTT_ERROR_TYPE_CONNECTION_REFUSED) {
|
||||
result = (int)event->error_handle->connect_return_code;
|
||||
}
|
||||
if (event && event->error_handle) {
|
||||
tlsError = event->error_handle->esp_tls_last_esp_err;
|
||||
tlsStackError = event->error_handle->esp_tls_stack_err;
|
||||
socketError = event->error_handle->esp_transport_sock_errno;
|
||||
}
|
||||
self->setStatusPhase(MqttBridgePhase::RetryWait, result, now, true,
|
||||
now + RECONNECT_INTERVAL_MS);
|
||||
portENTER_CRITICAL(&self->_statusMux);
|
||||
self->_status.lastTlsError = tlsError;
|
||||
self->_status.lastTlsStackError = tlsStackError;
|
||||
self->_status.lastSocketError = socketError;
|
||||
portEXIT_CRITICAL(&self->_statusMux);
|
||||
break;
|
||||
}
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
bool MqttBridge::observerCloudStart() {
|
||||
if (_observerCloudClient || !_identity) return _observerCloudClient != nullptr;
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
const uint32_t epoch = (uint32_t)time(nullptr);
|
||||
if (preset.auth == ObserverAuth::Jwt && epoch < 1700000000UL) return false;
|
||||
|
||||
static char token[768], username[72], clientId[80], statusTopic[192], offlinePayload[512];
|
||||
const char* connectUsername = nullptr;
|
||||
const char* connectPassword = nullptr;
|
||||
if (preset.auth == ObserverAuth::Jwt) {
|
||||
if (!createObserverJwt(token, sizeof(token))) return false;
|
||||
snprintf(username, sizeof(username), "v1_%s", _nodeId);
|
||||
connectUsername = username;
|
||||
connectPassword = token;
|
||||
} else if (preset.auth == ObserverAuth::UserPass) {
|
||||
connectUsername = preset.username == OBSERVER_PUBKEY_USERNAME
|
||||
? _nodeId : (preset.username ? preset.username : _user);
|
||||
connectPassword = preset.password ? preset.password : _pwd;
|
||||
if (!connectUsername[0] || !connectPassword[0]) return false;
|
||||
}
|
||||
snprintf(clientId, sizeof(clientId), "mqtt_%s-%.6s", preset.name, _nodeId);
|
||||
observerTopic(statusTopic, sizeof(statusTopic), "status");
|
||||
char timestamp[40], safeOrigin[72];
|
||||
formatObserverTimestamp(epoch, timestamp, sizeof(timestamp));
|
||||
escapeJson(_observerOrigin, safeOrigin, sizeof(safeOrigin));
|
||||
snprintf(offlinePayload, sizeof(offlinePayload),
|
||||
"{\"status\":\"offline\",\"timestamp\":\"%s\",\"origin\":\"%s\","
|
||||
"\"origin_id\":\"%s\",\"firmware_version\":\"%s\","
|
||||
"\"radio\":\"%s\",\"client_version\":\"wadamesh-observer/1\"}",
|
||||
timestamp, safeOrigin, _nodeId, _firmwareVersion, _observerRadio);
|
||||
|
||||
esp_mqtt_client_config_t config{};
|
||||
config.uri = preset.uri;
|
||||
config.client_id = clientId;
|
||||
config.username = connectUsername;
|
||||
config.password = connectPassword;
|
||||
config.lwt_topic = statusTopic;
|
||||
config.lwt_msg = offlinePayload;
|
||||
config.lwt_qos = 1;
|
||||
config.lwt_retain = preset.allowRetain ? 1 : 0;
|
||||
config.keepalive = preset.keepalive ? preset.keepalive : 60;
|
||||
config.buffer_size = 1536;
|
||||
config.out_buffer_size = 1536;
|
||||
config.network_timeout_ms = 10000;
|
||||
config.reconnect_timeout_ms = RECONNECT_INTERVAL_MS;
|
||||
if (strncmp(preset.uri, "wss://", 6) == 0 || strncmp(preset.uri, "mqtts://", 8) == 0)
|
||||
config.crt_bundle_attach = esp_crt_bundle_attach;
|
||||
config.user_context = this;
|
||||
|
||||
esp_mqtt_client_handle_t client = esp_mqtt_client_init(&config);
|
||||
if (!client) return false;
|
||||
esp_mqtt_client_register_event(client, MQTT_EVENT_ANY, observerCloudEvent, this);
|
||||
_observerCloudClient = client;
|
||||
_observerCloudTokenExpires = preset.auth == ObserverAuth::Jwt
|
||||
? epoch + preset.tokenLifetime : 0;
|
||||
setStatusPhase(MqttBridgePhase::Connecting, -1, millis(), false);
|
||||
if (esp_mqtt_client_start(client) != ESP_OK) {
|
||||
esp_mqtt_client_destroy(client);
|
||||
_observerCloudClient = nullptr;
|
||||
_observerCloudTokenExpires = 0;
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
void MqttBridge::observerCloudStop() {
|
||||
if (!_observerCloudClient) return;
|
||||
esp_mqtt_client_handle_t client = static_cast<esp_mqtt_client_handle_t>(_observerCloudClient);
|
||||
esp_mqtt_client_stop(client);
|
||||
esp_mqtt_client_destroy(client);
|
||||
_observerCloudClient = nullptr;
|
||||
_observerCloudConnected = false;
|
||||
_observerCloudTokenExpires = 0;
|
||||
}
|
||||
|
||||
void MqttBridge::observerCloudLoop(uint32_t now) {
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
if (WiFi.status() != WL_CONNECTED) {
|
||||
if (_observerCloudClient) observerCloudStop();
|
||||
setStatusPhase(MqttBridgePhase::WaitingForWifi, -1, now, false);
|
||||
return;
|
||||
}
|
||||
const uint32_t epoch = (uint32_t)time(nullptr);
|
||||
if (preset.auth == ObserverAuth::Jwt && epoch < 1700000000UL) {
|
||||
if (_observerCloudClient) observerCloudStop();
|
||||
setStatusPhase(MqttBridgePhase::WaitingForTime, -1, now, false);
|
||||
return;
|
||||
}
|
||||
if (preset.auth == ObserverAuth::Jwt && _observerCloudClient &&
|
||||
_observerCloudTokenExpires > 0 &&
|
||||
epoch + OBSERVER_TOKEN_REFRESH_SECS >= _observerCloudTokenExpires) {
|
||||
observerCloudStop();
|
||||
}
|
||||
if (!_observerCloudClient) {
|
||||
if (_observerCloudRetryAtMs && (int32_t)(now - _observerCloudRetryAtMs) < 0) return;
|
||||
if (!observerCloudStart()) {
|
||||
_observerCloudRetryAtMs = now + RECONNECT_INTERVAL_MS;
|
||||
setStatusPhase(MqttBridgePhase::RetryWait, -2, now, true,
|
||||
_observerCloudRetryAtMs);
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (!_observerCloudConnected) return;
|
||||
if (_lastObserverStatusMs == 0 ||
|
||||
(uint32_t)(now - _lastObserverStatusMs) >= OBSERVER_STATUS_INTERVAL_MS) {
|
||||
publishObserverStatus(true);
|
||||
}
|
||||
ObserverPacket observed;
|
||||
if (popObserverPacket(observed)) publishObserverPacket(observed);
|
||||
}
|
||||
|
||||
void MqttBridge::loop() {
|
||||
if (!_enabled || _connecting) return;
|
||||
if (!_enabled) return;
|
||||
uint32_t now = millis();
|
||||
if (isBuiltInObserverProfile(_observerProfile)) {
|
||||
observerCloudLoop(now);
|
||||
return;
|
||||
}
|
||||
if (!_plainEnabled || _connecting) return;
|
||||
if (WiFi.status() != WL_CONNECTED) {
|
||||
if (_mqtt.connected()) _mqtt.loop();
|
||||
setStatusPhase(MqttBridgePhase::WaitingForWifi, MQTT_DISCONNECTED,
|
||||
@@ -209,6 +674,15 @@ void MqttBridge::loop() {
|
||||
if (!_mqtt.loop()) {
|
||||
setStatusPhase(MqttBridgePhase::RetryWait, _mqtt.state(), now, true,
|
||||
now);
|
||||
return;
|
||||
}
|
||||
if (_pubObserver) {
|
||||
if (_lastObserverStatusMs == 0 ||
|
||||
(uint32_t)(now - _lastObserverStatusMs) >= OBSERVER_STATUS_INTERVAL_MS) {
|
||||
publishObserverStatus(true);
|
||||
}
|
||||
ObserverPacket observed;
|
||||
if (popObserverPacket(observed)) publishObserverPacket(observed);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -241,7 +715,7 @@ bool MqttBridge::sealToB64(const char* plain, char* out, size_t outCap) {
|
||||
}
|
||||
|
||||
void MqttBridge::pub(const char* subtopic, const char* json) {
|
||||
if (!_enabled || _connecting || !_mqtt.connected()) return;
|
||||
if (!_plainEnabled || _connecting || !_mqtt.connected()) return;
|
||||
char topic[80];
|
||||
snprintf(topic, sizeof(topic), "wadamesh/%s/%s", _nodeHex, subtopic);
|
||||
if (_encOn) {
|
||||
@@ -252,6 +726,191 @@ void MqttBridge::pub(const char* subtopic, const char* json) {
|
||||
}
|
||||
}
|
||||
|
||||
void MqttBridge::observerTopic(char* out, size_t outCap, const char* type) const {
|
||||
if (!out || outCap == 0) return;
|
||||
out[0] = '\0';
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
if (preset.topicStyle == ObserverTopicStyle::MeshRank) {
|
||||
snprintf(out, outCap, "meshrank/uplink/%s/%s/%s",
|
||||
_observerToken, _nodeId, type);
|
||||
return;
|
||||
}
|
||||
const char* pattern = (isBuiltInObserverProfile(_observerProfile) || !_observerTopic[0])
|
||||
? OBSERVER_DEFAULT_TOPIC : _observerTopic;
|
||||
const bool hasType = strstr(pattern, "{type}") != nullptr;
|
||||
size_t pos = 0;
|
||||
while (*pattern && pos + 1 < outCap) {
|
||||
if (strncmp(pattern, "{iata}", 6) == 0) {
|
||||
appendTopicText(out, outCap, pos, _observerIata[0] ? _observerIata : "UNSET");
|
||||
pattern += 6;
|
||||
} else if (strncmp(pattern, "{device}", 8) == 0) {
|
||||
appendTopicText(out, outCap, pos, _nodeId);
|
||||
pattern += 8;
|
||||
} else if (strncmp(pattern, "{type}", 6) == 0) {
|
||||
appendTopicText(out, outCap, pos, type);
|
||||
pattern += 6;
|
||||
} else {
|
||||
out[pos++] = *pattern++;
|
||||
out[pos] = '\0';
|
||||
}
|
||||
}
|
||||
if (!hasType) {
|
||||
if (pos > 0 && out[pos - 1] != '/') appendTopicText(out, outCap, pos, "/");
|
||||
appendTopicText(out, outCap, pos, type);
|
||||
}
|
||||
}
|
||||
|
||||
void MqttBridge::observeRx(float snr, float rssi, const uint8_t* raw, int len) {
|
||||
if (!_requestedEnabled || !_pubObserver || !raw || len <= 0 || len > OBSERVER_MAX_PACKET) return;
|
||||
|
||||
portENTER_CRITICAL(&_observerQueueMux);
|
||||
ObserverPacket& observed = _observerQueue[_observerQueueHead];
|
||||
observed.timestamp = (uint32_t)time(nullptr);
|
||||
observed.rssi = (int8_t)rssi;
|
||||
observed.snrQ4 = (int8_t)(snr * 4.0f);
|
||||
observed.len = (uint8_t)len;
|
||||
memcpy(observed.raw, raw, (size_t)len);
|
||||
_observerQueueHead = (uint8_t)((_observerQueueHead + 1) % OBSERVER_QUEUE_SIZE);
|
||||
if (_observerQueueCount < OBSERVER_QUEUE_SIZE) {
|
||||
++_observerQueueCount;
|
||||
} else {
|
||||
++_observerDropped; // full queue: overwrite the oldest observation
|
||||
}
|
||||
const uint8_t queued = _observerQueueCount;
|
||||
const uint32_t dropped = _observerDropped;
|
||||
portEXIT_CRITICAL(&_observerQueueMux);
|
||||
|
||||
portENTER_CRITICAL(&_statusMux);
|
||||
_status.observerQueued = queued;
|
||||
_status.observerDropped = dropped;
|
||||
portEXIT_CRITICAL(&_statusMux);
|
||||
}
|
||||
|
||||
bool MqttBridge::popObserverPacket(ObserverPacket& observed) {
|
||||
portENTER_CRITICAL(&_observerQueueMux);
|
||||
if (_observerQueueCount == 0) {
|
||||
portEXIT_CRITICAL(&_observerQueueMux);
|
||||
return false;
|
||||
}
|
||||
const uint8_t tail = (uint8_t)((_observerQueueHead + OBSERVER_QUEUE_SIZE -
|
||||
_observerQueueCount) % OBSERVER_QUEUE_SIZE);
|
||||
observed = _observerQueue[tail];
|
||||
--_observerQueueCount;
|
||||
const uint8_t queued = _observerQueueCount;
|
||||
portEXIT_CRITICAL(&_observerQueueMux);
|
||||
|
||||
portENTER_CRITICAL(&_statusMux);
|
||||
_status.observerQueued = queued;
|
||||
portEXIT_CRITICAL(&_statusMux);
|
||||
return true;
|
||||
}
|
||||
|
||||
void MqttBridge::publishObserverStatus(bool online) {
|
||||
if (!_pubObserver) return;
|
||||
char topic[192], timestamp[40], safeOrigin[72];
|
||||
observerTopic(topic, sizeof(topic), "status");
|
||||
formatObserverTimestamp((uint32_t)time(nullptr), timestamp, sizeof(timestamp));
|
||||
escapeJson(_observerOrigin, safeOrigin, sizeof(safeOrigin));
|
||||
|
||||
char json[640];
|
||||
snprintf(json, sizeof(json),
|
||||
"{\"status\":\"%s\",\"timestamp\":\"%s\",\"origin\":\"%s\","
|
||||
"\"origin_id\":\"%s\",\"model\":\"WADAMESH\","
|
||||
"\"firmware_version\":\"%s\",\"radio\":\"%s\","
|
||||
"\"client_version\":\"wadamesh-observer/1\","
|
||||
"\"location\":\"%s\","
|
||||
"\"stats\":{\"uptime_secs\":%lu,\"queue_len\":%u,"
|
||||
"\"packets_received\":%lu,\"observer_dropped\":%lu}}",
|
||||
online ? "online" : "offline", timestamp, safeOrigin, _nodeId,
|
||||
_firmwareVersion, _observerRadio,
|
||||
_observerIata[0] ? _observerIata : "UNSET",
|
||||
(unsigned long)(millis() / 1000u), (unsigned)_observerQueueCount,
|
||||
(unsigned long)(_observerPublished + _observerDropped + _observerQueueCount),
|
||||
(unsigned long)_observerDropped);
|
||||
publishObserverPayload(topic, json, true);
|
||||
_lastObserverStatusMs = millis();
|
||||
}
|
||||
|
||||
bool MqttBridge::publishObserverPayload(const char* topic, const char* json, bool retain) {
|
||||
if (!topic || !json) return false;
|
||||
if (isBuiltInObserverProfile(_observerProfile)) {
|
||||
if (!_observerCloudConnected || !_observerCloudClient) return false;
|
||||
const ObserverPreset& preset = observerPreset(_observerProfile);
|
||||
const bool effectiveRetain = retain && preset.allowRetain;
|
||||
esp_mqtt_client_handle_t client = static_cast<esp_mqtt_client_handle_t>(_observerCloudClient);
|
||||
if (esp_mqtt_client_get_outbox_size(client) > 16384) return false;
|
||||
return esp_mqtt_client_enqueue(client, topic, json, 0,
|
||||
effectiveRetain ? 1 : 0,
|
||||
effectiveRetain ? 1 : 0, true) >= 0;
|
||||
}
|
||||
return _plainEnabled && !_connecting && _mqtt.connected() &&
|
||||
_mqtt.publish(topic, json, retain);
|
||||
}
|
||||
|
||||
void MqttBridge::publishObserverPacket(const ObserverPacket& observed) {
|
||||
static mesh::Packet packet;
|
||||
static char rawHex[OBSERVER_MAX_PACKET * 2 + 1];
|
||||
static char pathJson[520];
|
||||
static char json[1280];
|
||||
|
||||
if (!packet.readFrom(observed.raw, observed.len)) {
|
||||
++_observerDropped;
|
||||
return;
|
||||
}
|
||||
|
||||
bytesToHex(observed.raw, observed.len, rawHex, sizeof(rawHex), true);
|
||||
uint8_t hash[MAX_HASH_SIZE];
|
||||
char hashHex[MAX_HASH_SIZE * 2 + 1];
|
||||
packet.calculatePacketHash(hash);
|
||||
bytesToHex(hash, sizeof(hash), hashHex, sizeof(hashHex), true);
|
||||
|
||||
pathJson[0] = '\0';
|
||||
if (packet.isRouteDirect() && packet.getPathHashCount() > 0) {
|
||||
size_t pos = 0;
|
||||
pos += snprintf(pathJson + pos, sizeof(pathJson) - pos, ",\"path\":[");
|
||||
const uint8_t hashSize = packet.getPathHashSize();
|
||||
const uint8_t hashCount = packet.getPathHashCount();
|
||||
for (uint8_t hop = 0; hop < hashCount && pos + hashSize * 2 + 4 < sizeof(pathJson); ++hop) {
|
||||
if (hop) pathJson[pos++] = ',';
|
||||
pathJson[pos++] = '\"';
|
||||
char hopHex[7];
|
||||
bytesToHex(packet.path + hop * hashSize, hashSize, hopHex, sizeof(hopHex), false);
|
||||
pos += snprintf(pathJson + pos, sizeof(pathJson) - pos, "%s\"", hopHex);
|
||||
}
|
||||
if (pos + 2 < sizeof(pathJson)) {
|
||||
pathJson[pos++] = ']';
|
||||
pathJson[pos] = '\0';
|
||||
}
|
||||
}
|
||||
|
||||
char timestamp[40], timeText[16], dateText[16], safeOrigin[72];
|
||||
formatObserverTimestamp(observed.timestamp, timestamp, sizeof(timestamp),
|
||||
timeText, sizeof(timeText), dateText, sizeof(dateText));
|
||||
escapeJson(_observerOrigin, safeOrigin, sizeof(safeOrigin));
|
||||
const int written = snprintf(
|
||||
json, sizeof(json),
|
||||
"{\"origin\":\"%s\",\"origin_id\":\"%s\",\"timestamp\":\"%s\","
|
||||
"\"type\":\"PACKET\",\"direction\":\"rx\",\"time\":\"%s\",\"date\":\"%s\","
|
||||
"\"len\":\"%u\",\"packet_type\":\"%u\",\"route\":\"%s\","
|
||||
"\"payload_len\":\"%u\",\"raw\":\"%s\",\"SNR\":\"%.1f\","
|
||||
"\"RSSI\":\"%d\",\"hash\":\"%s\"%s}",
|
||||
safeOrigin, _nodeId, timestamp, timeText, dateText,
|
||||
(unsigned)observed.len, (unsigned)packet.getPayloadType(),
|
||||
packet.isRouteDirect() ? "D" : "F", (unsigned)packet.payload_len,
|
||||
rawHex, (double)observed.snrQ4 / 4.0, (int)observed.rssi, hashHex, pathJson);
|
||||
|
||||
char topic[192];
|
||||
observerTopic(topic, sizeof(topic), "packets");
|
||||
const bool published = written > 0 && (size_t)written < sizeof(json) &&
|
||||
publishObserverPayload(topic, json, false);
|
||||
if (published) ++_observerPublished;
|
||||
else ++_observerDropped;
|
||||
portENTER_CRITICAL(&_statusMux);
|
||||
_status.observerPublished = _observerPublished;
|
||||
_status.observerDropped = _observerDropped;
|
||||
portEXIT_CRITICAL(&_statusMux);
|
||||
}
|
||||
|
||||
void MqttBridge::escapeJson(const char* src, char* dst, size_t dstLen) {
|
||||
static const char* hexd = "0123456789abcdef";
|
||||
size_t j = 0;
|
||||
@@ -271,7 +930,7 @@ void MqttBridge::escapeJson(const char* src, char* dst, size_t dstLen) {
|
||||
|
||||
void MqttBridge::publishDM(const char* senderName, const uint8_t* senderKey32,
|
||||
float snr, uint8_t hops, uint32_t ts, const char* text) {
|
||||
if (!_enabled || !_pubDm || _connecting || !_mqtt.connected()) return; // DMs are opt-in
|
||||
if (!_plainEnabled || !_pubDm || _connecting || !_mqtt.connected()) return; // DMs are opt-in
|
||||
char keyHex[13] = {};
|
||||
for (int i = 0; i < 6; ++i) snprintf(keyHex + i * 2, 3, "%02x", senderKey32[i]);
|
||||
|
||||
@@ -288,7 +947,7 @@ void MqttBridge::publishDM(const char* senderName, const uint8_t* senderKey32,
|
||||
|
||||
void MqttBridge::publishChannel(int channelIdx, const char* channelName,
|
||||
float snr, uint8_t hops, uint32_t ts, const char* text) {
|
||||
if (!_enabled || !_pubChannel || _connecting || !_mqtt.connected()) return; // channel publish toggle
|
||||
if (!_plainEnabled || !_pubChannel || _connecting || !_mqtt.connected()) return; // channel publish toggle
|
||||
char safeName[48], safeText[300];
|
||||
escapeJson(channelName, safeName, sizeof(safeName));
|
||||
escapeJson(text, safeText, sizeof(safeText));
|
||||
@@ -316,21 +975,65 @@ void MqttBridge::saveConfig(const char* host, uint16_t port,
|
||||
p.end();
|
||||
}
|
||||
|
||||
void MqttBridge::saveObserverConfig(bool enable, MqttObserverProfile profile, const char* origin,
|
||||
const char* iata, const char* topicTemplate,
|
||||
const char* profileToken) {
|
||||
char cleanIata[8] = {};
|
||||
size_t out = 0;
|
||||
for (size_t i = 0; iata && iata[i] && out + 1 < sizeof(cleanIata); ++i) {
|
||||
const unsigned char c = (unsigned char)iata[i];
|
||||
if (isalnum(c) || c == '-' || c == '_') cleanIata[out++] = (char)toupper(c);
|
||||
}
|
||||
|
||||
SdNvsPrefs p;
|
||||
if (!p.begin("mqtt", false)) return;
|
||||
p.putBool("obs", enable);
|
||||
p.putUInt("obs_profile", (uint32_t)profile);
|
||||
p.putString("obs_origin", origin ? origin : "");
|
||||
p.putString("obs_iata", cleanIata);
|
||||
p.putString("obs_topic", topicTemplate && topicTemplate[0]
|
||||
? topicTemplate : OBSERVER_DEFAULT_TOPIC);
|
||||
p.putString("obs_token", profileToken ? profileToken : "");
|
||||
p.end();
|
||||
}
|
||||
|
||||
void MqttBridge::reloadConfig() {
|
||||
// A connect attempt may be in flight on the one-shot task; PubSubClient is not
|
||||
// thread-safe, so wait it out (bounded: DNS + TCP + CONNACK <= ~10 s, and it
|
||||
// only overlaps when Save lands inside an attempt window on a dead broker).
|
||||
while (_connecting) delay(10);
|
||||
if (_pubObserver && ((isBuiltInObserverProfile(_observerProfile) && _observerCloudConnected) ||
|
||||
(_observerProfile == MqttObserverProfile::Custom && _mqtt.connected()))) {
|
||||
publishObserverStatus(false);
|
||||
}
|
||||
observerCloudStop();
|
||||
if (_mqtt.connected()) _mqtt.disconnect();
|
||||
loadConfig();
|
||||
portENTER_CRITICAL(&_observerQueueMux);
|
||||
_observerDropped += _observerQueueCount;
|
||||
_observerQueueHead = _observerQueueCount = 0;
|
||||
const uint32_t dropped = _observerDropped;
|
||||
portEXIT_CRITICAL(&_observerQueueMux);
|
||||
portENTER_CRITICAL(&_statusMux);
|
||||
_status.observerQueued = 0;
|
||||
_status.observerDropped = dropped;
|
||||
portEXIT_CRITICAL(&_statusMux);
|
||||
if (!_enabled) { Serial.println("[MQTT] disabled"); return; }
|
||||
_mqtt.setServer(_host, _port);
|
||||
_mqtt.setKeepAlive(60);
|
||||
_mqtt.setSocketTimeout(2);
|
||||
_mqtt.setBufferSize(_encOn ? 1024 : 512);
|
||||
if (_plainEnabled) {
|
||||
_mqtt.setServer(_host, _port);
|
||||
_mqtt.setKeepAlive(60);
|
||||
_mqtt.setSocketTimeout(2);
|
||||
_mqtt.setBufferSize(_pubObserver ? 1536 : (_encOn ? 1024 : 512));
|
||||
}
|
||||
_lastReconnectMs = 0; // reconnect on next loop() tick
|
||||
Serial.printf("[MQTT] reloaded -> %s:%u en=%d dm=%d ch=%d enc=%d\n",
|
||||
_host, _port, (int)_enabled, (int)_pubDm, (int)_pubChannel, (int)_encOn);
|
||||
_observerCloudRetryAtMs = 0;
|
||||
_lastObserverStatusMs = 0;
|
||||
Serial.printf("[MQTT] reloaded -> profile=%s host=%s:%u en=%d dm=%d ch=%d obs=%d enc=%d\n",
|
||||
observerPreset(_observerProfile).name,
|
||||
isBuiltInObserverProfile(_observerProfile) ? observerPreset(_observerProfile).host : _host,
|
||||
isBuiltInObserverProfile(_observerProfile) ?
|
||||
(unsigned)observerPreset(_observerProfile).port : (unsigned)_port,
|
||||
(int)_enabled, (int)_pubDm, (int)_pubChannel, (int)_pubObserver, (int)_encOn);
|
||||
}
|
||||
|
||||
#endif // !HAS_TANMATSU (real implementation)
|
||||
|
||||
+135
-12
@@ -3,11 +3,57 @@
|
||||
|
||||
#include <Arduino.h>
|
||||
|
||||
namespace mesh { class LocalIdentity; }
|
||||
|
||||
enum class MqttObserverProfile : uint8_t {
|
||||
Custom = 0,
|
||||
NebraskaMesh = 1,
|
||||
MichMesh = 2,
|
||||
AnalyzerUs,
|
||||
AnalyzerEu,
|
||||
NzAnalyzer,
|
||||
MeshMapper,
|
||||
MeshRank,
|
||||
Waev,
|
||||
Meshomatic,
|
||||
CascadiaMesh,
|
||||
TennMesh,
|
||||
NashMesh,
|
||||
CtMesh,
|
||||
ChiMesh,
|
||||
MeshatSe,
|
||||
EastIdahoMesh,
|
||||
ColoradoMesh,
|
||||
DutchMeshcore1,
|
||||
DutchMeshcore2,
|
||||
MeshcoreCa1,
|
||||
MeshcoreCa2,
|
||||
MeshcoreFi,
|
||||
OkiMesh1,
|
||||
OkiMesh2,
|
||||
InwMesh,
|
||||
BostonMesh,
|
||||
RfLab,
|
||||
IpntUk,
|
||||
FlMesh,
|
||||
CoreComms,
|
||||
MeshTexas,
|
||||
MeshChaun14,
|
||||
WcMesh,
|
||||
AtvirasTinklas,
|
||||
GoMesh,
|
||||
IdahoMesh,
|
||||
NtxMesh,
|
||||
BsMesh,
|
||||
Count,
|
||||
};
|
||||
|
||||
enum class MqttBridgePhase : uint8_t {
|
||||
Unavailable,
|
||||
Disabled,
|
||||
Misconfigured,
|
||||
WaitingForWifi,
|
||||
WaitingForTime,
|
||||
Connecting,
|
||||
Connected,
|
||||
RetryWait,
|
||||
@@ -19,14 +65,22 @@ struct MqttBridgeStatus {
|
||||
bool requestedEnabled = false;
|
||||
bool publishDm = false;
|
||||
bool publishChannel = false;
|
||||
bool publishObserver = false;
|
||||
MqttObserverProfile observerProfile = MqttObserverProfile::Custom;
|
||||
bool encrypted = false;
|
||||
char host[64] = {};
|
||||
uint16_t port = 1883;
|
||||
int8_t lastResult = -1;
|
||||
int32_t lastResult = -1;
|
||||
int32_t lastTlsError = 0;
|
||||
int32_t lastTlsStackError = 0;
|
||||
int32_t lastSocketError = 0;
|
||||
uint32_t lastAttemptMs = 0;
|
||||
uint32_t lastConnectedMs = 0;
|
||||
uint32_t lastResultMs = 0;
|
||||
uint32_t retryAtMs = 0;
|
||||
uint32_t observerPublished = 0;
|
||||
uint32_t observerDropped = 0;
|
||||
uint8_t observerQueued = 0;
|
||||
};
|
||||
|
||||
#if defined(HAS_TANMATSU) || defined(HAS_TDISPLAY_P4)
|
||||
@@ -37,14 +91,25 @@ struct MqttBridgeStatus {
|
||||
// the P4 build. Re-enable once esp-hosted TCP is proven on the P4.
|
||||
class MqttBridge {
|
||||
public:
|
||||
void begin(const char* /*nodeHex*/) {}
|
||||
void begin(const mesh::LocalIdentity*, const char*, float, float, uint8_t, uint8_t, const char*) {}
|
||||
void loop() {}
|
||||
void observeRx(float, float, const uint8_t*, int) {}
|
||||
void publishDM(const char*, const uint8_t*, float, uint8_t, uint32_t, const char*) {}
|
||||
void publishChannel(int, const char*, float, uint8_t, uint32_t, const char*) {}
|
||||
bool enabled() const { return false; }
|
||||
void getStatus(MqttBridgeStatus& out) const { out = MqttBridgeStatus{}; }
|
||||
static void saveConfig(const char*, uint16_t, const char*, const char*,
|
||||
bool, bool, const char*, bool) {}
|
||||
static void saveObserverConfig(bool, MqttObserverProfile, const char*, const char*, const char*, const char*) {}
|
||||
static uint8_t observerProfileCount();
|
||||
static const char* observerProfileOptions();
|
||||
static const char* observerProfileLabel(MqttObserverProfile profile);
|
||||
static MqttObserverProfile observerProfileFromIndex(uint16_t index);
|
||||
static uint16_t observerProfileIndex(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsIata(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsToken(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsUsername(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsPassword(MqttObserverProfile profile);
|
||||
void reloadConfig() {}
|
||||
};
|
||||
#else
|
||||
@@ -69,12 +134,11 @@ public:
|
||||
// - Channel messages publish by default; DIRECT MESSAGES are off by default and
|
||||
// must be explicitly enabled — they are private 1:1 traffic, and may be from
|
||||
// someone else who never consented to being bridged.
|
||||
// - If an encryption key (PSK) is set, every JSON payload is sealed with
|
||||
// - If an encryption key (PSK) is set, decoded message JSON is sealed with
|
||||
// AES-256-GCM (fresh random 12-byte nonce per message) before it leaves the
|
||||
// device, so the broker only ever sees opaque base64. TLS is deliberately NOT
|
||||
// used: the mbedTLS handshake (~30 KB heap) does not fit the ESP32-S3 budget
|
||||
// alongside Wi-Fi + BLE + LVGL (same reason map tiles are HTTP-only). App-layer
|
||||
// GCM is the lightweight, hardware-accelerated equivalent.
|
||||
// device. Observer payloads remain standard JSON. The custom profile uses
|
||||
// plain MQTT; community profiles use their WSS endpoints with CA-bundle
|
||||
// verification and an Ed25519-signed identity token.
|
||||
//
|
||||
// Topics (QoS 0, retained where noted):
|
||||
// wadamesh/{node_hex}/msg/dm — direct / signed messages (only if DM publish ON)
|
||||
@@ -87,14 +151,18 @@ public:
|
||||
// same passphrase. With no PSK, the payload is the plain JSON (use a private broker).
|
||||
//
|
||||
// Config persisted in Preferences namespace "mqtt" (file-backed via SdNvsPrefs):
|
||||
// en bool · host str · port u16 · user str · pwd str · dm bool · ch bool · psk str
|
||||
// en · host · port · user · pwd · dm · ch · psk, plus obs · obs_profile ·
|
||||
// obs_origin · obs_iata · obs_topic · obs_token.
|
||||
//
|
||||
// Call begin() once after the_mesh.begin() and SdNvsPrefs::useFile().
|
||||
// Call loop() every iteration of the Arduino loop().
|
||||
class MqttBridge {
|
||||
public:
|
||||
void begin(const char* nodeHex);
|
||||
void begin(const mesh::LocalIdentity* identity, const char* nodeName,
|
||||
float frequencyMhz, float bandwidthKhz, uint8_t spreadingFactor,
|
||||
uint8_t codingRate, const char* firmwareVersion);
|
||||
void loop();
|
||||
void observeRx(float snr, float rssi, const uint8_t* raw, int len);
|
||||
|
||||
void publishDM(const char* senderName, const uint8_t* senderKey32,
|
||||
float snr, uint8_t hops, uint32_t ts, const char* text);
|
||||
@@ -108,25 +176,70 @@ public:
|
||||
static void saveConfig(const char* host, uint16_t port,
|
||||
const char* user, const char* pwd,
|
||||
bool pubDm, bool pubChannel, const char* psk, bool enable);
|
||||
static void saveObserverConfig(bool enable, MqttObserverProfile profile, const char* origin,
|
||||
const char* iata, const char* topicTemplate,
|
||||
const char* profileToken);
|
||||
static uint8_t observerProfileCount();
|
||||
static const char* observerProfileOptions();
|
||||
static const char* observerProfileLabel(MqttObserverProfile profile);
|
||||
static MqttObserverProfile observerProfileFromIndex(uint16_t index);
|
||||
static uint16_t observerProfileIndex(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsIata(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsToken(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsUsername(MqttObserverProfile profile);
|
||||
static bool observerProfileNeedsPassword(MqttObserverProfile profile);
|
||||
// Re-read config from Preferences and reconnect (call after saveConfig).
|
||||
void reloadConfig();
|
||||
|
||||
private:
|
||||
MqttNbClient _wc;
|
||||
PubSubClient _mqtt{_wc};
|
||||
char _nodeHex[13] = {}; // 6-byte key → 12 hex chars + '\0'
|
||||
char _nodeHex[13] = {}; // legacy topic: first 6 key bytes, lowercase
|
||||
char _nodeId[65] = {}; // Observer topic: complete public key, uppercase
|
||||
char _nodeName[33] = {};
|
||||
char _observerRadio[48] = {};
|
||||
char _firmwareVersion[24] = {};
|
||||
bool _requestedEnabled = false;
|
||||
bool _enabled = false;
|
||||
bool _plainEnabled = false;
|
||||
bool _pubDm = false; // DMs off by default (private 1:1 traffic)
|
||||
bool _pubChannel = true; // channel messages on by default
|
||||
bool _pubObserver = false; // raw RX packets in MeshCore Observer format
|
||||
MqttObserverProfile _observerProfile = MqttObserverProfile::Custom;
|
||||
const mesh::LocalIdentity* _identity = nullptr;
|
||||
char _host[64] = {};
|
||||
char _user[32] = {};
|
||||
char _pwd[32] = {};
|
||||
char _user[65] = {};
|
||||
char _pwd[97] = {};
|
||||
char _psk[33] = {}; // passphrase; empty = no payload encryption
|
||||
char _observerOrigin[33] = {};
|
||||
char _observerIata[8] = {};
|
||||
char _observerTopic[96] = {};
|
||||
char _observerToken[65] = {};
|
||||
uint8_t _key[32] = {}; // SHA-256(psk), valid when _encOn
|
||||
bool _encOn = false;
|
||||
uint16_t _port = 1883;
|
||||
uint32_t _lastReconnectMs = 0;
|
||||
uint32_t _lastObserverStatusMs = 0;
|
||||
uint32_t _observerPublished = 0;
|
||||
uint32_t _observerDropped = 0;
|
||||
void* _observerCloudClient = nullptr;
|
||||
volatile bool _observerCloudConnected = false;
|
||||
uint32_t _observerCloudRetryAtMs = 0;
|
||||
uint32_t _observerCloudTokenExpires = 0;
|
||||
|
||||
static const uint8_t OBSERVER_QUEUE_SIZE = 8;
|
||||
static const uint16_t OBSERVER_MAX_PACKET = 255;
|
||||
struct ObserverPacket {
|
||||
uint32_t timestamp;
|
||||
int8_t rssi;
|
||||
int8_t snrQ4;
|
||||
uint8_t len;
|
||||
uint8_t raw[OBSERVER_MAX_PACKET];
|
||||
};
|
||||
ObserverPacket _observerQueue[OBSERVER_QUEUE_SIZE];
|
||||
uint8_t _observerQueueHead = 0;
|
||||
uint8_t _observerQueueCount = 0;
|
||||
portMUX_TYPE _observerQueueMux = portMUX_INITIALIZER_UNLOCKED;
|
||||
// True while the one-shot connect task owns _mqtt/_wc. The loop thread must
|
||||
// not touch either until it clears (PubSubClient is not thread-safe).
|
||||
volatile bool _connecting = false;
|
||||
@@ -143,6 +256,16 @@ private:
|
||||
bool reconnect();
|
||||
static void reconnectTask(void* arg); // one-shot task body wrapping reconnect()
|
||||
void pub(const char* subtopic, const char* json); // seals if _encOn
|
||||
void publishObserverStatus(bool online);
|
||||
void publishObserverPacket(const ObserverPacket& observed);
|
||||
bool publishObserverPayload(const char* topic, const char* json, bool retain);
|
||||
bool popObserverPacket(ObserverPacket& observed);
|
||||
void observerTopic(char* out, size_t outCap, const char* type) const;
|
||||
void observerCloudLoop(uint32_t now);
|
||||
void observerCloudStop();
|
||||
bool observerCloudStart();
|
||||
bool createObserverJwt(char* token, size_t tokenCap) const;
|
||||
static void observerCloudEvent(void* arg, const char* base, int32_t eventId, void* eventData);
|
||||
bool sealToB64(const char* plain, char* out, size_t outCap);
|
||||
static void escapeJson(const char* src, char* dst, size_t dstLen);
|
||||
};
|
||||
|
||||
+5
-3
@@ -1376,9 +1376,11 @@ void setup() {
|
||||
|
||||
#if defined(ESP32) && defined(MULTI_TRANSPORT_COMPANION)
|
||||
{
|
||||
char nodeHex[13] = {};
|
||||
for (int i = 0; i < 6; ++i) snprintf(nodeHex + i * 2, 3, "%02x", the_mesh.self_id.pub_key[i]);
|
||||
mqtt_bridge.begin(nodeHex);
|
||||
NodePrefs* prefs = the_mesh.getNodePrefs();
|
||||
mqtt_bridge.begin(&the_mesh.self_id, prefs ? prefs->node_name : "WADAMESH",
|
||||
prefs ? prefs->freq : 0.0f, prefs ? prefs->bw : 0.0f,
|
||||
prefs ? prefs->sf : 0, prefs ? prefs->cr : 0,
|
||||
FIRMWARE_VERSION);
|
||||
}
|
||||
#endif
|
||||
|
||||
|
||||
+193
-7
@@ -2658,6 +2658,12 @@ struct SettingsModalState {
|
||||
lv_obj_t* mqtt_pwd_ta;
|
||||
lv_obj_t* mqtt_ch_sw;
|
||||
lv_obj_t* mqtt_dm_sw;
|
||||
lv_obj_t* mqtt_obs_sw;
|
||||
lv_obj_t* mqtt_obs_profile_dd;
|
||||
lv_obj_t* mqtt_obs_origin_ta;
|
||||
lv_obj_t* mqtt_obs_iata_ta;
|
||||
lv_obj_t* mqtt_obs_topic_ta;
|
||||
lv_obj_t* mqtt_obs_token_ta;
|
||||
lv_obj_t* mqtt_psk_ta;
|
||||
lv_obj_t* mqtt_consent_cb;
|
||||
};
|
||||
@@ -16604,7 +16610,41 @@ static void mqttSaveCb(lv_event_t* e) {
|
||||
bool en = consent && g_set_modal.mqtt_en_sw && lv_obj_has_state(g_set_modal.mqtt_en_sw, LV_STATE_CHECKED);
|
||||
bool pub_ch = !g_set_modal.mqtt_ch_sw || lv_obj_has_state(g_set_modal.mqtt_ch_sw, LV_STATE_CHECKED);
|
||||
bool pub_dm = g_set_modal.mqtt_dm_sw && lv_obj_has_state(g_set_modal.mqtt_dm_sw, LV_STATE_CHECKED);
|
||||
bool pub_obs = g_set_modal.mqtt_obs_sw && lv_obj_has_state(g_set_modal.mqtt_obs_sw, LV_STATE_CHECKED);
|
||||
const MqttObserverProfile obs_profile = g_set_modal.mqtt_obs_profile_dd
|
||||
? MqttBridge::observerProfileFromIndex(lv_dropdown_get_selected(g_set_modal.mqtt_obs_profile_dd))
|
||||
: MqttObserverProfile::Custom;
|
||||
const char* obs_origin = g_set_modal.mqtt_obs_origin_ta
|
||||
? lv_textarea_get_text(g_set_modal.mqtt_obs_origin_ta) : "";
|
||||
const char* obs_iata = g_set_modal.mqtt_obs_iata_ta
|
||||
? lv_textarea_get_text(g_set_modal.mqtt_obs_iata_ta) : "";
|
||||
const char* obs_topic = g_set_modal.mqtt_obs_topic_ta
|
||||
? lv_textarea_get_text(g_set_modal.mqtt_obs_topic_ta) : "";
|
||||
const char* obs_token = g_set_modal.mqtt_obs_token_ta
|
||||
? lv_textarea_get_text(g_set_modal.mqtt_obs_token_ta) : "";
|
||||
if (en && pub_obs && MqttBridge::observerProfileNeedsIata(obs_profile) &&
|
||||
(!obs_iata || !obs_iata[0])) {
|
||||
g_lv.task->showAlert(TR("Set a location / IATA code for this profile"), 2200);
|
||||
return;
|
||||
}
|
||||
if (en && pub_obs && MqttBridge::observerProfileNeedsToken(obs_profile) &&
|
||||
(!obs_token || !obs_token[0])) {
|
||||
g_lv.task->showAlert(TR("Set the profile token"), 1800);
|
||||
return;
|
||||
}
|
||||
if (en && pub_obs && MqttBridge::observerProfileNeedsUsername(obs_profile) &&
|
||||
(!user || !user[0])) {
|
||||
g_lv.task->showAlert(TR("Set the broker username"), 1800);
|
||||
return;
|
||||
}
|
||||
if (en && pub_obs && MqttBridge::observerProfileNeedsPassword(obs_profile) &&
|
||||
(!pwd || !pwd[0])) {
|
||||
g_lv.task->showAlert(TR("Set the broker password"), 1800);
|
||||
return;
|
||||
}
|
||||
MqttBridge::saveConfig(host, port, user, pwd, pub_dm, pub_ch, psk, en);
|
||||
MqttBridge::saveObserverConfig(pub_obs, obs_profile, obs_origin, obs_iata,
|
||||
obs_topic, obs_token);
|
||||
{ SdNvsPrefs p; if (p.begin("mqtt", false)) { p.putBool("consent", consent); p.end(); } } // file-backed, not NVS (GH #128)
|
||||
mqtt_bridge.reloadConfig();
|
||||
closeSettingsModal();
|
||||
@@ -16619,7 +16659,7 @@ static void buildMqttSettings() {
|
||||
|
||||
// ---- Privacy warning (read before enabling) ----
|
||||
lv_obj_t* warn = lv_label_create(body);
|
||||
lv_label_set_text(warn, TR("Highly experimental. This forwards the text, sender name and timestamp of every message your node receives to an MQTT broker, where anyone able to read the broker can read them. Direct messages are private messages from other people who never agreed to be shared. Use a broker you control, set an encryption key below, and never a public broker."));
|
||||
lv_label_set_text(warn, TR("Highly experimental. Decoded message forwarding exposes message text to the configured broker; direct messages are private and stay off by default. Observer mode sends raw RF packet bytes and radio metadata. Use Custom only with a broker you trust; community profiles publish to that community network."));
|
||||
lv_label_set_long_mode(warn, LV_LABEL_LONG_WRAP);
|
||||
lv_obj_set_width(warn, cw);
|
||||
lv_obj_set_style_text_color(warn, lv_color_hex(0xCC6A00), LV_PART_MAIN);
|
||||
@@ -16699,7 +16739,7 @@ static void buildMqttSettings() {
|
||||
lv_obj_set_pos(g_set_modal.mqtt_user_ta, 0, y);
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_user_ta, true);
|
||||
taSetPlaceholder(g_set_modal.mqtt_user_ta, TR("Leave empty if not required"));
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_user_ta, 31);
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_user_ta, 64);
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_user_ta);
|
||||
y += SC(36);
|
||||
|
||||
@@ -16716,7 +16756,7 @@ static void buildMqttSettings() {
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_pwd_ta, true);
|
||||
lv_textarea_set_password_mode(g_set_modal.mqtt_pwd_ta, true);
|
||||
taSetPlaceholder(g_set_modal.mqtt_pwd_ta, TR("Leave empty if not required"));
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_pwd_ta, 31);
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_pwd_ta, 96);
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_pwd_ta);
|
||||
m9AttachSymbolButton(g_set_modal.mqtt_pwd_ta);
|
||||
y += SC(36);
|
||||
@@ -16740,9 +16780,136 @@ static void buildMqttSettings() {
|
||||
lv_obj_align(g_set_modal.mqtt_dm_sw, LV_ALIGN_TOP_RIGHT, 0, y);
|
||||
y += SC(38);
|
||||
|
||||
// ---- Encryption key (PSK): seals payloads with AES-GCM; empty = plaintext ----
|
||||
// ---- MeshCore Observer: raw RX packets in the community analyzer format ----
|
||||
lv_obj_t* obs_title = lv_label_create(body);
|
||||
lv_label_set_text(obs_title, TR("MeshCore Observer v1"));
|
||||
lv_obj_set_style_text_color(obs_title, lv_color_hex(COLOR_TEXT), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(obs_title, &g_font_14, LV_PART_MAIN);
|
||||
lv_obj_set_pos(obs_title, 2, y);
|
||||
y += SC(20);
|
||||
|
||||
lv_obj_t* obs_help = lv_label_create(body);
|
||||
lv_label_set_text(obs_help, TR("Publish every received RF packet with RSSI, SNR, route, path and raw data."));
|
||||
lv_label_set_long_mode(obs_help, LV_LABEL_LONG_WRAP);
|
||||
lv_obj_set_width(obs_help, cw);
|
||||
lv_obj_set_style_text_color(obs_help, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(obs_help, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(obs_help, 2, y);
|
||||
lv_obj_update_layout(obs_help);
|
||||
y += lv_obj_get_height(obs_help) + SC(8);
|
||||
|
||||
lv_obj_t* profile_lbl = lv_label_create(body);
|
||||
lv_label_set_text(profile_lbl, TR("Observer profile"));
|
||||
lv_obj_set_style_text_color(profile_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(profile_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(profile_lbl, 2, y + 8);
|
||||
g_set_modal.mqtt_obs_profile_dd = lv_dropdown_create(body);
|
||||
lv_dropdown_set_options(g_set_modal.mqtt_obs_profile_dd,
|
||||
MqttBridge::observerProfileOptions());
|
||||
lv_obj_set_size(g_set_modal.mqtt_obs_profile_dd, 142, SC(32));
|
||||
lv_obj_align(g_set_modal.mqtt_obs_profile_dd, LV_ALIGN_TOP_RIGHT, 0, y);
|
||||
y += SC(38);
|
||||
|
||||
lv_obj_t* profile_help = lv_label_create(body);
|
||||
lv_label_set_text(profile_help, TR("Built-in profiles apply their endpoint, authentication and standard topic automatically. Username/password are used only by profiles that require them."));
|
||||
lv_label_set_long_mode(profile_help, LV_LABEL_LONG_WRAP);
|
||||
lv_obj_set_width(profile_help, cw);
|
||||
lv_obj_set_style_text_color(profile_help, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(profile_help, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(profile_help, 2, y);
|
||||
lv_obj_update_layout(profile_help);
|
||||
y += lv_obj_get_height(profile_help) + SC(8);
|
||||
|
||||
lv_obj_t* obs_lbl = lv_label_create(body);
|
||||
lv_label_set_text(obs_lbl, TR("Enable Observer reporting"));
|
||||
lv_obj_set_style_text_color(obs_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(obs_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(obs_lbl, 2, y + 8);
|
||||
g_set_modal.mqtt_obs_sw = lv_switch_create(body);
|
||||
lv_obj_align(g_set_modal.mqtt_obs_sw, LV_ALIGN_TOP_RIGHT, 0, y);
|
||||
y += SC(38);
|
||||
|
||||
MqttBridgeStatus obs_status;
|
||||
mqtt_bridge.getStatus(obs_status);
|
||||
lv_obj_t* obs_stats = lv_label_create(body);
|
||||
lv_label_set_text_fmt(obs_stats, TR("Observer: %lu sent, %u queued, %lu dropped"),
|
||||
(unsigned long)obs_status.observerPublished,
|
||||
(unsigned)obs_status.observerQueued,
|
||||
(unsigned long)obs_status.observerDropped);
|
||||
lv_label_set_long_mode(obs_stats, LV_LABEL_LONG_WRAP);
|
||||
lv_obj_set_width(obs_stats, cw);
|
||||
lv_obj_set_style_text_color(obs_stats, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(obs_stats, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(obs_stats, 2, y);
|
||||
lv_obj_update_layout(obs_stats);
|
||||
y += lv_obj_get_height(obs_stats) + SC(8);
|
||||
|
||||
lv_obj_t* origin_lbl = lv_label_create(body);
|
||||
lv_label_set_text(origin_lbl, TR("Observer name"));
|
||||
lv_obj_set_style_text_color(origin_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(origin_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(origin_lbl, 2, y);
|
||||
y += SC(16);
|
||||
g_set_modal.mqtt_obs_origin_ta = lv_textarea_create(body);
|
||||
lv_obj_set_size(g_set_modal.mqtt_obs_origin_ta, lv_pct(100), SC(30));
|
||||
lv_obj_set_pos(g_set_modal.mqtt_obs_origin_ta, 0, y);
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_obs_origin_ta, true);
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_obs_origin_ta, 32);
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_obs_origin_ta);
|
||||
y += SC(36);
|
||||
|
||||
lv_obj_t* iata_lbl = lv_label_create(body);
|
||||
lv_label_set_text(iata_lbl, TR("Location / IATA code"));
|
||||
lv_obj_set_style_text_color(iata_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(iata_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(iata_lbl, 2, y);
|
||||
y += SC(16);
|
||||
g_set_modal.mqtt_obs_iata_ta = lv_textarea_create(body);
|
||||
lv_obj_set_size(g_set_modal.mqtt_obs_iata_ta, lv_pct(100), SC(30));
|
||||
lv_obj_set_pos(g_set_modal.mqtt_obs_iata_ta, 0, y);
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_obs_iata_ta, true);
|
||||
lv_textarea_set_accepted_chars(g_set_modal.mqtt_obs_iata_ta,
|
||||
"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_");
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_obs_iata_ta, 7);
|
||||
taSetPlaceholder(g_set_modal.mqtt_obs_iata_ta, "OMA");
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_obs_iata_ta);
|
||||
y += SC(36);
|
||||
|
||||
lv_obj_t* topic_lbl = lv_label_create(body);
|
||||
lv_label_set_text(topic_lbl, TR("Observer topic template"));
|
||||
lv_obj_set_style_text_color(topic_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(topic_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(topic_lbl, 2, y);
|
||||
y += SC(16);
|
||||
g_set_modal.mqtt_obs_topic_ta = lv_textarea_create(body);
|
||||
lv_obj_set_size(g_set_modal.mqtt_obs_topic_ta, lv_pct(100), SC(30));
|
||||
lv_obj_set_pos(g_set_modal.mqtt_obs_topic_ta, 0, y);
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_obs_topic_ta, true);
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_obs_topic_ta, 95);
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_obs_topic_ta);
|
||||
m9AttachSymbolButton(g_set_modal.mqtt_obs_topic_ta);
|
||||
y += SC(36);
|
||||
|
||||
lv_obj_t* token_lbl = lv_label_create(body);
|
||||
lv_label_set_text(token_lbl, TR("Profile token (MeshRank)"));
|
||||
lv_obj_set_style_text_color(token_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(token_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(token_lbl, 2, y);
|
||||
y += SC(16);
|
||||
g_set_modal.mqtt_obs_token_ta = lv_textarea_create(body);
|
||||
lv_obj_set_size(g_set_modal.mqtt_obs_token_ta, lv_pct(100), SC(30));
|
||||
lv_obj_set_pos(g_set_modal.mqtt_obs_token_ta, 0, y);
|
||||
lv_textarea_set_one_line(g_set_modal.mqtt_obs_token_ta, true);
|
||||
lv_textarea_set_password_mode(g_set_modal.mqtt_obs_token_ta, true);
|
||||
lv_textarea_set_max_length(g_set_modal.mqtt_obs_token_ta, 64);
|
||||
taSetPlaceholder(g_set_modal.mqtt_obs_token_ta, TR("Only required by token profiles"));
|
||||
attachSettingsTaEvents(g_set_modal.mqtt_obs_token_ta);
|
||||
y += SC(36);
|
||||
|
||||
// ---- Encryption key (PSK): seals decoded message payloads only. Observer
|
||||
// packets stay standard JSON so analyzer networks can consume them. ----
|
||||
lv_obj_t* psk_lbl = lv_label_create(body);
|
||||
lv_label_set_text(psk_lbl, TR("Encryption key (optional)"));
|
||||
lv_label_set_text(psk_lbl, TR("Message encryption key (optional)"));
|
||||
lv_obj_set_style_text_color(psk_lbl, lv_color_hex(COLOR_SUB), LV_PART_MAIN);
|
||||
lv_obj_set_style_text_font(psk_lbl, &g_font_12, LV_PART_MAIN);
|
||||
lv_obj_set_pos(psk_lbl, 2, y);
|
||||
@@ -16761,12 +16928,19 @@ static void buildMqttSettings() {
|
||||
// ---- Load current config ----
|
||||
{
|
||||
SdNvsPrefs p; // file-backed, matches MqttBridge (GH #128)
|
||||
bool cur_en = false, cur_dm = false, cur_ch = true, cur_consent = false;
|
||||
char cur_host[64] = {}, cur_port_s[8] = "1883", cur_user[32] = {}, cur_pwd[32] = {}, cur_psk[33] = {};
|
||||
bool cur_en = false, cur_dm = false, cur_ch = true, cur_obs = false, cur_consent = false;
|
||||
uint32_t cur_obs_profile = 0;
|
||||
char cur_host[64] = {}, cur_port_s[8] = "1883", cur_user[65] = {}, cur_pwd[97] = {}, cur_psk[33] = {};
|
||||
char cur_obs_origin[33] = {}, cur_obs_iata[8] = {}, cur_obs_topic[96] = "meshcore/{iata}/{device}";
|
||||
char cur_obs_token[65] = {};
|
||||
const NodePrefs* node_prefs = the_mesh.getNodePrefs();
|
||||
if (node_prefs) snprintf(cur_obs_origin, sizeof(cur_obs_origin), "%s", node_prefs->node_name);
|
||||
if (p.begin("mqtt", true)) {
|
||||
cur_en = p.getBool("en", false);
|
||||
cur_dm = p.getBool("dm", false);
|
||||
cur_ch = p.getBool("ch", true);
|
||||
cur_obs = p.getBool("obs", false);
|
||||
cur_obs_profile = p.getUInt("obs_profile", 0);
|
||||
cur_consent = p.getBool("consent", false);
|
||||
uint16_t port = (uint16_t)p.getUInt("port", 1883);
|
||||
snprintf(cur_port_s, sizeof(cur_port_s), "%u", port);
|
||||
@@ -16774,6 +16948,10 @@ static void buildMqttSettings() {
|
||||
if (p.isKey("user")) p.getString("user", cur_user, sizeof(cur_user));
|
||||
if (p.isKey("pwd")) p.getString("pwd", cur_pwd, sizeof(cur_pwd));
|
||||
if (p.isKey("psk")) p.getString("psk", cur_psk, sizeof(cur_psk));
|
||||
if (p.isKey("obs_origin")) p.getString("obs_origin", cur_obs_origin, sizeof(cur_obs_origin));
|
||||
if (p.isKey("obs_iata")) p.getString("obs_iata", cur_obs_iata, sizeof(cur_obs_iata));
|
||||
if (p.isKey("obs_topic")) p.getString("obs_topic", cur_obs_topic, sizeof(cur_obs_topic));
|
||||
if (p.isKey("obs_token")) p.getString("obs_token", cur_obs_token, sizeof(cur_obs_token));
|
||||
p.end();
|
||||
}
|
||||
if (cur_consent) {
|
||||
@@ -16783,10 +16961,18 @@ static void buildMqttSettings() {
|
||||
if (cur_en && cur_consent) lv_obj_add_state(g_set_modal.mqtt_en_sw, LV_STATE_CHECKED);
|
||||
if (cur_ch) lv_obj_add_state(g_set_modal.mqtt_ch_sw, LV_STATE_CHECKED);
|
||||
if (cur_dm) lv_obj_add_state(g_set_modal.mqtt_dm_sw, LV_STATE_CHECKED);
|
||||
if (cur_obs) lv_obj_add_state(g_set_modal.mqtt_obs_sw, LV_STATE_CHECKED);
|
||||
lv_dropdown_set_selected(g_set_modal.mqtt_obs_profile_dd,
|
||||
cur_obs_profile < MqttBridge::observerProfileCount()
|
||||
? (uint16_t)cur_obs_profile : 0);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_host_ta, cur_host);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_port_ta, cur_port_s);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_user_ta, cur_user);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_pwd_ta, cur_pwd);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_obs_origin_ta, cur_obs_origin);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_obs_iata_ta, cur_obs_iata);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_obs_topic_ta, cur_obs_topic);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_obs_token_ta, cur_obs_token);
|
||||
lv_textarea_set_text(g_set_modal.mqtt_psk_ta, cur_psk);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user