fix(meshcore): reply to DMs via discovered DIRECT route, not flood (#23)

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: Matt Johnson <mj@k7zvx.com>
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
malice 2026-07-03 22:25:40 -06:00 committed by GitHub
commit 5c0f4b42f3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 273 additions and 61 deletions

View file

@ -122,9 +122,8 @@ class MeshCoreTransport(MeshTransport):
"""Resolve a pubkey prefix (or key) to the full MeshCore contact dict. """Resolve a pubkey prefix (or key) to the full MeshCore contact dict.
Refreshes the roster first (ensure_contacts) so the lib can upgrade the 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 6-byte prefix to the full 32-byte key. Returns None if the contact can't
None if the contact can't be resolved. This mirrors what every working be resolved.
meshcore project does before send_msg_with_retry (never send to a bare prefix).
""" """
if self._mc is None: if self._mc is None:
return 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) logger.debug("MeshCore: get_contact_by_key_prefix failed for %s", dest, exc_info=True)
return None 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 # Channel table enumeration
# ------------------------------------------------------------------ # ------------------------------------------------------------------
@ -494,21 +563,38 @@ class MeshCoreTransport(MeshTransport):
contact = self._resolve_contact(destination) contact = self._resolve_contact(destination)
if contact is None: if contact is None:
logger.warning( logger.warning(
"MeshCore: could not resolve a contact for DM dest %s; cannot address reply " "MeshCore: could not resolve contact for DM dest %s; cannot reply",
"(recipient not in roster)", destination, destination,
) )
return False return False
label = contact.get("adv_name") or contact.get("name") or destination # Establish a real route so the reply goes DIRECT — flood DMs are
logger.debug("MeshCore: sending DM to %s via resolved contact", label) # silently rejected on this mesh.
# Pass the CONTACT OBJECT (not the bare prefix) so the lib can upgrade to the self._establish_direct_path(contact, destination)
# full key and reset_path->flood works — the pattern used by all working projects. # 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( result = self._run_coro(
self._mc.commands.send_msg_with_retry(contact, text), self._mc.commands.send_msg(contact, text),
timeout=40, timeout=15,
) )
if result is None: 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 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() success = not result.is_error()
if not success: if not success:
logger.warning("MeshCoreTransport: DM send returned error event") logger.warning("MeshCoreTransport: DM send returned error event")

View file

@ -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 Verifies that send_message(..., destination=...) resolves the destination to
the full contact object before calling send_msg_with_retry (never passes a bare the full contact object, calls _establish_direct_path (path discovery via
prefix string), and correctly maps the return value to True/False: send_path_discovery_sync) BEFORE send_msg, uses plain send_msg (NOT
- non-error Event returned True (ACKed, delivered) send_msg_with_retry), and correctly maps the return value to True/False:
- None returned False (no ACK, delivery not confirmed) - 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) - contact not in roster False (logged warning, send never called)
The meshcore lib is mocked via sys.modules (same pattern as the existing 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 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 = { _CONTACT_DICT = {
"public_key": "a" * 64, "public_key": "a" * 64,
"adv_name": "K7ZVX Matt", "adv_name": "K7ZVX Matt",
@ -82,7 +84,26 @@ def _ensure_fake_meshcore():
return result return result
@staticmethod @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 = MagicMock()
result.is_error.return_value = False result.is_error.return_value = False
return result return result
@ -109,6 +130,28 @@ def _mc_config():
return ConnectionConfig(meshcore_host="127.0.0.1", meshcore_port=5050) 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): def _transport_with_mc_mock(contact=_CONTACT_DICT):
"""Return a MeshCoreTransport with _mc as a MagicMock (no loop thread). """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: The mock exposes:
- mc.get_contact_by_key_prefix(prefix) contact dict (or None when contact=None) - mc.get_contact_by_key_prefix(prefix) contact dict (or None when contact=None)
- mc.ensure_contacts is an AsyncMock (async, returns True) - 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() cfg = _mc_config()
t = MeshCoreTransport(cfg) t = MeshCoreTransport(cfg)
mc = MagicMock() mc = MagicMock()
mc.get_contact_by_key_prefix.return_value = contact mc.get_contact_by_key_prefix.return_value = contact
mc.ensure_contacts = AsyncMock(return_value=True) 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._mc = mc
t._connected = True t._connected = True
@ -144,79 +195,144 @@ def _transport_with_mc_mock(contact=_CONTACT_DICT):
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
class TestMeshCoreDMDelivery: 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): def test_successful_dm_returns_true(self):
"""ACK received (non-error Event) → returns True; send_msg_with_retry called with CONTACT OBJECT.""" """Non-error send_msg result → True."""
t, mc = _transport_with_mc_mock() 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") result = t.send_message("reply text", destination="aabbccdd1122")
assert result is True 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): def test_path_discovery_called_before_send_msg(self):
"""No ACK (send_msg_with_retry returns None) → returns False.""" """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() 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") result = t.send_message("reply text", destination="aabbccdd1122")
assert result is False 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): 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() t, mc = _transport_with_mc_mock()
mc.commands.send_msg = AsyncMock(return_value=_err_event())
err_event = MagicMock()
err_event.is_error.return_value = True
mc.commands.send_msg_with_retry = AsyncMock(return_value=err_event)
result = t.send_message("fail text", destination="deadbeef0011") result = t.send_message("fail text", destination="deadbeef0011")
assert result is False 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): 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) t, mc = _transport_with_mc_mock(contact=None)
mc.commands.send_msg_with_retry = AsyncMock()
result = t.send_message("hello", destination="deadbeef0011") result = t.send_message("hello", destination="deadbeef0011")
assert result is False 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): 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 import logging
t, mc = _transport_with_mc_mock(contact=None) t, mc = _transport_with_mc_mock(contact=None)
mc.commands.send_msg_with_retry = AsyncMock()
with caplog.at_level(logging.WARNING): with caplog.at_level(logging.WARNING):
t.send_message("hello", destination="deadbeef0011") 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)

View file

@ -334,30 +334,40 @@ _DM_CONTACT = {"public_key": "a" * 64, "adv_name": "TestContact", "out_path_len"
class TestSendMessageDM: class TestSendMessageDM:
def test_dispatches_send_msg_with_retry(self): 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() t, mc, _ = _transport_with_mock_mc()
try: try:
# Supply a resolved contact so _resolve_contact succeeds. # Supply a resolved contact so _resolve_contact succeeds.
mc.get_contact_by_key_prefix.return_value = _DM_CONTACT mc.get_contact_by_key_prefix.return_value = _DM_CONTACT
mc.ensure_contacts = AsyncMock(return_value=True) 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 = MagicMock()
ok.is_error.return_value = False 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") result = t.send_message("hi DM", destination="aabbcc")
assert result is True assert result is True
# Must pass the CONTACT OBJECT (not the bare prefix) so the lib can # Must use send_msg (not send_msg_with_retry) with the CONTACT OBJECT.
# upgrade to the full 32-byte key and attempt reset_path->flood. mc.commands.send_msg.assert_awaited_once_with(_DM_CONTACT, "hi DM")
mc.commands.send_msg_with_retry.assert_awaited_once_with(_DM_CONTACT, "hi DM")
finally: finally:
_cleanup(t) _cleanup(t)
def test_send_msg_with_retry_error_returns_false(self): def test_send_msg_with_retry_error_returns_false(self):
"""Error event from send_msg → False."""
t, mc, _ = _transport_with_mock_mc() t, mc, _ = _transport_with_mock_mc()
try: try:
mc.get_contact_by_key_prefix.return_value = _DM_CONTACT mc.get_contact_by_key_prefix.return_value = _DM_CONTACT
mc.ensure_contacts = AsyncMock(return_value=True) 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 = MagicMock()
err.is_error.return_value = True 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 assert t.send_message("hi", destination="deadbeef") is False
finally: finally:
_cleanup(t) _cleanup(t)