mirror of
https://github.com/zvx-echo6/meshai.git
synced 2026-08-26 17:31:34 +00:00
This reverts commit 34c6f336cf.
Co-authored-by: Matt Johnson <mj@k7zvx.com>
This commit is contained in:
parent
34c6f336cf
commit
b38c16b3c2
2 changed files with 23 additions and 145 deletions
|
|
@ -600,32 +600,6 @@ class MeshCoreTransport(MeshTransport):
|
|||
# Message I/O
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def _send_dm_once(self, contact: dict, text: str, destination: str):
|
||||
"""Send one DM to *contact* and log the route type.
|
||||
|
||||
Returns the result event from send_msg, or None if the coroutine
|
||||
timed out or failed at the coro layer. Never raises.
|
||||
"""
|
||||
result = self._run_coro(
|
||||
self._mc.commands.send_msg(contact, text),
|
||||
timeout=15,
|
||||
)
|
||||
# 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
|
||||
return result
|
||||
|
||||
def send_message(
|
||||
self,
|
||||
text: str,
|
||||
|
|
@ -666,29 +640,34 @@ class MeshCoreTransport(MeshTransport):
|
|||
destination,
|
||||
)
|
||||
return False
|
||||
had_route = contact.get("out_path_len", -1) >= 0
|
||||
if not had_route:
|
||||
# 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
|
||||
# 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._send_dm_once(contact, text, destination)
|
||||
# Self-heal: if we trusted a cached route but it failed, the path
|
||||
# may be stale — run discovery and retry once.
|
||||
if had_route and (result is None or result.is_error()):
|
||||
logger.info(
|
||||
"MeshCore: cached route to %s failed; running path discovery and retrying",
|
||||
destination,
|
||||
)
|
||||
self._establish_direct_path(contact, destination)
|
||||
contact = self._resolve_contact(destination) or contact
|
||||
result = self._send_dm_once(contact, text, destination)
|
||||
result = self._run_coro(
|
||||
self._mc.commands.send_msg(contact, text),
|
||||
timeout=15,
|
||||
)
|
||||
if result is None:
|
||||
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")
|
||||
|
|
|
|||
|
|
@ -380,107 +380,6 @@ class TestSendMessageDM:
|
|||
_cleanup(t)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 3b. send_message — DM path-discovery skip / self-heal logic
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestSendMessageDMPathSkip:
|
||||
"""Verify the out_path_len guard and stale-route self-heal retry."""
|
||||
|
||||
def _make_ok_event(self):
|
||||
ev = MagicMock()
|
||||
ev.is_error.return_value = False
|
||||
ev.payload = {"type": 0}
|
||||
return ev
|
||||
|
||||
def _make_err_event(self):
|
||||
ev = MagicMock()
|
||||
ev.is_error.return_value = True
|
||||
ev.payload = {"reason": "no ack"}
|
||||
return ev
|
||||
|
||||
def _transport_with_contact(self, out_path_len):
|
||||
"""Return (transport, mc) with mc.get_contact_by_key_prefix returning a
|
||||
contact carrying the given out_path_len. Path-discovery helpers are
|
||||
AsyncMocks so they never block. _establish_direct_path is left intact
|
||||
(caller can monkeypatch as needed)."""
|
||||
t, mc, _ = _transport_with_mock_mc()
|
||||
contact = {
|
||||
"public_key": "a" * 64,
|
||||
"adv_name": "TestNode",
|
||||
"out_path_len": out_path_len,
|
||||
}
|
||||
mc.get_contact_by_key_prefix.return_value = contact
|
||||
mc.ensure_contacts = AsyncMock(return_value=True)
|
||||
# Path-discovery support (used by _establish_direct_path).
|
||||
path_ev = MagicMock()
|
||||
path_ev.is_error.return_value = False
|
||||
mc.commands.send_path_discovery_sync = AsyncMock(return_value=path_ev)
|
||||
mc.commands.get_advert_path = AsyncMock(return_value=path_ev)
|
||||
return t, mc
|
||||
|
||||
def test_routed_contact_skips_discovery(self):
|
||||
"""out_path_len=0 (already routed): _establish_direct_path must NOT be called."""
|
||||
t, mc = self._transport_with_contact(out_path_len=0)
|
||||
try:
|
||||
# Replace _establish_direct_path with a spy so we can assert it's not called.
|
||||
from unittest.mock import Mock
|
||||
spy = Mock()
|
||||
t._establish_direct_path = spy
|
||||
|
||||
mc.commands.send_msg = AsyncMock(return_value=self._make_ok_event())
|
||||
|
||||
result = t.send_message("hi", destination="7d4e07237294")
|
||||
|
||||
spy.assert_not_called()
|
||||
mc.commands.send_msg.assert_awaited_once()
|
||||
assert result is True
|
||||
finally:
|
||||
_cleanup(t)
|
||||
|
||||
def test_flood_contact_runs_discovery(self):
|
||||
"""out_path_len=-1 (flood/unknown): _establish_direct_path MUST be called once."""
|
||||
t, mc = self._transport_with_contact(out_path_len=-1)
|
||||
try:
|
||||
from unittest.mock import Mock
|
||||
spy = Mock()
|
||||
t._establish_direct_path = spy
|
||||
|
||||
mc.commands.send_msg = AsyncMock(return_value=self._make_ok_event())
|
||||
|
||||
result = t.send_message("hi", destination="7d4e07237294")
|
||||
|
||||
spy.assert_called_once()
|
||||
mc.commands.send_msg.assert_awaited_once()
|
||||
assert result is True
|
||||
finally:
|
||||
_cleanup(t)
|
||||
|
||||
def test_stale_routed_contact_self_heals(self):
|
||||
"""out_path_len=0 but first send fails (stale route): discovery runs,
|
||||
send retried once, final return is True."""
|
||||
t, mc = self._transport_with_contact(out_path_len=0)
|
||||
try:
|
||||
from unittest.mock import Mock
|
||||
spy = Mock()
|
||||
t._establish_direct_path = spy
|
||||
|
||||
# First send fails, second succeeds.
|
||||
mc.commands.send_msg = AsyncMock(
|
||||
side_effect=[self._make_err_event(), self._make_ok_event()]
|
||||
)
|
||||
|
||||
result = t.send_message("hi", destination="7d4e07237294")
|
||||
|
||||
# Discovery was triggered exactly once (the self-heal retry).
|
||||
spy.assert_called_once()
|
||||
# send_msg was called twice (first attempt + retry).
|
||||
assert mc.commands.send_msg.await_count == 2
|
||||
assert result is True
|
||||
finally:
|
||||
_cleanup(t)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 4. Inbound normalization — direct method calls (hermetic, no threads)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue