mirror of
https://github.com/ipnet-mesh/meshcore-hub.git
synced 2026-08-07 09:23:07 +02:00
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.
This commit is contained in:
+1
-10
@@ -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
|
||||
# -------------------
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
-17
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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 "<empty>"
|
||||
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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user