mirror of
https://github.com/zvx-echo6/meshai.git
synced 2026-08-26 17:31:34 +00:00
Revert "fix(meshcore): skip 25s path discovery when the contact is already routed (#57)"
This reverts commit 34c6f336cf.
This commit is contained in:
parent
34c6f336cf
commit
ca7c24cf82
2 changed files with 23 additions and 145 deletions
|
|
@ -600,32 +600,6 @@ class MeshCoreTransport(MeshTransport):
|
||||||
# Message I/O
|
# 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(
|
def send_message(
|
||||||
self,
|
self,
|
||||||
text: str,
|
text: str,
|
||||||
|
|
@ -666,29 +640,34 @@ class MeshCoreTransport(MeshTransport):
|
||||||
destination,
|
destination,
|
||||||
)
|
)
|
||||||
return False
|
return False
|
||||||
had_route = contact.get("out_path_len", -1) >= 0
|
# Establish a real route so the reply goes DIRECT — flood DMs are
|
||||||
if not had_route:
|
# silently rejected on this mesh.
|
||||||
# Establish a real route so the reply goes DIRECT — flood DMs are
|
self._establish_direct_path(contact, destination)
|
||||||
# silently rejected on this mesh.
|
# Re-resolve so we send with the freshly-learned out_path.
|
||||||
self._establish_direct_path(contact, destination)
|
contact = self._resolve_contact(destination) or contact
|
||||||
# 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
|
# Plain send_msg (NOT send_msg_with_retry — that calls reset_path
|
||||||
# which forces flood, defeating the path we just established).
|
# which forces flood, defeating the path we just established).
|
||||||
result = self._send_dm_once(contact, text, destination)
|
result = self._run_coro(
|
||||||
# Self-heal: if we trusted a cached route but it failed, the path
|
self._mc.commands.send_msg(contact, text),
|
||||||
# may be stale — run discovery and retry once.
|
timeout=15,
|
||||||
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)
|
|
||||||
if result is None:
|
if result is None:
|
||||||
logger.warning("MeshCoreTransport: DM to %s — no send result", destination)
|
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")
|
||||||
|
|
|
||||||
|
|
@ -380,107 +380,6 @@ class TestSendMessageDM:
|
||||||
_cleanup(t)
|
_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)
|
# 4. Inbound normalization — direct method calls (hermetic, no threads)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue