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.
pull/29/head
Rodrigo Pérez 3 months ago
parent 8820d2ac99
commit 667f803fd9

@ -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()

@ -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

@ -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).

@ -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

@ -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,

@ -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)

@ -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)

@ -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",

@ -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

@ -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))

@ -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

@ -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

@ -1,48 +0,0 @@
# ADN DMR Peer Server - infrastructure twisted adapters report pickle legacy
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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)

@ -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)

@ -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

@ -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)

@ -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"),
[

@ -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():

@ -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)

@ -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)

Loading…
Cancel
Save

Powered by TurnKey Linux.