From bf6b86696c0d65d84edfee1dcfbff03caaa41184 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sun, 12 Apr 2026 18:03:20 +0100 Subject: [PATCH] Migrate from Node.js meshcore-decoder CLI to native Python meshcoredecoder library Replace subprocess-based packet decoding with native Python meshcoredecoder>=0.3.2. The decoder is now always enabled (no toggle), removing ENABLED/COMMAND/TIMEOUT config vars. Adds _enrich_payload_decoded() to compensate for payload classes that lack to_dict() overrides in the library. Removes Node.js/npm from Dockerfile and deletes the patches/ directory. Adds parenthesized exception rule to AGENTS.md. --- .env.example | 11 +- AGENTS.md | 1 + Dockerfile | 17 - README.md | 5 +- docker-compose.yml | 3 - .../@michaelhart+meshcore-decoder+0.2.7.patch | 58 -- pyproject.toml | 2 + src/meshcore_hub/collector/cli.py | 22 +- .../collector/letsmesh_decoder.py | 180 ++--- src/meshcore_hub/collector/subscriber.py | 27 - src/meshcore_hub/common/config.py | 17 - src/meshcore_hub/web/app.py | 1 - tests/test_collector/test_letsmesh_decoder.py | 195 ++++-- .../test_letsmesh_normalizer_integration.py | 624 ++++++++++++++++++ 14 files changed, 856 insertions(+), 307 deletions(-) delete mode 100644 patches/@michaelhart+meshcore-decoder+0.2.7.patch create mode 100644 tests/test_collector/test_letsmesh_normalizer_integration.py diff --git a/.env.example b/.env.example index 7ef9556..f9e46e0 100644 --- a/.env.example +++ b/.env.example @@ -163,13 +163,7 @@ PACKETCAPTURE_EXIT_ON_RECONNECT_FAIL=true # The collector subscribes to MQTT events and stores them in the database # LetsMesh decoder support -# Set to false to disable external packet decoding -COLLECTOR_LETSMESH_DECODER_ENABLED=true - -# Decoder command (must be available in container PATH) -# Examples: meshcore-decoder, /usr/local/bin/meshcore-decoder, npx meshcore-decoder -COLLECTOR_LETSMESH_DECODER_COMMAND=meshcore-decoder - +# The native Python decoder is always enabled. # Optional: channel secret keys (comma or space separated) used to decrypt GroupText # packets. This supports unlimited keys. # Note: Public + #test keys are built into the collector code by default. @@ -177,9 +171,6 @@ COLLECTOR_LETSMESH_DECODER_COMMAND=meshcore-decoder # Without keys, encrypted packets cannot be shown as plaintext. # COLLECTOR_LETSMESH_DECODER_KEYS= -# Timeout in seconds per decode invocation -COLLECTOR_LETSMESH_DECODER_TIMEOUT_SECONDS=2.0 - # ------------------- # Webhook Settings # ------------------- diff --git a/AGENTS.md b/AGENTS.md index 725a98b..d9c12aa 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -67,6 +67,7 @@ MeshCore Hub is a Python 3.13+ monorepo for managing and orchestrating MeshCore - Use type hints for all function signatures - Write docstrings for public modules, classes, and functions - Keep functions focused and under 50 lines where possible +- **Always use parenthesized exception tuples** — `except (ValueError, TypeError):` not `except ValueError, TypeError:`. The comma form is Python 2 syntax and will fail at import time in Python 3 ### Imports diff --git a/Dockerfile b/Dockerfile index 3fb65db..31eb721 100644 --- a/Dockerfile +++ b/Dockerfile @@ -65,26 +65,9 @@ ENV PYTHONDONTWRITEBYTECODE=1 \ RUN apt-get update && apt-get install -y --no-install-recommends \ # For serial port access udev \ - # LetsMesh decoder runtime - nodejs \ - npm \ && rm -rf /var/lib/apt/lists/* \ && mkdir -p /data -# Install meshcore-decoder CLI. -RUN mkdir -p /opt/letsmesh-decoder \ - && cd /opt/letsmesh-decoder \ - && npm init -y >/dev/null 2>&1 \ - && npm install --omit=dev @michaelhart/meshcore-decoder@0.2.7 patch-package - -# Apply maintained meshcore-decoder compatibility patch. -COPY patches/@michaelhart+meshcore-decoder+0.2.7.patch /opt/letsmesh-decoder/patches/@michaelhart+meshcore-decoder+0.2.7.patch -RUN cd /opt/letsmesh-decoder \ - && npx patch-package --error-on-fail \ - && npm uninstall patch-package \ - && npm prune --omit=dev -RUN ln -s /opt/letsmesh-decoder/node_modules/.bin/meshcore-decoder /usr/local/bin/meshcore-decoder - # Copy virtual environment from builder COPY --from=builder /opt/venv /opt/venv ENV PATH="/opt/venv/bin:$PATH" diff --git a/README.md b/README.md index 2821e63..72831b1 100644 --- a/README.md +++ b/README.md @@ -185,10 +185,7 @@ All components are configured via environment variables. Create a `.env` file or | Variable | Default | Description | |----------|---------|-------------| -| `COLLECTOR_LETSMESH_DECODER_ENABLED` | `true` | Enable external LetsMesh packet decoding | -| `COLLECTOR_LETSMESH_DECODER_COMMAND` | `meshcore-decoder` | Decoder CLI command | | `COLLECTOR_LETSMESH_DECODER_KEYS` | *(none)* | Additional decoder channel keys (`label=hex`, `label:hex`, or `hex`) | -| `COLLECTOR_LETSMESH_DECODER_TIMEOUT_SECONDS` | `2.0` | Timeout per decoder invocation | #### LetsMesh Packet Decoding @@ -216,7 +213,7 @@ Normalization behavior: - Packets without decryptable message text are kept as informational `letsmesh_packet` events and are not shown in the messages feed; when decode succeeds the decoded JSON is attached to those packet log events. - When decoder output includes a human sender (`payload.decoded.decrypted.sender`), message text is normalized to `Name: Message` before storage; receiver/observer names are never used as sender fallback. - The collector keeps built-in keys for `Public` and `#test`, and merges any additional keys from `COLLECTOR_LETSMESH_DECODER_KEYS`. -- Docker runtime installs `@michaelhart/meshcore-decoder@0.2.7` and applies `patches/@michaelhart+meshcore-decoder+0.2.7.patch` via `patch-package` for Node compatibility. +- Docker runtime uses the native Python `meshcoredecoder` library (no external Node.js dependency). ### Webhooks diff --git a/docker-compose.yml b/docker-compose.yml index a924b10..951172f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -132,10 +132,7 @@ services: - MQTT_TLS=${MQTT_TLS:-false} - MQTT_TRANSPORT=${MQTT_TRANSPORT:-tcp} - MQTT_WS_PATH=${MQTT_WS_PATH:-/mqtt} - - COLLECTOR_LETSMESH_DECODER_ENABLED=${COLLECTOR_LETSMESH_DECODER_ENABLED:-true} - - COLLECTOR_LETSMESH_DECODER_COMMAND=${COLLECTOR_LETSMESH_DECODER_COMMAND:-meshcore-decoder} - COLLECTOR_LETSMESH_DECODER_KEYS=${COLLECTOR_LETSMESH_DECODER_KEYS:-} - - COLLECTOR_LETSMESH_DECODER_TIMEOUT_SECONDS=${COLLECTOR_LETSMESH_DECODER_TIMEOUT_SECONDS:-2.0} - DATA_HOME=/data - SEED_HOME=/seed # Webhook configuration diff --git a/patches/@michaelhart+meshcore-decoder+0.2.7.patch b/patches/@michaelhart+meshcore-decoder+0.2.7.patch deleted file mode 100644 index 4b0a779..0000000 --- a/patches/@michaelhart+meshcore-decoder+0.2.7.patch +++ /dev/null @@ -1,58 +0,0 @@ -diff --git a/node_modules/@michaelhart/meshcore-decoder/dist/crypto/ed25519-verifier.js b/node_modules/@michaelhart/meshcore-decoder/dist/crypto/ed25519-verifier.js -index d33ffd6..8d040d0 100644 ---- a/node_modules/@michaelhart/meshcore-decoder/dist/crypto/ed25519-verifier.js -+++ b/node_modules/@michaelhart/meshcore-decoder/dist/crypto/ed25519-verifier.js -@@ -36,7 +36,27 @@ var __importStar = (this && this.__importStar) || (function () { - })(); - Object.defineProperty(exports, "__esModule", { value: true }); - exports.Ed25519SignatureVerifier = void 0; --const ed25519 = __importStar(require("@noble/ed25519")); -+let _ed25519 = null; -+async function getEd25519() { -+ if (_ed25519) { -+ return _ed25519; -+ } -+ const mod = await import("@noble/ed25519"); -+ _ed25519 = mod.default ? mod.default : mod; -+ try { -+ _ed25519.etc.sha512Async = sha512Hash; -+ } -+ catch (error) { -+ console.debug("Could not set async SHA-512:", error); -+ } -+ try { -+ _ed25519.etc.sha512Sync = sha512HashSync; -+ } -+ catch (error) { -+ console.debug("Could not set up synchronous SHA-512:", error); -+ } -+ return _ed25519; -+} - const hex_1 = require("../utils/hex"); - const orlp_ed25519_wasm_1 = require("./orlp-ed25519-wasm"); - // Cross-platform SHA-512 implementation -@@ -90,16 +110,6 @@ function sha512HashSync(data) { - throw new Error('No SHA-512 implementation available for synchronous operation'); - } - } --// Set up SHA-512 for @noble/ed25519 --ed25519.etc.sha512Async = sha512Hash; --// Always set up sync version - @noble/ed25519 requires it --// It will throw in browser environments, which @noble/ed25519 can handle --try { -- ed25519.etc.sha512Sync = sha512HashSync; --} --catch (error) { -- console.debug('Could not set up synchronous SHA-512:', error); --} - class Ed25519SignatureVerifier { - /** - * Verify an Ed25519 signature for MeshCore advertisement packets -@@ -116,6 +126,7 @@ class Ed25519SignatureVerifier { - // Construct the signed message according to MeshCore format - const message = this.constructAdvertSignedMessage(publicKeyHex, timestamp, appData); - // Verify the signature using noble-ed25519 -+ const ed25519 = await getEd25519(); - return await ed25519.verify(signature, message, publicKey); - } - catch (error) { diff --git a/pyproject.toml b/pyproject.toml index 3d544fc..e396d50 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -42,6 +42,7 @@ dependencies = [ "python-frontmatter>=1.0.0", "markdown>=3.5.0", "prometheus-client>=0.20.0", + "meshcoredecoder>=0.3.2", ] [project.optional-dependencies] @@ -117,6 +118,7 @@ module = [ "frontmatter.*", "markdown.*", "prometheus_client.*", + "meshcoredecoder.*", ] ignore_missing_imports = true diff --git a/src/meshcore_hub/collector/cli.py b/src/meshcore_hub/collector/cli.py index 0f531b4..03a95e4 100644 --- a/src/meshcore_hub/collector/cli.py +++ b/src/meshcore_hub/collector/cli.py @@ -259,20 +259,9 @@ def _run_collector_service( click.echo(f" Interval: {settings.data_retention_interval_hours} hours") click.echo("") - click.echo("LetsMesh decode configuration:") - if settings.collector_letsmesh_decoder_enabled: - builtin_keys = len(LetsMeshPacketDecoder.BUILTIN_CHANNEL_KEYS) - env_keys = len(settings.collector_letsmesh_decoder_keys_list) - click.echo( - f" Decoder: Enabled ({settings.collector_letsmesh_decoder_command})" - ) - click.echo(f" Built-in keys: {builtin_keys}") - click.echo(f" Additional keys from .env: {env_keys} configured") - click.echo( - f" Timeout: {settings.collector_letsmesh_decoder_timeout_seconds:.2f}s" - ) - else: - click.echo(" Decoder: Disabled") + builtin_keys = len(LetsMeshPacketDecoder.BUILTIN_CHANNEL_KEYS) + env_keys = len(settings.collector_letsmesh_decoder_keys_list) + click.echo(f"Packet decoder: {builtin_keys} built-in keys, {env_keys} from .env") click.echo("") click.echo("Starting MQTT subscriber...") @@ -292,12 +281,7 @@ def _run_collector_service( cleanup_interval_hours=settings.data_retention_interval_hours, node_cleanup_enabled=settings.node_cleanup_enabled, node_cleanup_days=settings.node_cleanup_days, - letsmesh_decoder_enabled=settings.collector_letsmesh_decoder_enabled, - letsmesh_decoder_command=settings.collector_letsmesh_decoder_command, letsmesh_decoder_channel_keys=settings.collector_letsmesh_decoder_keys_list, - letsmesh_decoder_timeout_seconds=( - settings.collector_letsmesh_decoder_timeout_seconds - ), ) diff --git a/src/meshcore_hub/collector/letsmesh_decoder.py b/src/meshcore_hub/collector/letsmesh_decoder.py index b50061d..9d26328 100644 --- a/src/meshcore_hub/collector/letsmesh_decoder.py +++ b/src/meshcore_hub/collector/letsmesh_decoder.py @@ -1,25 +1,26 @@ """LetsMesh packet decoder integration. -Provides an optional bridge to the external `meshcore-decoder` CLI so the -collector can turn LetsMesh upload `raw` packet hex into decoded message data. +Provides native Python packet decoding via the ``meshcoredecoder`` library +so the collector can turn LetsMesh upload ``raw`` packet hex into decoded +message data. """ from __future__ import annotations import hashlib -import json import logging -import shlex -import shutil import string -import subprocess from typing import Any, NamedTuple +from meshcoredecoder import MeshCoreDecoder +from meshcoredecoder.crypto import MeshCoreKeyStore +from meshcoredecoder.types import DecryptionOptions + logger = logging.getLogger(__name__) class LetsMeshPacketDecoder: - """Decode LetsMesh packet payloads with `meshcore-decoder` CLI.""" + """Decode LetsMesh packet payloads with the native Python meshcore-decoder.""" class ChannelKey(NamedTuple): """Channel key metadata for decryption and channel labeling.""" @@ -28,9 +29,6 @@ class LetsMeshPacketDecoder: key_hex: str channel_hash: str - # Built-in keys required by your deployment. - # - Public channel - # - #test channel BUILTIN_CHANNEL_KEYS: tuple[tuple[str, str], ...] = ( ("Public", "8B3387E9C5CDEA6AC9E5EDBAA115CD72"), ("test", "9CD8FCF22A47333B591D96A2B848B73F"), @@ -38,13 +36,8 @@ class LetsMeshPacketDecoder: def __init__( self, - enabled: bool = True, - command: str = "meshcore-decoder", channel_keys: list[str] | None = None, - timeout_seconds: float = 2.0, ) -> None: - self._enabled = enabled - self._command_tokens = shlex.split(command.strip()) if command.strip() else [] self._channel_key_infos = self._normalize_channel_keys(channel_keys or []) self._channel_keys = [info.key_hex for info in self._channel_key_infos] self._channel_names_by_hash = { @@ -54,10 +47,14 @@ class LetsMeshPacketDecoder: } self._decode_cache: dict[str, dict[str, Any] | None] = {} self._decode_cache_maxsize = 2048 - self._timeout_seconds = timeout_seconds - self._checked_command = False - self._command_available = False - self._warned_unavailable = False + self._key_store = self._build_key_store() + + def _build_key_store(self) -> MeshCoreKeyStore: + """Build a MeshCoreKeyStore from configured channel keys.""" + key_store = MeshCoreKeyStore() + if self._channel_keys: + key_store.add_channel_secrets(self._channel_keys) + return key_store @classmethod def _normalize_channel_keys(cls, values: list[str]) -> list[ChannelKey]: @@ -193,83 +190,100 @@ class LetsMeshPacketDecoder: decoded = self._decode_raw(clean_hex) self._decode_cache[clean_hex] = decoded if len(self._decode_cache) > self._decode_cache_maxsize: - # Drop oldest cached payload (insertion-order dict). self._decode_cache.pop(next(iter(self._decode_cache))) return decoded def _decode_raw(self, raw_hex: str) -> dict[str, Any] | None: - """Decode raw packet hex with decoder CLI (cached per packet hex).""" - if not self._enabled: - return None - if not self._is_command_available(): - return None - - command = [*self._command_tokens, "decode", raw_hex, "--json"] - if self._channel_keys: - command.append("--key") - command.extend(self._channel_keys) - + """Decode raw packet hex with native Python decoder (cached per packet hex).""" try: - result = subprocess.run( - command, - check=False, - capture_output=True, - text=True, - timeout=self._timeout_seconds, + options = DecryptionOptions( + key_store=self._key_store, + attempt_decryption=True, ) - except subprocess.TimeoutExpired: - logger.debug( - "LetsMesh decoder timed out after %.2fs", - self._timeout_seconds, - ) - return None - except OSError as exc: - logger.debug("LetsMesh decoder failed to execute: %s", exc) + result = MeshCoreDecoder.decode(raw_hex, options) + except Exception as exc: + logger.debug("LetsMesh decoder failed: %s", exc) return None - if result.returncode != 0: - stderr = result.stderr.strip() if result.stderr else "" - logger.debug( - "LetsMesh decoder exited with code %s%s", - result.returncode, - f": {stderr}" if stderr else "", - ) - return None - - output = result.stdout.strip() - if not output: + if not result.is_valid: + errors = getattr(result, "errors", None) + if errors: + logger.debug("LetsMesh decoder errors: %s", errors) return None + raw_payload_obj = None try: - decoded = json.loads(output) - except json.JSONDecodeError: - logger.debug("LetsMesh decoder returned non-JSON output") - return None + raw_payload_obj = result.payload.get("decoded") + except (AttributeError, TypeError): + pass - return decoded if isinstance(decoded, dict) else None + decoded_dict = result.to_dict() - def _is_command_available(self) -> bool: - """Check decoder command availability once.""" - if self._checked_command: - return self._command_available + if raw_payload_obj is not None: + self._enrich_payload_decoded(decoded_dict, raw_payload_obj) - self._checked_command = True - if not self._command_tokens: - self._command_available = False - else: - command = self._command_tokens[0] - if "/" in command: - self._command_available = shutil.which(command) is not None - else: - self._command_available = shutil.which(command) is not None + self._flatten_control_parsed(decoded_dict) - if not self._command_available and not self._warned_unavailable: - self._warned_unavailable = True - command_text = " ".join(self._command_tokens) or "" - logger.warning( - "LetsMesh decoder command not found (%s). " - "Messages will remain encrypted placeholders until decoder is installed.", - command_text, - ) + return decoded_dict if isinstance(decoded_dict, dict) else None - return self._command_available + _PAYLOAD_ATTR_MAP: tuple[tuple[str, str], ...] = ( + ("channel_hash", "channelHash"), + ("cipher_mac", "cipherMac"), + ("ciphertext", "ciphertext"), + ("ciphertext_length", "ciphertextLength"), + ("decrypted", "decrypted"), + ("destination_hash", "destinationHash"), + ("source_hash", "sourceHash"), + ("sender_public_key", "senderPublicKey"), + ("path_length", "pathLength"), + ("path_hashes", "pathHashes"), + ("extra_type", "extraType"), + ("extra_data", "extraData"), + ("checksum", "checksum"), + ) + + @classmethod + def _enrich_payload_decoded( + cls, + decoded_dict: dict[str, Any], + payload_obj: Any, + ) -> None: + """Enrich payload.decoded dict with fields the library's to_dict() omits. + + Several payload classes (GroupTextPayload, TextMessagePayload, etc.) + inherit BasePayload.to_dict() which only returns type/version/isValid. + This method reads the actual object attributes and merges them in so + the normalizer can find decrypted text, channel hashes, etc. + """ + payload_section = decoded_dict.get("payload") + if not isinstance(payload_section, dict): + return + decoded_section = payload_section.get("decoded") + if not isinstance(decoded_section, dict): + return + for attr_name, dict_key in cls._PAYLOAD_ATTR_MAP: + value = getattr(payload_obj, attr_name, None) + if value is None: + continue + decoded_section.setdefault(dict_key, value) + + @staticmethod + def _flatten_control_parsed(decoded_dict: dict[str, Any]) -> None: + """Flatten Control payload parsed sub-fields into decoded dict. + + The Python library nests sub-type-specific fields (publicKey, subType, + nodeType, etc.) under ``payload.decoded.parsed`` while the TS CLI + returns them flat at ``payload.decoded.*``. The normalizer expects + the flat layout. + """ + payload = decoded_dict.get("payload") + if not isinstance(payload, dict): + return + decoded = payload.get("decoded") + if not isinstance(decoded, dict): + return + parsed = decoded.get("parsed") + if not isinstance(parsed, dict): + return + for key, value in parsed.items(): + decoded.setdefault(key, value) diff --git a/src/meshcore_hub/collector/subscriber.py b/src/meshcore_hub/collector/subscriber.py index a0e7e8a..2424a1a 100644 --- a/src/meshcore_hub/collector/subscriber.py +++ b/src/meshcore_hub/collector/subscriber.py @@ -47,10 +47,7 @@ class Subscriber(LetsMeshNormalizer): cleanup_interval_hours: int = 24, node_cleanup_enabled: bool = False, node_cleanup_days: int = 90, - letsmesh_decoder_enabled: bool = True, - letsmesh_decoder_command: str = "meshcore-decoder", letsmesh_decoder_channel_keys: list[str] | None = None, - letsmesh_decoder_timeout_seconds: float = 2.0, ): """Initialize subscriber. @@ -63,10 +60,7 @@ class Subscriber(LetsMeshNormalizer): cleanup_interval_hours: Hours between cleanup runs node_cleanup_enabled: Enable automatic cleanup of inactive nodes node_cleanup_days: Remove nodes not seen for this many days - letsmesh_decoder_enabled: Enable external LetsMesh packet decoder - letsmesh_decoder_command: Decoder CLI command letsmesh_decoder_channel_keys: Optional channel keys for decrypting group text - letsmesh_decoder_timeout_seconds: Decoder CLI timeout """ self.mqtt = mqtt_client self.db = db_manager @@ -90,10 +84,7 @@ class Subscriber(LetsMeshNormalizer): self._cleanup_thread: Optional[threading.Thread] = None self._last_cleanup: Optional[datetime] = None self._letsmesh_decoder = LetsMeshPacketDecoder( - enabled=letsmesh_decoder_enabled, - command=letsmesh_decoder_command, channel_keys=letsmesh_decoder_channel_keys, - timeout_seconds=letsmesh_decoder_timeout_seconds, ) @property @@ -469,10 +460,7 @@ def create_subscriber( cleanup_interval_hours: int = 24, node_cleanup_enabled: bool = False, node_cleanup_days: int = 90, - letsmesh_decoder_enabled: bool = True, - letsmesh_decoder_command: str = "meshcore-decoder", letsmesh_decoder_channel_keys: list[str] | None = None, - letsmesh_decoder_timeout_seconds: float = 2.0, ) -> Subscriber: """Create a configured subscriber instance. @@ -492,10 +480,7 @@ def create_subscriber( cleanup_interval_hours: Hours between cleanup runs node_cleanup_enabled: Enable automatic cleanup of inactive nodes node_cleanup_days: Remove nodes not seen for this many days - letsmesh_decoder_enabled: Enable external LetsMesh packet decoder - letsmesh_decoder_command: Decoder CLI command letsmesh_decoder_channel_keys: Optional channel keys for decrypting group text - letsmesh_decoder_timeout_seconds: Decoder CLI timeout Returns: Configured Subscriber instance @@ -528,10 +513,7 @@ def create_subscriber( cleanup_interval_hours=cleanup_interval_hours, node_cleanup_enabled=node_cleanup_enabled, node_cleanup_days=node_cleanup_days, - letsmesh_decoder_enabled=letsmesh_decoder_enabled, - letsmesh_decoder_command=letsmesh_decoder_command, letsmesh_decoder_channel_keys=letsmesh_decoder_channel_keys, - letsmesh_decoder_timeout_seconds=letsmesh_decoder_timeout_seconds, ) # Register handlers @@ -558,10 +540,7 @@ def run_collector( cleanup_interval_hours: int = 24, node_cleanup_enabled: bool = False, node_cleanup_days: int = 90, - letsmesh_decoder_enabled: bool = True, - letsmesh_decoder_command: str = "meshcore-decoder", letsmesh_decoder_channel_keys: list[str] | None = None, - letsmesh_decoder_timeout_seconds: float = 2.0, ) -> None: """Run the collector (blocking). @@ -581,10 +560,7 @@ def run_collector( cleanup_interval_hours: Hours between cleanup runs node_cleanup_enabled: Enable automatic cleanup of inactive nodes node_cleanup_days: Remove nodes not seen for this many days - letsmesh_decoder_enabled: Enable external LetsMesh packet decoder - letsmesh_decoder_command: Decoder CLI command letsmesh_decoder_channel_keys: Optional channel keys for decrypting group text - letsmesh_decoder_timeout_seconds: Decoder CLI timeout """ subscriber = create_subscriber( mqtt_host=mqtt_host, @@ -602,10 +578,7 @@ def run_collector( cleanup_interval_hours=cleanup_interval_hours, node_cleanup_enabled=node_cleanup_enabled, node_cleanup_days=node_cleanup_days, - letsmesh_decoder_enabled=letsmesh_decoder_enabled, - letsmesh_decoder_command=letsmesh_decoder_command, letsmesh_decoder_channel_keys=letsmesh_decoder_channel_keys, - letsmesh_decoder_timeout_seconds=letsmesh_decoder_timeout_seconds, ) # Set up signal handlers diff --git a/src/meshcore_hub/common/config.py b/src/meshcore_hub/common/config.py index 9d4b6e3..9e5f464 100644 --- a/src/meshcore_hub/common/config.py +++ b/src/meshcore_hub/common/config.py @@ -138,18 +138,6 @@ class CollectorSettings(CommonSettings): description="Remove nodes not seen for this many days (last_seen)", ge=1, ) - collector_letsmesh_decoder_enabled: bool = Field( - default=True, - description=("Enable external LetsMesh packet decoding via meshcore-decoder."), - ) - collector_letsmesh_decoder_command: str = Field( - default="meshcore-decoder", - description=( - "Command used to run LetsMesh packet decoder CLI " - "(for example: meshcore-decoder, /usr/local/bin/meshcore-decoder, " - "or 'npx meshcore-decoder')." - ), - ) collector_letsmesh_decoder_keys: Optional[str] = Field( default=None, description=( @@ -157,11 +145,6 @@ class CollectorSettings(CommonSettings): "Provide as comma/space separated hex values." ), ) - collector_letsmesh_decoder_timeout_seconds: float = Field( - default=2.0, - description="Timeout in seconds for each decoder invocation.", - ge=0.1, - ) @property def collector_data_dir(self) -> str: diff --git a/src/meshcore_hub/web/app.py b/src/meshcore_hub/web/app.py index 5149a0b..af8a7ac 100644 --- a/src/meshcore_hub/web/app.py +++ b/src/meshcore_hub/web/app.py @@ -43,7 +43,6 @@ def _build_channel_labels() -> dict[str, str]: """Build UI channel labels from built-in + configured decoder keys.""" raw_keys = os.getenv("COLLECTOR_LETSMESH_DECODER_KEYS") decoder = LetsMeshPacketDecoder( - enabled=False, channel_keys=_parse_decoder_key_entries(raw_keys), ) labels = decoder.channel_labels_by_index() diff --git a/tests/test_collector/test_letsmesh_decoder.py b/tests/test_collector/test_letsmesh_decoder.py index 4498fbe..ea86f6b 100644 --- a/tests/test_collector/test_letsmesh_decoder.py +++ b/tests/test_collector/test_letsmesh_decoder.py @@ -1,100 +1,89 @@ """Tests for LetsMesh packet decoder integration.""" -import subprocess -from unittest.mock import patch +from unittest.mock import MagicMock, patch from meshcore_hub.collector.letsmesh_decoder import LetsMeshPacketDecoder def test_decode_payload_returns_none_without_raw() -> None: """Decoder returns None when packet has no raw hex.""" - decoder = LetsMeshPacketDecoder(enabled=True) + decoder = LetsMeshPacketDecoder() assert decoder.decode_payload({"packet_type": 5}) is None def test_decode_payload_rejects_non_hex_raw_without_invoking_decoder() -> None: - """Decoder returns None and does not execute subprocess for invalid raw hex.""" - decoder = LetsMeshPacketDecoder(enabled=True, command="meshcore-decoder") + """Decoder returns None and does not call MeshCoreDecoder for invalid hex.""" + decoder = LetsMeshPacketDecoder() - with ( - patch("meshcore_hub.collector.letsmesh_decoder.shutil.which", return_value="1"), - patch("meshcore_hub.collector.letsmesh_decoder.subprocess.run") as mock_run, - ): + with patch("meshcore_hub.collector.letsmesh_decoder.MeshCoreDecoder") as mock_lib: assert decoder.decode_payload({"raw": "ZZ-not-hex"}) is None - mock_run.assert_not_called() + mock_lib.decode.assert_not_called() -def test_decode_payload_invokes_decoder_with_keys() -> None: - """Decoder command includes channel keys and returns parsed JSON.""" +def test_decode_payload_calls_native_decoder() -> None: + """Decoder calls MeshCoreDecoder.decode and returns dict output.""" decoder = LetsMeshPacketDecoder( - enabled=True, - command="meshcore-decoder", channel_keys=["0xABCDEF", "name=012345", "abcDEF"], - timeout_seconds=1.5, - ) - completed = subprocess.CompletedProcess( - args=["meshcore-decoder"], - returncode=0, - stdout='{"payload":{"decoded":{"decrypted":{"message":"hello"}}}}', - stderr="", ) - with ( - patch("meshcore_hub.collector.letsmesh_decoder.shutil.which", return_value="1"), - patch( - "meshcore_hub.collector.letsmesh_decoder.subprocess.run", - return_value=completed, - ) as mock_run, - ): + mock_result = MagicMock() + mock_result.is_valid = True + mock_result.to_dict.return_value = { + "payloadType": 5, + "pathLength": 2, + "payload": { + "decoded": { + "channelHash": "11", + "decrypted": {"message": "hello"}, + } + }, + } + + with patch( + "meshcore_hub.collector.letsmesh_decoder.MeshCoreDecoder.decode", + return_value=mock_result, + ) as mock_decode: decoded = decoder.decode_payload({"raw": "A1B2C3"}) assert isinstance(decoded, dict) - payload = decoded.get("payload") - assert isinstance(payload, dict) - decoded_payload = payload.get("decoded") - assert isinstance(decoded_payload, dict) - decrypted = decoded_payload.get("decrypted") - assert isinstance(decrypted, dict) - assert decrypted.get("message") == "hello" - command = mock_run.call_args.args[0] - assert command == [ - "meshcore-decoder", - "decode", - "A1B2C3", - "--json", - "--key", - "8B3387E9C5CDEA6AC9E5EDBAA115CD72", - "9CD8FCF22A47333B591D96A2B848B73F", - "ABCDEF", - "012345", - ] - assert mock_run.call_args.kwargs["timeout"] == 1.5 + assert decoded["payload"]["decoded"]["decrypted"]["message"] == "hello" + mock_decode.assert_called_once() + + args, kwargs = mock_decode.call_args + assert args[0] == "A1B2C3" + assert kwargs.get("options") is not None or len(args) > 1 def test_decode_payload_returns_none_for_decoder_error() -> None: - """Decoder returns None when decoder exits with failure.""" - decoder = LetsMeshPacketDecoder(enabled=True, command="meshcore-decoder") - completed = subprocess.CompletedProcess( - args=["meshcore-decoder"], - returncode=1, - stdout="", - stderr="decode error", - ) + """Decoder returns None when native decoder raises an exception.""" + decoder = LetsMeshPacketDecoder() - with ( - patch("meshcore_hub.collector.letsmesh_decoder.shutil.which", return_value="1"), - patch( - "meshcore_hub.collector.letsmesh_decoder.subprocess.run", - return_value=completed, - ), + with patch( + "meshcore_hub.collector.letsmesh_decoder.MeshCoreDecoder.decode", + side_effect=ValueError("bad packet"), + ): + assert decoder.decode_payload({"raw": "A1B2C3"}) is None + + +def test_decode_payload_returns_none_for_invalid_result() -> None: + """Decoder returns None when native decoder returns invalid result.""" + decoder = LetsMeshPacketDecoder() + + mock_result = MagicMock() + mock_result.is_valid = False + mock_result.errors = ["unsupported packet"] + + with patch( + "meshcore_hub.collector.letsmesh_decoder.MeshCoreDecoder.decode", + return_value=mock_result, ): assert decoder.decode_payload({"raw": "A1B2C3"}) is None def test_builtin_channel_keys_present_by_default() -> None: """Public and #test keys are always present even without .env keys.""" - decoder = LetsMeshPacketDecoder(enabled=True, command="meshcore-decoder") + decoder = LetsMeshPacketDecoder() assert decoder._channel_keys == [ "8B3387E9C5CDEA6AC9E5EDBAA115CD72", "9CD8FCF22A47333B591D96A2B848B73F", @@ -104,10 +93,7 @@ def test_builtin_channel_keys_present_by_default() -> None: def test_channel_name_lookup_from_decoded_hash() -> None: """Decoder resolves channel names from configured label=key entries.""" key_hex = "EB50A1BCB3E4E5D7BF69A57C9DADA211" - decoder = LetsMeshPacketDecoder( - enabled=False, - channel_keys=[f"#bot={key_hex}"], - ) + decoder = LetsMeshPacketDecoder(channel_keys=[f"#bot={key_hex}"]) channel_hash = decoder._compute_channel_hash(key_hex) decoded_packet = { "payload": { @@ -123,7 +109,6 @@ def test_channel_name_lookup_from_decoded_hash() -> None: def test_channel_labels_by_index_includes_labeled_entries() -> None: """Channel labels map includes built-ins and label=key env entries.""" decoder = LetsMeshPacketDecoder( - enabled=False, channel_keys=[ "bot=EB50A1BCB3E4E5D7BF69A57C9DADA211", "chat=D0BDD6D71538138ED979EEC00D98AD97", @@ -136,3 +121,77 @@ def test_channel_labels_by_index_includes_labeled_entries() -> None: assert labels[217] == "#test" assert labels[202] == "#bot" assert labels[184] == "#chat" + + +def test_decode_payload_caches_results() -> None: + """Repeated decode calls for same hex use cached result.""" + decoder = LetsMeshPacketDecoder() + + mock_result = MagicMock() + mock_result.is_valid = True + mock_result.to_dict.return_value = {"payloadType": 4} + + with patch( + "meshcore_hub.collector.letsmesh_decoder.MeshCoreDecoder.decode", + return_value=mock_result, + ) as mock_decode: + first = decoder.decode_payload({"raw": "AABBCC"}) + second = decoder.decode_payload({"raw": "AABBCC"}) + + assert first is second + mock_decode.assert_called_once() + + +def test_flatten_control_parsed_merges_parsed_into_decoded() -> None: + """Control payload parsed fields are flattened into decoded dict.""" + decoded_dict = { + "payload": { + "decoded": { + "type": 11, + "subType": 144, + "flags": 1, + "parsed": { + "publicKey": "AA" * 32, + "nodeType": 2, + }, + } + } + } + LetsMeshPacketDecoder._flatten_control_parsed(decoded_dict) + decoded = decoded_dict["payload"]["decoded"] + assert decoded["publicKey"] == "AA" * 32 + assert decoded["nodeType"] == 2 + assert "parsed" in decoded + + +def test_flatten_control_parsed_does_not_overwrite_existing() -> None: + """Flattening uses setdefault so existing keys are preserved.""" + decoded_dict = { + "payload": { + "decoded": { + "subType": 144, + "parsed": { + "subType": 999, + "extra": "value", + }, + } + } + } + LetsMeshPacketDecoder._flatten_control_parsed(decoded_dict) + decoded = decoded_dict["payload"]["decoded"] + assert decoded["subType"] == 144 + assert decoded["extra"] == "value" + + +def test_decode_payload_returns_none_for_empty_raw() -> None: + """Decoder returns None for empty raw string.""" + decoder = LetsMeshPacketDecoder() + assert decoder.decode_payload({"raw": ""}) is None + assert decoder.decode_payload({"raw": " "}) is None + + +def test_decode_payload_returns_none_for_non_string_raw() -> None: + """Decoder returns None when raw is not a string.""" + decoder = LetsMeshPacketDecoder() + assert decoder.decode_payload({"raw": 12345}) is None + assert decoder.decode_payload({"raw": None}) is None diff --git a/tests/test_collector/test_letsmesh_normalizer_integration.py b/tests/test_collector/test_letsmesh_normalizer_integration.py new file mode 100644 index 0000000..20f2e7d --- /dev/null +++ b/tests/test_collector/test_letsmesh_normalizer_integration.py @@ -0,0 +1,624 @@ +"""Integration tests for LetsMesh normalizer packet type routing. + +Tests the full _normalize_letsmesh_event flow with a mock decoder +returning realistic to_dict() structures for each packet type. +""" + +from __future__ import annotations + +from typing import Any +from unittest.mock import MagicMock + +from meshcore_hub.collector.letsmesh_decoder import LetsMeshPacketDecoder +from meshcore_hub.collector.letsmesh_normalizer import LetsMeshNormalizer + + +class _TestNormalizer(LetsMeshNormalizer): + """Minimal normalizer subclass wired with a mock decoder for testing.""" + + def __init__(self, decoder: LetsMeshPacketDecoder) -> None: + self._letsmesh_decoder = decoder + self.mqtt = MagicMock() + self.mqtt.topic_builder.parse_letsmesh_upload_topic = MagicMock( + return_value=None, + ) + + +def _make_decoder( + decoded_packets: dict[str, dict[str, Any] | None], +) -> LetsMeshPacketDecoder: + """Create a decoder that returns pre-configured results by raw hex.""" + decoder = LetsMeshPacketDecoder() + + original = decoder.decode_payload + + def stub(payload: dict[str, Any]) -> dict[str, Any] | None: + raw = payload.get("raw", "") + if raw in decoded_packets: + return decoded_packets[raw] + return original(payload) + + decoder.decode_payload = stub # type: ignore[method-assign] + return decoder + + +PUB_KEY = "AA" * 32 +OBSERVER_KEY = "BB" * 32 + + +class TestStatusFeed: + def test_status_feed_passes_through(self) -> None: + decoder = _make_decoder({}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "status", + ) + payload = {"uptime": 3600, "nodes_seen": 5} + result = norm._normalize_letsmesh_event("meshcore/BB.../status", payload) + assert result is not None + pk, event_type, pl = result + assert pk == OBSERVER_KEY + assert event_type == "letsmesh_status" + assert pl == payload + + +class TestInternalFeed: + def test_internal_feed_passes_through(self) -> None: + decoder = _make_decoder({}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "internal", + ) + payload = {"info": "restart"} + result = norm._normalize_letsmesh_event("meshcore/BB.../internal", payload) + assert result is not None + pk, event_type, pl = result + assert pk == OBSERVER_KEY + assert event_type == "letsmesh_internal" + assert pl == payload + + +class TestAdvertPacket: + def test_advert_with_name_and_location(self) -> None: + raw = "advert1" + decoded = { + "payloadType": 4, + "pathLength": 0, + "payload": { + "decoded": { + "type": 4, + "publicKey": PUB_KEY, + "timestamp": 1700000000, + "signature": "CC" * 64, + "appData": { + "flags": 0x90, + "deviceRole": 2, + "hasLocation": True, + "hasName": True, + "location": {"latitude": 47.5, "longitude": -122.1}, + "name": "TestRepeater", + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + pk, event_type, pl = result + assert pk == OBSERVER_KEY + assert event_type == "advertisement" + assert pl["public_key"] == PUB_KEY + assert pl["name"] == "TestRepeater" + assert pl["adv_type"] == "repeater" + assert pl["lat"] == 47.5 + assert pl["lon"] == -122.1 + + def test_advert_without_app_data_falls_back(self) -> None: + raw = "advert2" + decoded = { + "payloadType": 4, + "pathLength": 0, + "payload": { + "decoded": { + "type": 4, + "publicKey": PUB_KEY, + "timestamp": 1700000000, + "signature": "CC" * 64, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event( + "t", {"raw": raw, "origin": "FallbackName"} + ) + assert result is not None + _, event_type, pl = result + assert event_type == "advertisement" + assert pl["public_key"] == PUB_KEY + assert pl["name"] == "FallbackName" + + +class TestGroupTextPacket: + def test_grouptext_decrypted_routes_to_channel_msg(self) -> None: + raw = "gt1" + decoded = { + "payloadType": 5, + "pathLength": 1, + "payload": { + "decoded": { + "type": 5, + "channelHash": "11", + "sourceHash": "AB", + "decrypted": { + "message": "Hello world", + "sender": "Alice", + "timestamp": 1700000000, + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw, "hash": "abc123"}) + assert result is not None + _, event_type, pl = result + assert event_type == "channel_msg_recv" + assert "Alice" in pl["text"] + assert "Hello world" in pl["text"] + assert pl.get("channel_name") == "Public" + + def test_grouptext_not_decrypted_falls_through(self) -> None: + raw = "gt2" + decoded = { + "payloadType": 5, + "pathLength": 0, + "payload": { + "decoded": { + "type": 5, + "channelHash": "FF", + "ciphertext": "AABBCC", + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "letsmesh_packet" + assert pl.get("decoded_payload_type") == 5 + + +class TestTextMessagePacket: + def test_textmessage_decrypted_routes_to_contact_msg(self) -> None: + raw = "tm1" + decoded = { + "payloadType": 2, + "pathLength": 2, + "payload": { + "decoded": { + "type": 2, + "destinationHash": "CD", + "sourceHash": "AB", + "decrypted": { + "message": "Direct hello", + "sender": None, + "timestamp": 1700000000, + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "contact_msg_recv" + assert "Direct hello" in pl["text"] + + +class TestAnonRequestPacket: + def test_anonrequest_decrypted_routes_to_contact_msg(self) -> None: + raw = "ar1" + decoded = { + "payloadType": 7, + "pathLength": 1, + "payload": { + "decoded": { + "type": 7, + "destinationHash": "CD", + "senderPublicKey": "DD" * 32, + "decrypted": { + "message": "anon login message", + "timestamp": 1700000000, + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "contact_msg_recv" + assert "anon login message" in pl["text"] + + +class TestTracePacket: + def test_trace_with_snr_routes_to_trace_data(self) -> None: + raw = "tr1" + decoded = { + "payloadType": 9, + "pathLength": 3, + "payload": { + "decoded": { + "type": 9, + "traceTag": "A1B2C3D4", + "authCode": 42, + "flags": 0, + "pathHashes": ["AA", "BB", "CC"], + "snrValues": [5.25, -3.5, 12.0], + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "trace_data" + assert pl["initiator_tag"] == 0xA1B2C3D4 + assert pl["path_hashes"] == ["AA", "BB", "CC"] + assert pl["snr_values"] == [5.25, -3.5, 12.0] + assert pl["hop_count"] == 3 + + def test_trace_without_decoded_payload_falls_through(self) -> None: + raw = "tr2" + decoded = { + "payloadType": 9, + "pathLength": 0, + "payload": { + "decoded": None, + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "letsmesh_packet" + + +class TestControlPacket: + def test_control_discover_response_routes_to_contact(self) -> None: + raw = "ctrl1" + decoded = { + "payloadType": 11, + "pathLength": 0, + "payload": { + "decoded": { + "type": 11, + "flags": 1, + "subType": 144, + "dataHex": "", + "publicKey": PUB_KEY, + "nodeType": 2, + "parsed": { + "publicKey": PUB_KEY, + "nodeType": 2, + "snr": 10, + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "contact" + assert pl["public_key"] == PUB_KEY + assert pl["type"] == 2 + + def test_control_non_discover_routes_to_status_response(self) -> None: + raw = "ctrl2" + decoded = { + "payloadType": 11, + "pathLength": 0, + "payload": { + "decoded": { + "type": 11, + "flags": 0, + "subType": 128, + "dataHex": "", + "parsed": {}, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "status_response" + assert pl["control_subtype"] == 128 + + +class TestPathPacket: + def test_path_update_routes_correctly(self) -> None: + raw = "path1" + decoded = { + "payloadType": 8, + "pathLength": 0, + "payload": { + "decoded": { + "type": 8, + "pathLength": 4, + "pathHashes": ["AA", "BB", "CC", "DD"], + "extraType": 1, + "extraData": PUB_KEY, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "path_updated" + assert pl["hop_count"] == 4 + assert pl["path_hashes"] == ["AA", "BB", "CC", "DD"] + assert pl["node_public_key"] == PUB_KEY + + +class TestResponsePacket: + def test_response_with_battery_routes_to_battery(self) -> None: + raw = "resp1" + decoded = { + "payloadType": 1, + "pathLength": 1, + "payload": { + "decoded": { + "type": 1, + "destinationHash": "CD", + "sourceHash": "AB", + "cipherMac": "0000", + "ciphertext": "", + "tag": 0, + "decrypted": { + "content": { + "node_public_key": PUB_KEY, + "battery_voltage": 3.7, + "battery_percentage": 85, + } + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "battery" + assert pl["battery_voltage"] == 3.7 + assert pl["battery_percentage"] == 85 + + def test_response_with_telemetry_routes_to_telemetry(self) -> None: + raw = "resp2" + decoded = { + "payloadType": 1, + "pathLength": 1, + "payload": { + "decoded": { + "type": 1, + "destinationHash": "CD", + "sourceHash": "AB", + "cipherMac": "0000", + "ciphertext": "", + "tag": 0, + "decrypted": { + "content": { + "node_public_key": PUB_KEY, + "parsed_data": {"temperature": 22.5, "humidity": 60}, + } + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "telemetry_response" + assert pl["node_public_key"] == PUB_KEY + assert pl["parsed_data"]["temperature"] == 22.5 + + def test_response_with_status_routes_to_status_response(self) -> None: + raw = "resp3" + decoded = { + "payloadType": 1, + "pathLength": 0, + "payload": { + "decoded": { + "type": 1, + "destinationHash": "CD", + "sourceHash": "AB", + "cipherMac": "0000", + "ciphertext": "", + "tag": 0, + "decrypted": { + "content": { + "status": "running", + "uptime": 86400, + "message_count": 42, + } + }, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "status_response" + assert pl["status"] == "running" + assert pl["uptime"] == 86400 + assert pl["message_count"] == 42 + + +class TestAckPacket: + def test_ack_falls_through_to_letsmesh_packet(self) -> None: + raw = "ack1" + decoded = { + "payloadType": 3, + "pathLength": 0, + "payload": { + "decoded": { + "type": 3, + "checksum": "AABBCCDD", + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "letsmesh_packet" + assert pl["decoded_payload_type"] == 3 + assert pl["decoded_packet"]["payloadType"] == 3 + + +class TestRequestPacket: + def test_request_without_decrypted_content_falls_through(self) -> None: + raw = "req1" + decoded = { + "payloadType": 0, + "pathLength": 1, + "payload": { + "decoded": { + "type": 0, + "destinationHash": "CD", + "sourceHash": "AB", + "cipherMac": "0000", + "ciphertext": "AABB", + "timestamp": 1700000000, + "requestType": 1, + } + }, + } + decoder = _make_decoder({raw: decoded}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "packets", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": raw}) + assert result is not None + _, event_type, pl = result + assert event_type == "letsmesh_packet" + assert pl["decoded_payload_type"] == 0 + + +class TestUnrecognizedFeed: + def test_unknown_feed_type_returns_none(self) -> None: + decoder = _make_decoder({}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = ( + OBSERVER_KEY, + "unknown", + ) + + result = norm._normalize_letsmesh_event("t", {"raw": "FF"}) + assert result is None + + def test_unparseable_topic_returns_none(self) -> None: + decoder = _make_decoder({}) + norm = _TestNormalizer(decoder) + norm.mqtt.topic_builder.parse_letsmesh_upload_topic.return_value = None + + result = norm._normalize_letsmesh_event("garbage", {}) + assert result is None