From 44d13ace10cae704c30267ba53cdb10a7d4d4bdc Mon Sep 17 00:00:00 2001 From: Matt Johnson Date: Sat, 4 Jul 2026 04:25:35 +0000 Subject: [PATCH] fix(meshcore): reply to DMs via discovered DIRECT route, not flood MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pyMC-companion FLOOD direct-messages are silently rejected by recipient MeshCore nodes (43 sent / ~1 ack in 24h); DIRECT-routed packets deliver. A bare inbound DM does NOT populate the contact's out_path on pyMC, so the contact stays out_path_len=-1 (flood), and send_msg_with_retry actively reset_path→flood, guaranteeing the broken route. New behavior on a DM reply: - _establish_direct_path(): path discovery (CMD 52 send_path_discovery_sync) so the recipient returns a PATH packet → pyMC writes a real out_path → works at ANY hop count. Fallback: seed from the sender's cached advert (get_advert_path CMD 42 → update_contact CMD 9) if discovery is empty. - send via plain send_msg (CMD 2) — uses the learned path → DIRECT. Drops send_msg_with_retry (which forced flood). Logs the RESP_CODE_SENT route (direct/flood) so we can confirm. Tests reworked for the discover-then-direct-send path; 0 new failures. Co-Authored-By: Claude Opus 4.8 (1M context) --- work/meshai/transport/meshcore_transport.py | 110 ++++++++-- work/tests/test_meshcore_dm_delivery.py | 212 +++++++++++++++----- work/tests/test_meshcore_transport.py | 20 +- 3 files changed, 277 insertions(+), 65 deletions(-) diff --git a/work/meshai/transport/meshcore_transport.py b/work/meshai/transport/meshcore_transport.py index afcf329..1909c44 100644 --- a/work/meshai/transport/meshcore_transport.py +++ b/work/meshai/transport/meshcore_transport.py @@ -122,9 +122,8 @@ class MeshCoreTransport(MeshTransport): """Resolve a pubkey prefix (or key) to the full MeshCore contact dict. Refreshes the roster first (ensure_contacts) so the lib can upgrade the - 6-byte prefix to the full 32-byte key and run reset_path->flood. Returns - None if the contact can't be resolved. This mirrors what every working - meshcore project does before send_msg_with_retry (never send to a bare prefix). + 6-byte prefix to the full 32-byte key. Returns None if the contact can't + be resolved. """ if self._mc is None: return None @@ -140,6 +139,76 @@ class MeshCoreTransport(MeshTransport): logger.debug("MeshCore: get_contact_by_key_prefix failed for %s", dest, exc_info=True) return None + def _establish_direct_path(self, contact: dict, dst: str) -> None: + """Best-effort: establish a DIRECT route to *contact* before replying. + + Runs path discovery (CMD 52) so the recipient returns a PATH packet; the + lib updates the contact's out_path and subsequent send_msg goes DIRECT + instead of flood (flood DMs are silently rejected on this mesh). + + Falls back to seeding the path from a cached advert path if discovery + times out and the contact is still marked as flood (out_path_len == -1). + + Never raises; all steps are individually guarded. + """ + label = contact.get("adv_name") or dst + + # Step 1: trigger path discovery (CMD 52 → 0x34); wait for PATH_RESPONSE. + # The lib updates the contact's out_path internally when it receives the + # response, so a subsequent send_msg will use the direct route. + try: + path_event = self._run_coro( + self._mc.commands.send_path_discovery_sync(contact, timeout=25), + timeout=30, + ) + if path_event is not None and not path_event.is_error(): + logger.info( + "MeshCore: path discovery to %s succeeded (PATH_RESPONSE received)", label + ) + else: + logger.info( + "MeshCore: path discovery to %s — no PATH_RESPONSE; trying advert-path fallback", + label, + ) + except Exception as exc: + logger.debug("MeshCore: send_path_discovery_sync to %s failed: %s", dst, exc) + + # Step 2 (fallback): if the contact is still flood after discovery, + # seed the route from a cached advert path. + try: + fresh_contact = self._resolve_contact(dst) or contact + if fresh_contact.get("out_path_len", -1) < 0: + try: + ap_event = self._run_coro( + self._mc.commands.get_advert_path(fresh_contact), + timeout=10, + ) + if ( + ap_event is not None + and not ap_event.is_error() + and isinstance(getattr(ap_event, "payload", None), dict) + ): + path_hex = ap_event.payload.get("path", "") + path_len = ap_event.payload.get("path_len", -1) + if path_len > 0 and path_hex: + self._run_coro( + self._mc.commands.update_contact(fresh_contact, path=path_hex), + timeout=10, + ) + logger.info( + "MeshCore: seeded direct path for %s from advert-path cache (len=%d)", + fresh_contact.get("adv_name") or dst, + path_len, + ) + else: + logger.debug( + "MeshCore: advert-path for %s is flood/empty — leaving as flood", dst + ) + except Exception as exc: + logger.debug("MeshCore: advert-path fallback for %s failed: %s", dst, exc) + except Exception as exc: + logger.debug("MeshCore: _establish_direct_path step-2 error for %s: %s", dst, exc) + # ------------------------------------------------------------------ # Channel table enumeration # ------------------------------------------------------------------ @@ -494,21 +563,38 @@ class MeshCoreTransport(MeshTransport): contact = self._resolve_contact(destination) if contact is None: logger.warning( - "MeshCore: could not resolve a contact for DM dest %s; cannot address reply " - "(recipient not in roster)", destination, + "MeshCore: could not resolve contact for DM dest %s; cannot reply", + destination, ) return False - label = contact.get("adv_name") or contact.get("name") or destination - logger.debug("MeshCore: sending DM to %s via resolved contact", label) - # Pass the CONTACT OBJECT (not the bare prefix) so the lib can upgrade to the - # full key and reset_path->flood works — the pattern used by all working projects. + # Establish a real route so the reply goes DIRECT — flood DMs are + # silently rejected on this mesh. + self._establish_direct_path(contact, destination) + # Re-resolve so we send with the freshly-learned out_path. + contact = self._resolve_contact(destination) or contact + # Plain send_msg (NOT send_msg_with_retry — that calls reset_path + # which forces flood, defeating the path we just established). result = self._run_coro( - self._mc.commands.send_msg_with_retry(contact, text), - timeout=40, + self._mc.commands.send_msg(contact, text), + timeout=15, ) if result is None: - logger.warning("MeshCoreTransport: DM to %s not ACKed (delivery not confirmed)", label) + logger.warning("MeshCoreTransport: DM to %s — no send result", destination) return False + # Log whether the radio sent DIRECT or FLOOD from RESP_CODE_SENT type field. + try: + sent_type = ( + result.payload.get("type") + if hasattr(result, "payload") and isinstance(result.payload, dict) + else None + ) + logger.info( + "MeshCore: DM to %s sent (route=%s)", + contact.get("adv_name") or destination, + "flood" if sent_type == 1 else ("direct" if sent_type == 0 else "?"), + ) + except Exception: + pass success = not result.is_error() if not success: logger.warning("MeshCoreTransport: DM send returned error event") diff --git a/work/tests/test_meshcore_dm_delivery.py b/work/tests/test_meshcore_dm_delivery.py index 6b03cfe..7fb80fd 100644 --- a/work/tests/test_meshcore_dm_delivery.py +++ b/work/tests/test_meshcore_dm_delivery.py @@ -1,10 +1,12 @@ -"""Focused tests for the MeshCore DM delivery fix. +"""Focused tests for the MeshCore DM direct-route delivery fix. Verifies that send_message(..., destination=...) resolves the destination to -the full contact object before calling send_msg_with_retry (never passes a bare -prefix string), and correctly maps the return value to True/False: - - non-error Event returned → True (ACKed, delivered) - - None returned → False (no ACK, delivery not confirmed) +the full contact object, calls _establish_direct_path (path discovery via +send_path_discovery_sync) BEFORE send_msg, uses plain send_msg (NOT +send_msg_with_retry), and correctly maps the return value to True/False: + - non-error Event returned → True (sent, route logged) + - None returned → False (no send result) + - error Event returned → False - contact not in roster → False (logged warning, send never called) The meshcore lib is mocked via sys.modules (same pattern as the existing @@ -19,7 +21,7 @@ from unittest.mock import AsyncMock, MagicMock, patch import pytest -# Contact dict returned by the fake roster — 64-hex public_key, adv_name, out_path_len. +# Contact dict returned by the fake roster — 64-hex public_key, adv_name, out_path_len=-1 (flood). _CONTACT_DICT = { "public_key": "a" * 64, "adv_name": "K7ZVX Matt", @@ -82,7 +84,26 @@ def _ensure_fake_meshcore(): return result @staticmethod - async def send_msg(dst, text): + async def send_msg(dst, text, timestamp=None, attempt=0): + result = MagicMock() + result.is_error.return_value = False + result.payload = {"type": 0, "expected_ack": "00000000"} + return result + + @staticmethod + async def send_path_discovery_sync(dst, timeout=0, min_timeout=0): + result = MagicMock() + result.is_error.return_value = False + return result + + @staticmethod + async def get_advert_path(key): + result = MagicMock() + result.is_error.return_value = True # default: not found + return result + + @staticmethod + async def update_contact(contact, path=None, flags=None, path_hash_mode=None): result = MagicMock() result.is_error.return_value = False return result @@ -109,6 +130,28 @@ def _mc_config(): return ConnectionConfig(meshcore_host="127.0.0.1", meshcore_port=5050) +def _ok_event(route_type: int = 0): + """Non-error Event-like mock with MSG_SENT payload (0=direct, 1=flood).""" + ev = MagicMock() + ev.is_error.return_value = False + ev.payload = {"type": route_type, "expected_ack": "deadbeef"} + return ev + + +def _err_event(): + ev = MagicMock() + ev.is_error.return_value = True + ev.payload = {"reason": "test error"} + return ev + + +def _path_event(): + """Non-error PATH_RESPONSE-like Event mock.""" + ev = MagicMock() + ev.is_error.return_value = False + return ev + + def _transport_with_mc_mock(contact=_CONTACT_DICT): """Return a MeshCoreTransport with _mc as a MagicMock (no loop thread). @@ -119,12 +162,20 @@ def _transport_with_mc_mock(contact=_CONTACT_DICT): The mock exposes: - mc.get_contact_by_key_prefix(prefix) → contact dict (or None when contact=None) - mc.ensure_contacts is an AsyncMock (async, returns True) + - mc.commands.send_path_discovery_sync: AsyncMock → _path_event() + - mc.commands.send_msg: AsyncMock → _ok_event() + - mc.commands.get_advert_path: AsyncMock → _err_event() (not found by default) + - mc.commands.update_contact: AsyncMock → _ok_event() """ cfg = _mc_config() t = MeshCoreTransport(cfg) mc = MagicMock() mc.get_contact_by_key_prefix.return_value = contact mc.ensure_contacts = AsyncMock(return_value=True) + mc.commands.send_path_discovery_sync = AsyncMock(return_value=_path_event()) + mc.commands.send_msg = AsyncMock(return_value=_ok_event()) + mc.commands.get_advert_path = AsyncMock(return_value=_err_event()) + mc.commands.update_contact = AsyncMock(return_value=_ok_event()) t._mc = mc t._connected = True @@ -144,79 +195,144 @@ def _transport_with_mc_mock(contact=_CONTACT_DICT): # --------------------------------------------------------------------------- class TestMeshCoreDMDelivery: - """send_message with destination= must resolve the contact and use send_msg_with_retry.""" + """send_message with destination= must call path discovery then plain send_msg.""" - def test_acked_dm_returns_true_and_uses_send_msg_with_retry(self): - """ACK received (non-error Event) → returns True; send_msg_with_retry called with CONTACT OBJECT.""" + def test_successful_dm_returns_true(self): + """Non-error send_msg result → True.""" t, mc = _transport_with_mc_mock() - ok_event = MagicMock() - ok_event.is_error.return_value = False - mc.commands.send_msg_with_retry = AsyncMock(return_value=ok_event) - mc.commands.send_msg = AsyncMock() # must NOT be called - result = t.send_message("reply text", destination="aabbccdd1122") assert result is True - # Must pass the CONTACT OBJECT (dict), not the bare prefix string. - mc.commands.send_msg_with_retry.assert_awaited_once_with(_CONTACT_DICT, "reply text") - mc.commands.send_msg.assert_not_awaited() - def test_no_ack_returns_false(self): - """No ACK (send_msg_with_retry returns None) → returns False.""" + def test_path_discovery_called_before_send_msg(self): + """send_path_discovery_sync must be called before send_msg.""" + t, mc = _transport_with_mc_mock() + call_order = [] + + async def fake_pds(dst, timeout=0, min_timeout=0): + call_order.append("discovery") + return _path_event() + + async def fake_send(dst, text, timestamp=None, attempt=0): + call_order.append("send_msg") + return _ok_event() + + mc.commands.send_path_discovery_sync = fake_pds + mc.commands.send_msg = fake_send + + t.send_message("hello", destination="aabbccdd1122") + + assert "discovery" in call_order, "send_path_discovery_sync was never called" + assert "send_msg" in call_order, "send_msg was never called" + assert call_order.index("discovery") < call_order.index("send_msg"), ( + "send_path_discovery_sync must be called before send_msg" + ) + + def test_send_msg_used_not_send_msg_with_retry(self): + """Plain send_msg (not send_msg_with_retry) must be used for DMs.""" + t, mc = _transport_with_mc_mock() + mc.commands.send_msg_with_retry = AsyncMock() # must NOT be called + + t.send_message("reply text", destination="aabbccdd1122") + + mc.commands.send_msg.assert_awaited() + mc.commands.send_msg_with_retry.assert_not_awaited() + + def test_send_msg_receives_contact_dict(self): + """send_msg must receive the resolved contact dict as its first argument.""" t, mc = _transport_with_mc_mock() - mc.commands.send_msg_with_retry = AsyncMock(return_value=None) + t.send_message("reply text", destination="aabbccdd1122") + + args, _ = mc.commands.send_msg.await_args + assert args[0] == _CONTACT_DICT, ( + f"send_msg first arg should be the contact dict, got {args[0]!r}" + ) + + def test_none_send_result_returns_false(self): + """send_msg returning None → False.""" + t, mc = _transport_with_mc_mock() + mc.commands.send_msg = AsyncMock(return_value=None) result = t.send_message("reply text", destination="aabbccdd1122") assert result is False - mc.commands.send_msg_with_retry.assert_awaited_once_with(_CONTACT_DICT, "reply text") + + def test_none_send_result_logs_warning(self, caplog): + """send_msg returning None → warning is logged.""" + import logging + t, mc = _transport_with_mc_mock() + mc.commands.send_msg = AsyncMock(return_value=None) + + with caplog.at_level(logging.WARNING): + t.send_message("msg", destination="deadbeef0011") + + assert any("no send result" in r.getMessage() for r in caplog.records) def test_error_event_returns_false(self): - """Error event returned by send_msg_with_retry → returns False.""" + """Error event returned by send_msg → False.""" t, mc = _transport_with_mc_mock() - - err_event = MagicMock() - err_event.is_error.return_value = True - mc.commands.send_msg_with_retry = AsyncMock(return_value=err_event) + mc.commands.send_msg = AsyncMock(return_value=_err_event()) result = t.send_message("fail text", destination="deadbeef0011") assert result is False - def test_no_ack_logs_warning(self, caplog): - """No ACK → a warning mentioning the destination is logged.""" - import logging - t, mc = _transport_with_mc_mock() - mc.commands.send_msg_with_retry = AsyncMock(return_value=None) - - with caplog.at_level(logging.WARNING): - t.send_message("msg", destination="deadbeef0011") - - assert any( - "not ACKed" in r.getMessage() or "delivery not confirmed" in r.getMessage() - for r in caplog.records - ) - def test_unresolved_contact_returns_false_without_calling_send(self): - """When get_contact_by_key_prefix returns None, send_message returns False and never calls send_msg_with_retry.""" + """When get_contact_by_key_prefix returns None, returns False and never calls send_msg.""" t, mc = _transport_with_mc_mock(contact=None) - mc.commands.send_msg_with_retry = AsyncMock() - result = t.send_message("hello", destination="deadbeef0011") assert result is False - mc.commands.send_msg_with_retry.assert_not_awaited() + mc.commands.send_msg.assert_not_awaited() def test_unresolved_contact_logs_warning(self, caplog): - """Unresolved contact → a warning about 'not in roster' is logged.""" + """Unresolved contact → a warning about inability to reply is logged.""" import logging t, mc = _transport_with_mc_mock(contact=None) - mc.commands.send_msg_with_retry = AsyncMock() with caplog.at_level(logging.WARNING): t.send_message("hello", destination="deadbeef0011") - assert any("not in roster" in r.getMessage() for r in caplog.records) + assert any("cannot reply" in r.getMessage() for r in caplog.records) + + def test_path_discovery_failure_does_not_prevent_send(self): + """If path discovery raises, send_msg is still attempted (best-effort).""" + t, mc = _transport_with_mc_mock() + mc.commands.send_path_discovery_sync = AsyncMock( + side_effect=RuntimeError("discovery timed out") + ) + + result = t.send_message("hello", destination="aabbccdd1122") + + mc.commands.send_msg.assert_awaited() + assert result is True + + def test_advert_path_fallback_invoked_when_still_flood(self): + """When contact is still flood after path discovery, get_advert_path + update_contact are called.""" + t, mc = _transport_with_mc_mock() + # Path discovery returns None (no PATH_RESPONSE) — contact stays flood. + mc.commands.send_path_discovery_sync = AsyncMock(return_value=None) + # Advert path returns a real direct path. + ap_ev = MagicMock() + ap_ev.is_error.return_value = False + ap_ev.payload = {"path": "aabbcc1122dd", "path_len": 2, "path_hash_mode": 0} + mc.commands.get_advert_path = AsyncMock(return_value=ap_ev) + + t.send_message("hello", destination="aabbccdd1122") + + mc.commands.get_advert_path.assert_awaited() + mc.commands.update_contact.assert_awaited() + + def test_direct_route_type_logged(self, caplog): + """When send_msg reports route type 0 (direct), 'direct' is logged.""" + import logging + t, mc = _transport_with_mc_mock() + mc.commands.send_msg = AsyncMock(return_value=_ok_event(route_type=0)) + + with caplog.at_level(logging.INFO): + t.send_message("hi", destination="aabbccdd1122") + + assert any("direct" in r.getMessage() for r in caplog.records) diff --git a/work/tests/test_meshcore_transport.py b/work/tests/test_meshcore_transport.py index dd9e716..b3dc617 100644 --- a/work/tests/test_meshcore_transport.py +++ b/work/tests/test_meshcore_transport.py @@ -334,30 +334,40 @@ _DM_CONTACT = {"public_key": "a" * 64, "adv_name": "TestContact", "out_path_len" class TestSendMessageDM: def test_dispatches_send_msg_with_retry(self): + """DM send: path discovery then plain send_msg (not send_msg_with_retry).""" t, mc, _ = _transport_with_mock_mc() try: # Supply a resolved contact so _resolve_contact succeeds. mc.get_contact_by_key_prefix.return_value = _DM_CONTACT mc.ensure_contacts = AsyncMock(return_value=True) + # Path discovery must be awaitable. + path_ev = MagicMock() + path_ev.is_error.return_value = False + mc.commands.send_path_discovery_sync = AsyncMock(return_value=path_ev) ok = MagicMock() ok.is_error.return_value = False - mc.commands.send_msg_with_retry = AsyncMock(return_value=ok) + ok.payload = {"type": 0, "expected_ack": "00000000"} + mc.commands.send_msg = AsyncMock(return_value=ok) result = t.send_message("hi DM", destination="aabbcc") assert result is True - # Must pass the CONTACT OBJECT (not the bare prefix) so the lib can - # upgrade to the full 32-byte key and attempt reset_path->flood. - mc.commands.send_msg_with_retry.assert_awaited_once_with(_DM_CONTACT, "hi DM") + # Must use send_msg (not send_msg_with_retry) with the CONTACT OBJECT. + mc.commands.send_msg.assert_awaited_once_with(_DM_CONTACT, "hi DM") finally: _cleanup(t) def test_send_msg_with_retry_error_returns_false(self): + """Error event from send_msg → False.""" t, mc, _ = _transport_with_mock_mc() try: mc.get_contact_by_key_prefix.return_value = _DM_CONTACT mc.ensure_contacts = AsyncMock(return_value=True) + path_ev = MagicMock() + path_ev.is_error.return_value = False + mc.commands.send_path_discovery_sync = AsyncMock(return_value=path_ev) err = MagicMock() err.is_error.return_value = True - mc.commands.send_msg_with_retry = AsyncMock(return_value=err) + err.payload = {"reason": "test"} + mc.commands.send_msg = AsyncMock(return_value=err) assert t.send_message("hi", destination="deadbeef") is False finally: _cleanup(t)