From 667f803fd9bcf6e826f3cb766e925ca000e7c844 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 29 Jun 2026 22:15:43 -0400 Subject: [PATCH] chore: remove dead code (phase 1) Remove 24 functions/methods with zero production callers, including helpers only exercised by tests. Delete the standalone pickle_legacy module. Adapt affected tests to use production equivalents or inline constructions. 569 tests pass. --- src/adn_server/application/proxy/use_cases.py | 19 ---- .../application/report/monitor_topology.py | 12 --- src/adn_server/application/report/payloads.py | 8 -- src/adn_server/application/report/queue.py | 8 -- src/adn_server/application/routing/helpers.py | 90 ------------------- .../application/routing/subscription_table.py | 34 ------- .../subscription/subscription_queries.py | 58 ------------ src/adn_server/domain/__init__.py | 4 +- src/adn_server/domain/result.py | 10 --- .../infrastructure/config_reload.py | 10 --- src/adn_server/infrastructure/mesh/obp_v1.py | 4 - .../security/password_crypto.py | 21 ----- .../twisted_adapters/report/pickle_legacy.py | 48 ---------- .../twisted_adapters/report_server.py | 24 ----- tests/application/test_peer_rf_mode.py | 11 ++- .../application/test_peer_single_downlink.py | 2 - tests/application/test_proxy_use_cases.py | 28 ++---- tests/application/test_report_queue.py | 3 +- .../test_routing_table_legacy_view.py | 18 +--- tests/infrastructure/mesh/test_obp_v1.py | 4 +- 20 files changed, 19 insertions(+), 397 deletions(-) delete mode 100644 src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py diff --git a/src/adn_server/application/proxy/use_cases.py b/src/adn_server/application/proxy/use_cases.py index 7b4096d..ada6aaf 100644 --- a/src/adn_server/application/proxy/use_cases.py +++ b/src/adn_server/application/proxy/use_cases.py @@ -139,25 +139,6 @@ class ProxyUseCases: slot = self._slots.get_by_peer(peer_id) return slot.client if slot else None - def schedule_rpto(self, peer_id: bytes, payload: bytes) -> bool: - """Queue RPTO body for a connected peer (self-service / login options).""" - slot = self._slots.get_by_peer(peer_id) - if slot is None: - return False - self._rpto_queue.enqueue(peer_id, payload) - return True - - def next_pending_rpto(self) -> PendingRpto | None: - """Dequeue one pending RPTO with its client endpoint (for master inject loop).""" - item = self._rpto_queue.dequeue() - if item is None: - return None - peer_id, payload = item - slot = self._slots.get_by_peer(peer_id) - if slot is None: - return None - return PendingRpto(peer_id=peer_id, payload=payload, client=slot.client) - def list_slots(self) -> tuple[ClientSlot, ...]: return self._slots.list_slots() diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py index c3dc06c..007a39b 100644 --- a/src/adn_server/application/report/monitor_topology.py +++ b/src/adn_server/application/report/monitor_topology.py @@ -119,18 +119,6 @@ def expand_inject_proxy_systems( return out -def _slot_for_voice_peer( - peer_key: bytes, - *, - peers: dict[Any, dict[str, Any]], - peer_slots: dict[bytes, int] | None, - max_slots: int, -) -> int | None: - connected = _connected_peers(peers) - slot_map = _resolve_slot_map(connected, peer_slots, max_slots=max_slots) - return slot_map.get(peer_key) - - def _peer_key_from_int(peer_key: Any) -> bytes: if isinstance(peer_key, bytes): return peer_key diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index 1a126c7..3b89c03 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -148,14 +148,6 @@ def peer_options_pass_only(options: Any) -> bool: return list(parsed.keys()) == ["PASS"] and bool(parsed["PASS"].strip()) -def peer_options_pass_valid(options: Any) -> bool: - """True when ``PASS=`` is absent or is the only key with a non-empty value.""" - parsed = _parse_options_kv(options) - if "PASS" not in parsed: - return True - return peer_options_pass_only(options) - - def peer_options_static_valid(options: Any) -> bool: """False when TS1/TS2 static lists contain non-numeric tokens (legacy bridge_master parity). diff --git a/src/adn_server/application/report/queue.py b/src/adn_server/application/report/queue.py index 05dc125..ba403fe 100644 --- a/src/adn_server/application/report/queue.py +++ b/src/adn_server/application/report/queue.py @@ -58,14 +58,6 @@ class BoundedReportQueue: def enqueue_bridge(self, bridges: dict[str, Any], *, incremental: bool = False) -> None: self._pending_bridge = (bridges, incremental) - def pending_count(self) -> int: - n = len(self._events) - if self._pending_config is not None: - n += 1 - if self._pending_bridge is not None: - n += 1 - return n - def drain(self, sender: ReportSender) -> int: """Flush pending work to ``sender``; at most ``max_drain_per_tick`` voice events per call.""" sent = 0 diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 797c785..74abe74 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -98,13 +98,6 @@ def derive_peer_rf_mode(peer: dict[str, Any]) -> str: return RF_MODE_DUPLEX -def apply_peer_rf_mode(peer: dict[str, Any]) -> str: - """Store derived ``RF_MODE`` on the peer after RPTC (or test harness setup).""" - mode = derive_peer_rf_mode(peer) - peer["RF_MODE"] = mode - return mode - - def peer_rf_mode(peer: dict[str, Any]) -> str: """Cached or derived simplex/duplex mode for downlink and monitor.""" cached = peer.get("RF_MODE") @@ -458,26 +451,6 @@ def inject_only_defer_obp_hbp_slot_contention( return source_is_hbp and connected_count > 1 -def _downlink_same_stream_for_peer( - slot_st: dict[str, Any], - peer_id: bytes, - stream_id: bytes, -) -> bool: - """True when ``stream_id`` continues an RF leg owned by ``peer_id`` (downlink fan-out).""" - if not stream_id: - return False - pid = bytes_4(int_id(peer_id)) - if stream_id == slot_st.get("RX_STREAM_ID"): - rx_peer = slot_st.get("RX_PEER") - if rx_peer is not None and bytes_4(int_id(rx_peer)) == pid: - return True - if stream_id == slot_st.get("TX_STREAM_ID"): - tx_peer = slot_st.get("TX_PEER") - if tx_peer is not None and bytes_4(int_id(tx_peer)) == pid: - return True - return False - - def hbp_slot_blocks_group_voice_for_peer( slot_st: dict[str, Any], peer_id: bytes, @@ -1263,24 +1236,6 @@ def export_peer_ua_multi_tgs( return out -def sync_peer_ua_memory_from_store( - peer: dict[str, Any], - peer_id: bytes, - sys_cfg: dict[str, Any], -) -> None: - """Copy persisted UA rows from ``sys_cfg`` onto the live peer dict after DB restore.""" - pk = bytes_4(int_id(peer_id)) - store = sys_cfg.get("_PEER_UA_SESSIONS") - if isinstance(store, dict): - per_peer = store.get(pk) - if isinstance(per_peer, dict) and per_peer: - ua = peer.setdefault("_UA_SESSION", {}) - ua.clear() - for slot, entry in per_peer.items(): - if isinstance(entry, dict): - ua[slot] = dict(entry) - - def restore_peer_ua_entries_to_memory( sys_cfg: dict[str, Any], peer_id: bytes, @@ -1606,51 +1561,6 @@ def remap_dmrd_to_peer_static_slot( return packet[:15] + bytes([new_bits]) + packet[16:] -def repeat_downlink_report_slot( - wire_slot: int, - tgid: int, - peers: dict[Any, Any], - downlink_peer_ids: tuple[bytes, ...], - sys_cfg: dict[str, Any] | None, -) -> int: - """Monitor timeslot for REPEAT downlink START/TX (OBP bridge uses target TS, not wire slot). - - When every downlink peer maps the TG to the same OPTIONS or UA slot, use that slot so - CTABLE chips match RF (cross-slot static/dynamic). Otherwise fall back to wire slot. - """ - display_slots: set[int] = set() - tgid_i = int(tgid) - for peer_id in downlink_peer_ids: - peer = peers.get(peer_id) - if not isinstance(peer, dict): - continue - display_slots.add( - peer_downlink_voice_slot( - peer, wire_slot, tgid_i, sys_cfg, peer_id=peer_id, - ) - ) - if len(display_slots) == 1: - return display_slots.pop() - return int(wire_slot) - - -def peer_single_blocks_uplink( - peer: dict[str, Any], - peer_id: bytes, - slot: int, - tgid: int, - sys_cfg: dict[str, Any] | None, - *, - now: float | None = None, -) -> bool: - """SINGLE=1 never blocks local TX; a new TG replaces the session (see ``register_peer_ua_session``). - - Downlink exclusivity is enforced by :func:`peer_single_blocks_group_voice` only. - """ - del peer, peer_id, slot, tgid, sys_cfg, now - return False - - def _peer_owns_dynamic_ua( peer: dict[str, Any], slot: int, diff --git a/src/adn_server/application/routing/subscription_table.py b/src/adn_server/application/routing/subscription_table.py index 65cab36..5efb772 100644 --- a/src/adn_server/application/routing/subscription_table.py +++ b/src/adn_server/application/routing/subscription_table.py @@ -107,24 +107,6 @@ class SubscriptionTableMixin: time.time(), ) - def make_single_reflector(self, _tgid: bytes | int, _tmout: float, _sourcesystem: str) -> None: - """Legacy make_single_reflector: create reflector bridge #tgid with MASTERs and OBP.""" - _tgid_s = str(int_id(_tgid) if not isinstance(_tgid, int) else _tgid) - _bridge = "#" + _tgid_s - _tgid_b = _tgid if isinstance(_tgid, bytes) and len(_tgid) >= 3 else bytes_3(int(_tgid_s)) - if _tgid_s in ("9990", "9991", "9992", "9993", "9994", "9995", "9996", "9997", "9998", "9999"): - _tmout = 1.0 / 6.0 - from ..subscription.subscription_table_ops import make_single_reflector_store - - make_single_reflector_store( - self._subscription_store, - int(_tgid_s), - float(_tmout), - _sourcesystem, - self._config.get("SYSTEMS", {}), - time.time(), - ) - def make_default_reflector(self, reflector: int, _tmout: float, system: str) -> None: """Legacy make_default_reflector: ensure #reflector bridge exists and set system TS2 to ACTIVE/OFF.""" from ..subscription.subscription_table_ops import make_default_reflector_store @@ -179,11 +161,6 @@ class SubscriptionTableMixin: time.time(), ) - def remove_bridge_system(self, system: str) -> None: - """Deactivate all legs for one system (legacy remove_bridge_system).""" - from ..subscription.subscription_reset_ops import deactivate_system_legs_store - - deactivate_system_legs_store(self._subscription_store, system, time.time()) def ensure_stat_relay(self, _tgid: bytes) -> None: """Legacy ensure_stat_relay: on-the-fly relay bridges for OBP traffic when GEN_STAT_BRIDGES is True.""" _tgid_s = str(int_id(_tgid)) @@ -200,17 +177,6 @@ class SubscriptionTableMixin: from ..subscription.subscription_table_ops import deactivate_all_dynamic_relays_store deactivate_all_dynamic_relays_store(self._subscription_store, system_name) - def _readd_system_after_ua_timer_change(self, system: str, _tmout: float) -> None: - """After remove_bridge_system, re-add system to bridges that no longer have ts1/ts2 (legacy 1624-1639).""" - from ..subscription.subscription_table_ops import readd_system_after_ua_timer_change_store - - readd_system_after_ua_timer_change_store( - self._subscription_store, - system, - float(_tmout), - time.time(), - ) - def apply_startup_subscriptions(self) -> None: """Legacy startup: set default reflectors and static TGs for each MASTER system.""" prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) diff --git a/src/adn_server/application/subscription/subscription_queries.py b/src/adn_server/application/subscription/subscription_queries.py index 7daea34..f1d6f32 100644 --- a/src/adn_server/application/subscription/subscription_queries.py +++ b/src/adn_server/application/subscription/subscription_queries.py @@ -23,7 +23,6 @@ from __future__ import annotations from adn_server.application.ports import SubscriptionStore -from adn_server.domain.subscription import Subscription def store_has_table(store: SubscriptionStore, table_key: str) -> bool: @@ -72,60 +71,3 @@ def active_system_slots_for_tg_in_store( if slot not in slots: slots.append(slot) return tuple(sorted(slots)) - - -def system_has_other_active_bridge_on_slot( - store: SubscriptionStore, - system: str, - slot: int, - incoming_tgid: int, -) -> bool: - """True when another ACTIVE bridge leg occupies ``(system, slot)`` on a different TG. - - Diagnostic / table export only — static OPTIONS bridges stay ACTIVE idle and must - not gate per-peer downlink (see ``peer_voice_slots`` / hangtime in ``downlink.py``). - """ - incoming = int(incoming_tgid) - for sub in store.snapshot(): - if sub.system.value != system: - continue - if int(sub.channel.slot) != int(slot): - continue - if not sub.is_active(): - continue - if int(sub.target_tgid) != incoming: - return True - return False - - -def bridge_timer_active_on_slot( - store: SubscriptionStore, - system: str, - slot: int, - tgid: int, - *, - now: float, -) -> bool: - """True when an ACTIVE bridge leg for ``tgid`` on ``slot`` still has a live timer.""" - tg = int(tgid) - for sub in store.snapshot(): - if sub.system.value != system: - continue - if int(sub.channel.slot) != int(slot): - continue - if int(sub.target_tgid) != tg: - continue - if not sub.is_active(): - continue - exp = sub.state.timer_expires_at - if exp is not None and float(exp) > float(now): - return True - return False - - -def store_legs_for_table( - store: SubscriptionStore, - table_key: str, -) -> tuple[Subscription, ...]: - """All subscription legs in a bridge table.""" - return tuple(sub for sub in store.snapshot() if sub.table_key() == table_key) diff --git a/src/adn_server/domain/__init__.py b/src/adn_server/domain/__init__.py index 63fe32e..d3de20c 100644 --- a/src/adn_server/domain/__init__.py +++ b/src/adn_server/domain/__init__.py @@ -33,7 +33,7 @@ from .hbp_protocol import ( STREAM_TO, VER, ) -from .result import Fail, Result, Success, is_fail, is_ok, unwrap_or +from .result import Fail, Result, Success, is_fail from .subscription import ( ActivationPolicy, AudioChannel, @@ -78,8 +78,6 @@ __all__ = [ "Success", "Fail", "is_fail", - "is_ok", - "unwrap_or", "HBPF_VOICE", "HBPF_VOICE_SYNC", "HBPF_DATA_SYNC", diff --git a/src/adn_server/domain/result.py b/src/adn_server/domain/result.py index 1ef89ed..9bb0879 100644 --- a/src/adn_server/domain/result.py +++ b/src/adn_server/domain/result.py @@ -52,13 +52,3 @@ Result = Success[T] | Fail[E] def is_fail(r: Result[T, E]) -> bool: """Return True if result is Fail.""" return isinstance(r, Fail) - - -def is_ok(r: Result[T, E]) -> bool: - """Return True if result is Success.""" - return isinstance(r, Success) - - -def unwrap_or(r: Result[T, E], default: T) -> T: - """Return value if Success, else default.""" - return r.value if isinstance(r, Success) else default diff --git a/src/adn_server/infrastructure/config_reload.py b/src/adn_server/infrastructure/config_reload.py index 4833522..1ef9112 100644 --- a/src/adn_server/infrastructure/config_reload.py +++ b/src/adn_server/infrastructure/config_reload.py @@ -237,16 +237,6 @@ def reload_server_config( pending_stops: list[defer.Deferred] = [] deferred_starts: list[tuple[str, dict[str, Any], Any, BindSpec | None]] = [] - def _start_listener(name: str, sys_cfg: dict[str, Any], proto: Any) -> None: - if should_bind_udp is not None and not should_bind_udp(name, sys_cfg): - protocols[name] = proto - transports.pop(name, None) - log.info("(CONFIG-RELOAD) %s inject-only (no UDP bind)", name) - return - bind = bind_spec(sys_cfg) - transports[name] = listen_udp(name, bind, proto) - protocols[name] = proto - def _schedule_start(name: str, sys_cfg: dict[str, Any], proto: Any, bind: BindSpec | None) -> None: deferred_starts.append((name, sys_cfg, proto, bind)) diff --git a/src/adn_server/infrastructure/mesh/obp_v1.py b/src/adn_server/infrastructure/mesh/obp_v1.py index b54da20..3cda256 100644 --- a/src/adn_server/infrastructure/mesh/obp_v1.py +++ b/src/adn_server/infrastructure/mesh/obp_v1.py @@ -102,10 +102,6 @@ def verify_bcsq(packet: bytes, passphrase: bytes) -> VerifiedBcsq | None: return VerifiedBcsq(tgid=packet[4:7], stream_id=packet[7:11]) -def build_bcst(passphrase: bytes) -> bytes: - return BCST + obp_hmac_sha1(passphrase, BCST) - - def verify_bcst(packet: bytes, passphrase: bytes) -> bool: if packet[:4] != BCST or len(packet) < 4 + OBP_HMAC_LEN: return False diff --git a/src/adn_server/infrastructure/security/password_crypto.py b/src/adn_server/infrastructure/security/password_crypto.py index fb2aba7..ff73b35 100644 --- a/src/adn_server/infrastructure/security/password_crypto.py +++ b/src/adn_server/infrastructure/security/password_crypto.py @@ -66,24 +66,3 @@ def decrypt_password(encrypted_password: Optional[str], key_path: str = "config/ return decrypted.decode("utf-8") except Exception: return encrypted_password - - -def encrypt_password(password: Optional[str], key_path: str = "config/encryption_key.secret") -> Optional[str]: - """Encrypt a password for storage.""" - if not password: - return password - fernet = get_fernet(key_path) - encrypted = fernet.encrypt(password.encode("utf-8")) - return encrypted.decode("utf-8") - - -def is_encrypted(value: Optional[str], key_path: str = "config/encryption_key.secret") -> bool: - """Return True if value looks like Fernet-encrypted data.""" - if not value: - return False - try: - fernet = get_fernet(key_path) - fernet.decrypt(value.encode("utf-8")) - return True - except Exception: - return False diff --git a/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py b/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py deleted file mode 100644 index c68e415..0000000 --- a/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py +++ /dev/null @@ -1,48 +0,0 @@ -# ADN DMR Peer Server - infrastructure twisted adapters report pickle legacy -# -# Copyright (C) 2026 Rodrigo Pérez, CE5RPY -# -############################################################################### -# This program is free software; you can redistribute it and/or modify -# it under the terms of the GNU General Public License as published by -# the Free Software Foundation; either version 3 of the License, or -# (at your option) any later version. -# -# This program is distributed in the hope that it will be useful, -# but WITHOUT ANY WARRANTY; without even the implied warranty of -# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the -# GNU General Public License for more details. -# -# You should have received a copy of the GNU General Public License -# along with this program; if not, write to the Free Software Foundation, -# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA -############################################################################### - -"""Pickle wire helpers for legacy report v1 monitor shim.""" - -from __future__ import annotations - -import pickle -from typing import Any - -from .opcodes import REPORT_OPCODES - -_PICKLE_PROTOCOL = 2 - - -def encode_config_snd_frame( - systems: dict[str, Any], - *, - protocol: int = _PICKLE_PROTOCOL, -) -> bytes: - """``CONFIG_SND`` opcode + ``pickle.dumps(SYSTEMS, protocol=2)`` (legacy ``hblink.send_config``).""" - return REPORT_OPCODES["CONFIG_SND"] + pickle.dumps(systems, protocol=protocol) - - -def encode_bridge_snd_frame( - bridges: dict[str, Any], - *, - protocol: int = _PICKLE_PROTOCOL, -) -> bytes: - """``BRIDGE_SND`` opcode + ``pickle.dumps(BRIDGES, protocol=2)`` (legacy parity).""" - return REPORT_OPCODES["BRIDGE_SND"] + pickle.dumps(bridges, protocol=protocol) diff --git a/src/adn_server/infrastructure/twisted_adapters/report_server.py b/src/adn_server/infrastructure/twisted_adapters/report_server.py index 6c82428..e6cffe0 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report_server.py +++ b/src/adn_server/infrastructure/twisted_adapters/report_server.py @@ -35,7 +35,6 @@ from twisted.protocols.basic import NetstringReceiver from adn_server.application.ports import ReportMqttPublisher, ReportWireEncoder from adn_server.application.report.monitor_topology import ( expand_inject_proxy_systems, - remap_inject_proxy_voice_event_for_peer, remap_inject_proxy_voice_events, ) from adn_server.application.routing.downlink import DownlinkContext @@ -189,26 +188,3 @@ class ReportServerFactory(Factory): if self._mqtt is not None: self._mqtt.publish_frames(frames) - def send_routing_event_for_peer(self, event: str, peer_key: bytes) -> None: - """Send one remapped voice event for a single hotspot (no fan-out).""" - peer_slots = self._peer_slot_map() if self._peer_slot_map is not None else None - downlink_ctx = None - if self._downlink_ctx_for_system is not None: - target = self._config.get("PROXY", {}).get("TARGET_SYSTEM") - if isinstance(target, str) and target: - downlink_ctx = self._downlink_ctx_for_system(target) - mapped = remap_inject_proxy_voice_event_for_peer( - event, - self._config, - self._systems, - peer_key, - peer_slots, - self._bridges, - downlink_ctx, - ) - if not mapped: - return - frames = self._wire.bridge_event_frames(mapped) - self._broadcast_frames(frames) - if self._mqtt is not None: - self._mqtt.publish_frames(frames) diff --git a/tests/application/test_peer_rf_mode.py b/tests/application/test_peer_rf_mode.py index 012f183..b424142 100644 --- a/tests/application/test_peer_rf_mode.py +++ b/tests/application/test_peer_rf_mode.py @@ -11,11 +11,11 @@ from adn_server.application.routing.helpers import ( RF_MODE_DUPLEX, RF_MODE_SIMPLEX, SIMPLEX_VOICE_SLOT, - apply_peer_rf_mode, derive_peer_rf_mode, peer_downlink_voice_slot, peer_is_simplex, peer_options_static_tg_slot, + peer_rf_mode, remap_dmrd_to_peer_static_slot, ) @@ -41,8 +41,7 @@ def _duplex_peer() -> dict: def test_derive_simplex_from_slots_byte() -> None: peer = {"SLOTS": b"4", "RX_FREQ": b"", "TX_FREQ": b""} assert derive_peer_rf_mode(peer) == RF_MODE_SIMPLEX - assert apply_peer_rf_mode(peer) == RF_MODE_SIMPLEX - assert peer["RF_MODE"] == RF_MODE_SIMPLEX + assert peer_rf_mode(peer) == RF_MODE_SIMPLEX def test_derive_simplex_from_matching_frequencies() -> None: @@ -58,7 +57,7 @@ def test_derive_duplex_from_slots_and_split_frequencies() -> None: def test_simplex_forces_ts2_for_static_and_downlink_remap() -> None: peer = _simplex_peer() - apply_peer_rf_mode(peer) + peer_rf_mode(peer) assert peer_options_static_tg_slot(peer, 7144) == SIMPLEX_VOICE_SLOT assert peer_options_static_tg_slot(peer, 730444) == SIMPLEX_VOICE_SLOT assert peer_downlink_voice_slot(peer, 1, 7144) == SIMPLEX_VOICE_SLOT @@ -74,7 +73,7 @@ def test_simplex_forces_ts2_for_static_and_downlink_remap() -> None: def test_duplex_keeps_cross_slot_static_remap() -> None: peer = _duplex_peer() - apply_peer_rf_mode(peer) + peer_rf_mode(peer) assert peer_downlink_voice_slot(peer, 2, 7144) == 1 burst = DeterministicScenario.voice_burst_spec( PacketSpec(dst_id=7144, slot=2, peer_id=730001, rf_src=730001), @@ -87,6 +86,6 @@ def test_duplex_keeps_cross_slot_static_remap() -> None: def test_topology_peer_row_includes_rf_mode() -> None: peer = _simplex_peer() - apply_peer_rf_mode(peer) + peer_rf_mode(peer) row = _topology_peer_row(730002, peer) assert row["rf_mode"] == RF_MODE_SIMPLEX diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index f1acb7c..0c4fae3 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -30,7 +30,6 @@ from adn_server.application.routing.helpers import ( peer_receives_group_tgid, peer_should_receive_group_voice, peer_single_blocks_group_voice, - peer_single_blocks_uplink, peer_single_exclusive_tgid, register_peer_ua_multi_tg, register_peer_ua_session, @@ -223,7 +222,6 @@ def test_single_never_blocks_uplink_tx_switches_indigo() -> None: peer_id = _peer_id() now = 1_000_000.0 register_peer_ua_session(peer, peer_id, 2, 7305, sys_cfg, now=now) - assert not peer_single_blocks_uplink(peer, peer_id, 2, 730, sys_cfg, now=now + 10) assert peer_single_blocks_group_voice(peer, 2, 730, sys_cfg, peer_id=peer_id, now=now + 10) register_peer_ua_session(peer, peer_id, 2, 730, sys_cfg, now=now + 120) assert not peer_single_blocks_group_voice(peer, 2, 730, sys_cfg, peer_id=peer_id, now=now + 121) diff --git a/tests/application/test_proxy_use_cases.py b/tests/application/test_proxy_use_cases.py index dfaeef7..37fe8d0 100644 --- a/tests/application/test_proxy_use_cases.py +++ b/tests/application/test_proxy_use_cases.py @@ -28,7 +28,8 @@ import pytest from adn_server.application.proxy import ProxyUseCases, peer_id_from_packet from adn_server.domain.proxy import ClientEndpoint -from adn_server.domain.result import is_fail, is_ok +from adn_server.domain.result import is_fail +from adn_server.domain.result import Success from adn_server.domain.value_objects import bytes_4 from adn_server.infrastructure.proxy import InMemoryPendingRptoQueue, InMemoryProxySlotStore @@ -49,12 +50,12 @@ def _service(*, max_peers: int = _MAX_PEERS, black_list: tuple[int, ...] = ()) - def test_attach_creates_session_and_refreshes_client() -> None: svc = _service() first = svc.attach_client(_PEER_A, "10.0.0.1", 62031) - assert is_ok(first) + assert isinstance(first, Success) slot = first.value assert slot.client == ClientEndpoint(host="10.0.0.1", port=62031) again = svc.attach_client(_PEER_A, "10.0.0.2", 62032) - assert is_ok(again) + assert isinstance(again, Success) assert again.value.client == ClientEndpoint(host="10.0.0.2", port=62032) assert len(svc.list_slots()) == 1 @@ -69,7 +70,7 @@ def test_attach_fails_when_max_peers_exceeded() -> None: svc = _service(max_peers=4) for n in range(4): peer = bytes_4(1000 + n) - assert is_ok(svc.attach_client(peer, "10.0.0.1", 62031 + n)) + assert isinstance(svc.attach_client(peer, "10.0.0.1", 62031 + n), Success) assert is_fail(svc.attach_client(bytes_4(9999), "10.0.0.9", 62099)) @@ -79,7 +80,7 @@ def test_detach_removes_session() -> None: removed = svc.detach_client(_PEER_A) assert removed == slot assert svc.resolve_client(_PEER_A) is None - assert is_ok(svc.attach_client(_PEER_B, "10.0.0.5", 62035)) + assert isinstance(svc.attach_client(_PEER_B, "10.0.0.5", 62035), Success) def test_resolve_client_by_peer_id() -> None: @@ -118,25 +119,12 @@ def test_attach_client_assigns_report_slot() -> None: svc = _service(max_peers=4) first = svc.attach_client(_PEER_A, "10.0.0.1", 62031) second = svc.attach_client(_PEER_B, "10.0.0.2", 62032) - assert is_ok(first) - assert is_ok(second) + assert isinstance(first, Success) + assert isinstance(second, Success) assert first.value.report_slot == 0 assert second.value.report_slot == 1 -def test_schedule_and_dequeue_rpto() -> None: - svc = _service() - svc.attach_client(_PEER_A, "10.0.0.1", 62031) - payload = b"TS1=123;TS2=456;" - assert svc.schedule_rpto(_PEER_A, payload) is True - assert svc.schedule_rpto(bytes_4(999), payload) is False - pending = svc.next_pending_rpto() - assert pending is not None - assert pending.peer_id == _PEER_A - assert pending.payload == payload - assert pending.client == ClientEndpoint(host="10.0.0.1", port=62031) - - @pytest.mark.parametrize( ("data", "from_master", "expected"), [ diff --git a/tests/application/test_report_queue.py b/tests/application/test_report_queue.py index b58082c..cb4c97d 100644 --- a/tests/application/test_report_queue.py +++ b/tests/application/test_report_queue.py @@ -57,7 +57,8 @@ def test_enqueue_event_is_non_blocking(): sender = QueuedReportSender(queue, inner) sender.send_routing_event("GROUP VOICE,START,RX,SYS,1,2,3,1,4") assert inner.events == [] - assert queue.pending_count() == 1 + assert queue.drain(inner) == 1 + assert len(inner.events) == 1 def test_drain_delivers_events_before_snapshots(): diff --git a/tests/application/test_routing_table_legacy_view.py b/tests/application/test_routing_table_legacy_view.py index 77ab5a3..608430a 100644 --- a/tests/application/test_routing_table_legacy_view.py +++ b/tests/application/test_routing_table_legacy_view.py @@ -18,12 +18,10 @@ # Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA ############################################################################### -"""RoutingTableLegacyView: subscription store → pickle BRIDGE_SND shim.""" +"""RoutingTableLegacyView: subscription store → legacy BRIDGES shim.""" from __future__ import annotations -import pickle - from adn_server.application.subscription.routing_table_export import export_routing_table from adn_server.application.subscription.routing_table_legacy_view import RoutingTableLegacyView from adn_server.domain import bytes_3 @@ -38,8 +36,6 @@ from adn_server.domain.subscription import ( TgId, ) from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore -from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES -from adn_server.infrastructure.twisted_adapters.report.pickle_legacy import encode_bridge_snd_frame def _sample_store() -> InMemorySubscriptionStore: @@ -62,15 +58,3 @@ def test_generate_matches_export_routing_table(): view = RoutingTableLegacyView(store) now = 1_700_000_000.0 assert view.generate(now=now) == export_routing_table(store, now=now) - - -def test_bridge_snd_frame_is_pickle_protocol_2(): - store = _sample_store() - frame = encode_bridge_snd_frame(RoutingTableLegacyView(store).generate()) - assert frame[:1] == REPORT_OPCODES["BRIDGE_SND"] - bridges = pickle.loads(frame[1:], encoding="bytes") - assert "730444" in bridges - row = bridges["730444"][0] - assert row["SYSTEM"] == "MASTER-A" - assert row["TGID"] == bytes_3(730444) - assert isinstance(row["ACTIVE"], bool) diff --git a/tests/infrastructure/mesh/test_obp_v1.py b/tests/infrastructure/mesh/test_obp_v1.py index 804bc25..e05e5df 100644 --- a/tests/infrastructure/mesh/test_obp_v1.py +++ b/tests/infrastructure/mesh/test_obp_v1.py @@ -24,11 +24,11 @@ from __future__ import annotations from adn_server.domain import bytes_4 from adn_server.infrastructure.hbp_constants import BCKA, BCSQ, BCST, BCVE, DMRD +from adn_server.infrastructure.mesh.obp_v1 import obp_hmac_sha1 from adn_server.infrastructure.mesh.obp_v1 import ( DMRD_V1_WIRE_LEN, build_bcka, build_bcsq, - build_bcst, build_bcve, build_dmrd_v1, verify_bcka, @@ -73,7 +73,7 @@ def test_bcka_roundtrip() -> None: def test_bcst_roundtrip() -> None: - wire = build_bcst(_PASS) + wire = BCST + obp_hmac_sha1(_PASS, BCST) assert wire[:4] == BCST assert verify_bcst(wire, _PASS)