diff --git a/pyproject.toml b/pyproject.toml index 6a173a0..aa19992 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -7,7 +7,7 @@ build-backend = "setuptools.build_meta" [project] name = "adn-server" -version = "1.0.0" +version = "2.0.0-alpha1" description = "ADN DMR Peer Server" readme = "README.md" license = { text = "GPL-3.0-or-later" } diff --git a/src/adn_server/__init__.py b/src/adn_server/__init__.py index 1d8b677..0c3cd4e 100644 --- a/src/adn_server/__init__.py +++ b/src/adn_server/__init__.py @@ -37,4 +37,4 @@ """ADN DMR Peer Server — conference bridge (rewrite of bridge_master).""" -__version__ = "1.0.0" +__version__ = "2.0.0-alpha1" diff --git a/src/adn_server/application/bridge/bridge_table.py b/src/adn_server/application/bridge/bridge_table.py index 703d3a8..b211552 100644 --- a/src/adn_server/application/bridge/bridge_table.py +++ b/src/adn_server/application/bridge/bridge_table.py @@ -314,6 +314,7 @@ class BridgeTableMixin: if tg in prohibited_tgs: continue self.make_static_tg(tg, 2, tmout, system) + self._sync_subscription_store() def options_config_for_system(self, system_name: str) -> None: """Update static TG bridges for one system immediately (e.g. when RPTO received). So incoming OBP traffic reaches hotspots without waiting for the 26s options_config_loop.""" @@ -379,6 +380,7 @@ class BridgeTableMixin: _fp = f"{new_ts1}|{new_ts2}|{int(_tmout)}" if sys_cfg.get("_options_static_apply_fp") == _fp: self._restore_prohibited_static_bridge_legs(system_name) + self._sync_subscription_store() return # Legacy: reset TGs that were removed (bridge_master.py 1736-1767) old_ts1 = str(sys_cfg.get("TS1_STATIC") or "").strip() @@ -444,6 +446,7 @@ class BridgeTableMixin: if new_ts1 or new_ts2: logger.info("(OPTIONS) %s static TGs applied: TS1=%s TS2=%s", system_name, new_ts1 or "-", new_ts2 or "-") self._restore_prohibited_static_bridge_legs(system_name) + self._sync_subscription_store() except Exception as e: logger.debug("(OPTIONS) options_config_for_system %s: %s", system_name, e) @@ -827,4 +830,5 @@ class BridgeTableMixin: systems_cfg[_system]["DEFAULT_UA_TIMER"] = int(_options.get("DEFAULT_UA_TIMER", 10)) except Exception as e: logger.exception("(OPTIONS) caught exception: %s", e) + self._sync_subscription_store() diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 8363587..c19d184 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -40,7 +40,8 @@ from bitarray import bitarray from ..domain.dmr import bptc from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4, int_id -from .ports import BridgeRouter, DmrEmbeddedLcEncoder, TalkerAliasEmblcEncoder +from .ports import BridgeRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder +from .subscription.store_sync import replace_store_from_bridges from .talker_alias_use_cases import TalkerAliasUseCases from .bridge.helpers import obp_target_bcsq_quenches_stream, resolve_voice_peer_id from .bridge.timers import BridgeTimerMixin @@ -69,9 +70,11 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi call_later: Any = None, encode_emblc: DmrEmbeddedLcEncoder | None = None, ta_emblc_encoder: TalkerAliasEmblcEncoder | None = None, + subscription_store: SubscriptionStore | None = None, ) -> None: self._router = bridge_router self._config = config + self._subscription_store = subscription_store self._send_to_system = send_to_system # (system_name, packet, **kwargs) -> None self._get_protocols = get_protocols # () -> dict[str, protocol] self._reporting = reporting @@ -101,6 +104,12 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi except Exception: return False + def _sync_subscription_store(self) -> None: + """Mirror BRIDGES into the subscription store (OPTIONS/static paths only for now).""" + if self._subscription_store is None: + return + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + def get_bridges(self) -> dict[str, list[dict[str, Any]]]: """Return current BRIDGES.""" return self._router.get_bridges() diff --git a/src/adn_server/application/subscription/__init__.py b/src/adn_server/application/subscription/__init__.py index 683ba1f..17f4185 100644 --- a/src/adn_server/application/subscription/__init__.py +++ b/src/adn_server/application/subscription/__init__.py @@ -1,6 +1,8 @@ """Subscription application helpers.""" from .bridges_export import export_bridges, subscription_to_legacy_row +from .bridges_import import subscriptions_from_bridges +from .store_sync import replace_store_from_bridges from .bridges_legacy_view import BridgesLegacyView from .ingress import build_voice_ingress from .router import SubscriptionRouter @@ -10,5 +12,7 @@ __all__ = [ "SubscriptionRouter", "build_voice_ingress", "export_bridges", + "replace_store_from_bridges", "subscription_to_legacy_row", + "subscriptions_from_bridges", ] diff --git a/src/adn_server/application/subscription/bridges_import.py b/src/adn_server/application/subscription/bridges_import.py new file mode 100644 index 0000000..0814ded --- /dev/null +++ b/src/adn_server/application/subscription/bridges_import.py @@ -0,0 +1,75 @@ +"""Import legacy ``BRIDGES`` rows into domain subscriptions (mirror of ``bridges_export``).""" + +from __future__ import annotations + +from typing import Any + +from adn_server.domain import int_id +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + Subscription, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + + +def subscriptions_from_bridges(bridges: dict[str, list[dict[str, Any]]]) -> list[Subscription]: + """Build subscriptions from a legacy ``BRIDGES`` snapshot (OPTIONS/static TG / ECHO).""" + subs: list[Subscription] = [] + for table_key, rows in bridges.items(): + for row in rows: + if not isinstance(row, dict): + continue + subs.append(_subscription_from_row(table_key, row)) + return subs + + +def _subscription_from_row(table_key: str, row: dict[str, Any]) -> Subscription: + ts = int(row.get("TS") or 1) + if table_key.startswith("#"): + channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00") + bridge_key = table_key + else: + try: + channel_tgid = int(table_key) + except ValueError: + channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00") + bridge_key = None + to_type = str(row.get("TO_TYPE", "ON")) + timer = row.get("TIMER") + timer_at = float(timer) if isinstance(timer, (int, float)) else None + timeout = row.get("TIMEOUT") + timeout_sec = float(timeout) if isinstance(timeout, (int, float)) else None + return Subscription( + channel=AudioChannel(tgid=TgId(channel_tgid), slot=ts), # type: ignore[arg-type] + system=SystemId(str(row.get("SYSTEM", ""))), + target_tgid=TgId(int_id(row.get("TGID") or b"\x00\x00\x00")), + role=_role_from_to_type(to_type), + policy=_policy_from_to_type(to_type), + state=SubscriptionState( + phase=SubscriptionPhase.ACTIVE if row.get("ACTIVE") else SubscriptionPhase.IDLE, + timer_expires_at=timer_at, + ), + bridge_key=bridge_key, + timeout_seconds=timeout_sec, + ) + + +def _role_from_to_type(to_type: str) -> SubscriptionRole: + if to_type == "NONE": + return SubscriptionRole.ECHO + if to_type == "STAT": + return SubscriptionRole.PASSIVE_STAT + return SubscriptionRole.SINK + + +def _policy_from_to_type(to_type: str) -> ActivationPolicy: + if to_type == "STAT": + return ActivationPolicy.OPENBRIDGE_STAT + if to_type == "OFF": + return ActivationPolicy.STATIC + return ActivationPolicy.INBAND diff --git a/src/adn_server/application/subscription/store_sync.py b/src/adn_server/application/subscription/store_sync.py new file mode 100644 index 0000000..2a268f7 --- /dev/null +++ b/src/adn_server/application/subscription/store_sync.py @@ -0,0 +1,16 @@ +"""Keep ``SubscriptionStore`` aligned with legacy ``BRIDGES`` (read-only mirror; not routing authority yet).""" + +from __future__ import annotations + +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.application.subscription.bridges_import import subscriptions_from_bridges + + +def replace_store_from_bridges( + store: SubscriptionStore, + bridges: dict[str, list[dict[str, Any]]], +) -> None: + """Replace store contents from a full ``BRIDGES`` snapshot.""" + store.replace_all(subscriptions_from_bridges(bridges)) diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index da955ff..bdecbf5 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -36,6 +36,7 @@ from adn_server.application.runtime_context import ( from adn_server.domain import bytes_3 from adn_server.domain.dmr.bptc import encode_emblc from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore from adn_server.infrastructure.config_normalizer import ( ensure_system_runtime_config as _ensure_system_runtime_config, ) @@ -303,6 +304,7 @@ def run_peer_server( call_from_reactor=reactor.callFromThread, ) + subscription_store = InMemorySubscriptionStore() bridge_use_cases = BridgeUseCases( bridge_router, config, @@ -316,6 +318,7 @@ def run_peer_server( call_later=reactor.callLater, encode_emblc=encode_emblc, ta_emblc_encoder=default_ta_emblc_encoder, + subscription_store=subscription_store, ) bridge_use_cases.apply_startup_bridges() report_factory.set_bridges(bridge_router.get_bridges()) diff --git a/tests/application/test_bridges_import.py b/tests/application/test_bridges_import.py new file mode 100644 index 0000000..b5da3f1 --- /dev/null +++ b/tests/application/test_bridges_import.py @@ -0,0 +1,90 @@ +"""BRIDGES → Subscription import (policy / role parity).""" + +from __future__ import annotations + +from adn_server.application.subscription.bridges_import import subscriptions_from_bridges +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import ( + ActivationPolicy, + SubscriptionPhase, + SubscriptionRole, +) + + +def test_import_echo_and_inband_rows(): + tgid_b = bytes_3(9990) + bridges = { + "9990": [ + { + "SYSTEM": "ECHO", + "TS": 2, + "TGID": tgid_b, + "ACTIVE": True, + "TIMEOUT": 120.0, + "TO_TYPE": "NONE", + }, + { + "SYSTEM": "MASTER-A", + "TS": 1, + "TGID": tgid_b, + "ACTIVE": False, + "TIMEOUT": 600.0, + "TO_TYPE": "ON", + }, + ] + } + subs = subscriptions_from_bridges(bridges) + assert len(subs) == 2 + echo, master = subs + assert echo.role == SubscriptionRole.ECHO + assert echo.policy == ActivationPolicy.INBAND + assert echo.state.phase == SubscriptionPhase.ACTIVE + assert master.role == SubscriptionRole.SINK + assert master.policy == ActivationPolicy.INBAND + assert master.state.phase == SubscriptionPhase.IDLE + + +def test_import_static_and_stat_rows(): + bridges = { + "12345": [ + { + "SYSTEM": "MASTER-A", + "TS": 1, + "TGID": bytes_3(12345), + "ACTIVE": True, + "TIMEOUT": 600.0, + "TO_TYPE": "OFF", + }, + { + "SYSTEM": "OBP-CL", + "TS": 1, + "TGID": bytes_3(12345), + "ACTIVE": True, + "TIMEOUT": "", + "TO_TYPE": "STAT", + }, + ] + } + subs = subscriptions_from_bridges(bridges) + static, stat = subs + assert static.policy == ActivationPolicy.STATIC + assert static.role == SubscriptionRole.SINK + assert stat.policy == ActivationPolicy.OPENBRIDGE_STAT + assert stat.role == SubscriptionRole.PASSIVE_STAT + + +def test_import_hash_table_key_sets_bridge_key(): + bridges = { + "#730444": [ + { + "SYSTEM": "OBP-CL", + "TS": 1, + "TGID": bytes_3(730444), + "ACTIVE": True, + "TIMEOUT": "", + "TO_TYPE": "STAT", + } + ] + } + (sub,) = subscriptions_from_bridges(bridges) + assert sub.bridge_key == "#730444" diff --git a/tests/application/test_store_sync.py b/tests/application/test_store_sync.py new file mode 100644 index 0000000..475a02f --- /dev/null +++ b/tests/application/test_store_sync.py @@ -0,0 +1,84 @@ +"""SubscriptionStore sync from BRIDGES (V2-P2-008).""" + +from __future__ import annotations + +from typing import Any + +from adn_server.application.subscription.bridges_export import export_bridges +from adn_server.application.subscription.store_sync import replace_store_from_bridges +from adn_server.domain import bytes_3, int_id +from adn_server.infrastructure.bootstrap.peer_server import _make_echo_bridges +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + +def _row_fingerprint(row: dict[str, Any]) -> tuple[Any, ...]: + return ( + row.get("SYSTEM"), + int(row.get("TS") or 1), + int_id(row.get("TGID") or b"\x00\x00\x00"), + bool(row.get("ACTIVE")), + str(row.get("TO_TYPE", "ON")), + row.get("TIMEOUT"), + ) + + +def _bridges_fingerprints(bridges: dict[str, list[dict[str, Any]]]) -> set[tuple[Any, ...]]: + fps: set[tuple[Any, ...]] = set() + for rows in bridges.values(): + for row in rows: + if isinstance(row, dict): + fps.add(_row_fingerprint(row)) + return fps + + +def test_replace_store_from_echo_bridges_round_trip(): + config = { + "SYSTEMS": { + "ECHO": {"MODE": "PEER"}, + "MASTER-A": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 10}, + "MASTER-B": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 15}, + } + } + bridges = _make_echo_bridges(config) + store = InMemorySubscriptionStore() + replace_store_from_bridges(store, bridges) + + assert len(store.snapshot()) == len(_bridges_fingerprints(bridges)) + + now = 1_700_000_000.0 + exported = export_bridges(store, now=now) + assert _bridges_fingerprints(bridges) == _bridges_fingerprints(exported) + + +def test_replace_store_clears_stale_entries(): + store = InMemorySubscriptionStore() + bridges_a = { + "100": [ + { + "SYSTEM": "SYS-A", + "TS": 1, + "TGID": bytes_3(100), + "ACTIVE": False, + "TIMEOUT": 600.0, + "TO_TYPE": "ON", + } + ] + } + bridges_b = { + "200": [ + { + "SYSTEM": "SYS-B", + "TS": 2, + "TGID": bytes_3(200), + "ACTIVE": True, + "TIMEOUT": 600.0, + "TO_TYPE": "ON", + } + ] + } + replace_store_from_bridges(store, bridges_a) + assert len(store.snapshot()) == 1 + replace_store_from_bridges(store, bridges_b) + assert len(store.snapshot()) == 1 + (sub,) = store.snapshot() + assert int(sub.channel.tgid) == 200 diff --git a/tests/application/test_subscription_router.py b/tests/application/test_subscription_router.py index 56fa139..98c920e 100644 --- a/tests/application/test_subscription_router.py +++ b/tests/application/test_subscription_router.py @@ -4,61 +4,15 @@ from __future__ import annotations from typing import Any +from adn_server.application.subscription.bridges_import import subscriptions_from_bridges from adn_server.application.subscription.router import SubscriptionRouter from adn_server.domain import bytes_3, int_id -from adn_server.domain.subscription import ( - ActivationPolicy, - AudioChannel, - Subscription, - SubscriptionPhase, - SubscriptionRole, - SubscriptionState, - SystemId, - TgId, -) +from adn_server.domain.subscription import TgId from adn_server.domain.voice_routing import ForwardLeg, VoiceIngress from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore -def _role_from_to_type(to_type: str) -> SubscriptionRole: - if to_type == "NONE": - return SubscriptionRole.ECHO - if to_type == "STAT": - return SubscriptionRole.PASSIVE_STAT - return SubscriptionRole.SINK - - -def subscriptions_from_bridges(bridges: dict[str, list[dict[str, Any]]]) -> list[Subscription]: - """Build subscriptions that round-trip ``export_bridges`` shape for parity tests.""" - subs: list[Subscription] = [] - for table_key, rows in bridges.items(): - for row in rows: - if not isinstance(row, dict): - continue - ts = int(row.get("TS") or 1) - if table_key.startswith("#"): - channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00") - else: - channel_tgid = int(table_key) - to_type = str(row.get("TO_TYPE", "ON")) - subs.append( - Subscription( - channel=AudioChannel(tgid=TgId(channel_tgid), slot=ts), # type: ignore[arg-type] - system=SystemId(str(row.get("SYSTEM", ""))), - target_tgid=TgId(int_id(row.get("TGID") or b"\x00\x00\x00")), - role=_role_from_to_type(to_type), - policy=ActivationPolicy.INBAND, - state=SubscriptionState( - phase=SubscriptionPhase.ACTIVE if row.get("ACTIVE") else SubscriptionPhase.IDLE - ), - bridge_key=table_key if table_key.startswith("#") else None, - timeout_seconds=row.get("TIMEOUT") if isinstance(row.get("TIMEOUT"), (int, float)) else None, - ) - ) - return subs - - def legacy_forward_targets( bridges: dict[str, list[dict[str, Any]]], *,