diff --git a/adn-server.example.yaml b/adn-server.example.yaml index 094425a..4c6960f 100644 --- a/adn-server.example.yaml +++ b/adn-server.example.yaml @@ -10,10 +10,6 @@ GLOBAL: TGID_TS1_ACL: PERMIT:ALL TGID_TS2_ACL: PERMIT:ALL GEN_STAT_BRIDGES: true - # Phase 2b: route voice via SubscriptionStore mirror (BRIDGES still mutates ACTIVE/timers). - USE_SUBSCRIPTION_ROUTER: false - # Phase 2b: subscription store is routing authority; BRIDGES is export shim (off until RF sign-off). - USE_SUBSCRIPTION_STORE_AUTHORITY: false SERVER_ID: 73010 VALIDATE_SERVER_IDS: true URL_SECURITY: "" diff --git a/src/adn_server/application/bridge/bridge_table.py b/src/adn_server/application/bridge/bridge_table.py index 1b2a504..cf665ff 100644 --- a/src/adn_server/application/bridge/bridge_table.py +++ b/src/adn_server/application/bridge/bridge_table.py @@ -189,6 +189,7 @@ class BridgeTableMixin: else: bridgetemp.append(bridgesystem) bridges[key] = list(bridgetemp) + self._finalize_bridges_state() def reset_static_tg(self, tg: int, ts: int, _tmout: float, system: str) -> None: """Legacy reset_static_tg: set system/ts entry to ACTIVE False, TO_TYPE ON.""" @@ -504,8 +505,8 @@ class BridgeTableMixin: def _apply_master_runtime_options(self, system_name: str, _options: dict[str, Any]) -> None: """Apply SINGLE/TIMER/VOICE/LANG from peer OPTIONS over YAML defaults (legacy options_config). - When ``USE_SUBSCRIPTION_STORE_AUTHORITY`` is off, BRIDGES mutates first and the store - mirrors via ``_finalize_bridges_state``. When on, exported BRIDGES is the store shim. + BRIDGES mutates first; ``_finalize_bridges_state`` mirrors into the subscription store + and re-exports the legacy shim for report/debug. """ systems_cfg = self._config.get("SYSTEMS", {}) sys_cfg = systems_cfg.get(system_name, {}) diff --git a/src/adn_server/application/bridge/store_authority_mixin.py b/src/adn_server/application/bridge/store_authority_mixin.py index c636446..b84621d 100644 --- a/src/adn_server/application/bridge/store_authority_mixin.py +++ b/src/adn_server/application/bridge/store_authority_mixin.py @@ -20,30 +20,25 @@ class StoreAuthorityMixin: _config: dict[str, Any] _bridges_legacy_view: BridgesLegacyView | None - def _use_subscription_store_authority(self) -> bool: - if self._subscription_store is None: - return False - return bool(self._config.get("GLOBAL", {}).get("USE_SUBSCRIPTION_STORE_AUTHORITY", False)) - def _bridges_for_report(self) -> dict[str, list[dict[str, Any]]]: - """BRIDGES snapshot for monitor/report (export shim when store authority is on).""" - if self._subscription_store is not None and self._use_subscription_store_authority(): - view = getattr(self, "_bridges_legacy_view", None) - if view is None: - view = BridgesLegacyView(self._subscription_store) - self._bridges_legacy_view = view - return view.generate() - return self._router.get_bridges() + """BRIDGES snapshot for monitor/report (export shim from subscription store).""" + if self._subscription_store is None: + return self._router.get_bridges() + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + view = getattr(self, "_bridges_legacy_view", None) + if view is None: + view = BridgesLegacyView(self._subscription_store) + self._bridges_legacy_view = view + return view.generate() def _sync_store_for_voice_lookup(self) -> None: """Import router BRIDGES into the store (e.g. OBP ``_ensure_obp_source_for_tg``) before resolve.""" replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) - if self._use_subscription_store_authority(): - self._router.set_bridges(export_bridges(self._subscription_store)) + self._router.set_bridges(export_bridges(self._subscription_store)) self._router.rebuild_source_index() def _finalize_bridges_state(self) -> None: - """Import BRIDGES mutations into the store; export back when store is authority.""" + """Import BRIDGES mutations into the store; export back as legacy shim.""" if self._subscription_store is None: self._router.rebuild_source_index() return diff --git a/src/adn_server/application/bridge/voice_subscription.py b/src/adn_server/application/bridge/voice_subscription.py index 04e583b..4e0ccfc 100644 --- a/src/adn_server/application/bridge/voice_subscription.py +++ b/src/adn_server/application/bridge/voice_subscription.py @@ -1,7 +1,8 @@ """Subscription router helpers for the voice hot path (P2-009 / P2-010). BRIDGES mutates ACTIVE/timers; ``_finalize_bridges_state`` keeps the store aligned -after timer/OPTIONS paths. Per-packet store sync is skipped when store authority is on. +after timer/OPTIONS paths. Voice resolve always uses ``SubscriptionRouter`` when a +store is wired (production default). """ from __future__ import annotations @@ -17,16 +18,11 @@ logger = logging.getLogger(__name__) class VoiceSubscriptionMixin: - """Wire ``SubscriptionRouter`` into ``dmrd_received`` when configured.""" + """Wire ``SubscriptionRouter`` into ``dmrd_received``.""" _subscription_store: Any _subscription_router: SubscriptionRouter | None - def _use_subscription_router(self) -> bool: - if self._subscription_store is None: - return False - return bool(self._config.get("GLOBAL", {}).get("USE_SUBSCRIPTION_ROUTER", False)) - def _subscription_router_instance(self) -> SubscriptionRouter | None: if self._subscription_store is None: return None @@ -102,8 +98,8 @@ class VoiceSubscriptionMixin: dst_int: int, ingress_required: bool = True, ) -> tuple[tuple[str, ...], frozenset[tuple[str, int, int]] | None]: - """Return bridge tables and optional forward-leg filter (``None`` = legacy row scan).""" - if not self._use_subscription_router(): + """Return bridge tables and optional forward-leg filter (``None`` = no subscription store).""" + if self._subscription_store is None: return ( tuple(self._router.bridge_tables_with_active_source(system_name, bridge_match_slot, dst_int)), None, diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index f364d96..bb46dda 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -391,7 +391,7 @@ class BridgeUseCases( source_lc = b"\x00\x00\x20" + dst_id_b + rf_src # Legacy bridge_master routerOBP: _sysIgnore accumulates across each to_target(BRIDGES[_bridge]) # pass; dedupe (SYSTEM, TS) for OpenBridge targets so the same leg is not sent twice per packet. - # When USE_SUBSCRIPTION_ROUTER is on, resolve() already applies OBP dedup — skip sys_ignore_obp. + # SubscriptionRouter.resolve() already applies OBP dedup — skip sys_ignore_obp when using leg filter. forward_tables, forward_leg_key_set = self._voice_forward_plan( system_name=system_name, peer_id=peer_id, diff --git a/tests/application/test_store_authority_timers.py b/tests/application/test_store_authority_timers.py index f1af760..cd0c1b9 100644 --- a/tests/application/test_store_authority_timers.py +++ b/tests/application/test_store_authority_timers.py @@ -24,13 +24,12 @@ def test_rule_timer_syncs_subscription_store() -> None: assert master_a.state.phase == SubscriptionPhase.IDLE -def test_finalize_exports_when_store_authority_enabled() -> None: +def test_finalize_exports_from_subscription_store() -> None: bridges = active_bridge(730, (("MASTER-A", 2),)) bridges["730"][0]["ACTIVE"] = True bridges["730"][0]["TO_TYPE"] = "ON" scenario = DeterministicScenario(bridges=bridges) - scenario.config["GLOBAL"]["USE_SUBSCRIPTION_STORE_AUTHORITY"] = True scenario.bridge._finalize_bridges_state() exported = scenario.bridge.get_bridges() diff --git a/tests/application/test_voice_subscription_plan.py b/tests/application/test_voice_subscription_plan.py index 59edd79..516e8c1 100644 --- a/tests/application/test_voice_subscription_plan.py +++ b/tests/application/test_voice_subscription_plan.py @@ -13,17 +13,13 @@ from tests.application.test_subscription_router import _row from tests.harness.deterministic import minimal_config -def _bridge( - bridges: dict, - *, - use_subscription_router: bool, -) -> BridgeUseCases: +def _bridge(bridges: dict, *, with_store: bool) -> BridgeUseCases: config = minimal_config(("MASTER-A", "MASTER-B")) - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = use_subscription_router router = InMemoryBridgeRouter() router.set_bridges(bridges) - store = InMemorySubscriptionStore() - replace_store_from_bridges(store, bridges) + store = InMemorySubscriptionStore() if with_store else None + if store is not None: + replace_store_from_bridges(store, bridges) return BridgeUseCases( router, config, @@ -33,14 +29,15 @@ def _bridge( ) -def test_forward_plan_disabled_uses_legacy_tables(): +def test_forward_plan_without_store_uses_router_tables_only(): + """Harness without a store falls back to BRIDGES index (unit-test shortcut).""" bridges = { "730444": [ _row(system="MASTER-A", ts=1, tgid=730444, active=True), _row(system="MASTER-B", ts=1, tgid=730444, active=True), ] } - bridge = _bridge(bridges, use_subscription_router=False) + bridge = _bridge(bridges, with_store=False) tables, leg_keys = bridge._voice_forward_plan( system_name="MASTER-A", peer_id=b"\x00\x00\x03\xe9", @@ -57,7 +54,7 @@ def test_forward_plan_disabled_uses_legacy_tables(): assert leg_keys is None -def test_forward_plan_enabled_returns_leg_filter(): +def test_forward_plan_with_store_returns_leg_filter(): bridges = { "730444": [ _row(system="MASTER-A", ts=1, tgid=730444, active=True), @@ -65,7 +62,7 @@ def test_forward_plan_enabled_returns_leg_filter(): _row(system="OBP-CL", ts=1, tgid=730444, active=False), ] } - bridge = _bridge(bridges, use_subscription_router=True) + bridge = _bridge(bridges, with_store=True) tables, leg_keys = bridge._voice_forward_plan( system_name="MASTER-A", peer_id=b"\x00\x00\x03\xe9", diff --git a/tests/bridge/test_peer_options_override.py b/tests/bridge/test_peer_options_override.py index 79fa8d9..55a255d 100644 --- a/tests/bridge/test_peer_options_override.py +++ b/tests/bridge/test_peer_options_override.py @@ -20,10 +20,8 @@ from adn_server.domain import bytes_3, bytes_4 def _proxy_system_scenario( *, single_mode_yaml: bool = False, - use_subscription_router: bool = False, ) -> DeterministicScenario: config = DeterministicScenario().config - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = use_subscription_router sys_cfg = config["SYSTEMS"]["MASTER-A"] sys_cfg["SINGLE_MODE"] = single_mode_yaml sys_cfg["DEFAULT_UA_TIMER"] = 60 @@ -46,12 +44,8 @@ def _proxy_system_scenario( return scenario -@pytest.mark.parametrize("use_subscription_router", [False, True]) -def test_rpto_single_and_timer_override_yaml(use_subscription_router: bool) -> None: - scenario = _proxy_system_scenario( - single_mode_yaml=False, - use_subscription_router=use_subscription_router, - ) +def test_rpto_single_and_timer_override_yaml() -> None: + scenario = _proxy_system_scenario(single_mode_yaml=False) scenario.bridge.options_config_for_system( "MASTER-A", @@ -64,14 +58,8 @@ def test_rpto_single_and_timer_override_yaml(use_subscription_router: bool) -> N assert scenario.bridge.get_bridges()["52090"][0]["TIMEOUT"] == 300.0 -@pytest.mark.parametrize("use_subscription_router", [False, True]) -def test_options_loop_reads_connected_peer_without_yaml_options( - use_subscription_router: bool, -) -> None: - scenario = _proxy_system_scenario( - single_mode_yaml=False, - use_subscription_router=use_subscription_router, - ) +def test_options_loop_reads_connected_peer_without_yaml_options() -> None: + scenario = _proxy_system_scenario(single_mode_yaml=False) bridges = active_bridge(730444, (("MASTER-A", 2), ("MASTER-B", 2))) scenario.router.set_bridges(bridges) @@ -80,14 +68,8 @@ def test_options_loop_reads_connected_peer_without_yaml_options( assert scenario.config["SYSTEMS"]["MASTER-A"]["SINGLE_MODE"] is True -@pytest.mark.parametrize("use_subscription_router", [False, True]) -def test_single_mode_deactivates_other_static_tg_after_rpto( - use_subscription_router: bool, -) -> None: - scenario = _proxy_system_scenario( - single_mode_yaml=False, - use_subscription_router=use_subscription_router, - ) +def test_single_mode_deactivates_other_static_tg_after_rpto() -> None: + scenario = _proxy_system_scenario(single_mode_yaml=False) bridges = { "730444": [ { @@ -153,16 +135,15 @@ def test_single_mode_deactivates_other_static_tg_after_rpto( assert scenario.bridge.get_bridges()["730444"][0]["ACTIVE"] is True assert scenario.bridge.get_bridges()["52090"][0]["ACTIVE"] is False - if use_subscription_router: - scenario.bridge._sync_subscription_store() - legs = SubscriptionRouter(scenario.subscription_store).resolve( - VoiceIngress( - source_system="MASTER-B", - slot=2, - dst_tgid=TgId(52090), - ) + scenario.bridge._sync_subscription_store() + legs = SubscriptionRouter(scenario.subscription_store).resolve( + VoiceIngress( + source_system="MASTER-B", + slot=2, + dst_tgid=TgId(52090), ) - assert all(leg.target_system != "MASTER-A" for leg in legs) + ) + assert all(leg.target_system != "MASTER-A" for leg in legs) def test_ua_bridge_creation_pushes_routing_snapshot() -> None: @@ -207,7 +188,6 @@ def test_ua_bridge_uses_transmitting_peer_timer_minutes() -> None: def test_peer_timer_does_not_override_other_peer_static_tg_timeout() -> None: """Each hotspot TIMER applies only to that peer's static TGs (no max() across peers).""" config = DeterministicScenario().config - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = False config["SYSTEMS"]["MASTER-A"]["DEFAULT_UA_TIMER"] = 60 scenario = DeterministicScenario(config=config, bridges={}) scenario.bridge._get_protocols = lambda: scenario.protocols # noqa: SLF001 diff --git a/tests/bridge/test_subscription_router_dmrd.py b/tests/bridge/test_subscription_router_dmrd.py index b5beaef..52e0560 100644 --- a/tests/bridge/test_subscription_router_dmrd.py +++ b/tests/bridge/test_subscription_router_dmrd.py @@ -1,4 +1,4 @@ -"""P2-009: dmrd_received with USE_SUBSCRIPTION_ROUTER enabled.""" +"""P2-009: dmrd_received routes voice via SubscriptionRouter.""" from __future__ import annotations @@ -17,9 +17,8 @@ from tests.harness.deterministic import ( @pytest.mark.behavior def test_subscription_router_startup_bridge_voice_parity() -> None: - """Same forwards as legacy BRIDGES scan when USE_SUBSCRIPTION_ROUTER is on.""" + """Startup static TG forwards across masters via subscription resolve.""" config = minimal_config(("MASTER-A", "MASTER-B")) - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = True config["SYSTEMS"]["MASTER-A"]["TS2_STATIC"] = "52090" config["SYSTEMS"]["MASTER-A"]["DEFAULT_UA_TIMER"] = 10 bridges = active_bridge(52090, (("MASTER-A", 2), ("MASTER-B", 2))) @@ -37,27 +36,9 @@ def test_subscription_router_startup_bridge_voice_parity() -> None: assert_forwarded(scenario, "MASTER-B", count=2, dst_id=52090) -@pytest.mark.behavior -def test_subscription_router_matches_legacy_flag_off() -> None: - """USE_SUBSCRIPTION_ROUTER=false keeps legacy BRIDGES scan (default).""" - config = minimal_config(("MASTER-A", "MASTER-B")) - bridges = active_bridge(91, (("MASTER-A", 2), ("MASTER-B", 2))) - scenario = DeterministicScenario(config=config, bridges=bridges) - - base = PacketSpec(dst_id=91, stream_id=0x01020304, slot=2) - scenario.inject_hbp("MASTER-A", DeterministicScenario.voice_head_spec(base)) - scenario.inject_hbp( - "MASTER-A", - DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), - ) - - assert_forwarded(scenario, "MASTER-B", count=2, dst_id=91) - - @pytest.mark.behavior def test_subscription_router_hbp_slot1() -> None: config = minimal_config(("MASTER-A", "MASTER-B")) - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = True bridges = active_bridge(730444, (("MASTER-A", 1), ("MASTER-B", 1))) scenario = DeterministicScenario(config=config, bridges=bridges) replace_store_from_bridges(scenario.subscription_store, scenario.bridge.get_bridges()) @@ -73,11 +54,9 @@ def test_subscription_router_hbp_slot1() -> None: @pytest.mark.behavior -def test_subscription_router_with_store_authority_parity() -> None: - """P2-012: router + store authority export same forwards as legacy.""" +def test_subscription_router_store_export_parity() -> None: + """Store authority export keeps forwards and ACTIVE visible in get_bridges().""" config = minimal_config(("MASTER-A", "MASTER-B")) - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = True - config["GLOBAL"]["USE_SUBSCRIPTION_STORE_AUTHORITY"] = True bridges = active_bridge(52090, (("MASTER-A", 2), ("MASTER-B", 2))) scenario = DeterministicScenario(config=config, bridges=bridges) scenario.bridge._finalize_bridges_state() @@ -94,11 +73,9 @@ def test_subscription_router_with_store_authority_parity() -> None: @pytest.mark.behavior -def test_obp_to_system_with_store_authority_and_router() -> None: - """OBP RX must forward after _ensure_obp_source when store authority skips stale mirror.""" +def test_obp_to_system_forwards_after_obp_source_sync() -> None: + """OBP RX must forward after _ensure_obp_source syncs into the subscription store.""" config = minimal_config(("SYSTEM",)) - config["GLOBAL"]["USE_SUBSCRIPTION_ROUTER"] = True - config["GLOBAL"]["USE_SUBSCRIPTION_STORE_AUTHORITY"] = True add_openbridge_system(config, "OBP-CL") config["SYSTEMS"]["SYSTEM"]["TS2_STATIC"] = "7305" bridges = active_bridge(7305, (("OBP-CL", 1), ("SYSTEM", 2)))