added receiver plugin
This commit is contained in:
@@ -0,0 +1,289 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Receive Meshtastic mesh text messages over MQTT and expose them in Home Assistant.
|
||||
|
||||
Flow:
|
||||
1. Subscribe to a Meshtastic gateway's MQTT topic (public server or local broker).
|
||||
2. Decrypt channel packets (AES-CTR) using the channel PSK.
|
||||
3. Publish the latest text message and parsed Temperature/Humidity values to
|
||||
Home Assistant via MQTT Discovery (so entities appear automatically).
|
||||
"""
|
||||
|
||||
import base64
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
|
||||
import paho.mqtt.client as mqtt
|
||||
import requests
|
||||
from cryptography.hazmat.backends import default_backend
|
||||
from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes
|
||||
from meshtastic.protobuf import mesh_pb2, mqtt_pb2, portnums_pb2
|
||||
|
||||
|
||||
def log(msg):
|
||||
print(msg, file=sys.stdout, flush=True)
|
||||
|
||||
|
||||
# --- Add-on options (written by the Supervisor) ------------------------------
|
||||
with open("/data/options.json", encoding="utf-8") as fh:
|
||||
OPTS = json.load(fh)
|
||||
|
||||
SRC_BROKER = OPTS["source_broker"]
|
||||
SRC_PORT = int(OPTS.get("source_port", 1883))
|
||||
SRC_USER = OPTS.get("source_username", "") or ""
|
||||
SRC_PASS = OPTS.get("source_password", "") or ""
|
||||
SRC_TLS = bool(OPTS.get("source_tls", False))
|
||||
REGION = OPTS.get("region", "+") or "+"
|
||||
CHANNEL_NAME = OPTS["channel_name"]
|
||||
CHANNEL_PSK = OPTS.get("channel_psk", "AQ==")
|
||||
MATCH_PREFIX = OPTS.get("match_prefix", "") or ""
|
||||
DEVICE_NAME = OPTS.get("device_name", "Meshtastic Mesh")
|
||||
DISCOVERY_PREFIX = OPTS.get("discovery_prefix", "homeassistant")
|
||||
SENSORS = OPTS.get("sensors", [])
|
||||
|
||||
DEVICE_UID = re.sub(r"[^a-z0-9_]", "_", DEVICE_NAME.lower()) or "meshtastic_mesh"
|
||||
BASE_TOPIC = f"meshtastic_receiver/{DEVICE_UID}"
|
||||
AVAIL_TOPIC = f"{BASE_TOPIC}/status"
|
||||
|
||||
# --- PSK handling ------------------------------------------------------------
|
||||
# "AQ==" (byte 0x01) is Meshtastic's shorthand for this well-known default key.
|
||||
_DEFAULT_KEY = base64.b64decode("1PG7OiApB1nwvP+rz05pAQ==")
|
||||
|
||||
|
||||
def expand_psk(psk_b64):
|
||||
raw = base64.b64decode(psk_b64)
|
||||
return _DEFAULT_KEY if raw == b"\x01" else raw
|
||||
|
||||
|
||||
CHANNEL_KEY = expand_psk(CHANNEL_PSK)
|
||||
|
||||
# Matches "T1:20.8", "H2:54.4", "P:1013" etc. (unit suffix is ignored here).
|
||||
PAIR_RE = re.compile(r"([A-Za-z][A-Za-z0-9_]*)\s*:\s*(-?\d+(?:\.\d+)?)")
|
||||
|
||||
# Set once the HA broker client is connected.
|
||||
ha_client = None
|
||||
|
||||
|
||||
# --- Decryption --------------------------------------------------------------
|
||||
def decrypt_packet(encrypted, key, from_id, packet_id):
|
||||
"""AES-CTR decrypt. Nonce = packetId(uint64 LE) | fromNode(uint32 LE) | 0*4."""
|
||||
try:
|
||||
nonce = (
|
||||
packet_id.to_bytes(4, "little")
|
||||
+ b"\x00" * 4
|
||||
+ from_id.to_bytes(4, "little")
|
||||
+ b"\x00" * 4
|
||||
)
|
||||
cipher = Cipher(
|
||||
algorithms.AES(key), modes.CTR(nonce), backend=default_backend()
|
||||
)
|
||||
d = cipher.decryptor()
|
||||
return d.update(encrypted) + d.finalize()
|
||||
except Exception as exc: # noqa: BLE001
|
||||
log(f"[warn] decryption error: {exc}")
|
||||
return None
|
||||
|
||||
|
||||
# --- Home Assistant MQTT broker resolution (via Supervisor service) ----------
|
||||
def resolve_ha_mqtt():
|
||||
token = os.environ.get("SUPERVISOR_TOKEN")
|
||||
if not token:
|
||||
return None
|
||||
try:
|
||||
r = requests.get(
|
||||
"http://supervisor/services/mqtt",
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=10,
|
||||
)
|
||||
if r.status_code == 200:
|
||||
d = r.json()["data"]
|
||||
return {
|
||||
"host": d["host"],
|
||||
"port": int(d["port"]),
|
||||
"username": d.get("username") or None,
|
||||
"password": d.get("password") or None,
|
||||
"ssl": bool(d.get("ssl", False)),
|
||||
}
|
||||
log(f"[warn] supervisor mqtt service returned HTTP {r.status_code}")
|
||||
except Exception as exc: # noqa: BLE001
|
||||
log(f"[warn] could not resolve HA MQTT via supervisor: {exc}")
|
||||
return None
|
||||
|
||||
|
||||
# --- MQTT Discovery ----------------------------------------------------------
|
||||
def device_block():
|
||||
return {
|
||||
"identifiers": [DEVICE_UID],
|
||||
"name": DEVICE_NAME,
|
||||
"manufacturer": "Meshtastic",
|
||||
"model": "MQTT Mesh Receiver",
|
||||
}
|
||||
|
||||
|
||||
def publish_discovery():
|
||||
# Latest text message (with from/to/time as attributes).
|
||||
cfg = {
|
||||
"name": "Last Message",
|
||||
"unique_id": f"{DEVICE_UID}_last_message",
|
||||
"state_topic": f"{BASE_TOPIC}/last_message",
|
||||
"json_attributes_topic": f"{BASE_TOPIC}/last_message/attributes",
|
||||
"icon": "mdi:message-text",
|
||||
"availability_topic": AVAIL_TOPIC,
|
||||
"device": device_block(),
|
||||
}
|
||||
ha_client.publish(
|
||||
f"{DISCOVERY_PREFIX}/sensor/{DEVICE_UID}/last_message/config",
|
||||
json.dumps(cfg),
|
||||
retain=True,
|
||||
)
|
||||
|
||||
for s in SENSORS:
|
||||
label = s["label"]
|
||||
cfg = {
|
||||
"name": s.get("name", label),
|
||||
"unique_id": f"{DEVICE_UID}_{label}",
|
||||
"state_topic": f"{BASE_TOPIC}/{label}",
|
||||
"availability_topic": AVAIL_TOPIC,
|
||||
"device": device_block(),
|
||||
}
|
||||
if s.get("unit"):
|
||||
cfg["unit_of_measurement"] = s["unit"]
|
||||
if s.get("device_class"):
|
||||
cfg["device_class"] = s["device_class"]
|
||||
cfg["state_class"] = "measurement"
|
||||
ha_client.publish(
|
||||
f"{DISCOVERY_PREFIX}/sensor/{DEVICE_UID}/{label}/config",
|
||||
json.dumps(cfg),
|
||||
retain=True,
|
||||
)
|
||||
log(f"[ha] published discovery for last_message + {len(SENSORS)} sensors")
|
||||
|
||||
|
||||
# --- Message handling --------------------------------------------------------
|
||||
def handle_text(text, from_id, to_id):
|
||||
# HA sensor states are limited to 255 chars.
|
||||
ha_client.publish(f"{BASE_TOPIC}/last_message", text[:255], retain=True)
|
||||
ha_client.publish(
|
||||
f"{BASE_TOPIC}/last_message/attributes",
|
||||
json.dumps(
|
||||
{
|
||||
"from": f"!{from_id:08x}",
|
||||
"to": f"!{to_id:08x}",
|
||||
"received": time.strftime("%Y-%m-%dT%H:%M:%S%z"),
|
||||
"full_text": text,
|
||||
}
|
||||
),
|
||||
retain=True,
|
||||
)
|
||||
|
||||
if MATCH_PREFIX and MATCH_PREFIX not in text:
|
||||
return
|
||||
|
||||
known = {s["label"] for s in SENSORS}
|
||||
published = []
|
||||
for label, value in PAIR_RE.findall(text):
|
||||
if label in known:
|
||||
ha_client.publish(f"{BASE_TOPIC}/{label}", value, retain=True)
|
||||
published.append(f"{label}={value}")
|
||||
if published:
|
||||
log(f"[parse] {' '.join(published)}")
|
||||
|
||||
|
||||
# --- Source MQTT callbacks ---------------------------------------------------
|
||||
def on_src_connect(client, userdata, flags, reason_code, properties=None):
|
||||
if reason_code == 0:
|
||||
topic = f"msh/{REGION}/2/e/{CHANNEL_NAME}/#"
|
||||
client.subscribe(topic)
|
||||
log(f"[source] connected; subscribed to {topic}")
|
||||
else:
|
||||
log(f"[source] connection refused: reason_code={reason_code}")
|
||||
|
||||
|
||||
def on_src_message(client, userdata, msg):
|
||||
try:
|
||||
env = mqtt_pb2.ServiceEnvelope()
|
||||
env.ParseFromString(msg.payload)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
log(f"[warn] could not parse ServiceEnvelope: {exc}")
|
||||
return
|
||||
|
||||
mp = env.packet
|
||||
from_id = getattr(mp, "from")
|
||||
|
||||
if mp.HasField("decoded"):
|
||||
data = mp.decoded
|
||||
elif mp.HasField("encrypted"):
|
||||
plain = decrypt_packet(mp.encrypted, CHANNEL_KEY, from_id, mp.id)
|
||||
if plain is None:
|
||||
return
|
||||
data = mesh_pb2.Data()
|
||||
try:
|
||||
data.ParseFromString(plain)
|
||||
except Exception: # noqa: BLE001
|
||||
return # wrong PSK or not for this channel
|
||||
else:
|
||||
return
|
||||
|
||||
if data.portnum == portnums_pb2.PortNum.TEXT_MESSAGE_APP:
|
||||
text = data.payload.decode("utf-8", errors="replace")
|
||||
log(f"[rx] !{from_id:08x}: {text}")
|
||||
handle_text(text, from_id, mp.to)
|
||||
|
||||
|
||||
# --- Entry point -------------------------------------------------------------
|
||||
def main():
|
||||
global ha_client
|
||||
|
||||
ha = resolve_ha_mqtt()
|
||||
if ha is None:
|
||||
log(
|
||||
"[fatal] No MQTT broker available. Install and start the Mosquitto "
|
||||
"broker add-on, then restart this add-on."
|
||||
)
|
||||
sys.exit(1)
|
||||
log(f"[ha] using MQTT broker {ha['host']}:{ha['port']}")
|
||||
|
||||
ha_client = mqtt.Client(
|
||||
mqtt.CallbackAPIVersion.VERSION2, client_id="meshtastic_receiver"
|
||||
)
|
||||
if ha["username"]:
|
||||
ha_client.username_pw_set(ha["username"], ha["password"])
|
||||
if ha["ssl"]:
|
||||
ha_client.tls_set()
|
||||
ha_client.will_set(AVAIL_TOPIC, "offline", retain=True)
|
||||
|
||||
def on_ha_connect(client, userdata, flags, reason_code, properties=None):
|
||||
if reason_code == 0:
|
||||
log("[ha] connected")
|
||||
client.publish(AVAIL_TOPIC, "online", retain=True)
|
||||
publish_discovery()
|
||||
else:
|
||||
log(f"[ha] connection refused: reason_code={reason_code}")
|
||||
|
||||
ha_client.on_connect = on_ha_connect
|
||||
ha_client.connect(ha["host"], ha["port"], keepalive=60)
|
||||
ha_client.loop_start()
|
||||
|
||||
src = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2)
|
||||
if SRC_USER:
|
||||
src.username_pw_set(SRC_USER, SRC_PASS)
|
||||
if SRC_TLS:
|
||||
src.tls_set()
|
||||
src.on_connect = on_src_connect
|
||||
src.on_message = on_src_message
|
||||
log(f"[source] connecting to {SRC_BROKER}:{SRC_PORT} …")
|
||||
src.connect(SRC_BROKER, SRC_PORT, keepalive=60)
|
||||
|
||||
try:
|
||||
src.loop_forever()
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
finally:
|
||||
ha_client.publish(AVAIL_TOPIC, "offline", retain=True)
|
||||
ha_client.loop_stop()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user