import base64 import ctypes import importlib.util import json import os import queue import subprocess import sys import tempfile import threading import time import unittest from unittest import mock from pathlib import Path import RNS BRIDGE_PATH = Path(__file__).with_name("presence_bridge.py") BASE58_ALPHABET = "123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz" def base58_encode(data): leading_zeroes = len(data) - len(data.lstrip(b"\x00")) value = int.from_bytes(data, "big") encoded = "" while value > 0: value, remainder = divmod(value, 58) encoded = BASE58_ALPHABET[remainder] + encoded return ("1" * leading_zeroes) + (encoded or ("1" if not data else "")) def load_bridge(): spec = importlib.util.spec_from_file_location("presence_bridge_under_test", BRIDGE_PATH) module = importlib.util.module_from_spec(spec) assert spec.loader is not None spec.loader.exec_module(module) return module class PresenceBridgeOwnerLifecycleTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() def tearDown(self): self.bridge._shutdown.clear() def test_owner_pid_environment_is_strictly_validated(self): with mock.patch.dict(os.environ, {"QORTAL_RETICULUM_OWNER_PID": "4321"}): self.assertEqual(self.bridge._owner_pid_from_environment(), 4321) for invalid in ("", "not-a-pid", "0", "1", "-5"): with self.subTest(invalid=invalid), mock.patch.dict( os.environ, {"QORTAL_RETICULUM_OWNER_PID": invalid}, ): self.assertEqual(self.bridge._owner_pid_from_environment(), 0) @unittest.skipIf(os.name == "nt", "POSIX parent-reparenting behavior") def test_owner_watchdog_forces_exit_after_owner_is_lost(self): owner_pid = 4321 self.bridge._shutdown.clear() close_monitor = mock.Mock() with mock.patch.object( self.bridge, "_open_posix_owner_monitor", return_value=(lambda: False, close_monitor, "test"), ), mock.patch.object(self.bridge.time, "sleep") as sleep_mock, mock.patch.object( self.bridge.os, "_exit", side_effect=SystemExit(0) ) as exit_mock: with self.assertRaises(SystemExit): self.bridge._owner_watchdog_loop(owner_pid) self.assertTrue(self.bridge._shutdown.is_set()) close_monitor.assert_called_once_with() sleep_mock.assert_called_once_with(self.bridge._OWNER_EXIT_GRACE_SECONDS) exit_mock.assert_called_once_with(0) @unittest.skipIf(os.name == "nt", "POSIX parent-reparenting behavior") def test_owner_loss_still_forces_exit_when_stdin_already_started_shutdown(self): owner_pid = 4321 self.bridge._shutdown.set() with mock.patch.object( self.bridge, "_open_posix_owner_monitor", return_value=(lambda: False, lambda: None, "test"), ), mock.patch.object(self.bridge.time, "sleep"), mock.patch.object( self.bridge.os, "_exit", side_effect=SystemExit(0) ) as exit_mock: with self.assertRaises(SystemExit): self.bridge._owner_watchdog_loop(owner_pid) exit_mock.assert_called_once_with(0) @unittest.skipIf(os.name == "nt", "POSIX parent-reparenting behavior") def test_owner_watchdog_leaves_a_live_owner_alone_during_shutdown(self): owner_pid = 4321 self.bridge._shutdown.set() close_monitor = mock.Mock() with mock.patch.object( self.bridge, "_open_posix_owner_monitor", return_value=(lambda: True, close_monitor, "test"), ), mock.patch.object(self.bridge.os, "_exit") as exit_mock: self.bridge._owner_watchdog_loop(owner_pid) close_monitor.assert_called_once_with() exit_mock.assert_not_called() @unittest.skipIf(os.name == "nt", "POSIX process monitor") def test_owner_monitor_accepts_a_live_process_that_is_not_the_direct_parent(self): owner_alive, close_monitor, _kind = self.bridge._open_posix_owner_monitor( os.getpid() ) try: self.assertTrue(owner_alive()) finally: close_monitor() @unittest.skipIf(os.name == "nt", "POSIX process monitor") def test_owner_monitor_detects_the_exact_process_exiting(self): owner = subprocess.Popen( [sys.executable, "-c", "import time; time.sleep(0.1)"] ) owner_alive, close_monitor, _kind = self.bridge._open_posix_owner_monitor( owner.pid ) try: self.assertTrue(owner_alive()) owner.wait(timeout=2.0) deadline = time.monotonic() + 1.0 while owner_alive() and time.monotonic() < deadline: time.sleep(0.01) self.assertFalse(owner_alive()) finally: if owner.poll() is None: owner.terminate() owner.wait(timeout=2.0) close_monitor() @unittest.skipUnless(sys.platform.startswith("linux"), "Linux fallback monitor") def test_linux_owner_monitor_fallback_rejects_a_reused_pid(self): with mock.patch.object( self.bridge.os, "pidfd_open", side_effect=OSError("unsupported") ), mock.patch.object( self.bridge, "_linux_process_start_token", side_effect=["original-start", "original-start", "replacement-start"], ): owner_alive, close_monitor, kind = self.bridge._open_posix_owner_monitor( 4321 ) try: self.assertEqual(kind, "proc-start") self.assertTrue(owner_alive()) self.assertFalse(owner_alive()) finally: close_monitor() def test_macos_owner_monitor_uses_process_exit_events(self): class FakeKqueue: def __init__(self): self.exited = False self.closed = False self.changes = [] def control(self, changes, _max_events, _timeout): if changes is not None: self.changes.extend(changes) return [] return ["exit"] if self.exited else [] def close(self): self.closed = True fake_kqueue = FakeKqueue() with mock.patch.object(self.bridge.sys, "platform", "darwin"), mock.patch.object( self.bridge.select, "kqueue", return_value=fake_kqueue, create=True ), mock.patch.object( self.bridge.select, "kevent", return_value="process-event", create=True ), mock.patch.object( self.bridge.select, "KQ_FILTER_PROC", 1, create=True ), mock.patch.object( self.bridge.select, "KQ_EV_ADD", 2, create=True ), mock.patch.object( self.bridge.select, "KQ_EV_ENABLE", 4, create=True ), mock.patch.object( self.bridge.select, "KQ_EV_CLEAR", 8, create=True ), mock.patch.object( self.bridge.select, "KQ_NOTE_EXIT", 16, create=True ): owner_alive, close_monitor, kind = self.bridge._open_posix_owner_monitor( 4321 ) self.assertEqual(kind, "kqueue") self.assertTrue(owner_alive()) fake_kqueue.exited = True self.assertFalse(owner_alive()) close_monitor() self.assertEqual(fake_kqueue.changes, ["process-event"]) self.assertTrue(fake_kqueue.closed) def test_windows_owner_watchdog_uses_the_exact_process_handle(self): kernel32 = mock.Mock() kernel32.OpenProcess.return_value = 99 kernel32.WaitForSingleObject.return_value = 0 fake_windll = mock.Mock(kernel32=kernel32) self.bridge._shutdown.clear() with mock.patch.object(self.bridge.os, "name", "nt"), mock.patch.object( ctypes, "windll", fake_windll, create=True ), mock.patch.object( self.bridge._shutdown, "wait", return_value=False ), mock.patch.object(self.bridge.time, "sleep"), mock.patch.object( self.bridge.os, "_exit", side_effect=SystemExit(0) ) as exit_mock: with self.assertRaises(SystemExit): self.bridge._owner_watchdog_loop(4321) kernel32.OpenProcess.assert_called_once_with(0x00100000, False, 4321) kernel32.WaitForSingleObject.assert_called_once() kernel32.CloseHandle.assert_called_once_with(99) exit_mock.assert_called_once_with(0) class ReticulumPathVisibilityTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.destination_hash = bytes.fromhex("11" * 16) def test_shared_daemon_path_is_authoritative_when_local_table_misses(self): daemon = mock.Mock() daemon.is_connected_to_shared_instance = True daemon.get_path_snapshot.return_value = { "hops": 2, "timestamp": time.time(), } self.bridge._reticulum = daemon with mock.patch.object(RNS.Transport, "has_path", return_value=False): self.assertTrue(self.bridge._reticulum_has_path(self.destination_hash)) self.assertTrue(self.bridge._reticulum_has_path(self.destination_hash)) daemon.get_path_snapshot.assert_called_once_with(self.destination_hash) def test_embedded_instance_uses_local_path_without_rpc(self): reticulum = mock.Mock() reticulum.is_connected_to_shared_instance = False self.bridge._reticulum = reticulum with mock.patch.object(RNS.Transport, "has_path", return_value=False): self.assertFalse(self.bridge._reticulum_has_path(self.destination_hash)) reticulum.get_path_snapshot.assert_not_called() def test_dropping_path_invalidates_shared_availability_cache(self): daemon = mock.Mock() daemon.is_connected_to_shared_instance = True daemon.get_path_snapshot.side_effect = [ {"hops": 2, "timestamp": time.time()}, None, ] daemon.drop_path.return_value = True self.bridge._reticulum = daemon with mock.patch.object(RNS.Transport, "has_path", return_value=False), mock.patch.object( RNS.Transport, "expire_path", return_value=False ), mock.patch.object(RNS.Transport, "mark_path_unresponsive"): self.assertTrue(self.bridge._reticulum_has_path(self.destination_hash)) self.assertTrue(self.bridge._drop_reticulum_path(self.destination_hash)) self.assertFalse(self.bridge._reticulum_has_path(self.destination_hash)) self.assertEqual(daemon.get_path_snapshot.call_count, 2) def test_game_discovery_miss_nudges_without_dropping_path(self): peer_hash = self.destination_hash.hex() with mock.patch.object( self.bridge, "_nudge_cached_reticulum_path", return_value=True, ) as nudge, mock.patch.object( self.bridge, "_force_overlay_peer_path_refresh", ) as force_refresh: self.assertTrue( self.bridge._refresh_qortalland_game_path( peer_hash, "game_link_no_path", ) ) nudge.assert_called_once() force_refresh.assert_not_called() def test_game_failed_link_can_still_replace_bad_path(self): peer_hash = self.destination_hash.hex() with mock.patch.object( self.bridge, "_nudge_cached_reticulum_path", ) as nudge, mock.patch.object( self.bridge, "_force_overlay_peer_path_refresh", return_value=True, ) as force_refresh: self.assertTrue( self.bridge._refresh_qortalland_game_path( peer_hash, "game_link_attempt_closed", ) ) force_refresh.assert_called_once_with( peer_hash, target="qortalland-game", reason="game_link_attempt_closed", await_seconds=0.0, ) nudge.assert_not_called() def test_recent_media_success_avoids_route_rpc(self): peer_hash = self.destination_hash.hex() state = self.bridge._get_call_media_state(peer_hash) state.update({ "destination_hash_hex": peer_hash, "path_state": "fresh", "last_send_ok": time.time(), "last_send_fail": None, "last_inbound_at": None, }) with mock.patch.object( self.bridge, "_reticulum_has_path", side_effect=AssertionError("route lookup should not run for recent traffic"), ): self.assertEqual( self.bridge._classify_call_media_path_state( peer_hash, self.destination_hash, ), "fresh", ) class FakeLink: def __init__(self): self.closed_callback = None self.packet_callback = None self.remote_identified_callback = None self.resource_strategy = None self.resource_callback = None self.resource_started_callback = None self.resource_concluded_callback = None self.teardown_called = False self.remote_identity = None def get_mdu(self): return 4096 def set_link_closed_callback(self, callback): self.closed_callback = callback def set_packet_callback(self, callback): self.packet_callback = callback def set_remote_identified_callback(self, callback): self.remote_identified_callback = callback def get_remote_identity(self): return self.remote_identity def set_resource_strategy(self, strategy): self.resource_strategy = strategy def set_resource_callback(self, callback): self.resource_callback = callback def set_resource_started_callback(self, callback): self.resource_started_callback = callback def set_resource_concluded_callback(self, callback): self.resource_concluded_callback = callback def teardown(self): self.teardown_called = True class FakeSessionReceipt: FAILED = 0 SENT = 1 def __init__(self): self.progress = 0.0 self.metadata = None self.response = None self.cancelled = False self.status = self.SENT self.concluded_at = None self.resource = None self.link = None self.request_id = None def get_progress(self): return self.progress def get_response(self): return self.response def cancel(self): self.cancelled = True class FakeSessionLink(FakeLink): def __init__(self, link_id=None): super().__init__() self.link_id = link_id or bytes.fromhex("77" * 16) self.requests = [] self.identified_with = None self.remote_identity = None self.outgoing_resources = [] def identify(self, identity): self.identified_with = identity def get_remote_identity(self): return self.remote_identity def request(self, path, data=None, **callbacks): receipt = FakeSessionReceipt() self.requests.append((path, data, callbacks, receipt)) return receipt class FakePacket: def __init__(self, link): self.link = link class FakeRnsPacket: MDU = 500 ENCRYPTED_MDU = 500 sent_links = [] sent_payloads = [] def __init__(self, link, data, create_receipt=False): self.link = link self.data = data self.create_receipt = create_receipt def send(self): self.__class__.sent_links.append(self.link) self.__class__.sent_payloads.append(self.data) return True class FakeDestination: def __init__(self): self.hash = bytes.fromhex("44" * 16) class PresenceBridgeWireEncodingTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.bridge._destination = FakeDestination() def test_group_signal_stamps_transport_sender_without_overwriting_payload(self): read_state = { "y": "d", "q": "Q-peer", "u": 123, "n": 124, "p": "public-key", "z": "signature", } encoded = self.bridge._encode_group_signal_wire( {"t": "RCHAT", "k": "read_sync", "w": read_state} ) self.assertTrue(encoded["ok"]) wire = json.loads(encoded["wire_bytes"].decode("utf-8")) self.assertEqual(wire["r"], "44" * 16) self.assertEqual(wire["w"], read_state) def test_group_signal_relay_preserves_original_sender(self): encoded = self.bridge._encode_group_signal_wire( {"t": "GJ", "r": "11" * 16} ) self.assertTrue(encoded["ok"]) wire = json.loads(encoded["wire_bytes"].decode("utf-8")) self.assertEqual(wire["r"], "11" * 16) def test_call_signal_relay_preserves_original_sender(self): encoded = self.bridge._encode_call_signal_wire( {"t": "CR", "r": "22" * 16} ) self.assertTrue(encoded["ok"]) wire = json.loads(encoded["wire_bytes"].decode("utf-8")) self.assertEqual(wire["r"], "22" * 16) def test_route_bound_presence_fits_encrypted_mdu_direct_and_relayed(self): origin_raw = bytes.fromhex("55" * 16) origin_route = base64.urlsafe_b64encode(origin_raw).decode("ascii").rstrip("=") local_route = base64.urlsafe_b64encode(bytes.fromhex("44" * 16)).decode( "ascii" ).rstrip("=") envelope = { "type": "PRESENCE_ANNOUNCE", "id": "x" * 16, "timestamp": 9_999_999_999_999, "signature": "z" * 88, "payload": { "address": "Q" + ("a" * 33), "publicKey": "k" * 44, "sessionId": "P" + local_route + ("e" * 13), "status": "online", "clientVersion": "1.0.0", }, } direct = self.bridge.make_presence_wire(envelope, 4) relayed_envelope = { **envelope, "payload": { **envelope["payload"], "sessionId": "P" + origin_route + ("e" * 13), }, } relayed = self.bridge.make_presence_wire( relayed_envelope, 3, origin_sender_hash="55" * 16, ) self.assertLessEqual(len(direct), self.bridge._MAX_ENCRYPTED_WIRE_BYTES) self.assertLessEqual(len(relayed), self.bridge._MAX_ENCRYPTED_WIRE_BYTES) self.assertNotIn("o", json.loads(relayed.decode("utf-8"))) def test_oversized_legacy_presence_is_rejected_before_packet_send(self): envelope = { "type": "PRESENCE_ANNOUNCE", "id": "x" * 36, "timestamp": 9_999_999_999_999, "signature": "z" * 88, "payload": { "address": "Q" + ("a" * 33), "publicKey": "k" * 44, "sessionId": "legacy-session-id".ljust(36, "x"), "status": "online", "clientVersion": "1.0.0", }, } with self.assertRaisesRegex(RuntimeError, "exceeds encrypted MDU"): self.bridge.make_presence_wire(envelope, 4) def test_stale_route_bound_presence_is_not_republished_after_route_change(self): stale_route = base64.urlsafe_b64encode(bytes.fromhex("55" * 16)).decode( "ascii" ).rstrip("=") envelope = { "type": "PRESENCE_HEARTBEAT", "id": "x" * 16, "timestamp": 9_999_999_999_999, "signature": "z" * 88, "payload": { "address": "Q" + ("a" * 33), "publicKey": "k" * 44, "sessionId": "P" + stale_route + ("e" * 13), "status": "online", }, } with self.assertRaisesRegex(RuntimeError, "local destination"): self.bridge.make_presence_wire(envelope, 4) def test_presence_route_binding_decoder_rejects_noncanonical_ids(self): route = base64.urlsafe_b64encode(bytes.fromhex("55" * 16)).decode("ascii").rstrip("=") session_id = "P" + route + ("e" * 13) self.assertEqual( self.bridge._presence_route_bound_destination_hash(session_id), "55" * 16, ) self.assertIsNone( self.bridge._presence_route_bound_destination_hash("P" + ("!" * 35)) ) def test_route_bound_presence_recovers_signed_origin_from_relay(self): origin_raw = bytes.fromhex("55" * 16) route = base64.urlsafe_b64encode(origin_raw).decode("ascii").rstrip("=") emitted = [] message = { "t": "PRESENCE_HEARTBEAT", "i": "x" * 16, "a": "Q" + ("a" * 33), "k": "k" * 44, "n": "P" + route + ("e" * 13), "m": 9_999_999_999_999, "g": "z" * 88, "r": "44" * 16, "s": "online", "q": 3, } with mock.patch.object( self.bridge, "emit_event", side_effect=lambda event, payload: emitted.append((event, payload)), ): accepted = self.bridge._emit_presence_message(message, "relay-link") self.assertTrue(accepted) self.assertEqual(emitted[0][0], "presence_message") self.assertEqual( emitted[0][1]["route"], { "kind": "reticulum", "destinationHash": "55" * 16, "viaDestinationHash": "44" * 16, "overlayHopsRemaining": 3, "linkId": "relay-link", }, ) class PresenceBridgeReticulumChatInboundDedupTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() def test_identity_request_dedup_ignores_route_fields(self): request = { "t": "RCHAT", "k": "identity_req", "d": "dd" * 16, "rid": "11" * 12, "h": 0, "m": 5, "x": int((time.time() + 30) * 1000), } self.assertFalse( self.bridge._should_drop_duplicate_reticulum_chat_inbound(request) ) self.assertTrue( self.bridge._should_drop_duplicate_reticulum_chat_inbound( {**request, "h": 3, "r": "aa" * 16} ) ) def test_typing_dedup_ignores_origin_hops_and_ingress_sender(self): typing = { "t": "RCHAT", "k": "typing", "g": 73, "c": "general", "a": "Qsender", "ts": 123_456, "active": True, "o": "aa" * 16, "h": 1, "r": "bb" * 16, } self.assertFalse( self.bridge._should_drop_duplicate_reticulum_chat_inbound(typing) ) self.assertTrue( self.bridge._should_drop_duplicate_reticulum_chat_inbound( {**typing, "o": "cc" * 16, "h": 4, "r": "dd" * 16} ) ) self.assertFalse( self.bridge._should_drop_duplicate_reticulum_chat_inbound( {**typing, "ts": 123_457} ) ) class PresenceBridgeLandStateFastPathTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.bridge._destination = FakeDestination() self.group_id = 73 self.author = "Q-land-author" self.session_id = "land-session" self.source_hash = "11" * 16 self.target_hash = "22" * 16 self.origin = "AQIDBAUGBwgJCgsMDQ4PEA" self.private_key = RNS.Cryptography.Ed25519PrivateKey.generate() public_key = self.private_key.public_key().public_bytes() expires_at = int((time.time() + 60) * 1000) self.bridge._configure_land_state_forwarding( [ { "groupId": self.group_id, "targets": [ { "peerPresenceHash": self.target_hash, "expiresAt": expires_at, } ], } ], [ { "groupId": self.group_id, "authorAddress": self.author, "sessionId": self.session_id, "ephemeralPublicKey": base58_encode(public_key), "expiresAt": expires_at, } ], 7, ) def state(self, sequence=1, signature=None): timestamp = int(time.time() * 1000) fields = { "afk": True, "authorAddress": self.author, "direction": "r", "dnd": True, "groupId": self.group_id, "movement": "walk", "roomId": "room", "sequence": sequence, "sessionId": self.session_id, "skinId": 4, "timestamp": timestamp, "type": "QORTAL_LAND_STATE", "x": 50, "y": 60, } signed_bytes = json.dumps( fields, ensure_ascii=False, separators=(",", ":"), sort_keys=True, ).encode("utf-8") encoded_signature = signature if encoded_signature is None: encoded_signature = base58_encode( self.private_key.sign(signed_bytes) ) return { "t": "RCHAT", "k": "land_state", "g": self.group_id, "a": self.author, "s": self.session_id, "q": sequence, "x": 50, "y": 60, "u": "room", "d": "r", "m": "walk", "v": 3, "i": 4, "ts": timestamp, "z": encoded_signature, "o": self.origin, "h": 1, } def queue_without_scheduler(self, message): with mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True, ): queued = self.bridge._queue_land_state_fast_path( message, self.source_hash, "source-link", ) self.assertTrue(queued) return next(reversed(self.bridge._land_state_forward_pending)) def test_preserves_full_origin_and_reports_the_forwarding_revision(self): message = self.state() pending_key = self.queue_without_scheduler(message) emitted = [] sent = [] with mock.patch.object( self.bridge, "_send_wire_to_established_overlay_peer", side_effect=lambda peer, wire, _traffic: sent.append((peer, wire)) or True, ), mock.patch.object( self.bridge, "_emit_call_bridge_message", side_effect=lambda *args, **kwargs: emitted.append((args, kwargs)) or True, ): self.bridge._process_land_state_fast_path(pending_key) self.assertEqual(len(sent), 1) forwarded = json.loads(sent[0][1].decode("utf-8")) self.assertEqual(forwarded["o"], self.origin) self.assertEqual(forwarded["v"], 3) self.assertEqual(forwarded["i"], 4) self.assertTrue(emitted[0][1]["land_state_fast_forwarded"]) self.assertEqual(emitted[0][1]["land_state_forwarding_revision"], 7) def test_destination_hash_matching_requires_exact_hash_equivalence(self): full_hash = "0102030405060708090a0b0c0d0e0f10" same_prefix = "01020304" + ("ff" * 12) self.assertTrue( self.bridge._land_state_hash_matches(self.origin, full_hash) ) self.assertFalse( self.bridge._land_state_hash_matches(full_hash, same_prefix) ) def test_invalid_origin_is_replaced_with_the_verified_ingress_peer(self): message = self.state() message["o"] = "invalid" pending_key = self.queue_without_scheduler(message) sent = [] with mock.patch.object( self.bridge, "_send_wire_to_established_overlay_peer", side_effect=lambda peer, wire, _traffic: sent.append((peer, wire)) or True, ), mock.patch.object( self.bridge, "_emit_call_bridge_message", return_value=True, ): self.bridge._process_land_state_fast_path(pending_key) self.assertEqual(len(sent), 1) forwarded = json.loads(sent[0][1].decode("utf-8")) self.assertEqual(forwarded["o"], self.source_hash) def test_route_revision_change_forces_electron_fallback(self): pending_key = self.queue_without_scheduler(self.state()) emitted = [] def send_and_change_revision(_peer, _wire, _traffic): self.bridge._land_state_forwarding_revision = 8 return True with mock.patch.object( self.bridge, "_send_wire_to_established_overlay_peer", side_effect=send_and_change_revision, ), mock.patch.object( self.bridge, "_emit_call_bridge_message", side_effect=lambda *args, **kwargs: emitted.append((args, kwargs)) or True, ): self.bridge._process_land_state_fast_path(pending_key) self.assertFalse(emitted[0][1]["land_state_fast_forwarded"]) self.assertIsNone(emitted[0][1]["land_state_forwarding_revision"]) def test_unverified_higher_sequence_does_not_replace_pending_valid_state(self): valid_key = self.queue_without_scheduler(self.state(sequence=1)) invalid_key = self.queue_without_scheduler( self.state(sequence=999, signature="invalid") ) self.assertEqual(len(self.bridge._land_state_forward_pending), 2) emitted = [] with mock.patch.object( self.bridge, "_send_wire_to_established_overlay_peer", return_value=True, ) as send, mock.patch.object( self.bridge, "_emit_call_bridge_message", side_effect=lambda *args, **kwargs: emitted.append((args, kwargs)) or True, ): self.bridge._process_land_state_fast_path(valid_key) self.bridge._process_land_state_fast_path(invalid_key) self.assertEqual(send.call_count, 1) self.assertEqual(len(emitted), 1) def test_oversized_target_plan_falls_back_instead_of_installing_part_of_it(self): expires_at = int((time.time() + 60) * 1000) targets = [ { "peerPresenceHash": f"{index:032x}", "expiresAt": expires_at, } for index in range( 1, self.bridge._LAND_STATE_FORWARDING_MAX_TARGETS_PER_GROUP + 2, ) ] self.bridge._configure_land_state_forwarding( [{"groupId": self.group_id, "targets": targets}], [], 8, ) self.assertNotIn( self.group_id, self.bridge._land_state_forwarding_plans, ) def test_sequence_zero_is_forwarded_only_once(self): message = self.state(sequence=0) emitted = [] with mock.patch.object( self.bridge, "_send_wire_to_established_overlay_peer", return_value=True, ) as send, mock.patch.object( self.bridge, "_emit_call_bridge_message", side_effect=lambda *args, **kwargs: emitted.append((args, kwargs)) or True, ): self.bridge._process_land_state_fast_path( self.queue_without_scheduler(message) ) self.bridge._process_land_state_fast_path( self.queue_without_scheduler(message) ) self.assertEqual(send.call_count, 1) self.assertEqual(len(emitted), 1) class PresenceBridgeAudioForwardFastPathTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.room_id = "gcall-qortal-716" self.source_hash = "11" * 16 self.target_hash = "22" * 16 def plan(self, ingress_link="source-link", target_link="target-link"): return { "roomId": self.room_id, "topologyEpoch": 7, "rules": [ { "sourceAddress": "Q-source", "ingress": { "address": "Q-source", "transport": "link", "linkId": ingress_link, "peerPresenceHash": self.source_hash, "peerDestinationHash": self.source_hash, }, "targets": [ { "address": "Q-target", "transport": "link", "linkId": target_link, "peerPresenceHash": self.target_hash, "peerDestinationHash": self.target_hash, } ], } ], } def test_exact_verified_ingress_forwards_unchanged_media(self): rooms, rules = self.bridge._configure_audio_forwarding_plans([self.plan()]) self.assertEqual((rooms, rules), (1, 1)) captured = [] with mock.patch.object( self.bridge, "_audio_data_plane_broadcast_inbound_audio", return_value=True, ), mock.patch.object( self.bridge, "_put_audio_decoded_batch_keep_newest", side_effect=lambda frames: captured.extend(frames) or True, ): handled = self.bridge._try_group_audio_forward_fast_path( self.room_id, "source-link", self.source_hash, self.source_hash, 1234, b"encrypted-media", b"inbound-batch", ) self.assertTrue(handled) self.assertEqual(len(captured), 1) self.assertEqual(captured[0][0], "target-link") self.assertEqual(captured[0][1], self.room_id) self.assertEqual(captured[0][5], b"encrypted-media") def test_link_media_does_not_fall_back_to_hash_matching(self): self.bridge._configure_audio_forwarding_plans([self.plan()]) with mock.patch.object( self.bridge, "_audio_data_plane_broadcast_inbound_audio", ) as broadcast: handled = self.bridge._try_group_audio_forward_fast_path( self.room_id, "unverified-link", self.source_hash, self.source_hash, 1234, b"encrypted-media", b"inbound-batch", ) self.assertFalse(handled) broadcast.assert_not_called() def test_packet_media_matches_only_the_configured_peer_hash(self): plan = self.plan() ingress = plan["rules"][0]["ingress"] ingress["transport"] = "packet" ingress["linkId"] = "" self.bridge._configure_audio_forwarding_plans([plan]) captured = [] with mock.patch.object( self.bridge, "_audio_data_plane_broadcast_inbound_audio", return_value=True, ), mock.patch.object( self.bridge, "_put_audio_decoded_batch_keep_newest", side_effect=lambda frames: captured.extend(frames) or True, ): handled = self.bridge._try_group_audio_forward_fast_path( self.room_id, "", self.source_hash, "", 1234, b"encrypted-media", b"inbound-batch", ) rejected = self.bridge._try_group_audio_forward_fast_path( self.room_id, "", "33" * 16, "", 1235, b"other-media", b"other-batch", ) self.assertTrue(handled) self.assertFalse(rejected) self.assertEqual(len(captured), 1) def test_forward_queue_rejection_does_not_duplicate_local_delivery(self): self.bridge._configure_audio_forwarding_plans([self.plan()]) with mock.patch.object( self.bridge, "_audio_data_plane_broadcast_inbound_audio", return_value=True, ), mock.patch.object( self.bridge, "_put_audio_decoded_batch_keep_newest", return_value=False, ): handled = self.bridge._try_group_audio_forward_fast_path( self.room_id, "source-link", self.source_hash, self.source_hash, 1234, b"encrypted-media", b"inbound-batch", ) self.assertTrue(handled) def test_plan_replacement_removes_stale_room_and_loop_target(self): loop_plan = self.plan(target_link="source-link") rooms, rules = self.bridge._configure_audio_forwarding_plans([loop_plan]) self.assertEqual((rooms, rules), (1, 1)) stored_rule = self.bridge._audio_forwarding_plans_by_room[self.room_id][ "rules" ][0] self.assertEqual(stored_rule["targets"], []) rooms, rules = self.bridge._configure_audio_forwarding_plans([]) self.assertEqual((rooms, rules), (0, 0)) self.assertNotIn(self.room_id, self.bridge._audio_forwarding_plans_by_room) class PresenceBridgeOverlayAudioPromotionTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.sender_peer_hash = "22" * 16 self.original_rns_packet = RNS.Packet def tearDown(self): RNS.Packet = self.original_rns_packet def drain_audio_queue(self): while True: try: self.bridge._audio_binary_out_queue.get_nowait() except queue.Empty: return def group_audio_wire(self): room = b"gcall-qortal-1" sender_hash = bytes.fromhex(self.sender_peer_hash) payload = b"opus" return ( self.bridge._GROUP_AUDIO_BINARY_MAGIC + bytes( ( self.bridge._GROUP_AUDIO_BINARY_VERSION, len(room), len(sender_hash), ) ) + len(payload).to_bytes(2, "big") + room + sender_hash + payload ) def group_audio_heartbeat_wire(self): return self.bridge.json.dumps( { "t": self.bridge._GROUP_AUDIO_HEARTBEAT_WIRE_TYPE, "R": "gcall-qortal-1", "c": "PING", "m": int(time.time() * 1000), "r": self.sender_peer_hash, } ).encode("utf-8") def group_audio_heartbeat_wire_without_sender(self): return self.bridge.json.dumps( { "t": self.bridge._GROUP_AUDIO_HEARTBEAT_WIRE_TYPE, "R": "gcall-qortal-1", "c": "PING", "m": int(time.time() * 1000), } ).encode("utf-8") def qchat_file_auth_wire(self, transfer_id="transfer-1", peer_hash=None): return self.bridge.json.dumps( { "type": "QCHAT_FILE_LINK_AUTH", "transferId": transfer_id, "senderAddress": "Q-sender", "downloaderAddress": "Q-downloader", "downloaderPublicKey": "pub-downloader", "downloaderReticulumDestinationHash": peer_hash or self.sender_peer_hash, "downloaderReticulumIdentityPublicKeyBase64": "identity", "timestamp": int(time.time() * 1000), "signature": "sig", } ).encode("utf-8") def install_overlay_state(self, incoming=True): link = FakeLink() link_id = "overlay-test-link" peer_hash = "11" * 16 now = time.time() self.bridge._overlay_links_by_id[link_id] = { "link": link, "peerPresenceHash": peer_hash, "incoming": incoming, "established": True, "established_at": now, "created_at": now, "pending_packets": self.bridge.deque(maxlen=4), "last_activity_at": now, "last_rx_at": None, } self.bridge._overlay_link_ids_by_object[id(link)] = link_id self.bridge._active_overlay_link_id_by_peer_hash[peer_hash] = link_id if incoming: self.bridge._inbound_overlay_neighbors[peer_hash] = now else: self.bridge._active_overlay_neighbors[peer_hash] = now return link, link_id, peer_hash def install_audio_state( self, link_id, peer_hash=None, established=True, link=None, last_activity_at=None, ): peer_hash = peer_hash or self.sender_peer_hash link = link or FakeLink() now = time.time() self.bridge._audio_links_by_id[link_id] = { "link": link, "peerPresenceHash": peer_hash, "peerDestinationHash": peer_hash, "incoming": False, "established": established, "established_at": now if established else None, "created_at": now - 10, "last_activity_at": last_activity_at if last_activity_at is not None else now, "last_rx_at": None, "last_send_ok_at": None, "send_lock": self.bridge.threading.RLock(), "generation": 0, "closing": False, } self.bridge._audio_link_ids_by_object[id(link)] = link_id return link def drain_json_events(self): events = [] for event_queue in ( self.bridge._json_priority_event_queue, self.bridge._json_event_queue, ): while True: try: frame = event_queue.get_nowait() if frame is not None: events.append(frame) except queue.Empty: break while True: frame = self.bridge._pop_coalesced_json_event_line() if frame is None: return events events.append(frame) def drain_json_responses(self): responses = [] while True: try: responses.append(self.bridge._json_resp_queue.get_nowait()) except queue.Empty: return responses def install_fake_rns_packet(self): FakeRnsPacket.sent_links = [] FakeRnsPacket.sent_payloads = [] RNS.Packet = FakeRnsPacket self.bridge.RNS.Packet = FakeRnsPacket self.bridge._destination = FakeDestination() def test_incoming_overlay_group_audio_promotes_link_without_teardown(self): self.drain_audio_queue() link, overlay_link_id, overlay_peer_hash = self.install_overlay_state( incoming=True ) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, } self.bridge.on_overlay_link_packet(self.group_audio_wire(), packet) self.assertNotIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertNotIn(id(link), self.bridge._overlay_link_ids_by_object) self.assertNotIn( overlay_peer_hash, self.bridge._active_overlay_link_id_by_peer_hash, ) self.assertNotIn(overlay_peer_hash, self.bridge._inbound_overlay_neighbors) self.assertFalse(link.teardown_called) audio_link_id = self.bridge.get_audio_link_id(link) self.assertIsInstance(audio_link_id, str) audio_state = self.bridge.get_audio_link_state(audio_link_id) self.assertIsNotNone(audio_state) self.assertTrue(audio_state["incoming"]) self.assertEqual(audio_state["peerPresenceHash"], self.sender_peer_hash) self.assertEqual(audio_state["peerDestinationHash"], self.sender_peer_hash) self.assertEqual(audio_state["promoted_from_overlay_link_id"], overlay_link_id) self.assertIs(link.packet_callback, self.bridge.on_audio_link_packet) self.assertGreater(self.bridge._audio_binary_out_queue.qsize(), 0) def test_incoming_overlay_group_audio_without_desired_audio_is_not_promoted(self): self.drain_audio_queue() link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge.on_overlay_link_packet(self.group_audio_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_audio_link_id(link)) self.assertFalse(link.teardown_called) self.assertEqual(self.bridge._audio_binary_out_queue.qsize(), 0) def test_incoming_overlay_gac_promotes_link_when_audio_is_desired(self): self.drain_audio_queue() link, overlay_link_id, _overlay_peer_hash = self.install_overlay_state( incoming=True ) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, } self.bridge.on_overlay_link_packet(self.group_audio_heartbeat_wire(), packet) self.assertNotIn(overlay_link_id, self.bridge._overlay_links_by_id) audio_link_id = self.bridge.get_audio_link_id(link) self.assertIsInstance(audio_link_id, str) audio_state = self.bridge.get_audio_link_state(audio_link_id) self.assertIsNotNone(audio_state) self.assertTrue(audio_state["incoming"]) self.assertEqual(audio_state["peerPresenceHash"], self.sender_peer_hash) self.assertEqual(audio_state["peerDestinationHash"], self.sender_peer_hash) self.assertIs(link.packet_callback, self.bridge.on_audio_link_packet) def test_incoming_overlay_gac_without_desired_audio_is_not_promoted(self): self.drain_audio_queue() link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge.on_overlay_link_packet(self.group_audio_heartbeat_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_audio_link_id(link)) self.assertFalse(link.teardown_called) def test_incoming_overlay_gac_without_sender_is_not_promoted(self): self.drain_audio_queue() link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, } self.bridge.on_overlay_link_packet( self.group_audio_heartbeat_wire_without_sender(), packet, ) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_audio_link_id(link)) self.assertFalse(link.teardown_called) def test_stale_audio_mapping_does_not_allow_overlay_promotion(self): self.drain_audio_queue() link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge._active_audio_link_id_by_peer_hash[self.sender_peer_hash] = "stale-link" self.bridge.on_overlay_link_packet(self.group_audio_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_audio_link_id(link)) self.assertFalse(link.teardown_called) self.assertEqual(self.bridge._audio_binary_out_queue.qsize(), 0) def test_audio_send_with_stale_link_id_uses_established_peer_link(self): self.install_fake_rns_packet() current_link = self.install_audio_state("current-audio-link") self.bridge._active_audio_link_id_by_peer_hash[self.sender_peer_hash] = "stale-audio-link" self.bridge._outgoing_audio_link_id_by_peer_hash[self.sender_peer_hash] = "stale-audio-link" self.bridge._process_audio_batch( [ ( "stale-audio-link", "gcall-qortal-1", self.sender_peer_hash, "", int(time.time() * 1000), b"opus", ) ] ) self.assertEqual(FakeRnsPacket.sent_links, [current_link]) self.assertEqual( self.bridge._active_audio_link_id_by_peer_hash[self.sender_peer_hash], "current-audio-link", ) failures = [ frame for frame in self.drain_json_events() if frame.get("event") == "group_audio_send_failed" ] self.assertEqual(failures, []) def test_audio_heartbeat_with_stale_link_id_uses_established_peer_link(self): self.install_fake_rns_packet() current_link = self.install_audio_state("current-audio-link") self.bridge._active_audio_link_id_by_peer_hash[self.sender_peer_hash] = "stale-audio-link" self.bridge._outgoing_audio_link_id_by_peer_hash[self.sender_peer_hash] = "stale-audio-link" self.bridge.handle_send_group_audio_link_heartbeat( "req-1", { "linkId": "stale-audio-link", "peerPresenceHash": self.sender_peer_hash, "roomId": "gcall-qortal-1", "command": "PING", }, ) self.assertEqual(FakeRnsPacket.sent_links, [current_link]) responses = self.drain_json_responses() self.assertEqual(len(responses), 1) self.assertTrue(responses[0].get("ok")) self.assertEqual( responses[0].get("payload", {}).get("linkId"), "current-audio-link", ) def test_audio_rtt_probe_uses_established_audio_link(self): self.install_fake_rns_packet() current_link = self.install_audio_state("current-audio-link") state = self.bridge.get_audio_link_state("current-audio-link") self.bridge._ensure_audio_link_lifecycle_fields(state) self.bridge._process_audio_rtt_probe("current-audio-link", 0) self.assertEqual(FakeRnsPacket.sent_links, [current_link]) wire = json.loads(FakeRnsPacket.sent_payloads[0].decode("utf-8")) self.assertEqual(wire.get("t"), self.bridge._GROUP_AUDIO_RTT_WIRE_TYPE) self.assertEqual(wire.get("c"), self.bridge._GROUP_AUDIO_RTT_PROBE_COMMAND) self.assertRegex(str(wire.get("q") or ""), r"^[0-9a-f]{16}$") self.assertIn(wire["q"], state["rtt_pending"]) def test_audio_rtt_ack_resolves_probe_with_monotonic_clock(self): self.install_audio_state("current-audio-link") state = self.bridge.get_audio_link_state("current-audio-link") self.bridge._ensure_audio_link_lifecycle_fields(state) probe_id = "a1" * 8 state["rtt_pending"][probe_id] = {"sent_ns": 1_000_000_000} with mock.patch.object( self.bridge.time, "monotonic_ns", return_value=1_025_000_000, ): rtt_ms = self.bridge._resolve_audio_rtt_probe( "current-audio-link", state, probe_id, ) self.assertEqual(rtt_ms, 25.0) self.assertEqual(state.get("rtt_latest_ms"), 25.0) self.assertEqual(state.get("rtt_median_ms"), 25.0) self.assertNotIn(probe_id, state["rtt_pending"]) def test_audio_rtt_probe_is_acknowledged_inside_bridge(self): self.install_fake_rns_packet() current_link = self.install_audio_state("current-audio-link") probe_id = "b2" * 8 wire = json.dumps( { "t": self.bridge._GROUP_AUDIO_RTT_WIRE_TYPE, "c": self.bridge._GROUP_AUDIO_RTT_PROBE_COMMAND, "q": probe_id, "r": self.sender_peer_hash, } ).encode("utf-8") def run_immediately(_lane, _name, func, *args, **kwargs): func(*args, **kwargs) return True with mock.patch.object( self.bridge, "_enqueue_scheduler_task", side_effect=run_immediately, ): self.bridge.on_audio_link_packet(wire, FakePacket(current_link)) self.assertEqual(FakeRnsPacket.sent_links, [current_link]) ack = json.loads(FakeRnsPacket.sent_payloads[0].decode("utf-8")) self.assertEqual(ack.get("t"), self.bridge._GROUP_AUDIO_RTT_WIRE_TYPE) self.assertEqual(ack.get("c"), self.bridge._GROUP_AUDIO_RTT_ACK_COMMAND) self.assertEqual(ack.get("q"), probe_id) def test_audio_rtt_commands_do_not_intercept_call_heartbeat(self): current_link = self.install_audio_state("current-audio-link") heartbeat = json.dumps( { "t": self.bridge._GROUP_AUDIO_HEARTBEAT_WIRE_TYPE, "R": "gcall-qortal-1", "c": "PING", "m": int(time.time() * 1000), "r": self.sender_peer_hash, } ).encode("utf-8") with mock.patch.object( self.bridge, "_emit_call_bridge_message", return_value=True, ) as emit_call_message: self.bridge.on_audio_link_packet(heartbeat, FakePacket(current_link)) emit_call_message.assert_called_once() def test_removing_audio_link_discards_pending_rtt_probe(self): self.install_audio_state("current-audio-link") state = self.bridge.get_audio_link_state("current-audio-link") self.bridge._ensure_audio_link_lifecycle_fields(state) state["rtt_probe_queued"] = True state["rtt_pending"]["c3" * 8] = {"sent_ns": time.monotonic_ns()} removed = self.bridge.remove_audio_link("current-audio-link") self.assertIsNotNone(removed) self.assertFalse(removed.get("rtt_probe_queued")) self.assertEqual(removed.get("rtt_pending"), {}) def test_audio_open_stops_after_max_establish_attempts(self): self.bridge._destination = FakeDestination() self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, "attempts": self.bridge._AUDIO_LINK_MAX_ESTABLISH_ATTEMPTS, "retry_delay": self.bridge._AUDIO_LINK_RETRY_MIN_SECONDS, "retry_timer": None, "last_failure_reason": "establish_timeout", } ok, payload, error = self.bridge._open_group_audio_link_for_peer( self.sender_peer_hash, retry_reason="establish_timeout", ) self.assertFalse(ok) self.assertEqual(payload.get("code"), "max_establish_attempts") self.assertEqual(error, "Max group audio link establish attempts reached") events = [ frame for frame in self.drain_json_events() if frame.get("event") == "group_audio_send_failed" ] self.assertEqual(len(events), 1) self.assertEqual(events[0].get("payload", {}).get("reason"), "max_establish_attempts") self.assertEqual(events[0].get("payload", {}).get("code"), "max_establish_attempts") def test_audio_retry_not_scheduled_after_max_establish_attempts(self): self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, "attempts": self.bridge._AUDIO_LINK_MAX_ESTABLISH_ATTEMPTS, "retry_delay": self.bridge._AUDIO_LINK_RETRY_MIN_SECONDS, "retry_timer": None, "last_failure_reason": "establish_timeout", } self.bridge._schedule_audio_link_retry(self.sender_peer_hash, "establish_timeout") self.bridge._schedule_audio_link_retry(self.sender_peer_hash, "establish_timeout") desired = self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] self.assertIsNone(desired.get("retry_timer")) self.assertEqual(desired.get("last_failure_reason"), "max_establish_attempts") events = [ frame for frame in self.drain_json_events() if frame.get("event") == "group_audio_send_failed" ] self.assertEqual(len(events), 1) self.assertEqual(events[0].get("payload", {}).get("reason"), "max_establish_attempts") self.assertEqual(events[0].get("payload", {}).get("code"), "max_establish_attempts") def test_audio_retry_timer_callback_does_not_enqueue_after_max_attempts(self): enqueued = [] original_enqueue = self.bridge._enqueue_scheduler_task original_timer = self.bridge.threading.Timer class FakeTimer: def __init__(self, delay, function): self.delay = delay self.function = function self.daemon = False self.started = False def start(self): self.started = True def cancel(self): self.started = False try: self.bridge._enqueue_scheduler_task = lambda *args, **kwargs: enqueued.append( (args, kwargs) ) self.bridge.threading.Timer = FakeTimer self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, "attempts": self.bridge._AUDIO_LINK_MAX_ESTABLISH_ATTEMPTS - 1, "retry_delay": self.bridge._AUDIO_LINK_RETRY_MIN_SECONDS, "retry_timer": None, "last_failure_reason": "establish_timeout", } self.bridge._schedule_audio_link_retry( self.sender_peer_hash, "establish_timeout", immediate=True, ) desired = self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] timer = desired.get("retry_timer") self.assertIsNotNone(timer) desired["attempts"] = self.bridge._AUDIO_LINK_MAX_ESTABLISH_ATTEMPTS timer.function() self.assertEqual(enqueued, []) events = [ frame for frame in self.drain_json_events() if frame.get("event") == "group_audio_send_failed" ] self.assertEqual(len(events), 1) self.assertEqual( events[0].get("payload", {}).get("code"), "max_establish_attempts", ) finally: self.bridge._enqueue_scheduler_task = original_enqueue self.bridge.threading.Timer = original_timer def test_audio_established_resets_establish_attempts(self): link = self.install_audio_state("current-audio-link", established=False) self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, "attempts": self.bridge._AUDIO_LINK_MAX_ESTABLISH_ATTEMPTS, "retry_delay": self.bridge._AUDIO_LINK_RETRY_MAX_SECONDS, "retry_timer": None, "last_failure_reason": "establish_timeout", "max_attempts_emitted": True, } self.bridge.on_outgoing_audio_link_established(link) desired = self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] self.assertEqual(desired.get("attempts"), 0) self.assertEqual(desired.get("retry_delay"), self.bridge._AUDIO_LINK_RETRY_MIN_SECONDS) self.assertEqual(desired.get("last_failure_reason"), "") self.assertFalse(desired.get("max_attempts_emitted")) def test_outbound_overlay_group_audio_is_not_promoted(self): self.drain_audio_queue() link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=False) packet = FakePacket(link) self.bridge._known_peers[self.sender_peer_hash] = object() self.bridge._audio_link_desired_by_peer_hash[self.sender_peer_hash] = { "desired": True, } self.bridge.on_overlay_link_packet(self.group_audio_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_audio_link_id(link)) self.assertFalse(link.teardown_called) self.assertEqual(self.bridge._audio_binary_out_queue.qsize(), 0) def test_incoming_overlay_qchat_auth_promotes_when_transfer_is_pending(self): link, overlay_link_id, overlay_peer_hash = self.install_overlay_state( incoming=True ) packet = FakePacket(link) self.bridge._qchat_file_pending_sends_by_transfer["transfer-1"] = { "expires_at": time.time() + 60, } self.bridge.on_overlay_link_packet(self.qchat_file_auth_wire(), packet) self.assertNotIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertNotIn(id(link), self.bridge._overlay_link_ids_by_object) self.assertNotIn( overlay_peer_hash, self.bridge._active_overlay_link_id_by_peer_hash, ) self.assertNotIn(overlay_peer_hash, self.bridge._inbound_overlay_neighbors) self.assertFalse(link.teardown_called) file_link_id = self.bridge.get_qchat_file_link_id(link) self.assertIsInstance(file_link_id, str) file_state = self.bridge.get_qchat_file_link_state(file_link_id) self.assertIsNotNone(file_state) self.assertTrue(file_state["incoming"]) self.assertEqual(file_state["peerPresenceHash"], self.sender_peer_hash) self.assertEqual(file_state["transferId"], "transfer-1") self.assertIs(link.packet_callback, self.bridge.on_qchat_file_link_packet) def test_incoming_overlay_qchat_auth_without_pending_transfer_is_not_promoted(self): link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge.on_overlay_link_packet(self.qchat_file_auth_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_qchat_file_link_id(link)) self.assertFalse(link.teardown_called) def test_incoming_overlay_qchat_auth_with_expired_transfer_is_not_promoted(self): link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._qchat_file_pending_sends_by_transfer["transfer-1"] = { "expires_at": time.time() - 1, } self.bridge.on_overlay_link_packet(self.qchat_file_auth_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_qchat_file_link_id(link)) self.assertFalse(link.teardown_called) def test_incoming_overlay_qchat_auth_with_invalid_peer_hash_is_not_promoted(self): link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=True) packet = FakePacket(link) self.bridge._qchat_file_pending_sends_by_transfer["transfer-1"] = { "expires_at": time.time() + 60, } self.bridge.on_overlay_link_packet( self.qchat_file_auth_wire(peer_hash="not-a-hash"), packet, ) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_qchat_file_link_id(link)) self.assertFalse(link.teardown_called) def test_outbound_overlay_qchat_auth_is_not_promoted(self): link, overlay_link_id, _peer_hash = self.install_overlay_state(incoming=False) packet = FakePacket(link) self.bridge._qchat_file_pending_sends_by_transfer["transfer-1"] = { "expires_at": time.time() + 60, } self.bridge.on_overlay_link_packet(self.qchat_file_auth_wire(), packet) self.assertIn(overlay_link_id, self.bridge._overlay_links_by_id) self.assertIsNone(self.bridge.get_qchat_file_link_id(link)) self.assertFalse(link.teardown_called) class PresenceBridgeOverlayRouteMigrationTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.bridge._destination = FakeDestination() self.peer_hash = "ab" * 16 self.active_link = FakeLink() self.active_link_id = "active-overlay" self.active_state = { "linkId": self.active_link_id, "link": self.active_link, "peerPresenceHash": self.peer_hash, "incoming": False, "established": True, "established_at": time.time() - 120, "created_at": time.time() - 120, "overlay_transport_admitted": True, "manager_kind": "overlay", "manager_state": self.bridge._LINK_STATE_ESTABLISHED, "generation": 0, "peer_capabilities": { self.bridge._OVERLAY_ROUTE_MIGRATION_CAPABILITY, }, } self.bridge._overlay_links_by_id[self.active_link_id] = self.active_state self.bridge._overlay_link_ids_by_object[id(self.active_link)] = self.active_link_id self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash] = self.active_link_id def test_migration_uses_real_bounded_scheduler_lanes(self): for shard in range(self.bridge._SCHEDULER_OVERLAY_MIGRATION_SHARDS): lane = f"overlay-migration-{shard}" self.assertIn(lane, self.bridge._SCHEDULER_QUEUE_MAX_BY_LANE) self.assertGreater(self.bridge._SCHEDULER_QUEUE_MAX_BY_LANE[lane], 0) def test_hello_advertises_route_migration_and_marks_only_candidates(self): regular = json.loads( self.bridge._make_overlay_transport_wire( self.bridge._OVERLAY_HELLO_WIRE_TYPE, ).decode("utf-8") ) candidate = json.loads( self.bridge._make_overlay_transport_wire( self.bridge._OVERLAY_HELLO_WIRE_TYPE, migration_candidate=True, ).decode("utf-8") ) self.assertIn( self.bridge._OVERLAY_ROUTE_MIGRATION_CAPABILITY, regular["c"], ) self.assertNotIn("m", regular) self.assertEqual( candidate["m"], self.bridge._OVERLAY_ROUTE_MIGRATION_MARKER, ) def test_rtt_probe_resolution_requires_the_exact_nonce(self): event = threading.Event() state = { "rtt_pending": { "11" * 8: { "sent_ns": 1_000_000_000, "event": event, "rtt_ms": None, } } } with mock.patch.object( self.bridge.time, "monotonic_ns", return_value=1_025_000_000, ): self.assertIsNone( self.bridge._resolve_overlay_rtt_probe(state, "22" * 8) ) self.assertEqual( self.bridge._resolve_overlay_rtt_probe(state, "11" * 8), 25.0, ) self.assertTrue(event.is_set()) self.assertEqual(list(state["rtt_samples_ms"]), [25.0]) def test_ping_echoes_the_same_rtt_nonce(self): nonce = "33" * 8 wire = { "t": self.bridge._OVERLAY_PING_WIRE_TYPE, "r": self.peer_hash, "q": nonce, } with mock.patch.object( self.bridge, "_admit_overlay_peer_from_transport", return_value=True, ), mock.patch.object( self.bridge, "_send_overlay_transport_control", return_value=True, ) as send, mock.patch.object( self.bridge, "_overlay_enqueue_dedup", ): handled = self.bridge._handle_overlay_transport_control( wire, self.active_link, self.active_link_id, self.active_state, ) self.assertTrue(handled) self.assertEqual( send.call_args.kwargs["correlation_id"], nonce, ) def test_quality_gate_requires_lower_hops_and_clear_rtt_gain(self): accepted, result = self.bridge._overlay_migration_quality_acceptable( [120, 125, 130, 135, 140], [55, 60, 65, 70, 75], 8, 3, ) self.assertTrue(accepted) self.assertEqual(result["active_median_ms"], 130.0) self.assertEqual(result["candidate_median_ms"], 65.0) cases = ( ([120] * 5, [60] * 5, 4, 4), ([100] * 5, [85] * 5, 4, 2), ([120] * 5, [60] * 3, 4, 2), ) for active, candidate, active_hops, candidate_hops in cases: with self.subTest( active_hops=active_hops, candidate_hops=candidate_hops, candidate_samples=len(candidate), ): accepted, _result = self.bridge._overlay_migration_quality_acceptable( active, candidate, active_hops, candidate_hops, ) self.assertFalse(accepted) def test_rtt_rounds_alternate_probe_order(self): calls = [] def probe(link_id, purpose): calls.append((link_id, purpose)) event = threading.Event() event.set() return { "event": event, "rtt_ms": 100.0 if link_id == self.active_link_id else 40.0, } with mock.patch.object( self.bridge, "_send_overlay_rtt_probe", side_effect=probe, ): active, candidate = self.bridge._collect_overlay_migration_rtt_samples( self.active_link_id, "candidate-overlay", ) self.assertEqual( [link_id for link_id, _purpose in calls[:4]], [ self.active_link_id, "candidate-overlay", "candidate-overlay", self.active_link_id, ], ) self.assertEqual( len(active), self.bridge._OVERLAY_ROUTE_MIGRATION_PROBE_SAMPLES, ) self.assertEqual( len(candidate), self.bridge._OVERLAY_ROUTE_MIGRATION_PROBE_SAMPLES, ) def test_only_active_link_initiator_schedules_a_better_route_probe(self): with mock.patch.object( self.bridge, "_reticulum_path_snapshot", return_value={"has_path": True, "hops": 2}, ) as path_snapshot, mock.patch.object( self.bridge, "_reticulum_link_route_snapshot", return_value={"remote_hops": 7}, ) as route_snapshot, mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True, ) as enqueue: self.assertTrue( self.bridge._maybe_schedule_overlay_route_migration( self.peer_hash, "test_announce", ) ) self.assertEqual(enqueue.call_count, 1) self.assertIs( enqueue.call_args.args[2], self.bridge._overlay_route_migration_inspection_job, ) path_snapshot.assert_not_called() route_snapshot.assert_not_called() self.assertIn( self.peer_hash, self.bridge._overlay_route_migration_pending_by_peer_hash, ) self.bridge._overlay_route_migration_pending_by_peer_hash.clear() self.bridge._overlay_route_migration_last_attempt_at_by_peer_hash.clear() self.active_state["incoming"] = True with mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True, ) as enqueue: self.assertFalse( self.bridge._maybe_schedule_overlay_route_migration( self.peer_hash, "test_announce", ) ) enqueue.assert_not_called() def test_route_inspection_starts_migration_only_for_a_better_path(self): self.bridge._overlay_route_migration_pending_by_peer_hash.add(self.peer_hash) target_path = {"has_path": True, "hops": 2} active_route = {"remote_hops": 7} with mock.patch.object( self.bridge, "_reticulum_path_snapshot", return_value=target_path, ), mock.patch.object( self.bridge, "_reticulum_link_route_snapshot", return_value=active_route, ), mock.patch.object( self.bridge, "_overlay_route_migration_job", ) as migrate: self.bridge._overlay_route_migration_inspection_job( self.peer_hash, self.active_link_id, "test_announce", ) migrate.assert_called_once_with( self.peer_hash, self.active_link_id, target_path, active_route, ) self.assertIn( self.peer_hash, self.bridge._overlay_route_migration_last_attempt_at_by_peer_hash, ) self.assertNotIn( self.peer_hash, self.bridge._overlay_route_migration_pending_by_peer_hash, ) def test_dedup_does_not_select_or_close_a_migration_candidate(self): candidate_link_id = "candidate-overlay" candidate_state = dict(self.active_state) candidate_state.update( { "linkId": candidate_link_id, "link": FakeLink(), "migration_candidate": True, } ) self.bridge._overlay_links_by_id[candidate_link_id] = candidate_state with mock.patch.object( self.bridge, "_schedule_overlay_duplicate_close", ) as close: kept = self.bridge._dedup_overlay_links_for_peer(self.peer_hash) self.assertIs(kept, self.active_state) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], self.active_link_id, ) close.assert_not_called() def test_promoted_candidate_keeps_old_link_as_non_reclaiming_backup(self): candidate_link_id = "candidate-overlay" transaction_id = "44" * 8 candidate_state = dict(self.active_state) candidate_state.update( { "linkId": candidate_link_id, "link": FakeLink(), "migration_candidate": True, "migration_source_link_id": self.active_link_id, "overlay_transport_admitted": True, "migration_transaction_id": transaction_id, } ) self.bridge._overlay_links_by_id[candidate_link_id] = candidate_state with mock.patch.object( self.bridge, "emit_overlay_link_state", ), mock.patch.object( self.bridge, "_schedule_delayed_presence_announce_replay", ), mock.patch.object( self.bridge, "_flush_overlay_link_pending", ), mock.patch.object( self.bridge, "_schedule_overlay_duplicate_close", ) as delayed_close: self.assertTrue( self.bridge._promote_overlay_migration_candidate( self.peer_hash, candidate_link_id, "test_commit", ) ) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], candidate_link_id, ) self.assertTrue(self.active_state["migration_draining"]) delayed_close.assert_not_called() self.assertIs( self.bridge._dedup_overlay_links_for_peer(self.peer_hash), candidate_state, ) delayed_close.assert_not_called() self.assertTrue( self.bridge._finalize_overlay_migration( self.peer_hash, candidate_link_id, transaction_id, "test_finalize", ) ) delayed_close.assert_called_once_with( self.peer_hash, candidate_link_id, self.active_link_id, "route_migrated", ) kept = self.bridge._register_active_overlay_for_peer( self.peer_hash, self.active_link_id, ) self.assertIs(kept, candidate_state) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], candidate_link_id, ) removed = self.bridge.remove_overlay_link(candidate_link_id) self.assertIs(removed, candidate_state) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], self.active_link_id, ) self.assertNotIn("migration_draining", self.active_state) def test_candidate_send_failure_does_not_refresh_active_peer_path(self): state = { "linkId": "candidate-overlay", "link": FakeLink(), "peerPresenceHash": self.peer_hash, "migration_candidate": True, } with mock.patch.object( self.bridge, "_send_packet_on_link", return_value=False, ), mock.patch.object( self.bridge, "_overlay_enqueue_close", return_value=True, ) as close, mock.patch.object( self.bridge, "_force_overlay_peer_path_refresh", ) as refresh: self.assertFalse( self.bridge._send_overlay_transport_control( state["link"], state, self.bridge._OVERLAY_PING_WIRE_TYPE, "candidate_test", ) ) close.assert_called_once_with( "candidate-overlay", "overlay_transport_packet_send_false", ) refresh.assert_not_called() def test_candidate_teardown_leaves_active_link_and_peer_state_untouched(self): candidate_link = FakeLink() candidate_link_id = "candidate-overlay" self.bridge._overlay_links_by_id[candidate_link_id] = { "linkId": candidate_link_id, "link": candidate_link, "peerPresenceHash": self.peer_hash, "migration_candidate": True, "manager_kind": "overlay", "manager_state": self.bridge._LINK_STATE_CONNECTING, "generation": 0, } self.bridge._overlay_link_ids_by_object[id(candidate_link)] = candidate_link_id with mock.patch.object( self.bridge, "_run_with_timeout", return_value=(True, None, None), ), mock.patch.object( self.bridge, "emit_overlay_link_state", ), mock.patch.object( self.bridge, "_demote_overlay_fanout_peer", ) as demote, mock.patch.object( self.bridge, "_overlay_enqueue_peer_recovery", ) as recovery: self.bridge._teardown_overlay_link_id( candidate_link_id, "overlay_transport_packet_send_false", ) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], self.active_link_id, ) self.assertIn(self.active_link_id, self.bridge._overlay_links_by_id) demote.assert_not_called() recovery.assert_not_called() def test_unsolicited_incoming_candidate_is_rejected(self): self.bridge._overlay_links_by_id.clear() self.bridge._active_overlay_link_id_by_peer_hash.clear() incoming = FakeLink() with mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True, ) as enqueue: link_id = self.bridge._register_incoming_overlay_link( incoming, self.peer_hash, "test_candidate", migration_candidate=True, ) self.assertEqual(link_id, "") self.assertEqual(self.bridge._overlay_links_by_id, {}) self.assertEqual(enqueue.call_count, 1) def test_capable_active_peer_can_register_protected_incoming_candidate(self): self.active_state["incoming"] = True self.bridge._inbound_overlay_neighbors[self.peer_hash] = time.time() incoming = FakeLink() incoming.remote_identity = object() with mock.patch.object( self.bridge, "_send_overlay_hello_for_link", ) as send_hello, mock.patch.object( self.bridge, "derive_presence_destination_hash_for_identity", return_value=self.peer_hash, ): link_id = self.bridge._register_incoming_overlay_link( incoming, self.peer_hash, "test_candidate", migration_candidate=True, ) self.assertTrue(link_id) self.assertTrue( self.bridge._overlay_links_by_id[link_id]["migration_candidate"] ) self.assertTrue( self.bridge._overlay_links_by_id[link_id]["migration_peer_authenticated"] ) send_hello.assert_called_once_with(link_id, "incoming:test_candidate") self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], self.active_link_id, ) def test_unidentified_incoming_candidate_cannot_replace_active_link(self): self.active_state["incoming"] = True self.bridge._inbound_overlay_neighbors[self.peer_hash] = time.time() incoming = FakeLink() with mock.patch.object( self.bridge, "_send_overlay_hello_for_link", ) as send_hello, mock.patch.object( self.bridge, "_send_overlay_transport_control", return_value=True, ) as send_control: link_id = self.bridge._register_incoming_overlay_link( incoming, self.peer_hash, "test_candidate", migration_candidate=True, ) self.assertTrue(link_id) candidate = self.bridge._overlay_links_by_id[link_id] self.assertFalse(candidate["migration_peer_authenticated"]) self.assertFalse(candidate.get("overlay_transport_admitted") is True) send_hello.assert_not_called() self.assertTrue( self.bridge._handle_overlay_transport_control( { "t": self.bridge._OVERLAY_HELLO_WIRE_TYPE, "r": self.peer_hash, "c": [self.bridge._OVERLAY_ROUTE_MIGRATION_CAPABILITY], "m": self.bridge._OVERLAY_ROUTE_MIGRATION_MARKER, }, incoming, link_id, candidate, ) ) send_control.assert_not_called() self.assertFalse(candidate["migration_ready_event"].is_set()) self.assertFalse( self.bridge._promote_overlay_migration_candidate( self.peer_hash, link_id, "unauthenticated_commit", ) ) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], self.active_link_id, ) def test_mismatched_incoming_candidate_identity_is_rejected(self): self.active_state["incoming"] = True self.bridge._inbound_overlay_neighbors[self.peer_hash] = time.time() incoming = FakeLink() incoming.remote_identity = object() with mock.patch.object( self.bridge, "derive_presence_destination_hash_for_identity", return_value="cd" * 16, ), mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True, ) as enqueue: link_id = self.bridge._register_incoming_overlay_link( incoming, self.peer_hash, "test_candidate", migration_candidate=True, ) self.assertEqual(link_id, "") self.assertNotIn(id(incoming), self.bridge._overlay_link_ids_by_object) enqueue.assert_called_once() def test_remote_commit_retains_old_link_until_finalize(self): self.active_state["incoming"] = True self.bridge._inbound_overlay_neighbors[self.peer_hash] = time.time() candidate_link_id = "candidate-overlay" transaction_id = "55" * 8 candidate = dict(self.active_state) candidate.update( { "linkId": candidate_link_id, "link": FakeLink(), "incoming": True, "migration_candidate": True, "migration_source_link_id": self.active_link_id, "migration_peer_authenticated": True, "overlay_transport_admitted": True, } ) self.bridge._overlay_links_by_id[candidate_link_id] = candidate with mock.patch.object( self.bridge, "emit_overlay_link_state", ), mock.patch.object( self.bridge, "_schedule_delayed_presence_announce_replay", ), mock.patch.object( self.bridge, "_flush_overlay_link_pending", ), mock.patch.object( self.bridge, "_send_overlay_transport_control", return_value=True, ) as send, mock.patch.object( self.bridge, "_schedule_overlay_duplicate_close", ) as delayed_close, mock.patch.object( self.bridge, "_overlay_enqueue_dedup", ): commit = { "t": self.bridge._OVERLAY_MIGRATION_COMMIT_WIRE_TYPE, "r": self.peer_hash, "q": transaction_id, } self.assertTrue( self.bridge._handle_overlay_transport_control( commit, candidate["link"], candidate_link_id, candidate, ) ) self.assertEqual( self.bridge._active_overlay_link_id_by_peer_hash[self.peer_hash], candidate_link_id, ) delayed_close.assert_not_called() self.assertEqual( send.call_args.args[2], self.bridge._OVERLAY_MIGRATION_ACK_WIRE_TYPE, ) self.assertTrue( self.bridge._handle_overlay_transport_control( commit, candidate["link"], candidate_link_id, candidate, ) ) self.assertEqual(send.call_count, 2) delayed_close.assert_not_called() finalize = { "t": self.bridge._OVERLAY_MIGRATION_FINALIZE_WIRE_TYPE, "r": self.peer_hash, "q": transaction_id, } self.assertTrue( self.bridge._handle_overlay_transport_control( finalize, candidate["link"], candidate_link_id, candidate, ) ) delayed_close.assert_called_once_with( self.peer_hash, candidate_link_id, self.active_link_id, "route_migrated", ) def test_delayed_remote_identity_unlocks_incoming_candidate(self): self.active_state["incoming"] = True self.bridge._inbound_overlay_neighbors[self.peer_hash] = time.time() incoming = FakeLink() with mock.patch.object( self.bridge, "_send_overlay_hello_for_link", ) as send_hello: link_id = self.bridge._register_incoming_overlay_link( incoming, self.peer_hash, "test_candidate", migration_candidate=True, ) send_hello.assert_not_called() identity = object() incoming.remote_identity = identity with mock.patch.object( self.bridge, "derive_presence_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "find_peer_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "emit_overlay_link_state", ), mock.patch.object( self.bridge, "_note_overlay_peer_alive", ), mock.patch.object( self.bridge, "_register_active_overlay_for_peer", ), mock.patch.object( self.bridge, "_overlay_enqueue_dedup", ): self.bridge.on_overlay_link_remote_identified(incoming, identity) candidate = self.bridge._overlay_links_by_id[link_id] self.assertTrue(candidate["migration_peer_authenticated"]) self.assertTrue(candidate["overlay_transport_admitted"]) send_hello.assert_called_once_with(link_id, "migration_identity_verified") class PresenceBridgeResourceSchedulingTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.peer_hash = "ab" * 16 def test_resource_commands_use_fast_control_lane(self): actions = ( "accept_reticulum_chat_resource", "send_reticulum_chat_resource", "authorize_reticulum_chat_resource", "reject_reticulum_chat_resource", "cancel_reticulum_resource", ) for action in actions: with self.subTest(action=action): self.assertEqual( self.bridge._scheduler_lane_for_command(action), "resource-control", ) def test_resource_open_shard_is_stable_for_a_peer(self): lane = self.bridge._resource_open_scheduler_lane(self.peer_hash) self.assertEqual( lane, self.bridge._resource_open_scheduler_lane(self.peer_hash.upper()), ) self.assertIn(lane, self.bridge._SCHEDULER_QUEUE_MAX_BY_LANE) self.assertTrue(lane.startswith("resource-open-")) def test_only_one_link_handshake_starts_per_peer(self): first = {"peerPresenceHash": self.peer_hash, "transferId": "first"} second = {"peerPresenceHash": self.peer_hash, "transferId": "second"} scheduled = [] started = [] with mock.patch.object( self.bridge, "_enqueue_scheduler_task", side_effect=lambda lane, name, fn, *args, **kwargs: scheduled.append( (lane, name, fn, args, kwargs) ) or True, ), mock.patch.object( self.bridge, "_run_qchat_file_open_task", side_effect=lambda state: started.append(state), ): self.assertTrue(self.bridge._queue_qchat_file_open_state(first)) self.assertTrue(self.bridge._queue_qchat_file_open_state(second)) self.assertEqual(len(scheduled), 1) self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, [first]) self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, [first]) self.bridge._release_qchat_file_open_slot(first) self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, [first, second]) def test_chat_opens_are_queued_before_bulk_resources(self): bulk = { "peerPresenceHash": self.peer_hash, "transferId": "bulk", "resourceType": "reticulum_group_resource", } chat = { "peerPresenceHash": self.peer_hash, "transferId": "chat", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue(self.bridge._queue_qchat_file_open_state(bulk)) self.assertTrue(self.bridge._queue_qchat_file_open_state(chat)) queued = list(self.bridge._qchat_file_open_queue_by_peer[self.peer_hash]) self.assertEqual([item["transferId"] for item in queued], ["chat", "bulk"]) def test_active_resource_links_are_bounded_per_peer(self): waiting = { "peerPresenceHash": self.peer_hash, "transferId": "waiting", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } started = [] with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue(self.bridge._queue_qchat_file_open_state(waiting)) for index in range(self.bridge._QCHAT_FILE_ACTIVE_OUTGOING_MAX_PER_PEER): self.bridge._qchat_file_links_by_id[f"active-{index}"] = { "peerPresenceHash": self.peer_hash, "incoming": False, "transferId": f"active-{index}", } with mock.patch.object( self.bridge, "_run_qchat_file_open_task", side_effect=lambda state: started.append(state), ): self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, []) self.bridge._qchat_file_links_by_id.pop("active-0") self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, [waiting]) def test_chat_can_use_reserved_slot_while_bulk_is_at_its_limit(self): bulk_waiting = { "peerPresenceHash": self.peer_hash, "transferId": "bulk-waiting", "resourceType": "reticulum_group_resource", } chat_waiting = { "peerPresenceHash": self.peer_hash, "transferId": "chat-waiting", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue(self.bridge._queue_qchat_file_open_state(bulk_waiting)) self.assertTrue(self.bridge._queue_qchat_file_open_state(chat_waiting)) for index in range(self.bridge._QCHAT_FILE_BULK_ACTIVE_MAX_PER_PEER): self.bridge._qchat_file_links_by_id[f"bulk-active-{index}"] = { "peerPresenceHash": self.peer_hash, "incoming": False, "transferId": f"bulk-active-{index}", "resourceType": "reticulum_group_resource", } started = [] with mock.patch.object( self.bridge, "_run_qchat_file_open_task", side_effect=lambda state: started.append(state), ): self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) self.assertEqual(started, [chat_waiting]) self.assertIn( bulk_waiting, self.bridge._qchat_file_open_queue_by_peer[self.peer_hash], ) def test_successful_link_removal_wakes_waiting_peer_queue(self): waiting = {"peerPresenceHash": self.peer_hash, "transferId": "waiting"} with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue(self.bridge._queue_qchat_file_open_state(waiting)) link = FakeLink() active = { "peerPresenceHash": self.peer_hash, "transferId": "completed", "incoming": False, "link": link, } self.bridge._qchat_file_links_by_id["completed-link"] = active self.bridge._qchat_file_link_ids_by_object[id(link)] = "completed-link" with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ) as schedule, mock.patch.object( self.bridge, "_teardown_reticulum_link_bounded", ): self.bridge._qchat_file_receiver_transfer_done( self.peer_hash, "completed", ) schedule.assert_any_call(self.peer_hash) self.assertIn( waiting, self.bridge._qchat_file_open_queue_by_peer[self.peer_hash], ) def test_final_unestablished_link_failure_cleans_pending_transfer(self): transfer_id = "max-attempts" link = FakeLink() pending = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, } state = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "established": False, "open_attempts": self.bridge._QCHAT_FILE_LINK_MAX_OPEN_ATTEMPTS, "authMessage": {"type": "test"}, "receive_root": pending, "link": link, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) self.bridge._qchat_file_links_by_id["max-link"] = state self.bridge._qchat_file_link_ids_by_object[id(link)] = "max-link" with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge.on_qchat_file_link_closed(link) failed = [ call for call in emit.call_args_list if call.args and call.args[0] == "failed" ] self.assertEqual(len(failed), 1) self.assertEqual( failed[0].args[1]["reason"], "file_link_open_attempts_exhausted", ) self.assertNotIn( transfer_id, self.bridge._qchat_file_accepts_by_transfer, ) self.assertTrue(pending.get("cancelled")) def test_exhausted_parallel_link_keeps_viable_sibling_transfer(self): transfer_id = "parallel-transfer" failed_link = FakeLink() pending = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, } failed_state = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "established": False, "open_attempts": self.bridge._QCHAT_FILE_LINK_MAX_OPEN_ATTEMPTS, "authMessage": {"type": "test"}, "receive_root": pending, "link": failed_link, } sibling = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "established": True, "receive_root": pending, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) self.bridge._qchat_file_links_by_id["failed-link"] = failed_state self.bridge._qchat_file_link_ids_by_object[id(failed_link)] = "failed-link" self.bridge._qchat_file_links_by_id["sibling-link"] = sibling with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge.on_qchat_file_link_closed(failed_link) self.assertFalse( any(call.args and call.args[0] == "failed" for call in emit.call_args_list) ) self.assertIn( transfer_id, self.bridge._qchat_file_accepts_by_transfer, ) self.assertFalse(pending.get("cancelled", False)) def test_path_timeout_cleans_pending_transfer(self): transfer_id = "path-timeout" pending = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, } state = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "receive_root": pending, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object( self.bridge, "_open_qchat_file_link_for_state", side_effect=TimeoutError("no path"), ), mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._run_qchat_file_open_task(state) emit.assert_any_call( "failed", mock.ANY, ) self.assertNotIn( transfer_id, self.bridge._qchat_file_accepts_by_transfer, ) self.assertTrue(pending.get("cancelled")) def test_path_timeout_keeps_viable_parallel_sibling(self): transfer_id = "parallel-path-timeout" pending = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, } timed_out_state = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "receive_root": pending, } sibling = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "established": True, "receive_root": pending, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) self.bridge._qchat_file_links_by_id["path-sibling"] = sibling with mock.patch.object( self.bridge, "_open_qchat_file_link_for_state", side_effect=TimeoutError("no path"), ), mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._run_qchat_file_open_task(timed_out_state) self.assertFalse( any(call.args and call.args[0] == "failed" for call in emit.call_args_list) ) self.assertIn( transfer_id, self.bridge._qchat_file_accepts_by_transfer, ) self.assertFalse(pending.get("cancelled", False)) def test_path_timeout_fails_once_when_only_queued_siblings_remain(self): transfer_id = "queued-path-timeout" pending = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, } timed_out_state = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "receive_root": pending, } queued_sibling = { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "incoming": False, "receive_root": pending, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue( self.bridge._queue_qchat_file_open_state(queued_sibling) ) with mock.patch.object( self.bridge, "_open_qchat_file_link_for_state", side_effect=TimeoutError("no path"), ), mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._run_qchat_file_open_task(timed_out_state) failed = [ call for call in emit.call_args_list if call.args and call.args[0] == "failed" ] self.assertEqual(len(failed), 1) self.assertNotIn( transfer_id, self.bridge._qchat_file_accepts_by_transfer, ) self.assertTrue(pending.get("cancelled")) def test_terminal_cleanup_does_not_cancel_another_peer_transfer(self): current = { "peerPresenceHash": self.peer_hash, "transferId": "current-transfer", } stale_state = { "peerPresenceHash": self.peer_hash, "transferId": "stale-transfer", "incoming": False, } self.bridge._qchat_file_store_pending_receive(self.peer_hash, current) with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit"): self.bridge._qchat_file_fail_open_state( stale_state, "link_open_failed", force_transfer_failure=True, ) self.assertIn( "current-transfer", self.bridge._qchat_file_accepts_by_transfer, ) self.assertFalse(current.get("cancelled", False)) def test_cancelled_transfer_is_removed_from_open_queue(self): receive_root = {"cancelled": False} waiting = { "peerPresenceHash": self.peer_hash, "transferId": "cancelled", "receive_root": receive_root, } with mock.patch.object( self.bridge, "_schedule_qchat_file_peer_open_drain", return_value=True, ): self.assertTrue(self.bridge._queue_qchat_file_open_state(waiting)) receive_root["cancelled"] = True with mock.patch.object( self.bridge, "_run_qchat_file_open_task", ) as open_task: self.bridge._run_qchat_file_peer_open_queue(self.peer_hash) open_task.assert_not_called() self.assertNotIn( self.peer_hash, self.bridge._qchat_file_open_queue_by_peer, ) self.assertNotIn(id(waiting), self.bridge._qchat_file_open_queue_state_ids) def test_waiting_for_path_does_not_consume_link_attempt(self): state = { "peerPresenceHash": self.peer_hash, "peerIdentity": object(), "transferId": "waiting-for-path", } outbound = FakeDestination() outbound.hash = bytes.fromhex(self.peer_hash) with mock.patch.object( self.bridge, "build_outbound_destination", return_value=outbound, ), mock.patch.object( self.bridge, "_request_qchat_file_path", return_value=False, ): with self.assertRaises(self.bridge._QChatFilePathPending): self.bridge._open_qchat_file_link_for_state(state) self.assertEqual(int(state.get("open_attempts") or 0), 0) self.assertGreater(float(state.get("path_wait_started_at") or 0), 0) def test_failed_cached_path_is_refreshed_only_once_while_polling(self): state = { "peerPresenceHash": self.peer_hash, "peerIdentity": object(), "transferId": "refresh-once", } outbound = FakeDestination() outbound.hash = bytes.fromhex(self.peer_hash) refresh_permissions = [] def request_path(_destination_hash, _peer_hash, **kwargs): refresh_permissions.append(kwargs.get("allow_failed_path_refresh")) return False with mock.patch.object( self.bridge, "build_outbound_destination", return_value=outbound, ), mock.patch.object( self.bridge, "_peer_has_recent_unestablished_link_failure", return_value=True, ), mock.patch.object( self.bridge, "_request_qchat_file_path", side_effect=request_path, ): for _ in range(2): with self.assertRaises(self.bridge._QChatFilePathPending): self.bridge._open_qchat_file_link_for_state(state) self.assertEqual(refresh_permissions, [True, False]) self.assertTrue(state.get("failed_path_refresh_requested")) self.assertEqual(int(state.get("open_attempts") or 0), 0) class PresenceBridgeReusableResourceSessionTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.peer_hash = "ab" * 16 def pending( self, transfer_id, *, resource_type=None, event_id="", sha256="", timestamp=None, ): resource_type = resource_type or self.bridge._RETICULUM_CHAT_RESOURCE_TYPE auth = { "type": "RCR", "transferId": transfer_id, "ts": timestamp if timestamp is not None else int(time.time() * 1000), } if event_id: auth["eventId"] = event_id return { "peerPresenceHash": self.peer_hash, "transferId": transfer_id, "savePath": f"/tmp/{transfer_id}.recv", "fileName": f"{transfer_id}.bin", "size": 128, "sha256": sha256, "resourceType": resource_type, "metadata": { "groupId": 716, **({"eventId": event_id} if event_id else {}), }, "peerIdentity": object(), "authMessage": auth, } def session(self, lane="fast", established=True, slot=0): session_id = f"session-{lane}-{slot}" link = FakeSessionLink() state = { "linkId": session_id, "manager_kind": "resource_session", "sessionKey": self.bridge._resource_session_key(self.peer_hash, lane, slot), "sessionLane": lane, "sessionSlot": slot, "peerPresenceHash": self.peer_hash, "incoming": False, "established": established, "remote_ready": established, "provider_ready_sent": established, "created_at": time.time(), "last_used_at": time.time(), "pending_jobs": [], "active_requests": {}, "provider_active": 0, "link": link, "generation": 1, "activity_generation": 1, } self.bridge._qchat_file_links_by_id[session_id] = state self.bridge._qchat_file_link_ids_by_object[id(link)] = session_id self.bridge._resource_sessions_by_key[state["sessionKey"]] = session_id return state, link def test_parallel_peer_lookup_requires_an_exact_transfer_id(self): first = self.pending("first-range", resource_type="reticulum_group_resource_range") second = self.pending("second-range", resource_type="reticulum_group_resource_range") self.bridge._qchat_file_store_pending_receive(self.peer_hash, first) self.bridge._qchat_file_store_pending_receive(self.peer_hash, second) self.assertIs( self.bridge._qchat_file_get_pending_receive(self.peer_hash, "first-range"), first, ) self.assertIs( self.bridge._qchat_file_get_pending_receive(self.peer_hash, "second-range"), second, ) self.assertIsNone( self.bridge._qchat_file_get_pending_receive(self.peer_hash, "unknown-range") ) self.assertIsNone(self.bridge._qchat_file_get_pending_receive(self.peer_hash)) self.bridge._qchat_file_remove_pending_receive(self.peer_hash, "second-range") self.assertIs(self.bridge._qchat_file_get_pending_receive(self.peer_hash), first) def test_parallel_ranges_with_same_event_have_distinct_semantic_keys(self): first = self.pending( "first-range", resource_type="reticulum_group_resource_range", event_id="shared-event", ) second = self.pending( "second-range", resource_type="reticulum_group_resource_range", event_id="shared-event", ) file_hash = "ab" * 32 first["metadata"].update( {"fileHash": file_hash, "byteRanges": [[0, 1048576]]} ) second["metadata"].update( {"fileHash": file_hash, "byteRanges": [[1048576, 2097152]]} ) self.assertNotEqual( self.bridge._resource_session_semantic_key(first), self.bridge._resource_session_semantic_key(second), ) def test_session_response_skips_legacy_peer_fallback_callbacks(self): _state, link = self.session(lane="bulk") first = self.pending("first-range", resource_type="reticulum_group_resource_range") second = self.pending("second-range", resource_type="reticulum_group_resource_range") self.bridge._qchat_file_store_pending_receive(self.peer_hash, first) self.bridge._qchat_file_store_pending_receive(self.peer_hash, second) class ResponseResource: request_id = bytes.fromhex("33" * 16) def __init__(self, response_link): self.link = response_link resource = ResponseResource(link) with mock.patch.object( self.bridge, "_qchat_file_get_pending_receive", ) as get_pending, mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge.on_qchat_file_resource_started(resource) self.bridge.on_qchat_file_resource_concluded(resource) get_pending.assert_not_called() emit.assert_not_called() self.assertIn("first-range", self.bridge._qchat_file_accepts_by_transfer) self.assertIn("second-range", self.bridge._qchat_file_accepts_by_transfer) def test_range_response_is_bound_to_its_transfer_and_payload_hash(self): contents = b"verified range response" payload_hash = self.bridge.hashlib.sha256(contents).hexdigest() transfer_id = "verified-range" with tempfile.TemporaryDirectory() as directory: source_path = Path(directory) / "source.bin" save_path = Path(directory) / "received.bin" source_path.write_bytes(contents) response = open(source_path, "rb") receipt = FakeSessionReceipt() receipt.response = response receipt.metadata = { "transferId": transfer_id, "size": len(contents), "sha256": payload_hash, } pending = self.pending( transfer_id, resource_type="reticulum_group_resource_range", ) pending.update( { "savePath": str(save_path), "size": len(contents), } ) job = { "transferId": transfer_id, "pending": pending, "created_at": time.time(), "followers": [], } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_response_received(job, receipt) self.assertTrue(job["completed"]) self.assertEqual(save_path.read_bytes(), contents) received = [ call for call in emit.call_args_list if call.args and call.args[0] == "received" ] self.assertEqual(len(received), 1) self.assertEqual(received[0].args[1]["transferId"], transfer_id) self.assertEqual(received[0].args[1]["payloadHash"], payload_hash) def test_range_response_from_legacy_peer_uses_request_receipt_identity(self): contents = b"legacy range response" payload_hash = self.bridge.hashlib.sha256(contents).hexdigest() transfer_id = "expected-range" with tempfile.TemporaryDirectory() as directory: source_path = Path(directory) / "source.bin" save_path = Path(directory) / "received.bin" source_path.write_bytes(contents) receipt = FakeSessionReceipt() receipt.response = open(source_path, "rb") receipt.metadata = {"size": len(contents)} pending = self.pending( transfer_id, resource_type="reticulum_group_resource_range", ) pending.update( { "savePath": str(save_path), "size": len(contents), } ) job = { "transferId": transfer_id, "pending": pending, "created_at": time.time(), "followers": [], } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_response_received(job, receipt) self.assertTrue(job["completed"]) self.assertEqual(save_path.read_bytes(), contents) received = [ call for call in emit.call_args_list if call.args and call.args[0] == "received" ] self.assertEqual(len(received), 1) self.assertEqual(received[0].args[1]["transferId"], transfer_id) self.assertEqual(received[0].args[1]["payloadHash"], payload_hash) def test_range_response_rejects_mismatched_metadata_identity(self): contents = b"wrong transfer response" transfer_id = "expected-range" with tempfile.TemporaryDirectory() as directory: source_path = Path(directory) / "source.bin" save_path = Path(directory) / "received.bin" source_path.write_bytes(contents) receipt = FakeSessionReceipt() receipt.response = open(source_path, "rb") receipt.metadata = {"transferId": "different-range"} pending = self.pending( transfer_id, resource_type="reticulum_group_resource_range", ) pending.update({"savePath": str(save_path), "size": len(contents)}) job = { "transferId": transfer_id, "pending": pending, "created_at": time.time(), "followers": [], } self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_response_received(job, receipt) self.assertFalse(save_path.exists()) failures = [ call for call in emit.call_args_list if call.args and call.args[0] == "failed" ] self.assertEqual(len(failures), 1) self.assertEqual(failures[0].args[1]["reason"], "resource_response_invalid") self.assertIn("transfer id mismatch", failures[0].args[1]["error"]) def test_managed_accept_uses_reusable_session_instead_of_link_queue(self): payload = { "peerPresenceHash": self.peer_hash, "reticulumIdentityPublicKeyBase64": "identity", "authMessage": {"type": "RCR", "ts": int(time.time() * 1000)}, "transferId": "managed", "savePath": "/tmp/managed.recv", "fileName": "managed.bin", "size": 128, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "expires_at": time.time() + 60, } with mock.patch.object( self.bridge, "_parse_qchat_file_peer_identity", return_value=object(), ), mock.patch.object( self.bridge, "_resource_session_accept", ) as accept_session, mock.patch.object( self.bridge, "_open_qchat_file_link_async", ) as open_legacy: self.bridge.handle_accept_qchat_file_resource("req", payload) accept_session.assert_called_once() open_legacy.assert_not_called() def test_fast_and_bulk_resources_use_separate_peer_sessions(self): fast_job = { "pending": self.pending("fast"), "created_at": time.time(), } bulk_job = { "pending": self.pending( "bulk", resource_type="reticulum_group_resource_range", ), "created_at": time.time(), } with mock.patch.object(self.bridge, "_resource_session_poll_path"): self.assertTrue(self.bridge._resource_session_enqueue_job(fast_job)[0]) self.assertTrue(self.bridge._resource_session_enqueue_job(bulk_job)[0]) self.assertIn( self.bridge._resource_session_key(self.peer_hash, "fast"), self.bridge._resource_sessions_by_key, ) self.assertIn( self.bridge._resource_session_key(self.peer_hash, "bulk"), self.bridge._resource_sessions_by_key, ) self.assertEqual(len(self.bridge._resource_sessions_by_key), 2) def test_live_event_and_history_use_separate_priority_lanes(self): history_pending = self.pending("history") history_pending["metadata"].update( { "logicalResourceType": "reticulum_chat_history_page", "channelId": "general", "direction": "before", } ) history_pending["authMessage"].update({"before": {"id": "cursor"}}) live_pending = self.pending( "live", event_id="event-live", sha256="11" * 32, ) history_job = {"pending": history_pending, "created_at": time.time()} live_job = {"pending": live_pending, "created_at": time.time() + 0.01} with mock.patch.object(self.bridge, "_resource_session_poll_path"): self.assertTrue(self.bridge._resource_session_enqueue_job(history_job)[0]) self.assertTrue(self.bridge._resource_session_enqueue_job(live_job)[0]) fast_id = self.bridge._resource_sessions_by_key[ self.bridge._resource_session_key(self.peer_hash, "fast") ] bulk_id = self.bridge._resource_sessions_by_key[ self.bridge._resource_session_key(self.peer_hash, "bulk") ] self.assertEqual( [ job["pending"]["transferId"] for job in self.bridge._qchat_file_links_by_id[fast_id]["pending_jobs"] ], ["live"], ) self.assertEqual( [ job["pending"]["transferId"] for job in self.bridge._qchat_file_links_by_id[bulk_id]["pending_jobs"] ], ["history"], ) def test_dm_history_page_uses_fast_lane(self): self.assertEqual( self.bridge._resource_session_lane( self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "reticulum_chat_dm_page", ), "fast", ) def test_prepare_command_reports_fast_lane_for_dm_history(self): payload = { "peerPresenceHash": self.peer_hash, "reticulumIdentityPublicKeyBase64": "identity", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "logicalResourceType": "reticulum_chat_dm_page", } with mock.patch.object( self.bridge, "_resource_session_poll_path", ), mock.patch.object( self.bridge, "_parse_qchat_file_peer_identity", return_value=object(), ), mock.patch.object(self.bridge, "emit_resp") as emit_resp: self.bridge.handle_prepare_reticulum_resource_session("dm", payload) emit_resp.assert_called_once() self.assertTrue(emit_resp.call_args.args[1]) self.assertEqual(emit_resp.call_args.kwargs["payload"]["lane"], "fast") def test_prepare_command_reuses_pending_session_and_reports_state(self): peer_identity = object() payload = { "peerPresenceHash": self.peer_hash, "reticulumIdentityPublicKeyBase64": "identity", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "logicalResourceType": "reticulum_chat_history_page", } with mock.patch.object( self.bridge, "_resource_session_poll_path", ) as poll_path, mock.patch.object(self.bridge, "emit_resp") as emit_resp: with mock.patch.object( self.bridge, "_parse_qchat_file_peer_identity", return_value=peer_identity, ) as parse_identity: self.bridge.handle_prepare_reticulum_resource_session("one", payload) self.bridge.handle_prepare_reticulum_resource_session("two", payload) poll_path.assert_called_once() parse_identity.assert_has_calls( [ mock.call(self.peer_hash, "identity"), mock.call(self.peer_hash, "identity"), ] ) self.assertEqual(len(self.bridge._resource_sessions_by_key), 1) session_id = next(iter(self.bridge._resource_sessions_by_key.values())) self.assertIs( self.bridge._qchat_file_links_by_id[session_id]["peerIdentity"], peer_identity, ) self.assertEqual(emit_resp.call_count, 2) self.assertTrue(all(call.args[1] for call in emit_resp.call_args_list)) self.assertTrue( all( call.kwargs["payload"]["status"] == "pending" and call.kwargs["payload"]["lane"] == "bulk" for call in emit_resp.call_args_list ) ) def test_prepare_command_recalls_identity_when_public_key_is_omitted(self): peer_identity = object() payload = { "peerPresenceHash": self.peer_hash, "reticulumIdentityPublicKeyBase64": "", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } def recall(peer_hash, source): self.assertEqual(peer_hash, self.peer_hash) self.assertEqual(source, "ts_seed") self.bridge._known_peers[peer_hash] = peer_identity return True with mock.patch.object( self.bridge, "ensure_known_peer_from_recall", side_effect=recall, ), mock.patch.object( self.bridge, "_resource_session_poll_path", ) as poll_path, mock.patch.object(self.bridge, "emit_resp") as emit_resp: self.bridge.handle_prepare_reticulum_resource_session("req", payload) poll_path.assert_called_once() session_id = next(iter(self.bridge._resource_sessions_by_key.values())) self.assertIs( self.bridge._qchat_file_links_by_id[session_id]["peerIdentity"], peer_identity, ) emit_resp.assert_called_once_with( "req", True, payload=mock.ANY, ) def test_prepare_command_rejects_unverified_peer_identity(self): payload = { "peerPresenceHash": self.peer_hash, "reticulumIdentityPublicKeyBase64": "bad-identity", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } with mock.patch.object( self.bridge, "_parse_qchat_file_peer_identity", side_effect=ValueError("identity mismatch"), ), mock.patch.object( self.bridge, "_resource_session_poll_path", ) as poll_path, mock.patch.object(self.bridge, "emit_resp") as emit_resp: self.bridge.handle_prepare_reticulum_resource_session("req", payload) poll_path.assert_not_called() emit_resp.assert_called_once_with( "req", False, payload={"code": "bad_reticulum_identity"}, error="identity mismatch", ) def test_prepare_command_parses_real_reticulum_identity(self): peer_identity = RNS.Identity() peer_hash = self.bridge.destination_hash_hex( self.bridge.build_outbound_destination(peer_identity).hash ) public_key = base64.b64encode(peer_identity.get_public_key()).decode("ascii") payload = { "peerPresenceHash": peer_hash, "reticulumIdentityPublicKeyBase64": public_key, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } with mock.patch.object( self.bridge, "_resource_session_poll_path", ) as poll_path, mock.patch.object(self.bridge, "emit_resp") as emit_resp: self.bridge.handle_prepare_reticulum_resource_session("req", payload) poll_path.assert_called_once() session_id = next(iter(self.bridge._resource_sessions_by_key.values())) parsed_identity = self.bridge._qchat_file_links_by_id[session_id]["peerIdentity"] self.assertEqual(parsed_identity.get_public_key(), peer_identity.get_public_key()) emit_resp.assert_called_once_with("req", True, payload=mock.ANY) def test_sessions_are_bounded_per_peer_without_a_global_connection_cap(self): with mock.patch.object(self.bridge, "_resource_session_poll_path"): for index in range(24): peer_hash = f"{index + 1:032x}" state, reason = self.bridge._resource_session_get_or_create( peer_hash, object(), "fast", ) self.assertEqual(reason, "") self.assertIsInstance(state, dict) self.assertEqual(len(self.bridge._resource_sessions_by_key), 24) def test_bulk_pool_uses_one_resource_per_link_and_reserves_history(self): sessions = [ self.session(lane="bulk", slot=index) for index in range(self.bridge._RESOURCE_SESSION_BULK_POOL_SIZE) ] for index, (state, _link) in enumerate(sessions[:5]): state["pending_jobs"] = [ { "pending": self.pending( f"range-{index}", resource_type="reticulum_group_resource_range", ), "created_at": time.time() + index / 1000, "followers": [], } ] history = self.pending("visible-history") history["metadata"]["logicalResourceType"] = "reticulum_chat_history_page" sessions[5][0]["pending_jobs"] = [ { "pending": history, "created_at": time.time() + 1, "followers": [], } ] for state, _link in sessions: self.bridge._resource_session_dispatch_pending(state) active = [ job for state, _link in sessions for job in state["active_requests"].values() ] self.assertEqual(sum(len(link.requests) for _state, link in sessions), 6) self.assertTrue(all(len(link.requests) == 1 for _state, link in sessions)) self.assertEqual( sum( self.bridge._resource_session_job_class(job) == "attachment" for job in active ), 5, ) self.assertTrue( any( job["pending"]["transferId"] == "visible-history" for job in active ) ) def test_session_waits_for_provider_ready_before_dispatch(self): state, link = self.session(established=False) self.bridge._destination = FakeDestination() with mock.patch.object( self.bridge, "_send_packet_on_link", return_value=True, ), mock.patch.object( self.bridge, "_resource_session_dispatch_pending", ), mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ), mock.patch.object(self.bridge, "emit_event") as emit_event: self.bridge.on_outgoing_resource_session_established(link) self.assertTrue(state["established"]) self.assertFalse(state["remote_ready"]) self.assertFalse( any( call.args[0] == "reticulum_resource_session" and call.args[1].get("status") == "ready" for call in emit_event.call_args_list ) ) self.bridge._handle_qchat_file_link_packet( json.dumps( { "type": self.bridge._RESOURCE_SESSION_READY_TYPE, "r": self.peer_hash, "lane": "fast", "providerIdleMs": 180000, } ).encode("utf-8"), FakePacket(link), ) self.assertTrue(state["remote_ready"]) self.assertEqual(state["remote_provider_idle_seconds"], 180.0) emit_event.assert_any_call( "reticulum_resource_session", { "status": "ready", "peerPresenceHash": self.peer_hash, "lane": "fast", "linkId": state["linkId"], }, ) def test_first_packet_classifies_incoming_resource_session_without_overlay(self): link = FakeSessionLink() remote_identity = object() link.remote_identity = remote_identity self.bridge._destination = FakeDestination() hello = json.dumps( { "type": self.bridge._RESOURCE_SESSION_HELLO_TYPE, "r": self.peer_hash, "lane": "fast", } ).encode("utf-8") with mock.patch.object( self.bridge, "_schedule_inbound_classify_fallback", ), mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "_send_packet_on_link", return_value=True, ) as send_packet, mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ), mock.patch.object( self.bridge, "_teardown_reticulum_link_bounded", ): self.bridge.on_incoming_unified_link_established(link) link.packet_callback(hello, FakePacket(link)) link_id = self.bridge.get_qchat_file_link_id(link) self.assertIsInstance(link_id, str) self.assertIsNone(self.bridge.get_overlay_link_id(link)) state = self.bridge.get_qchat_file_link_state(link_id) self.assertEqual(state["linkId"], link_id) self.assertEqual(state["manager_kind"], "resource_session") self.assertEqual(state["peerPresenceHash"], self.peer_hash) self.assertTrue(state["provider_ready_sent"]) send_packet.assert_called_once() ready_wire = json.loads(send_packet.call_args.args[1].decode("utf-8")) self.assertEqual( ready_wire["providerIdleMs"], int( self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS * 1000 ), ) link.closed_callback(link) self.assertIsNone(self.bridge.get_qchat_file_link_id(link)) self.assertNotIn(link_id, self.bridge._qchat_file_links_by_id) def test_unidentified_incoming_resource_session_gets_idle_cleanup(self): link = FakeSessionLink() link.remote_identity = None with mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ) as schedule_idle, mock.patch.object( self.bridge, "_send_packet_on_link", ) as send_packet: link_id = self.bridge._register_incoming_resource_session( link, self.peer_hash, "bulk", ) self.assertIsInstance(link_id, str) state = self.bridge.get_qchat_file_link_state(link_id) self.assertIsInstance(state, dict) self.assertFalse(state.get("provider_ready_sent", False)) send_packet.assert_not_called() schedule_idle.assert_called_once_with(state) def test_parallel_dispatch_never_exceeds_lane_limit(self): for lane, limit in ( ("fast", self.bridge._RESOURCE_SESSION_FAST_CONCURRENCY), ("bulk", self.bridge._RESOURCE_SESSION_BULK_CONCURRENCY), ): with self.subTest(lane=lane): state, link = self.session(lane=lane) job_count = limit + 6 state["pending_jobs"] = [ { "pending": self.pending(f"{lane}-parallel-{index}"), "created_at": time.time(), "followers": [], } for index in range(job_count) ] threads = [ threading.Thread( target=self.bridge._resource_session_dispatch_pending, args=(state,), ) for _ in range(4) ] for thread in threads: thread.start() for thread in threads: thread.join(timeout=2) self.assertEqual(len(state["active_requests"]), limit) self.assertEqual(len(link.requests), limit) self.assertEqual(len(state["pending_jobs"]), job_count - limit) def test_bulk_jobs_are_balanced_across_independent_session_pool(self): selected_states = [] with mock.patch.object(self.bridge, "_resource_session_poll_path"): for index in range(self.bridge._RESOURCE_SESSION_BULK_POOL_SIZE * 2): state, reason = self.bridge._resource_session_get_or_create( self.peer_hash, object(), "bulk", ) self.assertEqual(reason, "") self.assertIsInstance(state, dict) state["pending_jobs"].append({"index": index}) selected_states.append(state) unique_states = {state["linkId"]: state for state in selected_states} self.assertEqual( len(unique_states), self.bridge._RESOURCE_SESSION_BULK_POOL_SIZE, ) self.assertEqual( sorted(len(state["pending_jobs"]) for state in unique_states.values()), [2] * self.bridge._RESOURCE_SESSION_BULK_POOL_SIZE, ) self.assertEqual( sorted(state["sessionSlot"] for state in unique_states.values()), list(range(self.bridge._RESOURCE_SESSION_BULK_POOL_SIZE)), ) def test_identical_event_downloads_share_one_network_job(self): first = self.pending("first", event_id="event-1", sha256="22" * 32) second = self.pending("second", event_id="event-1", sha256="22" * 32) self.bridge._qchat_file_store_pending_receive(self.peer_hash, first) self.bridge._qchat_file_store_pending_receive(self.peer_hash, second) with mock.patch.object( self.bridge, "_resource_session_enqueue_job", return_value=(True, ""), ) as enqueue, mock.patch.object(self.bridge, "emit_resp"): self.bridge._resource_session_accept("req-1", first) self.bridge._resource_session_accept("req-2", second) self.assertEqual(enqueue.call_count, 1) canonical = self.bridge._resource_session_jobs_by_transfer["first"] self.assertEqual(len(canonical["followers"]), 1) self.assertIs( self.bridge._resource_session_jobs_by_transfer["second"], canonical["followers"][0], ) def test_stale_authorization_is_not_dispatched(self): state, link = self.session() job = { "pending": self.pending( "stale", timestamp=int( (time.time() - self.bridge._RESOURCE_SESSION_AUTH_MAX_QUEUE_SECONDS - 1) * 1000 ), ), "created_at": time.time() - 100, "followers": [], } with mock.patch.object(self.bridge, "_qchat_file_emit") as emit: dispatched = self.bridge._resource_session_dispatch_job(state, job) self.assertFalse(dispatched) self.assertEqual(link.requests, []) self.assertTrue(job["completed"]) self.assertTrue( any( call.args[0] == "failed" and call.args[1]["reason"] == "resource_auth_refresh_required" for call in emit.call_args_list ) ) def test_request_exception_fails_job_instead_of_leaving_it_active(self): state, link = self.session() job = { "pending": self.pending("request-exception"), "created_at": time.time(), "followers": [], "session": state, } with mock.patch.object( link, "request", side_effect=RuntimeError("request failed"), ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: dispatched = self.bridge._resource_session_dispatch_job(state, job) self.assertFalse(dispatched) self.assertTrue(job["completed"]) self.assertNotIn("request-exception", state["active_requests"]) self.assertTrue( any( call.args[0] == "failed" and call.args[1]["reason"] == "resource_request_send_failed" for call in emit.call_args_list ) ) def test_late_response_after_cancellation_closes_response_file(self): job = { "completed": True, "pending": self.pending("late-response"), } receipt = FakeSessionReceipt() receipt.response = tempfile.TemporaryFile() self.bridge._resource_session_response_received(job, receipt) self.assertTrue(receipt.response.closed) def test_establishment_timeout_fails_jobs_once_and_refreshes_path(self): state, link = self.session(established=False) state["link_created_at"] = time.time() - 31 jobs = [ { "pending": self.pending(f"job-{index}"), "created_at": time.time(), "followers": [], "session": state, } for index in range(2) ] state["pending_jobs"] = jobs with mock.patch.object( self.bridge, "_force_overlay_peer_path_refresh", ) as refresh, mock.patch.object( self.bridge, "_teardown_reticulum_link_bounded", ) as teardown, mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_open_timeout(state) refresh.assert_called_once() teardown.assert_called_once_with(link, mock.ANY) self.assertTrue(all(job["completed"] for job in jobs)) self.assertEqual( len([call for call in emit.call_args_list if call.args[0] == "failed"]), 2, ) self.assertNotIn(state["sessionKey"], self.bridge._resource_sessions_by_key) def test_session_failure_requeues_jobs_that_were_never_dispatched(self): failed, failed_link = self.session(lane="bulk", slot=0) surviving, _surviving_link = self.session(lane="bulk", slot=1) job = { "pending": self.pending( "queued-range", resource_type="reticulum_group_resource_range", ), "created_at": time.time(), "followers": [], "session": failed, } failed["pending_jobs"] = [job] with mock.patch.object( self.bridge, "_teardown_reticulum_link_bounded", ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_fail_state( failed, "resource_session_link_closed", ) self.assertFalse(job.get("completed", False)) self.assertIs(job.get("session"), surviving) self.assertIs(surviving["active_requests"].get("queued-range"), job) self.assertFalse( any(call.args and call.args[0] == "failed" for call in emit.call_args_list) ) self.assertNotIn(failed["sessionKey"], self.bridge._resource_sessions_by_key) self.assertFalse(failed_link.teardown_called) def test_provider_session_outlives_requester_idle_leases(self): primary, _ = self.session(lane="bulk", slot=0) overflow, _ = self.session(lane="bulk", slot=1) incoming = dict(primary) incoming["incoming"] = True primary["remote_provider_idle_seconds"] = ( self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS ) overflow["remote_provider_idle_seconds"] = ( self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS ) self.assertEqual( self.bridge._resource_session_idle_timeout_seconds(primary), self.bridge._RESOURCE_SESSION_PRIMARY_IDLE_TIMEOUT_SECONDS, ) self.assertEqual( self.bridge._resource_session_idle_timeout_seconds(overflow), self.bridge._RESOURCE_SESSION_OVERFLOW_IDLE_TIMEOUT_SECONDS, ) self.assertEqual( self.bridge._resource_session_idle_timeout_seconds(incoming), self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS, ) self.assertGreater( self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS, self.bridge._RESOURCE_SESSION_PRIMARY_IDLE_TIMEOUT_SECONDS, ) self.assertGreater( self.bridge._RESOURCE_SESSION_INCOMING_IDLE_TIMEOUT_SECONDS, self.bridge._RESOURCE_SESSION_OVERFLOW_IDLE_TIMEOUT_SECONDS, ) def test_legacy_provider_session_closes_before_its_unadvertised_lease(self): primary, _ = self.session(lane="bulk", slot=0) overflow, _ = self.session(lane="bulk", slot=1) self.assertEqual( self.bridge._resource_session_idle_timeout_seconds(primary), self.bridge._RESOURCE_SESSION_LEGACY_PROVIDER_IDLE_TIMEOUT_SECONDS - self.bridge._RESOURCE_SESSION_PROVIDER_IDLE_GUARD_SECONDS, ) self.assertEqual( self.bridge._resource_session_idle_timeout_seconds(overflow), self.bridge._RESOURCE_SESSION_OVERFLOW_IDLE_TIMEOUT_SECONDS, ) def test_provider_request_waits_for_existing_electron_authorization(self): state, link = self.session() state["incoming"] = True transfer_id = "provider" with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"resource-response") file_path = temp_file.name pending = { "allowedRecipientAddress": "", "transferId": transfer_id, "filePath": file_path, "fileName": "provider.bin", "size": len(b"resource-response"), "sha256": "", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-provider"}, "expires_at": time.time() + 60, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending def authorize_on_event(status, payload): if status != "auth": return key = self.bridge._resource_session_waiter_key( state["linkId"], transfer_id, ) waiter = self.bridge._resource_session_provider_waiters[key] waiter["authorized"] = True waiter["event"].set() try: remote_identity = object() with mock.patch.object( self.bridge, "_qchat_file_emit", side_effect=authorize_on_event, ), mock.patch.object(self.bridge, "_resource_session_watch_provider_file"): with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ): response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": transfer_id, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-provider"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, remote_identity, time.time(), ) self.assertIsInstance(response, tuple) self.assertEqual(response[1]["transferId"], transfer_id) self.assertEqual(state["provider_active"], 1) self.assertEqual(response[0].read(), b"resource-response") response[0].close() finally: Path(file_path).unlink(missing_ok=True) def test_provider_rejects_unidentified_resource_session(self): _state, link = self.session() with mock.patch.object(self.bridge, "_qchat_file_emit") as emit: response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": "unidentified", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-unidentified"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, None, time.time(), ) self.assertEqual(response["reason"], "resource_peer_unidentified") emit.assert_not_called() def test_provider_rejects_request_after_idle_close_commits(self): state, link = self.session() state["incoming"] = True state["closing"] = True remote_identity = object() with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": "closing", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-closing"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, remote_identity, time.time(), ) self.assertEqual(response["reason"], "resource_session_unavailable") self.assertEqual(state["provider_active"], 0) emit.assert_not_called() def test_provider_waits_for_previous_resource_before_reusing_link(self): state, _link = self.session() state["incoming"] = True state["provider_active"] = 1 result = {} def acquire(): result["reason"] = ( self.bridge._resource_session_provider_acquire_link_slot( state, "next-resource", self.peer_hash, ) ) thread = threading.Thread(target=acquire) thread.start() deadline = time.time() + 1 while ( state.get("provider_link_waiter_transfer") != "next-resource" and time.time() < deadline ): time.sleep(0.01) self.assertEqual( state.get("provider_link_waiter_transfer"), "next-resource", ) with self.bridge._resource_session_provider_capacity_condition: state["provider_active"] = 0 self.bridge._resource_session_provider_capacity_condition.notify_all() thread.join(timeout=1) self.assertFalse(thread.is_alive()) self.assertEqual(result["reason"], "") self.assertEqual(state["provider_active"], 1) self.assertEqual(state["provider_link_handoffs"], 1) self.assertNotIn("provider_link_waiter_transfer", state) def test_provider_response_waits_for_link_handoff_then_starts_resource(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "handoff-response" with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"handoff-resource") file_path = temp_file.name self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = { "allowedRecipientAddress": self.peer_hash, "transferId": transfer_id, "filePath": file_path, "fileName": "handoff.bin", "size": len(b"handoff-resource"), "sha256": "", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-handoff"}, "expires_at": time.time() + 60, } auth_emitted = threading.Event() result = {} def authorize_on_event(status, _payload): if status != "auth": return waiter_key = self.bridge._resource_session_waiter_key( state["linkId"], transfer_id, ) waiter = self.bridge._resource_session_provider_waiters[waiter_key] waiter["authorized"] = True waiter["event"].set() auth_emitted.set() def request(): result["response"] = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": transfer_id, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-handoff"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) try: with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "_qchat_file_emit", side_effect=authorize_on_event, ), mock.patch.object( self.bridge, "_resource_session_watch_provider_file", ): thread = threading.Thread(target=request) thread.start() deadline = time.time() + 1 while ( state.get("provider_link_waiter_transfer") != transfer_id and time.time() < deadline ): time.sleep(0.01) self.assertEqual( state.get("provider_link_waiter_transfer"), transfer_id, ) self.assertFalse(auth_emitted.is_set()) with self.bridge._resource_session_provider_capacity_condition: state["provider_active"] = 0 self.bridge._resource_session_provider_capacity_condition.notify_all() thread.join(timeout=1) self.assertFalse(thread.is_alive()) self.assertTrue(auth_emitted.is_set()) response = result["response"] self.assertIsInstance(response, tuple) self.assertEqual(response[1]["transferId"], transfer_id) self.assertEqual(response[0].read(), b"handoff-resource") response[0].close() self.assertEqual(state["provider_active"], 1) self.assertEqual(state["provider_link_handoffs"], 1) finally: Path(file_path).unlink(missing_ok=True) def test_provider_allows_only_one_waiting_link_handoff(self): state, _link = self.session() state["incoming"] = True state["provider_active"] = 1 result = {} def acquire(): result["reason"] = ( self.bridge._resource_session_provider_acquire_link_slot( state, "first-waiter", self.peer_hash, ) ) thread = threading.Thread(target=acquire) thread.start() deadline = time.time() + 1 while ( state.get("provider_link_waiter_transfer") != "first-waiter" and time.time() < deadline ): time.sleep(0.01) self.assertEqual( self.bridge._resource_session_provider_acquire_link_slot( state, "first-waiter", self.peer_hash, ), "duplicate_resource_request", ) self.assertEqual( self.bridge._resource_session_provider_acquire_link_slot( state, "second-waiter", self.peer_hash, ), "resource_session_busy", ) with self.bridge._resource_session_provider_capacity_condition: state["provider_active"] = 0 self.bridge._resource_session_provider_capacity_condition.notify_all() thread.join(timeout=1) self.assertFalse(thread.is_alive()) self.assertEqual(result["reason"], "") self.assertEqual(state["provider_active"], 1) def test_provider_cancel_wakes_waiting_link_handoff(self): state, _link = self.session() state["incoming"] = True state["provider_active"] = 1 result = {} def acquire(): result["reason"] = ( self.bridge._resource_session_provider_acquire_link_slot( state, "cancelled-waiter", self.peer_hash, ) ) thread = threading.Thread(target=acquire) thread.start() deadline = time.time() + 1 while ( state.get("provider_link_waiter_transfer") != "cancelled-waiter" and time.time() < deadline ): time.sleep(0.01) self.bridge._resource_session_cancel_provider_transfer( state, "cancelled-waiter", ) thread.join(timeout=1) self.assertFalse(thread.is_alive()) self.assertEqual(result["reason"], "resource_requester_cancelled") self.assertEqual(state["provider_active"], 1) self.assertNotIn("provider_link_waiter_transfer", state) def test_provider_times_out_when_previous_resource_never_releases(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "_RESOURCE_SESSION_PROVIDER_LINK_HANDOFF_WAIT_SECONDS", 0.02, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": "same-link-second-resource", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-second-resource"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) self.assertEqual(response["reason"], "resource_session_busy") self.assertEqual(state["provider_active"], 1) self.assertNotIn("provider_link_waiter_transfer", state) emit.assert_not_called() def test_provider_pending_auth_limit_rejects_without_emitting_auth(self): state, link = self.session() state["incoming"] = True for index in range(self.bridge._RESOURCE_SESSION_PROVIDER_PENDING_AUTH_MAX): self.bridge._resource_session_provider_waiters[f"existing:{index}"] = { "peerPresenceHash": f"{index:032x}", } with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": "capacity", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-capacity"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) self.assertEqual(response["reason"], "resource_provider_busy") self.assertEqual(state["provider_active"], 0) emit.assert_not_called() def test_provider_handoff_shares_pending_authorization_limit(self): state, _link = self.session() state["incoming"] = True state["provider_active"] = 1 for index in range(self.bridge._RESOURCE_SESSION_PROVIDER_PENDING_AUTH_MAX): self.bridge._resource_session_provider_waiters[f"existing:{index}"] = { "peerPresenceHash": f"{index:032x}", } self.assertEqual( self.bridge._resource_session_provider_acquire_link_slot( state, "bounded-handoff", self.peer_hash, ), "resource_provider_busy", ) self.assertEqual(state["provider_active"], 1) self.assertNotIn("provider_link_waiter_transfer", state) def test_provider_waiting_for_auth_does_not_consume_transfer_capacity(self): state, link = self.session() state["incoming"] = True auth_emitted = threading.Event() result = {} def on_emit(status, _payload): if status == "auth": auth_emitted.set() def request(): result["response"] = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": "waiting-auth", "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "event-waiting-auth"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "_qchat_file_emit", side_effect=on_emit, ): thread = threading.Thread(target=request) thread.start() self.assertTrue(auth_emitted.wait(1.0)) self.assertEqual( sum(self.bridge._resource_session_provider_active_by_class.values()), 0, ) waiter = next( iter(self.bridge._resource_session_provider_waiters.values()) ) waiter["reason"] = "test_rejected" waiter["event"].set() thread.join(1.0) self.assertFalse(thread.is_alive()) self.assertEqual(result["response"]["reason"], "test_rejected") self.assertEqual(state["provider_active"], 0) self.assertFalse(self.bridge._resource_session_provider_waiters) self.assertFalse( self.bridge._resource_session_provider_pending_auth_by_peer ) def test_late_authorization_discards_abandoned_registered_send(self): state, _link = self.session() state["incoming"] = True transfer_id = "late-authorization" pending = { "allowedRecipientAddress": self.peer_hash, "transferId": transfer_id, "fileName": "late.bin", "size": 128, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "expires_at": time.time() + 60, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending with mock.patch.object(self.bridge, "_qchat_file_emit") as emit, mock.patch.object( self.bridge, "emit_resp", ) as emit_resp: self.bridge.handle_authorize_qchat_file_resource( "req", {"linkId": state["linkId"], "transferId": transfer_id}, ) self.assertNotIn( transfer_id, self.bridge._qchat_file_pending_sends_by_transfer, ) self.assertTrue(pending["cancelled"]) emit.assert_called_once_with("failed", mock.ANY) self.assertEqual( emit.call_args.args[1]["reason"], "resource_authorization_no_longer_active", ) emit_resp.assert_called_once_with( "req", False, payload={"code": "unknown_resource_request"}, error="Unknown resource session request", ) def test_late_authorization_cannot_discard_inflight_send(self): state, _link = self.session() state["incoming"] = True transfer_id = "inflight-authorization" pending = { "allowedRecipientAddress": self.peer_hash, "transferId": transfer_id, "fileName": "active.bin", "size": 128, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_inflight_transfers.add(transfer_id) with mock.patch.object(self.bridge, "_qchat_file_emit") as emit, mock.patch.object( self.bridge, "emit_resp", ): self.bridge.handle_authorize_qchat_file_resource( "req", {"linkId": state["linkId"], "transferId": transfer_id}, ) self.assertIs( self.bridge._qchat_file_pending_sends_by_transfer[transfer_id], pending, ) self.assertNotIn("cancelled", pending) emit.assert_not_called() def test_provider_capacity_preserves_live_and_sync_slots(self): active = self.bridge._resource_session_provider_active_by_class active["attachment"] = ( self.bridge._RESOURCE_SESSION_PROVIDER_ATTACHMENT_CONCURRENCY ) self.assertFalse( self.bridge._resource_session_provider_can_start_locked("attachment") ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked("history") ) active["history"] = 1 self.assertFalse( self.bridge._resource_session_provider_can_start_locked("history") ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked("live") ) active["attachment"] = 0 active["history"] = 0 active["metadata"] = ( self.bridge._RESOURCE_SESSION_PROVIDER_METADATA_CONCURRENCY ) self.assertFalse( self.bridge._resource_session_provider_can_start_locked("metadata") ) self.assertEqual( self.bridge._resource_session_provider_class( "reticulum_chat_calendar", "reticulum_chat_calendar", ), "metadata", ) def test_provider_capacity_waiters_are_prioritized(self): live_waiter = { "providerClass": "live", "peerPresenceHash": self.peer_hash, } metadata_waiter = { "providerClass": "metadata", "peerPresenceHash": "cd" * 16, } history_waiter = { "providerClass": "history", "peerPresenceHash": "ef" * 16, } queue = self.bridge._resource_session_provider_capacity_queue queue.extend([history_waiter, metadata_waiter, live_waiter]) self.assertFalse( self.bridge._resource_session_provider_can_start_locked( "metadata", "cd" * 16, metadata_waiter, ) ) self.assertFalse( self.bridge._resource_session_provider_can_start_locked( "history", "ef" * 16, history_waiter, ) ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked( "live", self.peer_hash, live_waiter, ) ) def test_ineligible_priority_waiter_does_not_block_other_peers(self): blocked_peer = self.peer_hash available_peer = "cd" * 16 live_waiter = { "providerClass": "live", "peerPresenceHash": blocked_peer, } history_waiter = { "providerClass": "history", "peerPresenceHash": available_peer, } self.bridge._resource_session_provider_capacity_queue.extend( [live_waiter, history_waiter] ) self.bridge._resource_session_provider_active_by_peer[blocked_peer] = ( self.bridge._RESOURCE_SESSION_PROVIDER_ACTIVE_MAX_PER_PEER ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked( "history", available_peer, history_waiter, ) ) def test_provider_active_capacity_is_bounded_per_peer(self): self.bridge._resource_session_provider_active_by_peer[self.peer_hash] = ( self.bridge._RESOURCE_SESSION_PROVIDER_ACTIVE_MAX_PER_PEER ) self.assertFalse( self.bridge._resource_session_provider_can_start_locked( "live", self.peer_hash, ) ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked( "live", "cd" * 16, ) ) def test_attachment_capacity_reserves_two_slots_for_same_peer_chat(self): self.bridge._resource_session_provider_active_by_class["attachment"] = ( self.bridge._RESOURCE_SESSION_PROVIDER_ATTACHMENT_MAX_PER_PEER ) self.bridge._resource_session_provider_active_by_peer[self.peer_hash] = ( self.bridge._RESOURCE_SESSION_PROVIDER_ATTACHMENT_MAX_PER_PEER ) self.bridge._resource_session_provider_active_attachments_by_peer[ self.peer_hash ] = self.bridge._RESOURCE_SESSION_PROVIDER_ATTACHMENT_MAX_PER_PEER self.assertFalse( self.bridge._resource_session_provider_can_start_locked( "attachment", self.peer_hash, ) ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked( "history", self.peer_hash, ) ) self.assertTrue( self.bridge._resource_session_provider_can_start_locked( "live", self.peer_hash, ) ) def test_provider_post_auth_wait_queue_is_bounded(self): state, _link = self.session() self.bridge._resource_session_provider_capacity_queue.extend( { "providerClass": "attachment", "peerPresenceHash": f"{index:032x}", "transferId": f"queued-{index}", } for index in range( self.bridge._RESOURCE_SESSION_PROVIDER_CAPACITY_QUEUE_MAX ) ) acquired = self.bridge._resource_session_provider_acquire_capacity( "live", "queue-full", {}, state, ) self.assertFalse(acquired) self.assertFalse( self.bridge._resource_session_provider_capacity_waiters_by_peer ) def test_provider_classifies_attachment_ranges_separately(self): self.assertEqual( self.bridge._resource_session_provider_class( "reticulum_group_resource_range", "reticulum_group_resource_range", ), "attachment", ) self.assertEqual( self.bridge._resource_session_provider_class( self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "reticulum_chat_history_page", ), "history", ) self.assertEqual( self.bridge._resource_session_provider_class( self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "reticulum_chat_metadata_snapshot", ), "metadata", ) self.assertEqual( self.bridge._resource_session_provider_class( self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "qortalland_chat", ), "live", ) def test_provider_uses_registered_resource_for_capacity_class(self): state, link = self.session() state["incoming"] = True transfer_id = "authoritative-class" with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"range-response") file_path = temp_file.name self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = { "allowedRecipientAddress": self.peer_hash, "transferId": transfer_id, "filePath": file_path, "fileName": "range.bin", "size": len(b"range-response"), "sha256": "", "resourceType": "reticulum_group_resource_range", "metadata": { "logicalResourceType": "reticulum_group_resource_range", }, "expires_at": time.time() + 60, } captured_classes = [] def authorize(status, _payload): if status != "auth": return waiter = next( iter(self.bridge._resource_session_provider_waiters.values()) ) waiter["authorized"] = True waiter["event"].set() def watch( _file, _transfer_id, _pending, _state, _request_id, provider_class, ): captured_classes.append(provider_class) try: with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object( self.bridge, "_qchat_file_emit", side_effect=authorize, ), mock.patch.object( self.bridge, "_resource_session_watch_provider_file", side_effect=watch, ): response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": transfer_id, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "mislabelled-range"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) self.assertIsInstance(response, tuple) self.assertEqual(captured_classes, ["attachment"]) response[0].close() finally: Path(file_path).unlink(missing_ok=True) def test_provider_slot_is_held_until_progressing_response_completes(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "provider-progress" request_id = bytes.fromhex("88" * 16) class ProviderResource: def __init__(self): self.request_id = request_id self.status = RNS.Resource.TRANSFERRING self.progress = 0.0 def get_progress(self): return self.progress resource = ProviderResource() link.outgoing_resources.append(resource) with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"progressing-resource") file_path = temp_file.name file_handle = open(file_path, "rb") pending = { "transferId": transfer_id, "fileName": "progress.bin", "size": len(b"progressing-resource"), "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_active_by_class["live"] = 1 try: with mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_STALL_TIMEOUT_SECONDS", 0.25, ), mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_INITIAL_PROGRESS_TIMEOUT_SECONDS", 0.25, ), mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_watch_provider_file( file_handle, transfer_id, pending, state, request_id, "live", ) for progress in (0.1, 0.2, 0.3, 0.4): time.sleep(0.11) resource.progress = progress resource.status = RNS.Resource.COMPLETE deadline = time.time() + 2 while state["provider_active"] > 0 and time.time() < deadline: time.sleep(0.02) self.assertEqual(state["provider_active"], 0) self.assertNotIn( transfer_id, self.bridge._qchat_file_pending_sends_by_transfer, ) self.assertTrue( any( call.args[0] == "sent" and call.args[1]["transferId"] == transfer_id for call in emit.call_args_list ) ) self.assertEqual( self.bridge._resource_session_provider_active_by_class["live"], 0, ) finally: if not file_handle.closed: file_handle.close() Path(file_path).unlink(missing_ok=True) def test_provider_zero_progress_uses_initial_progress_timeout(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "provider-no-progress" request_id = bytes.fromhex("89" * 16) resource = mock.Mock() resource.request_id = request_id resource.status = RNS.Resource.TRANSFERRING resource.get_progress.return_value = 0.0 link.outgoing_resources.append(resource) with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"no-progress-resource") file_path = temp_file.name file_handle = open(file_path, "rb") pending = { "transferId": transfer_id, "fileName": "no-progress.bin", "size": len(b"no-progress-resource"), "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_active_by_class["live"] = 1 try: with mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_INITIAL_PROGRESS_TIMEOUT_SECONDS", 0.12, ), mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_STALL_TIMEOUT_SECONDS", 1.0, ), mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_watch_provider_file( file_handle, transfer_id, pending, state, request_id, "live", ) deadline = time.time() + 1 while state["provider_active"] > 0 and time.time() < deadline: time.sleep(0.01) self.assertEqual(state["provider_active"], 0) resource.cancel.assert_called_once_with() self.assertTrue( any( call.args[0] == "failed" and call.args[1]["reason"] == "resource_response_not_started" for call in emit.call_args_list ) ) finally: if not file_handle.closed: file_handle.close() Path(file_path).unlink(missing_ok=True) def test_provider_watcher_follows_all_resource_segments(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "provider-segments" request_id = bytes.fromhex("99" * 16) class ProviderSegment: def __init__(self, index, status, progress): self.request_id = request_id self.segment_index = index self.total_segments = 2 self.status = status self.progress = progress self.next_segment = None def get_progress(self): return self.progress second = ProviderSegment(2, RNS.Resource.TRANSFERRING, 0.5) first = ProviderSegment(1, RNS.Resource.COMPLETE, 0.5) first.next_segment = second link.outgoing_resources.append(first) with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"segmented-resource") file_path = temp_file.name file_handle = open(file_path, "rb") pending = { "transferId": transfer_id, "fileName": "segments.bin", "size": len(b"segmented-resource"), "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_active_by_class["live"] = 1 try: with mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_STALL_TIMEOUT_SECONDS", 1.0, ), mock.patch.object( self.bridge, "_RESOURCE_SESSION_RESPONSE_INITIAL_PROGRESS_TIMEOUT_SECONDS", 1.0, ), mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: self.bridge._resource_session_watch_provider_file( file_handle, transfer_id, pending, state, request_id, "live", ) time.sleep(0.15) self.assertEqual(state["provider_active"], 1) self.assertFalse( any(call.args[0] == "sent" for call in emit.call_args_list) ) second.progress = 1.0 second.status = RNS.Resource.COMPLETE deadline = time.time() + 2 while time.time() < deadline: sent = any( call.args[0] == "sent" for call in emit.call_args_list ) if state["provider_active"] == 0 and sent: break time.sleep(0.02) self.assertEqual(state["provider_active"], 0) self.assertTrue( any(call.args[0] == "sent" for call in emit.call_args_list) ) finally: if not file_handle.closed: file_handle.close() Path(file_path).unlink(missing_ok=True) def test_provider_watcher_releases_inflight_admission_after_pending_replacement(self): state, link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "provider-replaced-pending" request_id = bytes.fromhex("aa" * 16) resource = mock.Mock() resource.request_id = request_id resource.status = RNS.Resource.TRANSFERRING resource.get_progress.return_value = 0.5 link.outgoing_resources.append(resource) with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"replacement-resource") file_path = temp_file.name file_handle = open(file_path, "rb") pending = { "transferId": transfer_id, "fileName": "original.bin", "size": len(b"replacement-resource"), "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, } replacement = {**pending, "fileName": "retry.bin"} self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_inflight_transfers.add(transfer_id) self.bridge._resource_session_provider_active_by_class["live"] = 1 try: with mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ): self.bridge._resource_session_watch_provider_file( file_handle, transfer_id, pending, state, request_id, "live", ) self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = ( replacement ) resource.status = RNS.Resource.COMPLETE deadline = time.time() + 1 while state["provider_active"] > 0 and time.time() < deadline: time.sleep(0.01) self.assertIs( self.bridge._qchat_file_pending_sends_by_transfer[transfer_id], replacement, ) self.assertNotIn( transfer_id, self.bridge._resource_session_provider_inflight_transfers, ) finally: if not file_handle.closed: file_handle.close() Path(file_path).unlink(missing_ok=True) def test_provider_watcher_preserves_file_during_response_handoff(self): state, _link = self.session() state["incoming"] = True state["provider_active"] = 1 transfer_id = "provider-handoff-cancel" with tempfile.NamedTemporaryFile(delete=False) as temp_file: temp_file.write(b"handoff-resource") file_path = temp_file.name file_handle = open(file_path, "rb") pending = { "transferId": transfer_id, "fileName": "handoff.bin", "size": len(b"handoff-resource"), "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "cancelled": True, } self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_provider_active_by_class["live"] = 1 try: with mock.patch.object( self.bridge, "_RESOURCE_SESSION_PROVIDER_RESPONSE_START_GRACE_SECONDS", 0.1, ), mock.patch.object( self.bridge, "_resource_session_schedule_idle_close", ): self.bridge._resource_session_watch_provider_file( file_handle, transfer_id, pending, state, b"request", "live", ) time.sleep(0.03) self.assertFalse(file_handle.closed) deadline = time.time() + 1 while state["provider_active"] > 0 and time.time() < deadline: time.sleep(0.01) self.assertTrue(file_handle.closed) self.assertEqual(state["provider_active"], 0) self.assertEqual( self.bridge._resource_session_provider_active_by_class["live"], 0, ) finally: if not file_handle.closed: file_handle.close() Path(file_path).unlink(missing_ok=True) def test_provider_cancel_wakes_auth_and_marks_registered_send(self): state, _link = self.session() state["incoming"] = True transfer_id = "provider-cancel" waiter = { "event": threading.Event(), "authorized": False, "reason": "resource_authorization_timeout", } waiter_key = self.bridge._resource_session_waiter_key( state["linkId"], transfer_id, ) pending = { "transferId": transfer_id, "allowedRecipientAddress": self.peer_hash, } self.bridge._resource_session_provider_waiters[waiter_key] = waiter self.bridge._qchat_file_pending_sends_by_transfer[transfer_id] = pending self.bridge._resource_session_cancel_provider_transfer( state, transfer_id, ) self.assertTrue(waiter["event"].is_set()) self.assertEqual(waiter["reason"], "resource_requester_cancelled") self.assertTrue(pending["cancelled"]) def test_provider_remembers_cancel_that_arrives_before_request(self): state, link = self.session() state["incoming"] = True transfer_id = "cancel-before-request" self.bridge._resource_session_cancel_provider_transfer( state, transfer_id, ) with mock.patch.object( self.bridge, "_destination_hash_for_identity", return_value=self.peer_hash, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: response = self.bridge._resource_session_response_generator( self.bridge._RESOURCE_SESSION_REQUEST_PATH, { "version": 1, "transferId": transfer_id, "resourceType": self.bridge._RETICULUM_CHAT_RESOURCE_TYPE, "metadata": {"eventId": "cancelled-event"}, "authMessage": {"type": "RCR"}, }, b"request", link.link_id, object(), time.time(), ) self.assertEqual(response["reason"], "resource_requester_cancelled") self.assertEqual(state["provider_active"], 0) emit.assert_not_called() def test_cancelled_request_does_not_close_reusable_session(self): state, link = self.session() pending = self.pending("cancel-me") receipt = FakeSessionReceipt() receipt.resource = mock.Mock() job = { "pending": pending, "created_at": time.time(), "followers": [], "session": state, "semanticKey": "cancel-key", "receipt": receipt, } state["active_requests"]["cancel-me"] = job self.bridge._resource_session_jobs_by_transfer["cancel-me"] = job self.bridge._resource_session_jobs_by_semantic_key["cancel-key"] = job self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object( self.bridge, "_qchat_file_emit", ), mock.patch.object( self.bridge, "_send_packet_on_link", return_value=True, ) as send_packet: closed = self.bridge._qchat_file_cancel_transfer( "cancel-me", self.peer_hash, "test-cancel", ) self.assertEqual(closed, 0) self.assertFalse(link.teardown_called) receipt.resource.cancel.assert_called_once_with() self.assertEqual(receipt.status, FakeSessionReceipt.FAILED) cancel_wire = json.loads(send_packet.call_args.args[1].decode("utf-8")) self.assertEqual( cancel_wire, { "type": self.bridge._RESOURCE_SESSION_CANCEL_TYPE, "transferId": "cancel-me", }, ) self.assertIn(state["sessionKey"], self.bridge._resource_sessions_by_key) self.assertTrue(job["completed"]) def test_authorization_timeout_retires_only_the_suspect_session(self): state, link = self.session(lane="bulk", slot=0) surviving, surviving_link = self.session(lane="bulk", slot=1) pending = self.pending( "authorization-timeout", resource_type="reticulum_group_resource_range", ) receipt = FakeSessionReceipt() receipt.resource = mock.Mock() job = { "pending": pending, "created_at": time.time(), "followers": [], "session": state, "semanticKey": "authorization-timeout-key", "receipt": receipt, } queued = { "pending": self.pending( "queued-after-timeout", resource_type="reticulum_group_resource_range", ), "created_at": time.time(), "followers": [], "session": state, "semanticKey": "queued-after-timeout-key", } state["active_requests"]["authorization-timeout"] = job state["pending_jobs"].append(queued) self.bridge._resource_session_jobs_by_transfer["authorization-timeout"] = job self.bridge._resource_session_jobs_by_transfer["queued-after-timeout"] = queued self.bridge._resource_session_jobs_by_semantic_key[ "authorization-timeout-key" ] = job self.bridge._resource_session_jobs_by_semantic_key[ "queued-after-timeout-key" ] = queued self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) with mock.patch.object( self.bridge, "_qchat_file_emit", ), mock.patch.object( self.bridge, "_send_packet_on_link", return_value=True, ), mock.patch.object( self.bridge, "_teardown_reticulum_link_bounded", ) as teardown: closed = self.bridge._qchat_file_cancel_transfer( "authorization-timeout", self.peer_hash, "resource_authorization_timeout", ) self.assertEqual(closed, 1) self.assertTrue(state["closing"]) self.assertNotIn(state["sessionKey"], self.bridge._resource_sessions_by_key) teardown.assert_called_once() self.assertIs(teardown.call_args.args[0], link) receipt.resource.cancel.assert_called_once_with() self.assertEqual(receipt.status, FakeSessionReceipt.FAILED) self.assertFalse(surviving.get("closing", False)) self.assertFalse(surviving_link.teardown_called) self.assertIs(queued.get("session"), surviving) self.assertIs( surviving["active_requests"].get("queued-after-timeout"), queued, ) def test_cancel_during_request_handoff_cancels_late_receipt(self): state, link = self.session() pending = self.pending("handoff-cancel") job = { "pending": pending, "created_at": time.time(), "followers": [], "session": state, "semanticKey": "handoff-key", } self.bridge._resource_session_jobs_by_transfer["handoff-cancel"] = job self.bridge._resource_session_jobs_by_semantic_key["handoff-key"] = job self.bridge._qchat_file_store_pending_receive(self.peer_hash, pending) receipt = FakeSessionReceipt() receipt.resource = mock.Mock() def request_then_cancel(*_args, **_kwargs): self.bridge._qchat_file_cancel_transfer( "handoff-cancel", self.peer_hash, "test-handoff-cancel", ) return receipt with mock.patch.object( link, "request", side_effect=request_then_cancel, ), mock.patch.object( self.bridge, "_send_packet_on_link", return_value=True, ), mock.patch.object(self.bridge, "_qchat_file_emit") as emit: dispatched = self.bridge._resource_session_dispatch_job(state, job) self.assertFalse(dispatched) self.assertTrue(job["completed"]) self.assertTrue(job["cancelled"]) self.assertEqual(receipt.status, FakeSessionReceipt.FAILED) receipt.resource.cancel.assert_called_once_with() self.assertNotIn("receipt", job) self.assertFalse( any(call.args[0] == "auth_sent" for call in emit.call_args_list) ) class PresenceBridgeOverlayGoodOutboundCacheTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.temp_dir = tempfile.TemporaryDirectory() self.bridge._reticulum_config_dir = self.temp_dir.name self.bridge._overlay_good_outbound_cache.clear() self.bridge._overlay_good_outbound_cache_loaded = False self.bridge._overlay_good_outbound_cache_dirty = False self.bridge._overlay_good_outbound_cache_last_write_at = 0.0 def tearDown(self): self.temp_dir.cleanup() self.bridge._shutdown.clear() def _write_cache(self, payload): path = self.bridge._overlay_good_outbound_cache_path() with open(path, "w", encoding="utf-8") as cache_file: json.dump(payload, cache_file) return path def test_namespace_fingerprint_covers_all_destination_name_components(self): baseline = self.bridge._overlay_good_outbound_cache_namespace_fingerprint() for field_name in ("APP_NAMESPACE", "PRESENCE_ASPECT", "PRESENCE_VERSION"): with self.subTest(field=field_name), mock.patch.object( self.bridge, field_name, f"{getattr(self.bridge, field_name)}-changed", ): self.assertNotEqual( self.bridge._overlay_good_outbound_cache_namespace_fingerprint(), baseline, ) def test_old_cache_is_replaced_without_seeding_old_namespace_peers(self): old_peer = "ab" * 16 path = self._write_cache( { "version": 1, "updatedAt": time.time(), "peers": [{"peerHash": old_peer, "lastRxAt": time.time()}], } ) self.bridge._load_overlay_good_outbound_cache() self.assertEqual(self.bridge._overlay_good_outbound_cache, {}) with open(path, "r", encoding="utf-8") as cache_file: replacement = json.load(cache_file) self.assertEqual( replacement["version"], self.bridge._OVERLAY_GOOD_OUTBOUND_CACHE_VERSION, ) self.assertEqual( replacement["namespaceFingerprint"], self.bridge._overlay_good_outbound_cache_namespace_fingerprint(), ) self.assertEqual(replacement["peers"], []) def test_current_version_cache_from_another_namespace_is_replaced(self): path = self._write_cache( { "version": self.bridge._OVERLAY_GOOD_OUTBOUND_CACHE_VERSION, "namespaceFingerprint": "wrong-namespace", "updatedAt": time.time(), "peers": [{"peerHash": "cd" * 16, "lastRxAt": time.time()}], } ) self.bridge._load_overlay_good_outbound_cache() self.assertEqual(self.bridge._overlay_good_outbound_cache, {}) with open(path, "r", encoding="utf-8") as cache_file: replacement = json.load(cache_file) self.assertEqual( replacement["namespaceFingerprint"], self.bridge._overlay_good_outbound_cache_namespace_fingerprint(), ) def test_failed_recall_removes_cached_peer_instead_of_marking_candidate(self): peer_hash = "ef" * 16 self.bridge._overlay_good_outbound_cache_loaded = True self.bridge._overlay_good_outbound_cache[peer_hash] = { "first_rx_at": time.time(), "last_rx_at": time.time(), "rx_count": 1, } with mock.patch.object( self.bridge, "_overlay_peer_available_for_new_outbound", return_value=True, ), mock.patch.object( self.bridge, "ensure_known_peer_from_recall", return_value=False, ), mock.patch.object( self.bridge, "_mark_candidate_peer", ) as mark_candidate, mock.patch.object( self.bridge, "_flush_overlay_good_outbound_cache", ) as flush_cache: self.bridge._seed_overlay_good_outbound_cache_candidates() mark_candidate.assert_not_called() flush_cache.assert_called_once_with(force=True) self.assertNotIn(peer_hash, self.bridge._overlay_good_outbound_cache) self.assertNotIn(peer_hash, self.bridge._candidate_peers) class PresenceBridgeAccountEndpointLeaseTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() self.bridge._destination = FakeDestination() def lease(self, address, destination, session, verification="direct-bound", offset=0): now_ms = int(time.time() * 1000) return { "address": address, "destinationHash": destination, "sessionId": session, "lastSeen": now_ms + offset, "expiresAt": now_ms + 45_000 + offset, "verification": verification, } def test_one_transport_can_serve_multiple_signed_account_leases(self): destination = "aa" * 16 self.bridge._set_verified_overlay_peers( [{"destinationHash": destination, "lastSeen": int(time.time() * 1000)}], [destination], [ self.lease("Q-account-a", destination, "session-a"), self.lease("Q-account-b", destination, "session-b"), ], ) self.assertNotIn("address", self.bridge._verified_overlay_peers[destination]) self.assertEqual( self.bridge._resolve_verified_game_peer("Q-account-a", destination), destination, ) self.assertEqual( self.bridge._resolve_verified_game_peer("Q-account-b", destination), destination, ) def test_account_switch_removes_only_the_ended_lease(self): destination = "aa" * 16 self.bridge._set_verified_overlay_peers( [{"destinationHash": destination, "lastSeen": int(time.time() * 1000)}], [destination], [ self.lease("Q-account-a", destination, "session-a"), self.lease("Q-account-b", destination, "session-b"), ], ) self.bridge._set_verified_overlay_peers( [{"destinationHash": destination, "lastSeen": int(time.time() * 1000)}], [destination], [self.lease("Q-account-b", destination, "session-b", offset=1)], ) self.assertIsNone( self.bridge._resolve_verified_game_peer("Q-account-a", destination) ) self.assertEqual( self.bridge._resolve_verified_game_peer("Q-account-b", destination), destination, ) self.assertIn(destination, self.bridge._verified_overlay_peers) def test_expired_account_lease_is_never_resolved(self): destination = "aa" * 16 expired = self.lease("Q-account-a", destination, "session-a") expired["expiresAt"] = int(time.time() * 1000) - 1 self.bridge._set_verified_overlay_peers( [{"destinationHash": destination, "lastSeen": int(time.time() * 1000)}], [destination], [expired], ) self.assertIsNone( self.bridge._resolve_verified_game_peer("Q-account-a", destination) ) self.assertIn(destination, self.bridge._verified_overlay_peers) def test_unpreferred_resolution_uses_strongest_fresh_lease(self): relayed = "aa" * 16 direct = "bb" * 16 self.bridge._set_verified_overlay_peers( [ {"destinationHash": relayed, "lastSeen": int(time.time() * 1000)}, {"destinationHash": direct, "lastSeen": int(time.time() * 1000)}, ], [relayed, direct], [ self.lease( "Q-account", relayed, "session-relayed", "relayed-bound", 10 ), self.lease( "Q-account", direct, "session-direct", "direct-bound", 0 ), ], ) self.assertEqual( self.bridge._resolve_verified_game_peer("Q-account"), direct ) class PresenceBridgePinnedChatPeersTest(unittest.TestCase): def setUp(self): self.bridge = load_bridge() def tearDown(self): self.bridge._shutdown.clear() def test_pins_only_current_account_endpoint_leases_and_clears_them(self): peer_hash = "aa" * 16 expires_at = time.time() + 60 self.bridge._account_endpoint_leases = { "Q-account": { peer_hash: { "destination_hash": peer_hash, "expires_at": expires_at, } } } responses = [] with mock.patch.object( self.bridge, "emit_resp", side_effect=lambda *args, **kwargs: responses.append((args, kwargs)), ), mock.patch.object( self.bridge, "_enqueue_scheduler_task", return_value=True ), mock.patch.object( self.bridge, "ensure_known_peer_from_recall", return_value=True ): self.bridge.handle_configure_reticulum_chat_pinned_peers( "pin", { "peers": [ { "accountAddress": "Q-account", "destinationHash": peer_hash, "expiresAt": int((expires_at + 30) * 1000), } ] }, ) self.assertEqual( self.bridge._pinned_chat_overlay_peers, {peer_hash: expires_at}, ) self.assertTrue(responses[-1][1]["payload"]["maintenanceQueued"]) self.bridge.handle_configure_reticulum_chat_pinned_peers( "clear", {"peers": []} ) self.assertEqual(self.bridge._pinned_chat_overlay_peers, {}) def test_rejects_a_pin_without_a_current_endpoint_lease(self): responses = [] with mock.patch.object( self.bridge, "emit_resp", side_effect=lambda *args, **kwargs: responses.append((args, kwargs)), ): self.bridge.handle_configure_reticulum_chat_pinned_peers( "pin", { "peers": [ { "accountAddress": "Q-account", "destinationHash": "bb" * 16, "expiresAt": int((time.time() + 60) * 1000), } ] }, ) self.assertFalse(responses[-1][0][1]) self.assertEqual(self.bridge._pinned_chat_overlay_peers, {}) def test_expired_pin_keeps_an_otherwise_valid_overlay_neighbor(self): peer_hash = "cc" * 16 self.bridge._pinned_chat_overlay_peers[peer_hash] = time.time() - 1 self.bridge._active_overlay_neighbors[peer_hash] = time.time() self.bridge._candidate_peers[peer_hash] = { "last_seen": time.time(), "source": "test", } self.assertEqual(self.bridge._prune_pinned_chat_overlay_peers(), set()) self.assertIn(peer_hash, self.bridge._active_overlay_neighbors) self.assertIn(peer_hash, self.bridge._candidate_peers) if __name__ == "__main__": unittest.main()