diff --git a/data/mesh_ingestor/__init__.py b/data/mesh_ingestor/__init__.py index 784faca..85b6d49 100644 --- a/data/mesh_ingestor/__init__.py +++ b/data/mesh_ingestor/__init__.py @@ -21,7 +21,7 @@ import threading as threading # re-exported for compatibility import sys import types -from . import config, daemon, handlers, interfaces, queue, serialization +from . import channels, config, daemon, handlers, interfaces, queue, serialization __all__: list[str] = [] @@ -40,7 +40,7 @@ def _export_constants() -> None: __all__.extend(["json", "urllib", "glob", "threading", "signal"]) -for _module in (daemon, handlers, interfaces, queue, serialization): +for _module in (channels, daemon, handlers, interfaces, queue, serialization): _reexport(_module) _export_constants() diff --git a/data/mesh_ingestor/channels.py b/data/mesh_ingestor/channels.py new file mode 100644 index 0000000..a2546a5 --- /dev/null +++ b/data/mesh_ingestor/channels.py @@ -0,0 +1,223 @@ +# Copyright (C) 2025 l5yth +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Helpers for capturing and exposing mesh channel metadata.""" + +from __future__ import annotations + +import os +from typing import Any, Iterable, Iterator + +from . import config + +try: # pragma: no cover - optional dependency for enum introspection + from meshtastic.protobuf import channel_pb2 +except Exception: # pragma: no cover - exercised in environments without protobufs + channel_pb2 = None # type: ignore[assignment] + +_ROLE_PRIMARY = 1 +_ROLE_SECONDARY = 2 + +if channel_pb2 is not None: # pragma: no branch - evaluated once at import time + try: + _ROLE_PRIMARY = int(channel_pb2.Channel.Role.PRIMARY) + _ROLE_SECONDARY = int(channel_pb2.Channel.Role.SECONDARY) + except Exception: # pragma: no cover - defensive, version specific + _ROLE_PRIMARY = 1 + _ROLE_SECONDARY = 2 + +_CHANNEL_MAPPINGS: tuple[tuple[int, str], ...] = () +_CHANNEL_LOOKUP: dict[int, str] = {} + + +def _iter_channel_objects(channels_obj: Any) -> Iterator[Any]: + """Yield channel descriptors from ``channels_obj``. + + The real Meshtastic API exposes channels via protobuf containers that are + list-like. This helper converts the container into a deterministic iterator + while avoiding runtime errors if an unexpected type is supplied. + """ + + if channels_obj is None: + return iter(()) + + if isinstance(channels_obj, dict): + return iter(channels_obj.values()) + + if isinstance(channels_obj, Iterable): + return iter(list(channels_obj)) + + length_fn = getattr(channels_obj, "__len__", None) + getitem = getattr(channels_obj, "__getitem__", None) + if callable(length_fn) and callable(getitem): + try: + length = int(length_fn()) + except Exception: # pragma: no cover - defensive only + length = None + if length is not None and length >= 0: + snapshot = [] + for index in range(length): + try: + snapshot.append(getitem(index)) + except Exception: # pragma: no cover - best effort copy + break + return iter(snapshot) + + return iter(()) + + +def _primary_channel_name() -> str | None: + """Return the name to use for the primary channel when available.""" + + preset = getattr(config, "MODEM_PRESET", None) + if isinstance(preset, str) and preset.strip(): + return preset + env_name = os.environ.get("CHANNEL", "").strip() + if env_name: + return env_name + return None + + +def _normalize_role(role: Any) -> int | None: + """Convert a channel role descriptor into an integer value.""" + + if isinstance(role, int): + return role + if isinstance(role, str): + value = role.strip().upper() + if value == "PRIMARY": + return _ROLE_PRIMARY + if value == "SECONDARY": + return _ROLE_SECONDARY + try: + return int(value) + except ValueError: + return None + name_attr = getattr(role, "name", None) + if isinstance(name_attr, str): + return _normalize_role(name_attr) + value_attr = getattr(role, "value", None) + if isinstance(value_attr, int): + return value_attr + try: + return int(role) # type: ignore[arg-type] + except Exception: + return None + + +def _channel_tuple(channel_obj: Any) -> tuple[int, str] | None: + """Return ``(index, name)`` for ``channel_obj`` when resolvable.""" + + role_value = _normalize_role(getattr(channel_obj, "role", None)) + if role_value == _ROLE_PRIMARY: + channel_index = 0 + channel_name = _primary_channel_name() + elif role_value == _ROLE_SECONDARY: + raw_index = getattr(channel_obj, "index", None) + try: + channel_index = int(raw_index) + except Exception: + channel_index = None + settings = getattr(channel_obj, "settings", None) + channel_name = getattr(settings, "name", None) if settings else None + else: + return None + + if not isinstance(channel_index, int): + return None + + if isinstance(channel_name, str): + channel_name = channel_name.strip() + else: + channel_name = None + + if not channel_name: + return None + + return channel_index, channel_name + + +def capture_from_interface(iface: Any) -> None: + """Populate the channel cache by inspecting ``iface`` when possible.""" + + global _CHANNEL_MAPPINGS, _CHANNEL_LOOKUP + + if iface is None or _CHANNEL_MAPPINGS: + return + + try: + wait_for_config = getattr(iface, "waitForConfig", None) + if callable(wait_for_config): + wait_for_config() + except Exception: # pragma: no cover - hardware dependent safeguard + pass + + local_node = getattr(iface, "localNode", None) + channels_obj = getattr(local_node, "channels", None) if local_node else None + + channel_entries: list[tuple[int, str]] = [] + seen_indices: set[int] = set() + for candidate in _iter_channel_objects(channels_obj): + result = _channel_tuple(candidate) + if result is None: + continue + index, name = result + if index in seen_indices: + continue + channel_entries.append((index, name)) + seen_indices.add(index) + + if not channel_entries: + return + + _CHANNEL_MAPPINGS = tuple(channel_entries) + _CHANNEL_LOOKUP = {index: name for index, name in _CHANNEL_MAPPINGS} + + config._debug_log( + "Captured channel metadata", + context="channels.capture", + severity="info", + always=True, + channels=_CHANNEL_MAPPINGS, + ) + + +def channel_mappings() -> tuple[tuple[int, str], ...]: + """Return the cached ``(index, name)`` channel tuples.""" + + return _CHANNEL_MAPPINGS + + +def channel_name(channel_index: int | None) -> str | None: + """Return the channel name for ``channel_index`` when known.""" + + if channel_index is None: + return None + return _CHANNEL_LOOKUP.get(int(channel_index)) + + +def _reset_channel_cache() -> None: + """Clear cached channel data. Intended for use in tests only.""" + + global _CHANNEL_MAPPINGS, _CHANNEL_LOOKUP + _CHANNEL_MAPPINGS = () + _CHANNEL_LOOKUP = {} + + +__all__ = [ + "capture_from_interface", + "channel_mappings", + "channel_name", + "_reset_channel_cache", +] diff --git a/data/mesh_ingestor/daemon.py b/data/mesh_ingestor/daemon.py index 06f782f..3fdc97a 100644 --- a/data/mesh_ingestor/daemon.py +++ b/data/mesh_ingestor/daemon.py @@ -242,6 +242,7 @@ def main() -> None: iface, resolved_target = interfaces._create_default_interface() active_candidate = resolved_target interfaces._ensure_radio_metadata(iface) + interfaces._ensure_channel_metadata(iface) retry_delay = max(0.0, config._RECONNECT_INITIAL_DELAY_SECS) initial_snapshot_sent = False if not announced_target and resolved_target: diff --git a/data/mesh_ingestor/handlers.py b/data/mesh_ingestor/handlers.py index 269de38..efba9f0 100644 --- a/data/mesh_ingestor/handlers.py +++ b/data/mesh_ingestor/handlers.py @@ -21,7 +21,7 @@ import json import time from collections.abc import Mapping -from . import config, queue +from . import channels, config, queue from .serialization import ( _canonical_node_id, _coerce_float, @@ -64,6 +64,21 @@ def _apply_radio_metadata(payload: dict) -> dict: return payload +def _is_encrypted_flag(value) -> bool: + """Return ``True`` when ``value`` represents an encrypted payload.""" + + if isinstance(value, bool): + return value + if isinstance(value, (int, float)): + return value != 0 + if isinstance(value, str): + normalized = value.strip().lower() + if normalized in {"", "0", "false", "no"}: + return False + return True + return bool(value) + + def _apply_radio_metadata_to_nodes(payload: dict) -> dict: """Attach radio metadata to each node entry stored in ``payload``.""" @@ -857,6 +872,8 @@ def store_packet_dict(packet: Mapping) -> None: rssi = _first(packet, "rssi", "rx_rssi", "rxRssi", default=None) hop = _first(packet, "hopLimit", "hop_limit", default=None) + encrypted_flag = _is_encrypted_flag(encrypted) + message_payload = { "id": int(pkt_id), "rx_time": rx_time, @@ -871,6 +888,12 @@ def store_packet_dict(packet: Mapping) -> None: "rssi": int(rssi) if rssi is not None else None, "hop_limit": int(hop) if hop is not None else None, } + + channel_name_value = None + if not encrypted_flag: + channel_name_value = channels.channel_name(channel) + if channel_name_value: + message_payload["channel_name"] = channel_name_value _queue_post_json( "/api/messages", _apply_radio_metadata(message_payload), @@ -881,14 +904,17 @@ def store_packet_dict(packet: Mapping) -> None: from_label = _canonical_node_id(from_id) or from_id to_label = _canonical_node_id(to_id) or to_id payload_desc = "Encrypted" if text is None and encrypted else text - config._debug_log( - "Queued message payload", - context="handlers.store_packet_dict", - from_id=from_label, - to_id=to_label, - channel=channel, - payload=payload_desc, - ) + log_kwargs = { + "context": "handlers.store_packet_dict", + "from_id": from_label, + "to_id": to_label, + "channel": channel, + "channel_display": channel_name_value or channel, + "payload": payload_desc, + } + if channel_name_value: + log_kwargs["channel_name"] = channel_name_value + config._debug_log("Queued message payload", **log_kwargs) _last_packet_monotonic: float | None = None diff --git a/data/mesh_ingestor/interfaces.py b/data/mesh_ingestor/interfaces.py index 00f0452..4ee08da 100644 --- a/data/mesh_ingestor/interfaces.py +++ b/data/mesh_ingestor/interfaces.py @@ -26,7 +26,7 @@ from typing import TYPE_CHECKING, Any from meshtastic.serial_interface import SerialInterface from meshtastic.tcp_interface import TCPInterface -from . import config, serialization +from . import channels, config, serialization if TYPE_CHECKING: # pragma: no cover - import only used for type checking from meshtastic.ble_interface import BLEInterface as _BLEInterface @@ -334,6 +334,24 @@ def _ensure_radio_metadata(iface: Any) -> None: ) +def _ensure_channel_metadata(iface: Any) -> None: + """Capture channel metadata by inspecting ``iface`` once per runtime.""" + + if iface is None: + return + + try: + channels.capture_from_interface(iface) + except Exception as exc: # pragma: no cover - defensive instrumentation + config._debug_log( + "Failed to capture channel metadata", + context="interfaces.ensure_channel_metadata", + severity="warn", + error_class=exc.__class__.__name__, + error_message=str(exc), + ) + + _DEFAULT_TCP_PORT = 4403 _DEFAULT_TCP_TARGET = "http://127.0.0.1" @@ -581,6 +599,7 @@ def _create_default_interface() -> tuple[object, str]: __all__ = [ "BLEInterface", "NoAvailableMeshInterface", + "_ensure_channel_metadata", "_ensure_radio_metadata", "_DummySerialInterface", "_DEFAULT_TCP_PORT", diff --git a/tests/test_mesh.py b/tests/test_mesh.py index 731f1f0..14d6f00 100644 --- a/tests/test_mesh.py +++ b/tests/test_mesh.py @@ -185,6 +185,7 @@ def mesh_module(monkeypatch): for attr in ("LORA_FREQ", "MODEM_PRESET"): if attr in module.__dict__: delattr(module, attr) + module.channels._reset_channel_cache() yield module @@ -402,6 +403,64 @@ def test_ensure_radio_metadata_extracts_config(mesh_module, capsys): assert second_log == "" +def test_capture_channels_from_interface_records_metadata(mesh_module, capsys): + mesh = mesh_module + + mesh.config.MODEM_PRESET = "MediumFast" + + class DummyInterface: + def __init__(self) -> None: + self.wait_calls = 0 + primary = SimpleNamespace(role=1, settings=SimpleNamespace()) + secondary = SimpleNamespace( + role="SECONDARY", + index="7", + settings=SimpleNamespace(name="TestChannel"), + ) + self.localNode = SimpleNamespace(channels=[primary, secondary]) + + def waitForConfig(self) -> None: # noqa: D401 - matches interface contract + self.wait_calls += 1 + + iface = DummyInterface() + + mesh.channels.capture_from_interface(iface) + log_output = capsys.readouterr().out + + assert iface.wait_calls == 1 + assert mesh.channels.channel_mappings() == ((0, "MediumFast"), (7, "TestChannel")) + assert mesh.channels.channel_name(7) == "TestChannel" + assert "Captured channel metadata" in log_output + assert "channels=((0, 'MediumFast'), (7, 'TestChannel'))" in log_output + + mesh.channels.capture_from_interface(SimpleNamespace(localNode=None)) + assert mesh.channels.channel_mappings() == ((0, "MediumFast"), (7, "TestChannel")) + + +def test_capture_channels_primary_falls_back_to_env(mesh_module, monkeypatch, capsys): + mesh = mesh_module + + mesh.config.MODEM_PRESET = None + monkeypatch.setenv("CHANNEL", "FallbackName") + + class DummyInterface: + def __init__(self) -> None: + self.localNode = SimpleNamespace( + channels={"primary": SimpleNamespace(role="PRIMARY")} + ) + + def waitForConfig(self) -> None: # noqa: D401 - placeholder + return None + + mesh.channels._reset_channel_cache() + mesh.channels.capture_from_interface(DummyInterface()) + log_output = capsys.readouterr().out + + assert mesh.channels.channel_mappings() == ((0, "FallbackName"),) + assert mesh.channels.channel_name(0) == "FallbackName" + assert "FallbackName" in log_output + + def test_create_default_interface_falls_back_to_tcp(mesh_module, monkeypatch): mesh = mesh_module attempts = [] @@ -1292,6 +1351,64 @@ def test_store_packet_dict_handles_invalid_channel(mesh_module, monkeypatch): assert priority == mesh._MESSAGE_POST_PRIORITY +def test_store_packet_dict_appends_channel_name(mesh_module, monkeypatch, capsys): + mesh = mesh_module + mesh.channels._reset_channel_cache() + mesh.config.MODEM_PRESET = "MediumFast" + + class DummyInterface: + def __init__(self) -> None: + self.localNode = SimpleNamespace( + channels=[ + SimpleNamespace(role=1, settings=SimpleNamespace()), + SimpleNamespace( + role=2, + index=5, + settings=SimpleNamespace(name="Chat"), + ), + ] + ) + + def waitForConfig(self) -> None: # noqa: D401 - matches interface contract + return None + + mesh.channels.capture_from_interface(DummyInterface()) + capsys.readouterr() + + captured = [] + monkeypatch.setattr( + mesh, + "_queue_post_json", + lambda path, payload, *, priority: captured.append((path, payload, priority)), + ) + + monkeypatch.setattr(mesh, "DEBUG", True) + + packet = { + "id": "789", + "rxTime": 123456, + "from": "!abc", + "to": "!def", + "channel": 5, + "decoded": {"text": "hi", "portnum": 1}, + } + + mesh.store_packet_dict(packet) + + assert captured, "Expected message to be stored" + path, payload, priority = captured[0] + assert path == "/api/messages" + assert payload["channel_name"] == "Chat" + assert payload["channel"] == 5 + assert payload["text"] == "hi" + assert payload["encrypted"] is None + assert priority == mesh._MESSAGE_POST_PRIORITY + + log_output = capsys.readouterr().out + assert "channel_name='Chat'" in log_output + assert "channel_display='Chat'" in log_output + + def test_store_packet_dict_includes_encrypted_payload(mesh_module, monkeypatch): mesh = mesh_module captured = [] @@ -1322,6 +1439,7 @@ def test_store_packet_dict_includes_encrypted_payload(mesh_module, monkeypatch): assert payload["text"] is None assert payload["from_id"] == 2988082812 assert payload["to_id"] == "!receiver" + assert "channel_name" not in payload assert payload["lora_freq"] == 868 assert payload["modem_preset"] == "MediumFast" assert priority == mesh._MESSAGE_POST_PRIORITY @@ -2056,6 +2174,8 @@ def test_store_packet_dict_debug_message(mesh_module, monkeypatch, capsys): out = capsys.readouterr().out assert "context=handlers.store_packet_dict" in out assert "Queued message payload" in out + assert "channel_display=0" in out + assert "channel_name=" not in out def test_on_receive_skips_seen_packets(mesh_module):