From 0a65e0566b29a288dc3eac798ca5ee9bee6be78b Mon Sep 17 00:00:00 2001 From: agessaman Date: Sun, 19 Jul 2026 07:12:24 -0700 Subject: [PATCH] fix(router): mark companion delivery only after a bridge receives the packet MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The PATH and protocol-response paths recorded the companion dedupe unconditionally after the bridge fan-out, so a delivery where every bridge raised was still suppressed for the full dedupe TTL — the client lost the packet even though later mesh copies arrived. Have the fan-out report delivery (at least one bridge completed without raising) next to authentication, and mark the dedupe only on delivery, which is the behaviour the marking helper's contract already documented. One healthy bridge still counts as delivered, so duplicate suppression of repeated mesh copies is unchanged. --- repeater/packet_router.py | 47 ++++++++++++++++++------------ tests/test_packet_router.py | 57 +++++++++++++++++++++++++++++++++++++ 2 files changed, 86 insertions(+), 18 deletions(-) diff --git a/repeater/packet_router.py b/repeater/packet_router.py index 197975d..122d379 100644 --- a/repeater/packet_router.py +++ b/repeater/packet_router.py @@ -273,14 +273,21 @@ class PacketRouter: return {} return companion_bridges - async def _fan_out_to_bridges(self, packet, bridges, *, context: str) -> bool: - """Offer packet to each bridge; True if any bridge authenticated it. + async def _fan_out_to_bridges(self, packet, bridges, *, context: str) -> tuple: + """Offer packet to each bridge; report ``(delivered, authenticated)``. Accepts a dict of bridges — pass a single-entry dict for targeted delivery to the bridge that owns ``dest_hash``. A bridge that raises is logged and skipped; ``result.authenticated`` is read directly (every bridge returns a HandlerResult) so a broken contract surfaces instead of being hidden. + + ``delivered`` is True when at least one bridge completed without raising — + the signal callers must use before ``_mark_delivered_to_companions`` so a + delivery where every bridge raised is retried on the next copy instead of + being suppressed for the dedupe TTL. ``authenticated`` is True when any + bridge authenticated the packet. """ + delivered = False authenticated = False for bridge in bridges.values(): try: @@ -288,9 +295,10 @@ class PacketRouter: except Exception as e: logger.debug("Companion bridge %s error: %s", context, e) continue + delivered = True if result.authenticated is True: authenticated = True - return authenticated + return delivered, authenticated async def _consume_via_local_candidates( self, packet, metadata: dict, dest_hash, helper, process_method_name: str @@ -638,21 +646,21 @@ class PacketRouter: companion_bridges = self._companion_bridges_for_packet(packet, metadata) if dest_hash is not None and dest_hash in companion_bridges: if not self._was_delivered_to_companions(packet): - consumed = ( - await self._fan_out_to_bridges( - packet, {dest_hash: companion_bridges[dest_hash]}, context="PATH" - ) - or consumed + delivered, authenticated = await self._fan_out_to_bridges( + packet, {dest_hash: companion_bridges[dest_hash]}, context="PATH" ) - self._mark_delivered_to_companions(packet) + consumed = authenticated or consumed + if delivered: + self._mark_delivered_to_companions(packet) elif companion_bridges and not self._was_delivered_to_companions(packet): # Dest not in bridges: path-return with ephemeral dest (e.g. multi-hop login). # Deliver to all bridges; each will try to decrypt and ignore if not relevant. - consumed = ( - await self._fan_out_to_bridges(packet, companion_bridges, context="PATH") - or consumed + delivered, authenticated = await self._fan_out_to_bridges( + packet, companion_bridges, context="PATH" ) - self._mark_delivered_to_companions(packet) + consumed = authenticated or consumed + if delivered: + self._mark_delivered_to_companions(packet) logger.debug( "PATH dest=0x%02x (anon) delivered to %d bridge(s) for matching", dest_hash or 0, @@ -674,13 +682,13 @@ class PacketRouter: companion_bridges = self._companion_bridges_for_packet(packet, metadata) local_hash = getattr(self.daemon, "local_hash", None) if dest_hash is not None and dest_hash in companion_bridges: - consumed = await self._fan_out_to_bridges( + _, consumed = await self._fan_out_to_bridges( packet, {dest_hash: companion_bridges[dest_hash]}, context="RESPONSE" ) logger.info("RESPONSE dest=0x%02x delivered to companion bridge", dest_hash) elif dest_hash == local_hash and companion_bridges: # Response addressed to this repeater (e.g. path-based reply to first hop) - consumed = await self._fan_out_to_bridges( + _, consumed = await self._fan_out_to_bridges( packet, companion_bridges, context="RESPONSE" ) logger.info( @@ -692,7 +700,7 @@ class PacketRouter: # Dest not in bridges and not local: likely ANON_REQ response (dest = ephemeral # sender hash). Deliver to all bridges; each will try to decrypt and ignore if # not relevant (firmware-like behavior, works with multiple companion bridges). - consumed = await self._fan_out_to_bridges( + _, consumed = await self._fan_out_to_bridges( packet, companion_bridges, context="RESPONSE" ) logger.debug( @@ -717,8 +725,11 @@ class PacketRouter: companion_bridges = self._companion_bridges_for_packet(packet, metadata) final_hop = _is_direct_final_hop(packet) if companion_bridges and (final_hop or not self._was_delivered_to_companions(packet)): - await self._fan_out_to_bridges(packet, companion_bridges, context="RESPONSE") - self._mark_delivered_to_companions(packet) + delivered, _ = await self._fan_out_to_bridges( + packet, companion_bridges, context="RESPONSE" + ) + if delivered: + self._mark_delivered_to_companions(packet) if companion_bridges and final_hop: # DIRECT with empty path: we're the final hop, so consume after delivery. processed_by_injection = True diff --git a/tests/test_packet_router.py b/tests/test_packet_router.py index 4d2456f..a1e46f5 100644 --- a/tests/test_packet_router.py +++ b/tests/test_packet_router.py @@ -1326,6 +1326,63 @@ class TestInjectedTxRawEcho(unittest.IsolatedAsyncioTestCase): class TestCompanionDeliveryFailureHandling(unittest.IsolatedAsyncioTestCase): """Dedupe marking and candidate-loop behaviour when companion bridges raise.""" + async def test_path_all_bridges_raising_is_retried_on_next_copy(self): + """A PATH delivery where every bridge raised must not be marked + delivered; the next mesh copy gets another delivery attempt.""" + daemon = _make_daemon() + bridge = _make_bridge() + bridge.process_received_packet = AsyncMock(side_effect=RuntimeError("bridge down")) + daemon.companion_bridges = {0x01: bridge} + router = PacketRouter(daemon) + pkt = _make_packet(PathHandler.payload_type()) + pkt.payload = bytes([0x01, 0xAA]) + + await router._route_packet(pkt) + self.assertEqual(bridge.process_received_packet.await_count, 1) + + # Bridge recovers; the second copy must be delivered, not TTL-suppressed. + recovered = AsyncMock(return_value=HandlerResult.not_for_us()) + bridge.process_received_packet = recovered + await router._route_packet(pkt) + recovered.assert_awaited_once() + + async def test_path_partial_bridge_failure_still_marks_delivered(self): + """One healthy bridge is a delivery: the duplicate copy stays suppressed.""" + daemon = _make_daemon() + raising = _make_bridge() + raising.process_received_packet = AsyncMock(side_effect=RuntimeError("boom")) + healthy = _make_bridge() + healthy.process_received_packet = AsyncMock(return_value=HandlerResult.not_for_us()) + daemon.companion_bridges = {0x01: raising, 0x02: healthy} + router = PacketRouter(daemon) + pkt = _make_packet(PathHandler.payload_type()) + # Dest not in bridges: anon path-return, delivered to all bridges. + pkt.payload = bytes([0xEE, 0xAA]) + + await router._route_packet(pkt) + await router._route_packet(pkt) + + healthy.process_received_packet.assert_awaited_once() + raising.process_received_packet.assert_awaited_once() + + async def test_protocol_response_all_bridges_raising_is_retried_on_next_copy(self): + daemon = _make_daemon() + bridge = _make_bridge() + bridge.process_received_packet = AsyncMock(side_effect=RuntimeError("bridge down")) + daemon.companion_bridges = {0x01: bridge} + router = PacketRouter(daemon) + pkt = _make_packet(ProtocolResponseHandler.payload_type()) + pkt.header = ROUTE_TYPE_FLOOD # not a final hop: dedupe decides delivery + + with patch("repeater.packet_router.PathHandler.payload_type", return_value=0x55): + await router._route_packet(pkt) + self.assertEqual(bridge.process_received_packet.await_count, 1) + + recovered = AsyncMock(return_value=HandlerResult.not_for_us()) + bridge.process_received_packet = recovered + await router._route_packet(pkt) + recovered.assert_awaited_once() + async def test_login_candidate_bridge_error_still_offers_local_identity(self): """A raising companion bridge must not abort the candidate loop: the hash-colliding room-server/repeater identity still gets the packet."""