feat(subscription): remove USE_SUBSCRIPTION_* flags (P2-014)

Subscription store and router are always active when wired at bootstrap;
drop YAML toggles and flag-off test paths. Sync router into store before
get_bridges() export and after make_static_tg mutations.
pull/1/head
Rodrigo Pérez 4 months ago
parent c10f7f899a
commit f680fe4e58

@ -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: "<set-in-adn-server.yaml>"

@ -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, {})

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

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

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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.