mirror of
https://github.com/pyMC-dev/pyMC_Repeater.git
synced 2026-08-08 09:52:55 +02:00
fix(companion): route group and raw data like MeshCore
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user