From 612a4c3f0e9fbd66e546b52ad9893362bc55319c Mon Sep 17 00:00:00 2001 From: agessaman Date: Tue, 14 Jul 2026 20:40:13 -0700 Subject: [PATCH] fix(companion): route group and raw data like MeshCore --- repeater/main.py | 34 ++++++++++ repeater/packet_router.py | 25 ++++++- tests/test_main_py_coverage.py | 118 ++++++++++++++++++++++++++++++++- tests/test_packet_router.py | 74 ++++++++++++++++++++- 4 files changed, 246 insertions(+), 5 deletions(-) diff --git a/repeater/main.py b/repeater/main.py index 3bcd9c8..e12245e 100644 --- a/repeater/main.py +++ b/repeater/main.py @@ -37,6 +37,8 @@ from repeater.sensors import SensorManager from repeater.utils_packet import create_scoped_advert_packet from repeater.web.http_server import HTTPStatsServer, _log_buffer +from openhop_core.protocol.constants import PAYLOAD_TYPE_RAW_CUSTOM + logger = logging.getLogger("RepeaterDaemon") _COMPANION_LOAD_RETRY_DELAY_SEC = 0.5 @@ -284,6 +286,10 @@ class RepeaterDaemon: self.dispatcher.register_fallback_handler(self._router_callback) logger.info("Packet router registered as fallback (catches all packets)") + # Final-hop RAW_CUSTOM is local-only. Direct packets with remaining + # hops are handed to the router; flood RAW_CUSTOM is discarded. + self._register_raw_custom_handler() + # Set default path hash mode for flood 0-hop packets (adverts, etc.) path_hash_mode = self.config.get("mesh", {}).get("path_hash_mode", 0) if path_hash_mode not in (0, 1, 2): @@ -952,6 +958,34 @@ class RepeaterDaemon: except Exception as e: logger.debug("Push RX raw to companion: %s", e) + def _register_raw_custom_handler(self) -> None: + """Register firmware-compatible RAW_CUSTOM handling ahead of fallback routing.""" + if self.dispatcher: + self.dispatcher.register_handler( + PAYLOAD_TYPE_RAW_CUSTOM, self._on_raw_data_for_companions + ) + + async def _on_raw_data_for_companions(self, packet) -> None: + """Deliver final direct RAW_CUSTOM packets and route direct intermediate hops.""" + if not packet.is_route_direct(): + return + + if getattr(packet, "path", None): + await self._router_callback(packet) + return + + handler = self.repeater_handler + if handler: + if handler.is_duplicate(packet): + return + handler.mark_seen(packet) + + for bridge in self.companion_bridges.values(): + try: + await bridge.process_received_packet(packet) + except Exception as e: + logger.debug("Companion bridge RAW_CUSTOM error: %s", e) + def _register_duplicate_logging_hook(self, dedupe_enabled: bool) -> None: """Register pre-dedup duplicate logging only when dispatcher dedupe is active.""" if not self.dispatcher or not dedupe_enabled: diff --git a/repeater/packet_router.py b/repeater/packet_router.py index 9f2f457..afc4254 100644 --- a/repeater/packet_router.py +++ b/repeater/packet_router.py @@ -15,6 +15,7 @@ from openhop_core.node.handlers.protocol_response import ProtocolResponseHandler from openhop_core.node.handlers.text import TextMessageHandler from openhop_core.node.handlers.trace import TraceHandler from openhop_core.protocol.constants import ( + PAYLOAD_TYPE_GRP_DATA, PH_ROUTE_MASK, ROUTE_TYPE_DIRECT, ROUTE_TYPE_TRANSPORT_DIRECT, @@ -67,6 +68,14 @@ def _is_direct_final_hop(packet) -> bool: return not path or len(path) == 0 +def _is_direct_intermediate_hop(packet) -> bool: + """True for a direct packet that still has one or more routing hops.""" + route = getattr(packet, "header", 0) & PH_ROUTE_MASK + return route in (ROUTE_TYPE_DIRECT, ROUTE_TYPE_TRANSPORT_DIRECT) and not _is_direct_final_hop( + packet + ) + + def _is_expected_drop_reason(reason: str | None) -> bool: if not isinstance(reason, str) or not reason: return False @@ -573,9 +582,7 @@ class PacketRouter: consumed = False if self.daemon.path_helper: try: - consumed = ( - await self.daemon.path_helper.process_path_packet(packet) - ) is True + consumed = (await self.daemon.path_helper.process_path_packet(packet)) is True except Exception as e: logger.debug(f"Path helper processing error: {e}") # The helper/bridge results decide ownership: a direct middle hop @@ -725,6 +732,18 @@ class PacketRouter: except Exception as e: logger.debug(f"Companion bridge GRP_TXT error: {e}") + elif payload_type == PAYLOAD_TYPE_GRP_DATA: + # MeshCore forwards direct packets with remaining hops before payload + # handling. Otherwise, companions authenticate and filter channels. + if not _is_direct_intermediate_hop(packet): + companion_bridges = self._companion_bridges_for_packet(packet, metadata) + if companion_bridges: + for bridge in companion_bridges.values(): + try: + await bridge.process_received_packet(packet) + except Exception as e: + logger.debug(f"Companion bridge GRP_DATA error: {e}") + # Only pass to repeater engine if not already processed by injection # Skip engine for packets we injected for TX (already sent; avoid double-send/double-count) if getattr(packet, "_injected_for_tx", False): diff --git a/tests/test_main_py_coverage.py b/tests/test_main_py_coverage.py index 89562fa..94c1774 100644 --- a/tests/test_main_py_coverage.py +++ b/tests/test_main_py_coverage.py @@ -1,11 +1,19 @@ from types import SimpleNamespace -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, call, patch import pytest from repeater.companion.constants import STATS_TYPE_CORE, STATS_TYPE_PACKETS, STATS_TYPE_RADIO from repeater.main import RepeaterDaemon from repeater.main import main as repeater_main +from openhop_core.node.dispatcher import Dispatcher +from openhop_core.protocol import PacketBuilder +from openhop_core.protocol.constants import ( + PAYLOAD_TYPE_RAW_CUSTOM, + ROUTE_TYPE_FLOOD, + ROUTE_TYPE_TRANSPORT_DIRECT, + ROUTE_TYPE_TRANSPORT_FLOOD, +) class _FakeIdentity: @@ -179,6 +187,114 @@ async def test_raw_rx_and_duplicate_logging_hooks(): engine.record_duplicate.assert_called_once_with(pkt, rssi=-77, snr=1.5) +@pytest.mark.asyncio +async def test_raw_custom_route_type_matrix_is_direct_only_and_first_seen(): + daemon = RepeaterDaemon(_base_config(), radio=object()) + dispatcher = SimpleNamespace(register_handler=MagicMock()) + daemon.dispatcher = dispatcher + daemon.router = SimpleNamespace(enqueue=AsyncMock()) + first_bridge = SimpleNamespace(process_received_packet=AsyncMock()) + second_bridge = SimpleNamespace(process_received_packet=AsyncMock()) + daemon.companion_bridges = {1: first_bridge, 2: second_bridge} + engine = SimpleNamespace( + is_duplicate=MagicMock(side_effect=[False, True, False]), mark_seen=MagicMock() + ) + daemon.repeater_handler = engine + + daemon._register_raw_custom_handler() + dispatcher.register_handler.assert_called_once_with( + PAYLOAD_TYPE_RAW_CUSTOM, daemon._on_raw_data_for_companions + ) + + # Firmware wire vector: RAW_CUSTOM (0x0f), version 0, DIRECT (0x02). + direct_packet = PacketBuilder.create_raw_data(b"\xa5") + assert direct_packet.header == 0x3E + assert direct_packet.write_to() == b"\x3e\x00\xa5" + await daemon._on_raw_data_for_companions(direct_packet) + await daemon._on_raw_data_for_companions(direct_packet) + + # Firmware also treats TRANSPORT_DIRECT as direct routing. + transport_direct_packet = PacketBuilder.create_raw_data(b"\xa6") + transport_direct_packet.header = (PAYLOAD_TYPE_RAW_CUSTOM << 2) | ROUTE_TYPE_TRANSPORT_DIRECT + transport_direct_packet.transport_codes = [0x1234, 0x5678] + assert transport_direct_packet.header == 0x3F + assert transport_direct_packet.write_to() == b"\x3f\x34\x12\x78\x56\x00\xa6" + await daemon._on_raw_data_for_companions(transport_direct_packet) + + first_bridge.process_received_packet.assert_has_awaits( + [call(direct_packet), call(transport_direct_packet)] + ) + second_bridge.process_received_packet.assert_has_awaits( + [call(direct_packet), call(transport_direct_packet)] + ) + engine.mark_seen.assert_has_calls([call(direct_packet), call(transport_direct_packet)]) + + # A direct packet with a remaining hop is router traffic, not local raw data. + intermediate_direct_packet = PacketBuilder.create_raw_data(b"\xa9") + intermediate_direct_packet.path_len = 1 + intermediate_direct_packet.path = bytearray([0x42]) + assert intermediate_direct_packet.write_to() == b"\x3e\x01\x42\xa9" + await daemon._on_raw_data_for_companions(intermediate_direct_packet) + daemon.router.enqueue.assert_awaited_once_with(intermediate_direct_packet) + + # The corresponding FLOOD vector is discarded rather than delivered or routed. + flood_packet = PacketBuilder.create_raw_data(b"\xa7") + flood_packet.header = (PAYLOAD_TYPE_RAW_CUSTOM << 2) | ROUTE_TYPE_FLOOD + assert flood_packet.header == 0x3D + assert flood_packet.write_to() == b"\x3d\x00\xa7" + await daemon._on_raw_data_for_companions(flood_packet) + + transport_flood_packet = PacketBuilder.create_raw_data(b"\xa8") + transport_flood_packet.header = (PAYLOAD_TYPE_RAW_CUSTOM << 2) | ROUTE_TYPE_TRANSPORT_FLOOD + transport_flood_packet.transport_codes = [0x1234, 0x5678] + assert transport_flood_packet.header == 0x3C + assert transport_flood_packet.write_to() == b"\x3c\x34\x12\x78\x56\x00\xa8" + await daemon._on_raw_data_for_companions(transport_flood_packet) + + assert engine.is_duplicate.call_count == 3 + daemon.router.enqueue.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_final_raw_custom_dispatcher_handler_bypasses_router_and_repeater_forwarding(): + daemon = RepeaterDaemon(_base_config(), radio=object()) + radio = SimpleNamespace( + set_rx_callback=MagicMock(), + get_last_rssi=lambda: -80, + get_last_snr=lambda: 2.0, + ) + dispatcher = Dispatcher(radio, dedupe_enabled=False) + daemon.dispatcher = dispatcher + daemon.router = SimpleNamespace(enqueue=AsyncMock()) + daemon.repeater_handler = SimpleNamespace( + is_duplicate=MagicMock(return_value=False), mark_seen=MagicMock() + ) + bridge = SimpleNamespace(process_received_packet=AsyncMock()) + daemon.companion_bridges = {1: bridge} + dispatcher.register_fallback_handler(daemon._router_callback) + daemon._register_raw_custom_handler() + + packet = PacketBuilder.create_raw_data(b"\xde\xad") + await dispatcher._process_received_packet(packet.write_to(), rssi=-81, snr=3.0) + + bridge.process_received_packet.assert_awaited_once() + delivered = bridge.process_received_packet.await_args.args[0] + assert delivered.header == 0x3E + assert bytes(delivered.payload) == b"\xde\xad" + daemon.router.enqueue.assert_not_awaited() + + intermediate = PacketBuilder.create_raw_data(b"\xbe") + intermediate.path_len = 1 + intermediate.path = bytearray([0x42]) + await dispatcher._process_received_packet(intermediate.write_to(), rssi=-81, snr=3.0) + daemon.router.enqueue.assert_awaited_once() + + flood_packet = PacketBuilder.create_raw_data(b"\xef") + flood_packet.header = (PAYLOAD_TYPE_RAW_CUSTOM << 2) | ROUTE_TYPE_FLOOD + await dispatcher._process_received_packet(flood_packet.write_to(), rssi=-81, snr=3.0) + daemon.router.enqueue.assert_awaited_once() + + @pytest.mark.asyncio async def test_deliver_control_data_filters_non_discovery_and_pushes_valid(): daemon = RepeaterDaemon(_base_config(), radio=object()) diff --git a/tests/test_packet_router.py b/tests/test_packet_router.py index d0218b3..f99ff75 100644 --- a/tests/test_packet_router.py +++ b/tests/test_packet_router.py @@ -15,10 +15,12 @@ or: """ import asyncio +import hashlib import time import unittest from unittest.mock import AsyncMock, MagicMock, patch +from openhop_core.companion import CompanionBridge from openhop_core.node.handlers.ack import AckHandler from openhop_core.node.handlers.advert import AdvertHandler from openhop_core.node.handlers.control import ControlHandler @@ -31,7 +33,12 @@ from openhop_core.node.handlers.protocol_response import ProtocolResponseHandler from openhop_core.node.handlers.result import HandlerResult from openhop_core.node.handlers.text import TextMessageHandler from openhop_core.node.handlers.trace import TraceHandler -from openhop_core.protocol.constants import ROUTE_TYPE_DIRECT +from openhop_core.protocol.constants import ( + PAYLOAD_TYPE_GRP_DATA, + ROUTE_TYPE_DIRECT, + ROUTE_TYPE_FLOOD, +) +from openhop_core.protocol import LocalIdentity, PacketBuilder from repeater.packet_router import ( PacketRouter, @@ -284,6 +291,71 @@ class TestInFlightCap(unittest.IsolatedAsyncioTestCase): bridge.process_received_packet.assert_awaited_once() daemon.repeater_handler.assert_awaited_once() + async def test_route_grp_data_fans_out_to_companions_and_forwards(self): + """GRP_DATA reaches each companion for channel matching, then reaches the engine.""" + daemon = _make_daemon() + first_bridge = _make_bridge() + second_bridge = _make_bridge() + daemon.companion_bridges = {0x01: first_bridge, 0x02: second_bridge} + router = PacketRouter(daemon) + pkt = _make_packet(PAYLOAD_TYPE_GRP_DATA) + # Firmware wire vector: GRP_DATA (0x06), version 0, FLOOD (0x01). + pkt.header = (PAYLOAD_TYPE_GRP_DATA << 2) | ROUTE_TYPE_FLOOD + self.assertEqual(pkt.header, 0x19) + + await router._route_packet(pkt) + + first_bridge.process_received_packet.assert_awaited_once_with(pkt) + second_bridge.process_received_packet.assert_awaited_once_with(pkt) + daemon.repeater_handler.assert_awaited_once() + + async def test_route_grp_data_decrypts_in_a_companion_and_forwards(self): + """A firmware-format GRP_DATA packet is both delivered locally and forwarded.""" + channel_secret = b"\x11" * 32 + + async def inject_packet(packet, wait_for_ack=False): + return True + + bridge = CompanionBridge(LocalIdentity(), inject_packet, node_name="Test") + assert bridge.set_channel(0, "Public", channel_secret) + packet = PacketBuilder.create_group_data_packet( + PAYLOAD_TYPE_GRP_DATA, + channel_hash=hashlib.sha256(channel_secret).digest()[0], + channel_secret=channel_secret, + plaintext=b"\x34\x12\x02\xaa\xbb", + secret=channel_secret, + ) + self.assertEqual(packet.header, 0x19) + + daemon = _make_daemon() + daemon.companion_bridges = {0x01: bridge} + router = PacketRouter(daemon) + + await router._route_packet(packet) + + queued = bridge.sync_next_message() + self.assertIsNotNone(queued) + self.assertTrue(queued.is_channel) + self.assertEqual(queued.channel_idx, 0) + self.assertEqual(queued.channel_data_type, 0x1234) + self.assertEqual(queued.channel_data_payload, b"\xaa\xbb") + daemon.repeater_handler.assert_awaited_once() + + async def test_route_direct_grp_data_intermediate_skips_companion_delivery(self): + """A direct GRP_DATA packet with a remaining hop is only offered to the engine.""" + daemon = _make_daemon() + bridge = _make_bridge() + daemon.companion_bridges = {0x01: bridge} + router = PacketRouter(daemon) + packet = _make_packet(PAYLOAD_TYPE_GRP_DATA) + packet.header = (PAYLOAD_TYPE_GRP_DATA << 2) | ROUTE_TYPE_DIRECT + packet.path = bytearray([0x42]) + + await router._route_packet(packet) + + bridge.process_received_packet.assert_not_awaited() + daemon.repeater_handler.assert_awaited_once() + async def test_non_injected_handler_false_is_logged(self): """Inbound packets should log when repeater_handler reports TX failure.""" daemon = _make_daemon()