From 7e7279875cd3bf96e428f8ca17b3c900e41dc9d2 Mon Sep 17 00:00:00 2001 From: agessaman Date: Wed, 8 Jul 2026 10:49:06 -0700 Subject: [PATCH] feat(packet_capture): decode packet payloads to MQTT and log Add optional payload decoding to the packet capture service. GRP_TXT channel messages are decrypted (sender/text), ADVERTs are parsed (name/role/lat-lon), and a nested "decoded" object is attached to each packet alongside the unchanged raw fields. - Comprehensive channel key store: bot's configured radio channels plus decode_hashtag_channels, [Channels_List], decode_channel_keys, and the built-in default Public key. - Publishing the decoded object to MQTT is off by default and configurable per broker via mqttN_include_decoded. - Configurable packet-log rotation (off/size/time) for historical dumps. The decoder lives in a standalone, dependency-free module (modules/meshcore_payload_decode.py) so it can be shared verbatim with the meshcore-packet-capture project. Closes #197 Closes #35 --- config.ini.example | 42 ++- docs/packet-capture.md | 69 ++++ modules/meshcore_payload_decode.py | 338 ++++++++++++++++++ .../service_plugins/packet_capture_service.py | 194 +++++++++- .../test_packet_capture_payload_decode.py | 322 +++++++++++++++++ 5 files changed, 957 insertions(+), 8 deletions(-) create mode 100644 modules/meshcore_payload_decode.py create mode 100644 tests/unit/test_packet_capture_payload_decode.py diff --git a/config.ini.example b/config.ini.example index 4ee9df7..c9cacaa 100644 --- a/config.ini.example +++ b/config.ini.example @@ -1608,6 +1608,42 @@ mqtt_skip_unparseable_packets = true # (MeshCore createAdvert signature over pubkey || timestamp || app_data). File output is unchanged. advert_require_valid_signature = false +# --- Payload decoding (issues #197 & #35) --- +# When enabled, a nested "decoded" object is added to each packet's JSON (file + MQTT) with +# plain-text / structured fields: GRP_TXT channel messages are decrypted (sender, text), +# ADVERTs are parsed (name, role, lat/lon), and human-readable type/route labels + path are +# included. Raw fields are unchanged, so this is backward compatible. Default off. +# NOTE: Direct messages (TXT_MSG) are ECDH-encrypted between two nodes and cannot be decrypted +# by a passive observer; they are marked {"encrypted": true} instead. +decode_payloads = false + +# Global default for whether the "decoded" object is published to MQTT (default false: opt in +# per broker with mqttN_include_decoded, or set this to true to publish decoded to all brokers). +# The log file always follows decode_payloads, independent of this setting. +include_decoded = false + +# Channels used to decrypt GRP_TXT messages (in addition to the bot's own configured radio +# channels, which are used automatically): +# decode_hashtag_channels - comma-separated public/hashtag channel names; keys are derived +# from the name (SHA256 of "#name"). '#' optional. +# decode_channel_keys - comma-separated name=key pairs for custom channels; key may be +# 32-char hex or base64 (16 bytes). +# decode_include_public - include the well-known default "Public" channel key (default true). +decode_hashtag_channels = +decode_channel_keys = +decode_include_public = true + +# --- Packet log rotation (issue #35) --- +# Roll the output_file so historical dumps stay manageable. Default off = single appended file. +# log_rotation - off | size | time +# log_max_bytes - roll when the file reaches this size (size mode); accepts 50MB, 10M, etc. +# log_rotation_when - interval for time mode (midnight, H, D, W0-W6) - see TimedRotatingFileHandler +# log_backup_count - number of rolled files to keep +log_rotation = off +log_max_bytes = 50MB +log_rotation_when = midnight +log_backup_count = 5 + # Owner information (for packet analyzer registration) # Owner public key (64-character hex string) owner_public_key = @@ -1654,6 +1690,7 @@ iata = XYZ # mqttN_topic_prefix = # Legacy topic prefix (fallback if topic_status/topic_packets not set) # mqttN_client_id = # MQTT client ID (optional, auto-generated from bot name) # mqttN_upload_packet_types = # Comma-separated packet types to upload (e.g. 2,4); empty = all +# mqttN_include_decoded = true/false # Publish the decoded object to this broker (default: include_decoded) # # Topic template placeholders: # {IATA} - Uppercase IATA code (e.g., SEA) @@ -1672,8 +1709,9 @@ mqtt1_token_audience = mqtt-us-v1.letsmesh.net mqtt1_topic_status = meshcore/{IATA}/{PUBLIC_KEY}/status mqtt1_topic_packets = meshcore/{IATA}/{PUBLIC_KEY}/packets mqtt1_websocket_path = /mqtt -mqtt1_client_id = -mqtt1_upload_packet_types = +mqtt1_client_id = +mqtt1_upload_packet_types = +# mqtt1_include_decoded = true # Publish decoded payloads to this broker (default: include_decoded) # Optional per-broker JWT (inherit globals if omitted): # mqtt1_jwt_ttl_seconds = 3600 # JWT exp claim: iat + this many seconds # mqtt1_jwt_renewal_interval = 1800 # Refresh password this often; use < ttl; 0 = no renewal loop diff --git a/docs/packet-capture.md b/docs/packet-capture.md index ba09965..e741bd3 100644 --- a/docs/packet-capture.md +++ b/docs/packet-capture.md @@ -173,6 +173,54 @@ jwt_renewal_interval = 43200 # Default proactive refresh cadence (12 hours) } ``` +### Decoded Payloads + +When `decode_payloads = true`, each packet gains a nested `decoded` object with plain-text / +structured fields, in addition to the unchanged raw fields above. This makes dumps easy to +process with tools like `jq` (e.g. `jq 'select(.decoded.kind=="GRP_TXT") | .decoded.text'`). + +```json +{ + "type": "PACKET", + "packet_type": "5", + "route": "F", + "raw": "1540CAB3...", + "decoded": { + "kind": "GRP_TXT", + "channel_hash": "ca", + "channel": "#bot", + "sender": "Alice", + "text": "hello mesh", + "msg_timestamp": "2026-07-08T21:22:31Z", + "decrypted": true, + "path": ["A1", "B2"] + } +} +``` + +The `decoded` object holds only payload-specific content — it does not restate header fields +(`packet_type`, `route`) that already exist at the top level. + +What can be decoded: + +- **GRP_TXT** (channel messages) are decrypted when a matching channel key is available. + Keys come from the bot's own configured radio channels automatically, plus + `decode_hashtag_channels` (keys derived from the `#name`), `decode_channel_keys` + (`name=hexOrBase64` pairs), and the built-in default **Public** channel key + (`decode_include_public = true`). +- **ADVERT** packets are parsed into `name`, `mode` (role), `lat`/`lon`, and `public_key`. +- The decoded **path** hop list is included in `decoded.path` when it isn't already present at the + top level (the top-level `path` is only added for `route=D`), so flood-route paths are captured + without duplication. +- **Direct messages (TXT_MSG)** are ECDH-encrypted between two nodes and **cannot** be decrypted + by a passive observer — they appear as `{"kind": "TXT_MSG", "encrypted": true}`. + +Publishing of the `decoded` object to MQTT is **off by default** (`include_decoded = false`) — opt +in per broker with `mqttN_include_decoded = true`, or set `include_decoded = true` to publish it to +all brokers. This lets you, e.g., send decoded text to a private broker while public brokers receive +only raw packets. The log file always includes the `decoded` object when `decode_payloads = true`, +independent of this setting. + ### Status Message ```json { @@ -238,6 +286,27 @@ health_check_interval = 30 # Check connection every 30s health_check_grace_period = 2 # Allow 2 failures before warning ``` +### Log Rotation + +By default `output_file` is a single file that is appended to forever. To keep historical dumps +manageable, enable rotation: + +```ini +# Size-based: roll at 50 MB, keep 5 backups (packets.jsonl.1 ... .5) +log_rotation = size +log_max_bytes = 50MB +log_backup_count = 5 + +# Or time-based: roll daily at midnight, keep 14 days +log_rotation = time +log_rotation_when = midnight +log_backup_count = 14 +``` + +`log_rotation = off` (default) keeps the original single-file behavior. `log_max_bytes` accepts +plain bytes or suffixes like `10M` / `1G`. `log_rotation_when` uses Python's +`TimedRotatingFileHandler` values (`midnight`, `H`, `D`, `W0`–`W6`). + ### JWT Authentication Tokens are valid for 24 hours and auto-renewed. The service tries on-device signing first (if `auth_token_method = device`), then falls back to Python signing. diff --git a/modules/meshcore_payload_decode.py b/modules/meshcore_payload_decode.py new file mode 100644 index 0000000..f90e90e --- /dev/null +++ b/modules/meshcore_payload_decode.py @@ -0,0 +1,338 @@ +#!/usr/bin/env python3 +"""Standalone MeshCore packet-payload decoder. + +Decodes the *application payload* of a MeshCore packet into plain-text / +structured fields: GRP_TXT (channel) message decryption, ADVERT parsing, and +light structured fields for other payload types. + +This module is intentionally **self-contained** — it depends only on the Python +standard library plus ``cryptography`` (already a project dependency). It does +NOT import any bot-specific modules so that it can be copied verbatim into the +parent project ``meshcore-packet-capture`` (canonical home), mirroring the +existing ``auth_token.py`` <-> ``packet_capture_utils.py`` lineage. + +Key sourcing (which channel keys to try) is the host's responsibility: build a +:class:`ChannelKeyStore` from your own config / database and hand it to +:func:`decode_payload`. + +GRP_TXT wire format and crypto follow the reference implementation +https://github.com/michaelhart/meshcore-decoder : + + payload = channel_hash(1) + cipher_mac(2) + ciphertext(...) + + channel_hash = first byte of SHA256(channel_key_16) + MAC = HMAC_SHA256(key32, ciphertext)[:2], key32 = key16 + 16 zero bytes + cipher = AES-128-ECB, NoPadding, key = key16 + plaintext = timestamp(4, LE u32) + flags(1) + text(UTF-8, NUL-terminated), + text usually "sender: message" +""" + +from __future__ import annotations + +import hashlib +import hmac +import logging +from datetime import datetime, timezone +from typing import Any, Optional + +from cryptography.hazmat.backends import default_backend +from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes + +logger = logging.getLogger("meshcore_payload_decode") + +# --- Protocol constants (mirror modules/enums.py; kept local for portability) --- + +# PayloadType values (header bits 2-5). See modules/enums.py:PayloadType. +PT_REQ = 0x00 +PT_RESPONSE = 0x01 +PT_TXT_MSG = 0x02 +PT_ACK = 0x03 +PT_ADVERT = 0x04 +PT_GRP_TXT = 0x05 +PT_GRP_DATA = 0x06 +PT_ANON_REQ = 0x07 +PT_PATH = 0x08 +PT_TRACE = 0x09 +PT_MULTIPART = 0x0A +PT_RAW_CUSTOM = 0x0F + +PAYLOAD_TYPE_NAMES = { + PT_REQ: "REQ", + PT_RESPONSE: "RESPONSE", + PT_TXT_MSG: "TXT_MSG", + PT_ACK: "ACK", + PT_ADVERT: "ADVERT", + PT_GRP_TXT: "GRP_TXT", + PT_GRP_DATA: "GRP_DATA", + PT_ANON_REQ: "ANON_REQ", + PT_PATH: "PATH", + PT_TRACE: "TRACE", + PT_MULTIPART: "MULTIPART", + 0x0B: "Type11", + 0x0C: "Type12", + 0x0D: "Type13", + 0x0E: "Type14", + PT_RAW_CUSTOM: "RAW_CUSTOM", +} + +# Advert flag bits (see modules/enums.py:AdvertFlags / C++ AdvertDataHelpers.h) +ADV_TYPE_MASK = 0x0F +ADV_TYPE_CHAT = 0x01 +ADV_TYPE_REPEATER = 0x02 +ADV_TYPE_ROOM = 0x03 +ADV_TYPE_SENSOR = 0x04 +ADV_LATLON_MASK = 0x10 +ADV_FEAT1_MASK = 0x20 +ADV_FEAT2_MASK = 0x40 +ADV_NAME_MASK = 0x80 + +_ADV_TYPE_NAMES = { + ADV_TYPE_CHAT: "Companion", + ADV_TYPE_REPEATER: "Repeater", + ADV_TYPE_ROOM: "RoomServer", + ADV_TYPE_SENSOR: "Sensor", +} + +# The well-known MeshCore default "Public" channel key (base64 izOH6cXN6mrJ5e26oRXNcg==). +# NOTE: this is a fixed constant, NOT the hashtag derivation of "#public" +# (SHA256("#public")[:16] = 8b4b705b... which is different). +DEFAULT_PUBLIC_CHANNEL_KEY = bytes.fromhex("8b3387e9c5cdea6ac9e5edbaa115cd72") + + +def derive_hashtag_key(name: str) -> bytes: + """Derive a public/hashtag channel key from its name. + + The key is the first 16 bytes of SHA256 of the lowercased ``#name``. + + NOTE: This duplicates ``modules/channel_manager.py:generate_hashtag_key`` on + purpose — this module stays dependency-free for portability. If the MeshCore + derivation ever changes, update both (and the shared test vector). + """ + if not name.startswith("#"): + name = "#" + name + return hashlib.sha256(name.lower().encode("utf-8")).digest()[:16] + + +def channel_hash_for_key(key16: bytes) -> str: + """Return the 2-hex channel hash (first byte of SHA256(key)) for a channel key.""" + return f"{hashlib.sha256(key16).digest()[0]:02x}" + + +class ChannelKeyStore: + """Maps a channel hash -> candidate 16-byte keys (handles hash collisions).""" + + def __init__(self) -> None: + # channel_hash (2-hex, lower) -> list of (key16, name) + self._by_hash: dict[str, list[tuple[bytes, Optional[str]]]] = {} + + def add_secret(self, key16: bytes, name: Optional[str] = None) -> None: + """Add a raw 16-byte channel key (optionally with a display name).""" + if not key16 or len(key16) != 16: + logger.debug("Ignoring channel key with invalid length: %r", key16) + return + h = channel_hash_for_key(key16) + bucket = self._by_hash.setdefault(h, []) + if any(existing == key16 for existing, _ in bucket): + return # de-dup identical keys + bucket.append((key16, name)) + + def add_hex(self, key_hex: str, name: Optional[str] = None) -> None: + """Add a channel key from a 32-char hex string.""" + try: + self.add_secret(bytes.fromhex(key_hex.strip()), name) + except ValueError: + logger.debug("Ignoring non-hex channel key: %r", key_hex) + + def add_hashtag(self, name: str) -> None: + """Add a public/hashtag channel by name (key derived from the name).""" + normalized = name if name.startswith("#") else "#" + name + self.add_secret(derive_hashtag_key(name), normalized.lower()) + + def has(self, channel_hash: str) -> bool: + return channel_hash.lower() in self._by_hash + + def keys_for(self, channel_hash: str) -> list[tuple[bytes, Optional[str]]]: + return self._by_hash.get(channel_hash.lower(), []) + + def __len__(self) -> int: + return sum(len(v) for v in self._by_hash.values()) + + +def decrypt_group_text(ciphertext: bytes, cipher_mac: bytes, key16: bytes) -> Optional[dict[str, Any]]: + """Verify+decrypt a GRP_TXT ciphertext with a single channel key. + + Returns ``{timestamp, flags, sender, text}`` on success, or ``None`` if the + MAC fails or the plaintext is malformed. + """ + if len(ciphertext) < 16 or len(ciphertext) % 16 != 0: + return None + + # MAC: HMAC-SHA256 over ciphertext with 32-byte secret (key16 + 16 zero bytes) + key32 = key16 + b"\x00" * 16 + calc_mac = hmac.new(key32, ciphertext, hashlib.sha256).digest() + if not hmac.compare_digest(calc_mac[:2], cipher_mac[:2]): + return None + + # Decrypt: AES-128-ECB, no padding + try: + decryptor = Cipher(algorithms.AES(key16), modes.ECB(), backend=default_backend()).decryptor() + plaintext = decryptor.update(ciphertext) + decryptor.finalize() + except Exception as e: # pragma: no cover - defensive + logger.debug("AES decrypt failed: %s", e) + return None + + if len(plaintext) < 5: + return None + + timestamp = int.from_bytes(plaintext[0:4], "little") + flags = plaintext[4] + + text = plaintext[5:].decode("utf-8", errors="ignore") + nul = text.find("\x00") + if nul >= 0: + text = text[:nul] + + # Split "sender: message" when the prefix looks like a name + sender: Optional[str] = None + content = text + colon = text.find(": ") + if 0 < colon < 50: + candidate = text[:colon] + if not any(c in candidate for c in ":[]"): + sender = candidate + content = text[colon + 2:] + + return {"timestamp": timestamp, "flags": flags, "sender": sender, "text": content} + + +def _iso_utc(unix_ts: int) -> Optional[str]: + try: + return datetime.fromtimestamp(unix_ts, tz=timezone.utc).isoformat().replace("+00:00", "Z") + except (OverflowError, OSError, ValueError): + return None + + +def decode_group_text(payload: bytes, key_store: Optional[ChannelKeyStore]) -> dict[str, Any]: + """Decode (and, if a key matches, decrypt) a GRP_TXT payload.""" + if len(payload) < 3: + return {"kind": "GRP_TXT", "decrypted": False, "error": "payload_too_short"} + + channel_hash = f"{payload[0]:02x}" + cipher_mac = payload[1:3] + ciphertext = payload[3:] + + result: dict[str, Any] = { + "kind": "GRP_TXT", + "channel_hash": channel_hash, + "cipher_mac": cipher_mac.hex(), + "ciphertext_len": len(ciphertext), + "decrypted": False, + } + + if key_store and key_store.has(channel_hash): + for key16, name in key_store.keys_for(channel_hash): + decrypted = decrypt_group_text(ciphertext, cipher_mac, key16) + if decrypted: + result["decrypted"] = True + result["channel"] = name + result["sender"] = decrypted["sender"] + result["text"] = decrypted["text"] + result["flags"] = decrypted["flags"] + result["msg_timestamp"] = _iso_utc(decrypted["timestamp"]) + break + + return result + + +def parse_advert(payload: bytes) -> dict[str, Any]: + """Parse an ADVERT payload (port of meshcore-packet-capture parse_advert). + + Layout: pub_key(32) + timestamp(4) + signature(64) + app_data(flags + optional + latlon/feat1/feat2/name). + """ + result: dict[str, Any] = {"kind": "ADVERT"} + try: + if len(payload) < 100: + result.update({"advert_parse_ok": False, "advert_error": "payload_too_short_header"}) + return result + + result.update( + { + "advert_parse_ok": True, + "public_key": payload[0:32].hex(), + "advert_time": int.from_bytes(payload[32:36], "little"), + "signature": payload[36:100].hex(), + } + ) + + app_data = payload[100:] + if not app_data: + return result + + flags_byte = app_data[0] + adv_type = flags_byte & ADV_TYPE_MASK + result["mode"] = _ADV_TYPE_NAMES.get(adv_type, f"Type{adv_type}") + + i = 1 + if flags_byte & ADV_LATLON_MASK: + if len(app_data) < i + 8: + return result + lat = int.from_bytes(app_data[i:i + 4], "little", signed=True) + lon = int.from_bytes(app_data[i + 4:i + 8], "little", signed=True) + result["lat"] = round(lat / 1000000.0, 6) + result["lon"] = round(lon / 1000000.0, 6) + i += 8 + + if flags_byte & ADV_FEAT1_MASK: + if len(app_data) < i + 2: + return result + result["feat1"] = int.from_bytes(app_data[i:i + 2], "little") + i += 2 + + if flags_byte & ADV_FEAT2_MASK: + if len(app_data) < i + 2: + return result + result["feat2"] = int.from_bytes(app_data[i:i + 2], "little") + i += 2 + + if flags_byte & ADV_NAME_MASK and len(app_data) > i: + result["name"] = app_data[i:].decode("utf-8", errors="ignore").rstrip("\x00") + + return result + except Exception as e: # pragma: no cover - defensive + logger.debug("Error parsing ADVERT: %s", e) + result.update({"advert_parse_ok": False, "advert_error": "exception", "advert_error_detail": str(e)}) + return result + + +def decode_payload( + payload_type_value: int, + payload: bytes, + key_store: Optional[ChannelKeyStore] = None, +) -> dict[str, Any]: + """Decode a packet's application payload into structured / plain-text fields. + + Args: + payload_type_value: PayloadType (header bits 2-5), e.g. 5 for GRP_TXT. + payload: The application payload bytes (after header/transport/path). + key_store: Optional channel keys used to decrypt GRP_TXT messages. + + Returns: + A dict describing the decoded payload. Always contains ``kind``. + """ + if payload_type_value == PT_GRP_TXT: + return decode_group_text(payload, key_store) + if payload_type_value == PT_ADVERT: + return parse_advert(payload) + if payload_type_value == PT_TXT_MSG: + # Direct messages are ECDH-encrypted between two nodes; a passive + # observer cannot decrypt them. + return { + "kind": "TXT_MSG", + "encrypted": True, + "note": "direct message; not decryptable by observer", + } + if payload_type_value == PT_ACK: + return {"kind": "ACK", "ack": payload.hex()} + + return {"kind": PAYLOAD_TYPE_NAMES.get(payload_type_value, f"Type{payload_type_value}")} diff --git a/modules/service_plugins/packet_capture_service.py b/modules/service_plugins/packet_capture_service.py index 165b25f..97098f9 100644 --- a/modules/service_plugins/packet_capture_service.py +++ b/modules/service_plugins/packet_capture_service.py @@ -19,6 +19,11 @@ from meshcore import EventType # Import bot's enums from ..enums import PayloadType, PayloadVersion, RouteType +from ..meshcore_payload_decode import ( + DEFAULT_PUBLIC_CHANNEL_KEY, + ChannelKeyStore, + decode_payload, +) # Import bot's utilities for packet hash from ..utils import ( @@ -43,6 +48,81 @@ from .base_service import BaseServicePlugin from .packet_capture_utils import create_auth_token_async, read_private_key_file +def _decode_key_str(key_str: str) -> Optional[bytes]: + """Decode a 16-byte channel key from a hex (32 chars) or base64 string.""" + key_str = key_str.strip() + try: + if len(key_str) == 32 and all(c in "0123456789abcdefABCDEF" for c in key_str): + return bytes.fromhex(key_str) + import base64 + + raw = base64.b64decode(key_str, validate=True) + return raw if len(raw) == 16 else None + except (ValueError, TypeError): + return None + + +def _parse_size(value: str) -> int: + """Parse a size string like '50MB', '10M', or '1048576' into bytes.""" + value = str(value).strip().upper().replace("B", "") + multipliers = {"K": 1024, "M": 1024**2, "G": 1024**3} + try: + if value and value[-1] in multipliers: + return int(float(value[:-1]) * multipliers[value[-1]]) + return int(value or 0) + except ValueError: + return 0 + + +class _RotatingPacketLog: + """Writes one JSON line per packet, with optional size/time rotation. + + ``rotation='off'`` appends to a single file (original behavior). ``'size'`` + and ``'time'`` use stdlib rotating handlers. Kept host-agnostic so the same + logic can be ported to meshcore-packet-capture. + """ + + def __init__( + self, + path: str, + rotation: str = "off", + max_bytes: int = 0, + backup_count: int = 5, + when: str = "midnight", + ) -> None: + self._handler = None + self._fh = None + if rotation == "size" and max_bytes > 0: + from logging.handlers import RotatingFileHandler + + self._handler = RotatingFileHandler( + path, maxBytes=max_bytes, backupCount=backup_count, encoding="utf-8" + ) + elif rotation == "time": + from logging.handlers import TimedRotatingFileHandler + + self._handler = TimedRotatingFileHandler( + path, when=when, backupCount=backup_count, encoding="utf-8" + ) + else: + self._fh = open(path, "a", encoding="utf-8") + + def write_line(self, line: str) -> None: + if self._handler is not None: + # Default logging formatter emits just the message plus a terminator. + record = logging.LogRecord("packetlog", logging.INFO, "(packet)", 0, line, None, None) + self._handler.emit(record) + else: + self._fh.write(line + "\n") + self._fh.flush() + + def close(self) -> None: + if self._handler is not None: + self._handler.close() + elif self._fh is not None: + self._fh.close() + + class PacketCaptureService(BaseServicePlugin): """Packet capture service using bot's meshcore connection. @@ -186,6 +266,19 @@ class PacketCaptureService(BaseServicePlugin): self.debug = config.getboolean("PacketCapture", "debug", fallback=False) self._apply_log_level() + # Packet log rotation (off|size|time). Default off preserves single-file behavior. + self.log_rotation = config.get("PacketCapture", "log_rotation", fallback="off").strip().lower() + self.log_max_bytes = _parse_size(config.get("PacketCapture", "log_max_bytes", fallback="0")) + self.log_backup_count = config.getint("PacketCapture", "log_backup_count", fallback=5) + self.log_rotation_when = config.get("PacketCapture", "log_rotation_when", fallback="midnight").strip() + + # Payload decoding (decode plain text / structured payload into a nested "decoded" object) + self.decode_payloads = config.getboolean("PacketCapture", "decode_payloads", fallback=False) + # Global default for the per-broker mqttN_include_decoded toggle (default off: + # opt in per broker, or set include_decoded = true to publish decoded everywhere) + self.include_decoded = config.getboolean("PacketCapture", "include_decoded", fallback=False) + self.channel_key_store = self._build_channel_key_store(config) if self.decode_payloads else None + # MQTT configuration self.mqtt_enabled = config.getboolean("PacketCapture", "mqtt_enabled", fallback=True) self.mqtt_brokers = self._parse_mqtt_brokers(config) @@ -228,6 +321,60 @@ class PacketCaptureService(BaseServicePlugin): # The create_auth_token_async function will automatically try to export the key # from the device if private_key_hex is None and meshcore_instance is available + def _build_channel_key_store(self, config) -> ChannelKeyStore: + """Build a comprehensive channel key store for GRP_TXT decryption. + + Sources (mirrors the web viewer's channel sourcing, + modules/web_viewer/app.py:_get_additional_decode_channels): + 1. The bot's live radio channels (channel_manager, with real keys). + 2. ``decode_hashtag_channels`` + the ``[Channels_List]`` section (keys derived). + 3. Explicit ``decode_channel_keys`` = name=hexkey list. + 4. The well-known default Public channel key (unless disabled). + """ + store = ChannelKeyStore() + + # 1. Bot's configured radio channels (have real key material). + try: + channel_manager = getattr(self.bot, "channel_manager", None) + if channel_manager: + for ch in channel_manager.get_configured_channels(): + key_hex = ch.get("channel_key_hex") + if key_hex: + store.add_hex(key_hex, ch.get("channel_name")) + except Exception as e: + self.logger.debug(f"Could not load channel_manager keys: {e}") + + # 2a. decode_hashtag_channels: comma list of names (keys derived). + hashtag_raw = config.get("PacketCapture", "decode_hashtag_channels", fallback="") + for name in (n.strip() for n in hashtag_raw.split(",")): + if name: + store.add_hashtag(name) + + # 2b. [Channels_List] section names (keys derived), matching web viewer behavior. + if config.has_section("Channels_List"): + for key in config.options("Channels_List"): + name = key.split(".")[-1] if "." in key else key + name = name.strip() + if name: + store.add_hashtag(name) + + # 3. Explicit name=hexkey pairs (hex or base64). + keys_raw = config.get("PacketCapture", "decode_channel_keys", fallback="") + for entry in (e.strip() for e in keys_raw.split(",")): + if not entry or "=" not in entry: + continue + name, _, key_str = entry.partition("=") + key_bytes = _decode_key_str(key_str.strip()) + if key_bytes: + store.add_secret(key_bytes, name.strip()) + + # 4. Built-in default Public channel key. + if config.getboolean("PacketCapture", "decode_include_public", fallback=True): + store.add_secret(DEFAULT_PUBLIC_CHANNEL_KEY, "public") + + self.logger.info(f"Payload decoding enabled with {len(store)} channel key(s)") + return store + def _prune_correlation_caches(self, current_time: Optional[float] = None) -> None: """Drop stale rf_data_cache and recent_rf_packets entries. @@ -308,6 +455,11 @@ class PacketCaptureService(BaseServicePlugin): "websocket_path": config.get("PacketCapture", f"mqtt{broker_num}_websocket_path", fallback="/mqtt"), "client_id": config.get("PacketCapture", f"mqtt{broker_num}_client_id", fallback=None), "upload_packet_types": upload_packet_types, + "include_decoded": config.getboolean( + "PacketCapture", + f"mqtt{broker_num}_include_decoded", + fallback=getattr(self, "include_decoded", False), + ), "jwt_renewal_interval": jwt_renewal_interval, "jwt_ttl_seconds": jwt_ttl_seconds, } @@ -448,11 +600,18 @@ class PacketCaptureService(BaseServicePlugin): self.logger.info("Starting packet capture service...") - # Open output file if specified + # Open output file if specified (with optional size/time rotation) if self.output_file: try: - self.output_handle = open(self.output_file, "a") - self.logger.info(f"Writing packets to: {self.output_file}") + self.output_handle = _RotatingPacketLog( + self.output_file, + rotation=self.log_rotation, + max_bytes=self.log_max_bytes, + backup_count=self.log_backup_count, + when=self.log_rotation_when, + ) + rotation_note = "" if self.log_rotation == "off" else f" (rotation: {self.log_rotation})" + self.logger.info(f"Writing packets to: {self.output_file}{rotation_note}") except Exception as e: self.logger.error(f"Failed to open output file: {e}") @@ -824,6 +983,25 @@ class PacketCaptureService(BaseServicePlugin): if route == "D" and packet_info.get("path"): packet_data["path"] = ",".join(packet_info["path"]) + # Attach decoded payload (issues #197 & #35): plain text / structured fields. + # Only payload-specific content goes here — header fields (packet_type, route) + # already exist at the top level, so we don't restate them. + if self.decode_payloads and self.channel_key_store is not None: + try: + payload_type_value = packet_info.get("payload_type_value", 0) + if hasattr(payload_type_value, "value"): + payload_type_value = payload_type_value.value + payload_bytes = bytes.fromhex(packet_info.get("payload_hex", "") or "") + decoded = decode_payload(int(payload_type_value), payload_bytes, self.channel_key_store) + # Include the decoded hop path only when it isn't already at the top level + # (top-level "path" is added for route=D) — captures flood paths without duplicating. + if packet_info.get("path") and "path" not in packet_data: + decoded["path"] = list(packet_info["path"]) + packet_data["decoded"] = decoded + except Exception as e: + if self.debug: + self.logger.debug(f"Payload decode failed: {e}") + return packet_data async def process_packet( @@ -901,8 +1079,7 @@ class PacketCaptureService(BaseServicePlugin): # Write to file if self.output_handle: - self.output_handle.write(json.dumps(formatted_packet, default=str) + "\n") - self.output_handle.flush() + self.output_handle.write_line(json.dumps(formatted_packet, default=str)) # Publish to MQTT if enabled # The publish function will check per-broker connection status @@ -1510,7 +1687,12 @@ class PacketCaptureService(BaseServicePlugin): if not topic: continue - payload = json.dumps(packet_info, default=str) + # Per-broker: strip the decoded object for brokers that opt out. + broker_packet = packet_info + if "decoded" in packet_info and not config.get("include_decoded", False): + broker_packet = {k: v for k, v in packet_info.items() if k != "decoded"} + + payload = json.dumps(broker_packet, default=str) # Log topic and payload size for debugging self.logger.debug(f"Publishing to topic '{topic}' on {config['host']} (payload: {len(payload)} bytes)") diff --git a/tests/unit/test_packet_capture_payload_decode.py b/tests/unit/test_packet_capture_payload_decode.py new file mode 100644 index 0000000..c0cc1a8 --- /dev/null +++ b/tests/unit/test_packet_capture_payload_decode.py @@ -0,0 +1,322 @@ +"""Tests for MeshCore payload decoding (issues #197 & #35). + +Covers the standalone decoder (GRP_TXT decryption, ADVERT parsing, key store), +the service's nested ``decoded`` object, the per-broker include_decoded toggle, +and packet-log rotation. +""" + +from __future__ import annotations + +import asyncio +import configparser +import json +import logging +from unittest.mock import MagicMock + +from modules.meshcore_payload_decode import ( + DEFAULT_PUBLIC_CHANNEL_KEY, + ChannelKeyStore, + channel_hash_for_key, + decode_payload, + decrypt_group_text, + derive_hashtag_key, + parse_advert, +) +from modules.service_plugins.packet_capture_service import ( + PacketCaptureService, + _decode_key_str, + _parse_size, + _RotatingPacketLog, +) + +# Known GRP_TXT vector from michaelhart/meshcore-decoder tests. +# raw = header(0x15) + path_len(0x40=0 hops) + payload; key is the "#bot" channel key. +BOT_KEY = bytes.fromhex("eb50a1bcb3e4e5d7bf69a57c9dada211") +GRP_RAW = "1540cab3b15626481a5ba64247ab25766e410b026e0678a32da9f0c3946fae5b714cab170f" +GRP_PAYLOAD = bytes.fromhex(GRP_RAW[4:]) # after header + path bytes + + +# --------------------------------------------------------------------------- # +# Standalone decoder +# --------------------------------------------------------------------------- # +def test_decode_group_text_decrypts_known_vector(): + store = ChannelKeyStore() + store.add_secret(BOT_KEY, "#bot") + result = decode_payload(5, GRP_PAYLOAD, store) + assert result["kind"] == "GRP_TXT" + assert result["decrypted"] is True + assert result["channel"] == "#bot" + assert result["sender"] == "Howl 👾" + assert result["text"] == "prefix 0101" + assert result["msg_timestamp"].endswith("Z") + + +def test_group_text_mac_reject_without_key(): + result = decode_payload(5, GRP_PAYLOAD, ChannelKeyStore()) + assert result["decrypted"] is False + assert "text" not in result + + +def test_group_text_wrong_key_fails_mac(): + # A key with the same channel hash bucket is not guaranteed; verify a + # deliberately wrong key does not decrypt. + store = ChannelKeyStore() + wrong = bytes(16) # all-zero key + # Force it into the same hash bucket so keys_for() returns it. + store._by_hash[f"{GRP_PAYLOAD[0]:02x}"] = [(wrong, "wrong")] + result = decode_payload(5, GRP_PAYLOAD, store) + assert result["decrypted"] is False + + +def test_channel_hash_and_keystore_collision(): + store = ChannelKeyStore() + store.add_secret(BOT_KEY, "#bot") + h = channel_hash_for_key(BOT_KEY) + assert store.has(h) + # Add a second (wrong) key under the same hash bucket -> both tried. + store._by_hash[h].append((bytes(16), "decoy")) + assert len(store.keys_for(h)) == 2 + # Decryption still succeeds because the real key is tried. + assert decode_payload(5, GRP_PAYLOAD, store)["decrypted"] is True + + +def test_derive_hashtag_key_matches_channel_manager(): + from modules.channel_manager import ChannelManager + + assert derive_hashtag_key("bot") == ChannelManager.generate_hashtag_key("#bot") + assert derive_hashtag_key("#bot") == derive_hashtag_key("bot") + + +def test_default_public_key_hash_is_fixed_constant(): + # The default Public key is NOT the hashtag derivation of "#public". + assert derive_hashtag_key("public") != DEFAULT_PUBLIC_CHANNEL_KEY + assert channel_hash_for_key(DEFAULT_PUBLIC_CHANNEL_KEY) == "11" + + +def test_decrypt_group_text_rejects_unaligned_ciphertext(): + assert decrypt_group_text(b"\x00" * 15, b"\x00\x00", BOT_KEY) is None + + +def test_txt_msg_marked_not_decryptable(): + result = decode_payload(2, b"\x00" * 20, None) + assert result["kind"] == "TXT_MSG" + assert result["encrypted"] is True + + +def test_parse_advert_extracts_fields(): + pub_key = bytes(range(32)) + timestamp = (1_700_000_000).to_bytes(4, "little") + signature = bytes(64) + # flags: name (0x80) | latlon (0x10) | chat (0x01) + flags = bytes([0x91]) + lat = (47_600_000).to_bytes(4, "little", signed=True) + lon = (-122_300_000 & 0xFFFFFFFF).to_bytes(4, "little") + name = b"TestNode" + payload = pub_key + timestamp + signature + flags + lat + lon + name + + advert = parse_advert(payload) + assert advert["kind"] == "ADVERT" + assert advert["advert_parse_ok"] is True + assert advert["mode"] == "Companion" + assert advert["name"] == "TestNode" + assert advert["lat"] == 47.6 + assert advert["public_key"] == pub_key.hex() + + +# --------------------------------------------------------------------------- # +# Service integration +# --------------------------------------------------------------------------- # +def _service_for_format(decode_payloads: bool) -> PacketCaptureService: + svc = object.__new__(PacketCaptureService) + svc.decode_payloads = decode_payloads + svc.debug = False + svc.logger = logging.getLogger("test-packet-capture") + svc.bot = None + svc._get_bot_name = lambda: "TestBot" # type: ignore[method-assign] + if decode_payloads: + svc.channel_key_store = ChannelKeyStore() + svc.channel_key_store.add_secret(BOT_KEY, "#bot") + else: + svc.channel_key_store = None + return svc + + +def _grp_packet_info() -> dict: + return { + "route_type": "FLOOD", + "payload_type": "GRP_TXT", + "payload_type_value": 5, + "path": [], + "path_byte_length": 0, + "payload_hex": GRP_PAYLOAD.hex(), + "payload_bytes": len(GRP_PAYLOAD), + "packet_hash": "ABCDEF0123456789", + } + + +def test_format_packet_data_attaches_decoded(): + svc = _service_for_format(True) + out = PacketCaptureService._format_packet_data( + svc, GRP_RAW, _grp_packet_info(), {"snr": 5, "rssi": -100}, None + ) + assert "decoded" in out + assert out["decoded"]["text"] == "prefix 0101" + assert out["decoded"]["sender"] == "Howl 👾" + # Header fields are not restated inside decoded (no redundancy with top level) + assert "type_label" not in out["decoded"] + assert "route_label" not in out["decoded"] + + +def test_decoded_path_only_when_not_at_top_level(): + svc = _service_for_format(True) + + # Flood route with hops: top level has no "path", so decoded carries it. + flood = _grp_packet_info() + flood["path"] = ["F8DA", "7A2A"] + out = PacketCaptureService._format_packet_data( + svc, GRP_RAW, flood, {"snr": 5, "rssi": -100}, None + ) + assert "path" not in out # top-level path is route=D only + assert out["decoded"]["path"] == ["F8DA", "7A2A"] + + # Direct route with hops: top level has "path", so decoded does NOT duplicate it. + direct = _grp_packet_info() + direct["route_type"] = "DIRECT" + direct["path"] = ["AB", "CD"] + out = PacketCaptureService._format_packet_data( + svc, GRP_RAW, direct, {"snr": 5, "rssi": -100}, None + ) + assert out["path"] == "AB,CD" + assert "path" not in out["decoded"] + + +def test_format_packet_data_no_decoded_when_disabled(): + svc = _service_for_format(False) + out = PacketCaptureService._format_packet_data( + svc, GRP_RAW, _grp_packet_info(), {"snr": 5, "rssi": -100}, None + ) + assert "decoded" not in out + + +def _bot_from_ini(ini: str) -> MagicMock: + cp = configparser.ConfigParser() + cp.read_string(ini.strip()) + bot = MagicMock() + bot.config = cp + return bot + + +def test_per_broker_include_decoded_flag_parsing(): + bot = _bot_from_ini( + """ + [PacketCapture] + enabled = false + include_decoded = true + mqtt1_server = full.example + mqtt2_server = raw-only.example + mqtt2_include_decoded = false + """ + ) + svc = object.__new__(PacketCaptureService) + svc.bot = bot + svc.include_decoded = True + brokers = PacketCaptureService._parse_mqtt_brokers(svc, bot.config) + assert brokers[0]["include_decoded"] is True + assert brokers[1]["include_decoded"] is False + + +def test_include_decoded_defaults_off(): + # With nothing configured, the global default is off and brokers inherit it. + bot = _bot_from_ini( + """ + [PacketCapture] + enabled = false + mqtt1_server = a.example + mqtt2_server = b.example + mqtt2_include_decoded = true + """ + ) + svc = object.__new__(PacketCaptureService) + svc.bot = bot + svc.include_decoded = False + brokers = PacketCaptureService._parse_mqtt_brokers(svc, bot.config) + assert brokers[0]["include_decoded"] is False # inherits global default (off) + assert brokers[1]["include_decoded"] is True # explicit opt-in + + +def test_publish_strips_decoded_per_broker(): + published: dict[str, dict] = {} + + def make_client(name): + client = MagicMock() + + def publish(topic, payload, qos=0): + published[name] = json.loads(payload) + result = MagicMock() + result.rc = 0 + return result + + client.publish.side_effect = publish + return client + + svc = object.__new__(PacketCaptureService) + svc.logger = logging.getLogger("test-publish") + svc.packet_count = 1 + svc.bot = None + svc.mqtt_clients = [ + { + "connected": True, + "client": make_client("full"), + "config": {"host": "full", "topic_prefix": "mc/full", "include_decoded": True}, + }, + { + "connected": True, + "client": make_client("raw"), + "config": {"host": "raw", "topic_prefix": "mc/raw", "include_decoded": False}, + }, + ] + + packet = {"type": "PACKET", "packet_type": "5", "raw": "1540AA", "decoded": {"text": "hi"}} + asyncio.run(PacketCaptureService.publish_packet_mqtt(svc, packet)) + + assert "decoded" in published["full"] + assert "decoded" not in published["raw"] + assert published["raw"]["raw"] == "1540AA" + + +# --------------------------------------------------------------------------- # +# Helpers & rotation +# --------------------------------------------------------------------------- # +def test_parse_size_units(): + assert _parse_size("50MB") == 50 * 1024 * 1024 + assert _parse_size("10M") == 10 * 1024 * 1024 + assert _parse_size("1048576") == 1048576 + assert _parse_size("0") == 0 + assert _parse_size("garbage") == 0 + + +def test_decode_key_str_hex_and_base64(): + assert _decode_key_str("eb50a1bcb3e4e5d7bf69a57c9dada211") == BOT_KEY + assert _decode_key_str("izOH6cXN6mrJ5e26oRXNcg==") == DEFAULT_PUBLIC_CHANNEL_KEY + assert _decode_key_str("not-a-key") is None + + +def test_rotating_log_off_single_file(tmp_path): + path = tmp_path / "packets.jsonl" + log = _RotatingPacketLog(str(path)) + log.write_line('{"a": 1}') + log.write_line('{"b": 2}') + log.close() + assert path.read_text().splitlines() == ['{"a": 1}', '{"b": 2}'] + + +def test_rotating_log_size_rolls(tmp_path): + path = tmp_path / "packets.jsonl" + log = _RotatingPacketLog(str(path), rotation="size", max_bytes=200, backup_count=3) + for i in range(100): + log.write_line(json.dumps({"packet": i, "pad": "x" * 20})) + log.close() + backups = list(tmp_path.glob("packets.jsonl.*")) + assert path.exists() + assert len(backups) >= 1 # rotation produced at least one backup + assert len(backups) <= 3 # backup_count respected