From 9e0b44128b2acd7b3f53426ca7a19f33defca79a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Thu, 11 Jun 2026 14:31:39 -0400 Subject: [PATCH] feat(subscription): P2-015 store-native timers, OPTIONS, and OBP source Move stat_trimmer, bridge_debug, bridge_reset, bridge_table (static TG, reflectors, UA), and OBP source ensure to mutate SubscriptionStore first and export to the BRIDGES shim; batch OPTIONS loops export once at the end. --- .../application/bridge/bridge_table.py | 268 ++++++-- .../application/bridge/obp_forward.py | 16 + src/adn_server/application/bridge/timers.py | 64 +- .../application/bridge_use_cases.py | 7 +- .../subscription/bridge_debug_ops.py | 122 ++++ .../subscription/bridge_reset_ops.py | 111 ++++ .../subscription/bridge_table_ops.py | 620 ++++++++++++++++++ .../subscription/obp_source_ops.py | 72 ++ .../subscription/stat_trimmer_ops.py | 30 + .../test_bridge_table_store_ops.py | 65 ++ .../test_timer_store_ops_slice3.py | 93 +++ 11 files changed, 1395 insertions(+), 73 deletions(-) create mode 100644 src/adn_server/application/subscription/bridge_debug_ops.py create mode 100644 src/adn_server/application/subscription/bridge_reset_ops.py create mode 100644 src/adn_server/application/subscription/bridge_table_ops.py create mode 100644 src/adn_server/application/subscription/obp_source_ops.py create mode 100644 src/adn_server/application/subscription/stat_trimmer_ops.py create mode 100644 tests/application/test_bridge_table_store_ops.py create mode 100644 tests/application/test_timer_store_ops_slice3.py diff --git a/src/adn_server/application/bridge/bridge_table.py b/src/adn_server/application/bridge/bridge_table.py index cf665ff..d639cd4 100644 --- a/src/adn_server/application/bridge/bridge_table.py +++ b/src/adn_server/application/bridge/bridge_table.py @@ -32,6 +32,22 @@ class BridgeTableMixin: _tgid_b = _tgid if isinstance(_tgid, bytes) and len(_tgid) >= 3 else bytes_3(tgid_int) if _tgid_s in ("9990", "9991", "9992", "9993", "9994", "9995", "9996", "9997", "9998", "9999"): _tmout = 1.0 / 6.0 + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import make_single_bridge_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + make_single_bridge_store( + self._subscription_store, + tgid_int, + _sourcesystem, + _slot, + float(_tmout), + self._config.get("SYSTEMS", {}), + time.time(), + ) + self._export_store_to_router() + return timeout_sec = _tmout * 60.0 now = time.time() bridges = self._router.get_bridges() @@ -63,6 +79,18 @@ class BridgeTableMixin: _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 + if self._subscription_store is not None: + from ..subscription.bridge_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(), + ) + return now = time.time() bridges = self._router.get_bridges() bridges[_bridge] = [] @@ -79,6 +107,18 @@ class BridgeTableMixin: 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.""" + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import make_default_reflector_store + + make_default_reflector_store( + self._subscription_store, + reflector, + float(_tmout), + system, + self._config.get("SYSTEMS", {}), + time.time(), + ) + return bridge = "#" + str(reflector) bridges = self._router.get_bridges() if bridge not in bridges: @@ -152,6 +192,23 @@ class BridgeTableMixin: def make_static_tg(self, tg: int, ts: int, _tmout: float, system: str) -> None: """Legacy make_static_tg: ensure bridge for tg exists and set system/ts to ACTIVE/OFF.""" + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import make_static_tg_store + + single_mode = bool( + self._config.get("SYSTEMS", {}).get(system, {}).get("SINGLE_MODE", False) + ) + make_static_tg_store( + self._subscription_store, + tg, + ts, + float(_tmout), + system, + self._config.get("SYSTEMS", {}), + time.time(), + single_mode=single_mode, + ) + return bridges = self._router.get_bridges() key = str(tg) if key not in bridges or not bridges.get(key): @@ -193,6 +250,18 @@ class BridgeTableMixin: 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.""" + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import reset_static_tg_store + + reset_static_tg_store( + self._subscription_store, + tg, + ts, + float(_tmout), + system, + time.time(), + ) + return bridges = self._router.get_bridges() key = str(tg) if key not in bridges: @@ -207,6 +276,16 @@ class BridgeTableMixin: def reset_all_reflector_system(self, _tmout: float, system: str) -> None: """Legacy reset_all_reflector_system: set system's TS2 entry to inactive in every # bridge.""" + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import reset_all_reflector_system_store + + reset_all_reflector_system_store( + self._subscription_store, + float(_tmout), + system, + time.time(), + ) + return bridges = self._router.get_bridges() timeout_sec = _tmout * 60.0 now = time.time() @@ -245,6 +324,19 @@ class BridgeTableMixin: def make_stat_bridge(self, _tgid: bytes) -> None: """Legacy make_stat_bridge: on-the-fly relay bridges for OBP traffic when GEN_STAT_BRIDGES is True.""" _tgid_s = str(int_id(_tgid)) + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import make_stat_bridge_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + make_stat_bridge_store( + self._subscription_store, + _tgid, + self._config.get("SYSTEMS", {}), + time.time(), + ) + self._export_store_to_router() + return bridges = self._router.get_bridges() bridges[_tgid_s] = [] systems_cfg = self._config.get("SYSTEMS", {}) @@ -262,6 +354,14 @@ class BridgeTableMixin: def deactivate_all_dynamic_bridges(self, system_name: str) -> None: """Legacy deactivate_all_dynamic_bridges: deactivate all non-STAT, non-reflector bridges for a system (TG 4000).""" + if self._subscription_store is not None: + from ..subscription.bridge_table_ops import deactivate_all_dynamic_bridges_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + deactivate_all_dynamic_bridges_store(self._subscription_store, system_name) + self._export_store_to_router() + return bridges = self._router.get_bridges() for _bridge in list(bridges): if _bridge not in bridges: @@ -280,6 +380,16 @@ class BridgeTableMixin: 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).""" + if self._subscription_store is not None: + from ..subscription.bridge_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(), + ) + return bridges = self._router.get_bridges() timeout_sec = _tmout * 60.0 now = time.time() @@ -308,6 +418,10 @@ class BridgeTableMixin: def apply_startup_bridges(self) -> None: """Legacy startup: set default reflectors and static TGs for each MASTER system.""" + if self._subscription_store is not None: + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) logger.debug("(ROUTER) Setting default reflectors") for system, sys_cfg in self._config.get("SYSTEMS", {}).items(): @@ -543,6 +657,10 @@ class BridgeTableMixin: ``peer_options`` from RPTO overrides YAML (inject-only proxy: OPTIONS live on each peer). """ + if self._subscription_store is not None and not getattr(self, "_options_store_batch", False): + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) systems_cfg = self._config.get("SYSTEMS", {}) sys_cfg = systems_cfg.get(system_name, {}) @@ -855,83 +973,93 @@ class BridgeTableMixin: def options_config_loop(self) -> None: """Legacy options_config: parse OPTIONS from MASTER systems and update bridges (default reflector, static TGs).""" + batch_store = self._subscription_store is not None + if batch_store: + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + self._options_store_batch = True prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) logger.debug("(OPTIONS) Running options parser") systems_cfg = self._config.get("SYSTEMS", {}) - for _system in list(systems_cfg.keys()): - try: - if systems_cfg.get(_system, {}).get("MODE") != "MASTER": - continue - if not systems_cfg.get(_system, {}).get("ENABLED", True): - continue - opt_str: bytes | str | None = systems_cfg.get(_system, {}).get("OPTIONS") - if opt_str is None: - protocols = self._get_protocols() if self._get_protocols else {} - proto = protocols.get(_system) - peers = getattr(proto, "_peers", {}) if proto is not None else {} - if isinstance(peers, dict): - for peer in peers.values(): - if isinstance(peer, dict) and peer.get("CONNECTION") == "YES" and peer.get("OPTIONS"): - opt_str = peer["OPTIONS"] + try: + for _system in list(systems_cfg.keys()): + try: + if systems_cfg.get(_system, {}).get("MODE") != "MASTER": + continue + if not systems_cfg.get(_system, {}).get("ENABLED", True): + continue + opt_str: bytes | str | None = systems_cfg.get(_system, {}).get("OPTIONS") if opt_str is None: + protocols = self._get_protocols() if self._get_protocols else {} + proto = protocols.get(_system) + peers = getattr(proto, "_peers", {}) if proto is not None else {} + if isinstance(peers, dict): + for peer in peers.values(): + if isinstance(peer, dict) and peer.get("CONNECTION") == "YES" and peer.get("OPTIONS"): + opt_str = peer["OPTIONS"] + if opt_str is None: + continue + _options = self._parse_options_string(opt_str) + if not _options: continue - _options = self._parse_options_string(opt_str) - if not _options: - continue - logger.debug("(OPTIONS) Options found for %s", _system) - if "_opt_key" in systems_cfg[_system] and systems_cfg[_system].get("_opt_key"): - if "KEY" not in _options: - logger.debug("(OPTIONS) %s, options key set but no key in options string, skipping", _system) + logger.debug("(OPTIONS) Options found for %s", _system) + if "_opt_key" in systems_cfg[_system] and systems_cfg[_system].get("_opt_key"): + if "KEY" not in _options: + logger.debug("(OPTIONS) %s, options key set but no key in options string, skipping", _system) + continue + if systems_cfg[_system]["_opt_key"] != _options.get("KEY"): + logger.debug("(OPTIONS) %s, options key set but key sent does not match, skipping", _system) + continue + elif _options.get("KEY"): + systems_cfg[_system]["_opt_key"] = _options["KEY"] + logger.debug("(OPTIONS) %s, _opt_key not set but key sent. Setting to sent key", _system) + else: + systems_cfg[_system]["_opt_key"] = False + logger.debug("(OPTIONS) %s, _opt_key not set and no key sent. Set to false", _system) + self._apply_master_runtime_options(_system, _options) + _options.setdefault("TS1_STATIC", False) + _options.setdefault("TS2_STATIC", False) + _options.setdefault("DEFAULT_REFLECTOR", 0) + _options.setdefault("OVERRIDE_IDENT_TG", False) + _options.setdefault("DEFAULT_UA_TIMER", systems_cfg[_system].get("DEFAULT_UA_TIMER", 10)) + if "TS1_STATIC" not in _options or "TS2_STATIC" not in _options or "DEFAULT_REFLECTOR" not in _options or "DEFAULT_UA_TIMER" not in _options: + logger.debug("(OPTIONS) %s - Required field missing, ignoring", _system) continue - if systems_cfg[_system]["_opt_key"] != _options.get("KEY"): - logger.debug("(OPTIONS) %s, options key set but key sent does not match, skipping", _system) + if _options["TS1_STATIC"] == "": + _options["TS1_STATIC"] = False + if _options["TS2_STATIC"] == "": + _options["TS2_STATIC"] = False + if _options.get("TS1_STATIC") and re.search(r"[^\d,]", str(_options["TS1_STATIC"])): + logger.debug("(OPTIONS) %s - TS1_STATIC contains characters other than numbers and comma, ignoring", _system) continue - elif _options.get("KEY"): - systems_cfg[_system]["_opt_key"] = _options["KEY"] - logger.debug("(OPTIONS) %s, _opt_key not set but key sent. Setting to sent key", _system) - else: - systems_cfg[_system]["_opt_key"] = False - logger.debug("(OPTIONS) %s, _opt_key not set and no key sent. Set to false", _system) - self._apply_master_runtime_options(_system, _options) - _options.setdefault("TS1_STATIC", False) - _options.setdefault("TS2_STATIC", False) - _options.setdefault("DEFAULT_REFLECTOR", 0) - _options.setdefault("OVERRIDE_IDENT_TG", False) - _options.setdefault("DEFAULT_UA_TIMER", systems_cfg[_system].get("DEFAULT_UA_TIMER", 10)) - if "TS1_STATIC" not in _options or "TS2_STATIC" not in _options or "DEFAULT_REFLECTOR" not in _options or "DEFAULT_UA_TIMER" not in _options: - logger.debug("(OPTIONS) %s - Required field missing, ignoring", _system) - continue - if _options["TS1_STATIC"] == "": - _options["TS1_STATIC"] = False - if _options["TS2_STATIC"] == "": - _options["TS2_STATIC"] = False - if _options.get("TS1_STATIC") and re.search(r"[^\d,]", str(_options["TS1_STATIC"])): - logger.debug("(OPTIONS) %s - TS1_STATIC contains characters other than numbers and comma, ignoring", _system) - continue - if _options.get("TS2_STATIC") and re.search(r"[^\d,]", str(_options["TS2_STATIC"])): - logger.debug("(OPTIONS) %s - TS2_STATIC contains characters other than numbers and comma, ignoring", _system) - continue - for key in ("DEFAULT_REFLECTOR", "OVERRIDE_IDENT_TG", "DEFAULT_UA_TIMER"): - if isinstance(_options.get(key), str) and not str(_options[key]).isdigit(): - logger.debug("(OPTIONS) %s - %s is not an integer, ignoring", _system, key) + if _options.get("TS2_STATIC") and re.search(r"[^\d,]", str(_options["TS2_STATIC"])): + logger.debug("(OPTIONS) %s - TS2_STATIC contains characters other than numbers and comma, ignoring", _system) continue - if int(_options.get("DEFAULT_UA_TIMER", 0)) == 0: - _options["DEFAULT_UA_TIMER"] = 35791394 - _tmout = float(int(_options["DEFAULT_UA_TIMER"])) - new_ref = int(_options.get("DEFAULT_REFLECTOR", 0)) - cur_ref = int(systems_cfg[_system].get("DEFAULT_REFLECTOR", 0)) - if new_ref != cur_ref: - if new_ref > 0: - logger.debug("(OPTIONS) %s default reflector changed, updating", _system) - self.reset_all_reflector_system(_tmout, _system) - self.make_default_reflector(new_ref, _tmout, _system) - elif new_ref in prohibited_tgs and not bool(new_ref): - logger.debug("(OPTIONS) %s default reflector is prohibited, ignoring change", _system) - else: - logger.debug("(OPTIONS) %s default reflector disabled, updating", _system) - self.reset_all_reflector_system(_tmout, _system) - self.options_config_for_system(_system) - except Exception as e: - logger.exception("(OPTIONS) caught exception: %s", e) + for key in ("DEFAULT_REFLECTOR", "OVERRIDE_IDENT_TG", "DEFAULT_UA_TIMER"): + if isinstance(_options.get(key), str) and not str(_options[key]).isdigit(): + logger.debug("(OPTIONS) %s - %s is not an integer, ignoring", _system, key) + continue + if int(_options.get("DEFAULT_UA_TIMER", 0)) == 0: + _options["DEFAULT_UA_TIMER"] = 35791394 + _tmout = float(int(_options["DEFAULT_UA_TIMER"])) + new_ref = int(_options.get("DEFAULT_REFLECTOR", 0)) + cur_ref = int(systems_cfg[_system].get("DEFAULT_REFLECTOR", 0)) + if new_ref != cur_ref: + if new_ref > 0: + logger.debug("(OPTIONS) %s default reflector changed, updating", _system) + self.reset_all_reflector_system(_tmout, _system) + self.make_default_reflector(new_ref, _tmout, _system) + elif new_ref in prohibited_tgs and not bool(new_ref): + logger.debug("(OPTIONS) %s default reflector is prohibited, ignoring change", _system) + else: + logger.debug("(OPTIONS) %s default reflector disabled, updating", _system) + self.reset_all_reflector_system(_tmout, _system) + self.options_config_for_system(_system) + except Exception as e: + logger.exception("(OPTIONS) caught exception: %s", e) + finally: + if batch_store: + self._options_store_batch = False self._sync_subscription_store() diff --git a/src/adn_server/application/bridge/obp_forward.py b/src/adn_server/application/bridge/obp_forward.py index 8e33e19..9ded1f4 100644 --- a/src/adn_server/application/bridge/obp_forward.py +++ b/src/adn_server/application/bridge/obp_forward.py @@ -42,6 +42,22 @@ class BridgeObpForwardMixin: return if not (79 <= dst_int < 9990 or dst_int > 9999): return + if self._subscription_store is not None: + from ..subscription.obp_source_ops import ensure_obp_source_for_tg_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + ensure_obp_source_for_tg_store( + self._subscription_store, + system_name, + bridge_key, + dst_id_b, + dst_int, + time.time(), + ) + self._export_store_to_router() + return + bridges = self._router.get_bridges() now = time.time() diff --git a/src/adn_server/application/bridge/timers.py b/src/adn_server/application/bridge/timers.py index fb2bccd..513c1ca 100644 --- a/src/adn_server/application/bridge/timers.py +++ b/src/adn_server/application/bridge/timers.py @@ -134,6 +134,19 @@ class BridgeTimerMixin: def bridge_debug_loop(self) -> None: """Legacy bridgeDebug (bridge_master.py 487-543): remove invalid bridges, fix >1 active dial per MASTER.""" logger.debug("(BRIDGEDEBUG) Running bridge debug") + if self._subscription_store is not None: + from ..subscription.bridge_debug_ops import apply_bridge_debug_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + apply_bridge_debug_store( + self._subscription_store, + self._config.get("SYSTEMS", {}), + time.time(), + ) + self._export_store_to_router() + return + bridges = self._router.get_bridges() systems_cfg = self._config.get("SYSTEMS", {}) now = time.time() @@ -595,6 +608,35 @@ class BridgeTimerMixin: def bridge_reset_loop(self) -> None: """Bridge reset iteration (legacy bridge_reset, 6s). Clear _reset and remove_bridge_system.""" systems_cfg = self._config.get("SYSTEMS", {}) + if self._subscription_store is not None: + from ..subscription.bridge_reset_ops import ( + deactivate_system_legs_store, + restore_prohibited_static_legs_store, + ) + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + now = time.time() + for system_name in list(systems_cfg.keys()): + sys_cfg = systems_cfg.get(system_name, {}) + if not sys_cfg.get("_reset"): + continue + logger.info("(BRIDGERESET) Bridge reset for %s - no peers", system_name) + deactivate_system_legs_store(self._subscription_store, system_name, now) + sys_cfg.pop("_opt_key", None) + sys_cfg.pop("_options_static_apply_fp", None) + restore_prohibited_static_legs_store( + self._subscription_store, + system_name, + sys_cfg, + self.acl_check, + now, + ) + sys_cfg["_reset"] = False + sys_cfg["_resetlog"] = False + self._export_store_to_router() + return + for system_name in list(systems_cfg.keys()): sys_cfg = systems_cfg.get(system_name, {}) if sys_cfg.get("_reset"): @@ -615,10 +657,21 @@ class BridgeTimerMixin: def _restore_prohibited_static_bridge_legs(self, system_name: str) -> None: """After BRIDGERESET / peer RPTO: restore static TGs in prohibited_tgs (parity with _make_echo_bridges).""" - prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) sys_cfg = self._config.get("SYSTEMS", {}).get(system_name, {}) if sys_cfg.get("MODE") != "MASTER" or not sys_cfg.get("ENABLED", True): return + if self._subscription_store is not None: + from ..subscription.bridge_reset_ops import restore_prohibited_static_legs_store + + restore_prohibited_static_legs_store( + self._subscription_store, + system_name, + sys_cfg, + self.acl_check, + time.time(), + ) + return + prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) bridges = self._router.get_bridges() now = time.time() for ts, static_key, acl_key in ( @@ -691,6 +744,15 @@ class BridgeTimerMixin: def stat_trimmer_loop(self) -> None: """Trim STAT-only bridges with no ON/OFF in use (legacy statTrimmer, 303s).""" logger.debug("(ROUTER) STAT trimmer loop started") + if self._subscription_store is not None: + from ..subscription.stat_trimmer_ops import apply_stat_trimmer_store + from ..subscription.store_sync import replace_store_from_bridges + + replace_store_from_bridges(self._subscription_store, self._router.get_bridges()) + apply_stat_trimmer_store(self._subscription_store) + self._export_store_to_router() + return + bridges = self._router.get_bridges() remove_bridges: deque = deque() for bridge_key, entries in list(bridges.items()): diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 96c8bd4..34a0238 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -129,8 +129,11 @@ class BridgeUseCases( pass def _sync_subscription_store(self) -> None: - """Mirror BRIDGES into the subscription store after OPTIONS/static mutations.""" - self._finalize_bridges_state() + """Export subscription store to BRIDGES shim after OPTIONS/static batch mutations.""" + if self._subscription_store is not None: + self._export_store_to_router() + else: + self._finalize_bridges_state() def get_bridges(self) -> dict[str, list[dict[str, Any]]]: """Return current BRIDGES (export shim when store authority is enabled).""" diff --git a/src/adn_server/application/subscription/bridge_debug_ops.py b/src/adn_server/application/subscription/bridge_debug_ops.py new file mode 100644 index 0000000..a7fe172 --- /dev/null +++ b/src/adn_server/application/subscription/bridge_debug_ops.py @@ -0,0 +1,122 @@ +"""Store-native bridge_debug_loop (P2-015).""" + +from __future__ import annotations + +import logging +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.application.subscription.bridges_export import _legacy_to_type +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + InbandTriggers, + Subscription, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + +logger = logging.getLogger(__name__) + +_PROHIBITED_TABLE_KEYS = tuple(str(b) for b in range(10)) + tuple(f"#{b}" for b in range(10)) + + +def apply_bridge_debug_store( + store: SubscriptionStore, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """Remove invalid bridge keys and fix >1 active dial (#) bridge per MASTER.""" + for key in _PROHIBITED_TABLE_KEYS: + for sub in [s for s in store.snapshot() if s.table_key() == key]: + store.remove(sub.subscription_id) + + statroll = sum(1 for sub in store.snapshot() if _legacy_to_type(sub) == "STAT") + + for system, sys_cfg in systems_cfg.items(): + bridgeroll = 0 + dialroll = 0 + activeroll = 0 + for sub in store.snapshot(): + if sub.system.value != system: + continue + bridgeroll += 1 + if sub.is_active(): + if sub.table_key().startswith("#"): + dialroll += 1 + activeroll += 1 + else: + activeroll += 1 + if bridgeroll: + logger.debug( + "(BRIDGEDEBUG) system %s has %s bridges of which %s are in an ACTIVE state", + system, + bridgeroll, + activeroll, + ) + if dialroll > 1 and sys_cfg.get("MODE") == "MASTER": + logger.warning( + "(BRIDGEDEBUG) system %s has more than one active dial bridge (%s) - fixing", + system, + dialroll, + ) + _fix_duplicate_dial_bridges(store, system, sys_cfg, now) + + logger.info("(BRIDGEDEBUG) The server currently has %s STATic bridges", statroll) + + +def _fix_duplicate_dial_bridges( + store: SubscriptionStore, + system: str, + sys_cfg: dict[str, Any], + now: float, +) -> None: + times: dict[float, str] = {} + for sub in store.snapshot(): + if sub.system.value != system or not sub.is_active(): + continue + bridge_key = sub.table_key() + if not bridge_key.startswith("#"): + continue + timer = sub.state.timer_expires_at + if isinstance(timer, (int, float)): + times[float(timer)] = bridge_key + + _tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10)) + timeout_sec = _tmout * 60.0 + system_id = SystemId(system) + + for bridge_key in set(times.values()): + logger.warning("(BRIDGEDEBUG) deactivating system: %s for bridge: %s", system, bridge_key) + try: + setbridge = int(bridge_key[1:]) if bridge_key.startswith("#") else int(bridge_key) + except ValueError: + setbridge = 9 + on_trigger = bytes_3(setbridge) + + for sub in list(store.snapshot()): + if sub.system != system_id or sub.table_key() != bridge_key: + continue + if int(sub.channel.slot) != 2: + continue + store.remove(sub.subscription_id) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(9), slot=2), + system=system_id, + target_tgid=TgId(9), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.INBAND, + state=SubscriptionState( + phase=SubscriptionPhase.IDLE, + timer_expires_at=now + timeout_sec, + ), + bridge_key=bridge_key, + timeout_seconds=timeout_sec, + triggers=InbandTriggers(on=(on_trigger,), off=(), reset=()), + ) + ) diff --git a/src/adn_server/application/subscription/bridge_reset_ops.py b/src/adn_server/application/subscription/bridge_reset_ops.py new file mode 100644 index 0000000..2a47dcd --- /dev/null +++ b/src/adn_server/application/subscription/bridge_reset_ops.py @@ -0,0 +1,111 @@ +"""Store-native bridge_reset_loop helpers (P2-015).""" + +from __future__ import annotations + +import logging +from collections.abc import Callable +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.application.subscription.bridges_export import _legacy_to_type +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + InbandTriggers, + Subscription, + SubscriptionId, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + +logger = logging.getLogger(__name__) + +_PROHIBITED_STATIC_TGS = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) + + +def deactivate_system_legs_store( + store: SubscriptionStore, + system_name: str, + now: float, +) -> None: + """Mirror ``remove_bridge_system``: deactivate all legs for one system.""" + system_id = SystemId(system_name) + for sub in list(store.list_by_system(system_id)): + timeout = sub.timeout_seconds + if timeout is None or isinstance(timeout, str): + timeout = 600.0 + timeout_sec = float(timeout) + target_b = bytes_3(int(sub.target_tgid)) + sub.state.phase = SubscriptionPhase.IDLE + sub.role = SubscriptionRole.SINK + sub.policy = ActivationPolicy.INBAND + sub.state.timer_expires_at = now + timeout_sec + sub.triggers = InbandTriggers(on=(target_b,), off=(), reset=()) + store.upsert(sub) + + +def restore_prohibited_static_legs_store( + store: SubscriptionStore, + system_name: str, + sys_cfg: dict[str, Any], + acl_check: Callable[[bytes, Any], bool], + now: float, +) -> None: + """Restore service (ECHO/NONE) legs for prohibited static TGs after BRIDGERESET.""" + if sys_cfg.get("MODE") != "MASTER" or not sys_cfg.get("ENABLED", True): + return + + system_id = SystemId(system_name) + timeout_sec = (1.0 / 6.0) * 60.0 + + for ts, static_key, acl_key in ( + (1, "TS1_STATIC", "TG1_ACL"), + (2, "TS2_STATIC", "TG2_ACL"), + ): + for tg_s in str(sys_cfg.get(static_key) or "").split(","): + tg_s = tg_s.strip() + if not tg_s: + continue + try: + tg = int(tg_s) + except ValueError: + continue + if tg not in _PROHIBITED_STATIC_TGS: + continue + if sys_cfg.get("USE_ACL") and not acl_check(bytes_3(tg), sys_cfg.get(acl_key, (True, []))): + continue + + channel = AudioChannel(tgid=TgId(tg), slot=ts) + sub_id = SubscriptionId(channel=channel, system=system_id) + existing = store.get(sub_id) + if existing is not None and existing.is_active() and _legacy_to_type(existing) == "NONE": + continue + + service_leg = Subscription( + channel=channel, + system=system_id, + target_tgid=TgId(tg), + role=SubscriptionRole.ECHO, + policy=ActivationPolicy.INBAND, + state=SubscriptionState( + phase=SubscriptionPhase.ACTIVE, + timer_expires_at=now + timeout_sec, + ), + timeout_seconds=timeout_sec, + triggers=InbandTriggers(), + ) + if existing is not None: + store.remove(sub_id) + store.upsert(service_leg) + action = "Restored" if existing is not None else "Re-added" + logger.info( + "(ROUTER) %s service bridge leg: %s bridge %s TS %s", + action, + system_name, + tg, + ts, + ) diff --git a/src/adn_server/application/subscription/bridge_table_ops.py b/src/adn_server/application/subscription/bridge_table_ops.py new file mode 100644 index 0000000..10b83d5 --- /dev/null +++ b/src/adn_server/application/subscription/bridge_table_ops.py @@ -0,0 +1,620 @@ +"""Store-native bridge table mutations (P2-015 slice 4): OPTIONS, static TG, UA, reflectors.""" + +from __future__ import annotations + +import logging +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.application.subscription.bridges_export import _legacy_to_type +from adn_server.domain import bytes_3, int_id +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + InbandTriggers, + Subscription, + SubscriptionId, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + +logger = logging.getLogger(__name__) + +_SERVICE_TG_STRS = frozenset( + str(t) for t in (9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) +) + + +def _effective_tmout_minutes(tgid_int: int, tmout: float) -> float: + if str(tgid_int) in _SERVICE_TG_STRS: + return 1.0 / 6.0 + return tmout + + +def _table_has_legs(store: SubscriptionStore, table_key: str) -> bool: + return any(sub.table_key() == table_key for sub in store.snapshot()) + + +def _remove_table(store: SubscriptionStore, table_key: str) -> None: + for sub in list(store.snapshot()): + if sub.table_key() == table_key: + store.remove(sub.subscription_id) + + +def _find_leg( + store: SubscriptionStore, + system: str, + table_key: str, + slot: int, +) -> Subscription | None: + system_id = SystemId(system) + for sub in store.snapshot(): + if sub.system == system_id and sub.table_key() == table_key and int(sub.channel.slot) == slot: + return sub + return None + + +def _upsert_inband_sink( + store: SubscriptionStore, + *, + system: str, + table_key: str, + channel_tgid: int, + slot: int, + target_tgid: int, + active: bool, + timeout_sec: float, + now: float, + timer_at: float | None = None, + on_trigger: bytes | None = None, + bridge_key: str | None = None, +) -> None: + trigger = on_trigger if on_trigger is not None else bytes_3(target_tgid) + timer = timer_at + if timer is None: + timer = now + timeout_sec if active else now + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(channel_tgid), slot=slot), # type: ignore[arg-type] + system=SystemId(system), + target_tgid=TgId(target_tgid), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.INBAND, + state=SubscriptionState( + phase=SubscriptionPhase.ACTIVE if active else SubscriptionPhase.IDLE, + timer_expires_at=timer, + ), + bridge_key=bridge_key, + timeout_seconds=timeout_sec, + triggers=InbandTriggers(on=(trigger,), off=(), reset=()), + ) + ) + + +def _upsert_static_off( + store: SubscriptionStore, + *, + system: str, + tg: int, + slot: int, + active: bool, + timeout_sec: float, + timer_at: float, +) -> None: + tgid_b = bytes_3(tg) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(tg), slot=slot), # type: ignore[arg-type] + system=SystemId(system), + target_tgid=TgId(tg), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.STATIC, + state=SubscriptionState( + phase=SubscriptionPhase.ACTIVE if active else SubscriptionPhase.IDLE, + timer_expires_at=timer_at, + ), + timeout_seconds=timeout_sec, + triggers=InbandTriggers(on=(tgid_b,), off=(), reset=()), + ) + ) + + +def _upsert_obp_none( + store: SubscriptionStore, + *, + system: str, + table_key: str, + channel_tgid: int, + target_tgid: int, + now: float, + bridge_key: str | None = None, +) -> None: + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(channel_tgid), slot=1), # type: ignore[arg-type] + system=SystemId(system), + target_tgid=TgId(target_tgid), + role=SubscriptionRole.ECHO, + policy=ActivationPolicy.INBAND, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE, timer_expires_at=now), + bridge_key=bridge_key, + timeout_seconds=None, + triggers=InbandTriggers(), + ) + ) + + +def _upsert_stat_obp( + store: SubscriptionStore, + *, + system: str, + tgid_b: bytes, + now: float, +) -> None: + tgid_int = int_id(tgid_b) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(tgid_int), slot=1), # type: ignore[arg-type] + system=SystemId(system), + target_tgid=TgId(tgid_int), + role=SubscriptionRole.PASSIVE_STAT, + policy=ActivationPolicy.OPENBRIDGE_STAT, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE, timer_expires_at=now), + timeout_seconds=None, + triggers=InbandTriggers(), + ) + ) + + +def make_single_bridge_store( + store: SubscriptionStore, + tgid_int: int, + source_system: str, + slot: int, + tmout: float, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """Create bridge table for TG with per-system legs (legacy make_single_bridge).""" + tmout_eff = _effective_tmout_minutes(tgid_int, tmout) + timeout_sec = tmout_eff * 60.0 + table_key = str(tgid_int) + tgid_b = bytes_3(tgid_int) + _remove_table(store, table_key) + + for system, sys_cfg in systems_cfg.items(): + mode = sys_cfg.get("MODE") + if mode == "OPENBRIDGE": + if 79 <= tgid_int < 9990 or tgid_int > 9999: + _upsert_obp_none( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + target_tgid=tgid_int, + now=now, + ) + continue + if system == source_system: + if slot == 1: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=1, + target_tgid=tgid_int, + active=True, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + on_trigger=tgid_b, + ) + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=2, + target_tgid=tgid_int, + active=False, + timeout_sec=timeout_sec, + now=now, + on_trigger=tgid_b, + ) + else: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=2, + target_tgid=tgid_int, + active=True, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + on_trigger=tgid_b, + ) + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=1, + target_tgid=tgid_int, + active=False, + timeout_sec=timeout_sec, + now=now, + on_trigger=tgid_b, + ) + else: + for ts in (1, 2): + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=ts, + target_tgid=tgid_int, + active=False, + timeout_sec=timeout_sec, + now=now, + on_trigger=tgid_b, + ) + + +def make_single_reflector_store( + store: SubscriptionStore, + tgid_int: int, + tmout: float, + source_system: str, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """Create #tgid reflector bridge (legacy make_single_reflector).""" + table_key = "#" + str(tgid_int) + tgid_b = bytes_3(tgid_int) + tmout_eff = _effective_tmout_minutes(tgid_int, tmout) + timeout_sec = tmout_eff * 60.0 + _remove_table(store, table_key) + + for system, sys_cfg in systems_cfg.items(): + mode = sys_cfg.get("MODE") + if mode == "MASTER": + def_ua = float(sys_cfg.get("DEFAULT_UA_TIMER", 10)) * 60.0 + if system == source_system: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=9, + slot=2, + target_tgid=9, + active=True, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + on_trigger=tgid_b, + bridge_key=table_key, + ) + else: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=9, + slot=2, + target_tgid=9, + active=False, + timeout_sec=def_ua, + now=now, + on_trigger=tgid_b, + bridge_key=table_key, + ) + elif mode == "OPENBRIDGE" and (79 <= tgid_int < 9990 or tgid_int > 9999): + _upsert_obp_none( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + target_tgid=tgid_int, + now=now, + bridge_key=table_key, + ) + + +def make_default_reflector_store( + store: SubscriptionStore, + reflector: int, + tmout: float, + system: str, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """Ensure #reflector exists and set system's TS2 leg ACTIVE/OFF (legacy make_default_reflector).""" + table_key = "#" + str(reflector) + if not _table_has_legs(store, table_key): + make_single_reflector_store(store, reflector, tmout, system, systems_cfg, now) + timeout_sec = tmout * 60.0 + reflector_b = bytes_3(reflector) + existing = _find_leg(store, system, table_key, 2) + if existing is not None: + store.remove(existing.subscription_id) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(9), slot=2), + system=SystemId(system), + target_tgid=TgId(9), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.STATIC, + state=SubscriptionState( + phase=SubscriptionPhase.ACTIVE, + timer_expires_at=now + timeout_sec, + ), + bridge_key=table_key, + timeout_seconds=timeout_sec, + triggers=InbandTriggers(on=(reflector_b,), off=(), reset=()), + ) + ) + + +def ensure_master_legs_in_tg_bridge_store( + store: SubscriptionStore, + tg: int, + system: str, + tmout: float, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """Append missing TS1/TS2 MASTER legs on an existing TG table.""" + sys_cfg = systems_cfg.get(system, {}) + if sys_cfg.get("MODE") != "MASTER": + return + table_key = str(tg) + if table_key.startswith("#"): + return + if not _table_has_legs(store, table_key): + return + tmout_eff = _effective_tmout_minutes(tg, tmout) + if tmout_eff <= 0: + tmout_eff = 35791394.0 + timeout_sec = tmout_eff * 60.0 + for ts in (1, 2): + if _find_leg(store, system, table_key, ts) is None: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tg, + slot=ts, + target_tgid=tg, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + ) + + +def make_static_tg_store( + store: SubscriptionStore, + tg: int, + ts: int, + tmout: float, + system: str, + systems_cfg: dict[str, Any], + now: float, + *, + single_mode: bool, +) -> None: + """Ensure TG bridge exists and mark system/ts STATIC active (legacy make_static_tg).""" + table_key = str(tg) + if not _table_has_legs(store, table_key): + make_single_bridge_store(store, tg, system, ts, tmout, systems_cfg, now) + ensure_master_legs_in_tg_bridge_store(store, tg, system, tmout, systems_cfg, now) + timeout_sec = tmout * 60.0 + timer_at = now + timeout_sec + active = True + existing = _find_leg(store, system, table_key, ts) + if existing is not None and single_mode and not existing.is_active(): + active = False + timer_at = float(existing.state.timer_expires_at or timer_at) + _upsert_static_off( + store, + system=system, + tg=tg, + slot=ts, + active=active, + timeout_sec=timeout_sec, + timer_at=timer_at, + ) + + +def reset_static_tg_store( + store: SubscriptionStore, + tg: int, + ts: int, + tmout: float, + system: str, + now: float, +) -> None: + """Deactivate static TG leg (legacy reset_static_tg).""" + table_key = str(tg) + if not _table_has_legs(store, table_key): + return + timeout_sec = tmout * 60.0 + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tg, + slot=ts, + target_tgid=tg, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + ) + + +def reset_all_reflector_system_store( + store: SubscriptionStore, + tmout: float, + system: str, + now: float, +) -> None: + """Deactivate system's TS2 legs in every # bridge (legacy reset_all_reflector_system).""" + timeout_sec = tmout * 60.0 + for sub in list(store.snapshot()): + if sub.system.value != system: + continue + table_key = sub.table_key() + if not table_key.startswith("#"): + continue + if int(sub.channel.slot) != 2: + continue + try: + on_tgid = int(table_key[1:]) + except ValueError: + on_tgid = 9 + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=9, + slot=2, + target_tgid=9, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + on_trigger=bytes_3(on_tgid), + bridge_key=table_key, + ) + + +def make_stat_bridge_store( + store: SubscriptionStore, + tgid_b: bytes, + systems_cfg: dict[str, Any], + now: float, +) -> None: + """On-the-fly STAT relay bridge for OBP (legacy make_stat_bridge).""" + tgid_int = int_id(tgid_b) + table_key = str(tgid_int) + _remove_table(store, table_key) + for system, sys_cfg in systems_cfg.items(): + if sys_cfg.get("MODE") != "OPENBRIDGE": + if sys_cfg.get("MODE") == "MASTER": + tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10)) + timeout_sec = tmout * 60.0 + for ts in (1, 2): + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tgid_int, + slot=ts, + target_tgid=tgid_int, + active=False, + timeout_sec=timeout_sec, + now=now, + on_trigger=tgid_b, + ) + else: + _upsert_stat_obp(store, system=system, tgid_b=tgid_b, now=now) + + +def deactivate_all_dynamic_bridges_store( + store: SubscriptionStore, + system_name: str, +) -> None: + """Deactivate non-STAT, non-reflector legs for a system (TG 4000 path).""" + for sub in list(store.snapshot()): + if sub.system.value != system_name: + continue + bridge_key = sub.table_key() + if bridge_key.startswith("#"): + continue + if _legacy_to_type(sub) == "STAT": + continue + if not sub.is_active(): + continue + sub.state.phase = SubscriptionPhase.IDLE + store.upsert(sub) + logger.info( + "(ROUTER) Deactivated dynamic bridge due to TG/ID 4000: System: %s, Bridge: %s, TS: %s, TGID: %s", + system_name, + bridge_key, + sub.channel.slot, + int(sub.target_tgid), + ) + + +def readd_system_after_ua_timer_change_store( + store: SubscriptionStore, + system: str, + tmout: float, + now: float, +) -> None: + """Re-add missing TS1/TS2 legs after UA timer change (legacy _readd_system_after_ua_timer_change).""" + timeout_sec = tmout * 60.0 + table_keys = {sub.table_key() for sub in store.snapshot()} + for table_key in table_keys: + if table_key.startswith("#"): + if _find_leg(store, system, table_key, 2) is None: + try: + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=9, + slot=2, + target_tgid=9, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + on_trigger=bytes_3(4000), + bridge_key=table_key, + ) + except ValueError: + pass + continue + if _find_leg(store, system, table_key, 1) is None: + try: + tg = int(table_key) + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tg, + slot=1, + target_tgid=tg, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + ) + except ValueError: + pass + if _find_leg(store, system, table_key, 2) is None: + try: + tg = int(table_key) + _upsert_inband_sink( + store, + system=system, + table_key=table_key, + channel_tgid=tg, + slot=2, + target_tgid=tg, + active=False, + timeout_sec=timeout_sec, + now=now, + timer_at=now + timeout_sec, + ) + except ValueError: + pass diff --git a/src/adn_server/application/subscription/obp_source_ops.py b/src/adn_server/application/subscription/obp_source_ops.py new file mode 100644 index 0000000..d437e20 --- /dev/null +++ b/src/adn_server/application/subscription/obp_source_ops.py @@ -0,0 +1,72 @@ +"""Store-native OBP source leg ensure (P2-015 slice 4).""" + +from __future__ import annotations + +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.domain import bytes_3, int_id +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + InbandTriggers, + Subscription, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + + +def _tgid_match(entry_tgid: Any, dst_id_b: bytes, dst_int: int) -> bool: + if entry_tgid == dst_id_b: + return True + try: + return int_id(entry_tgid or b"\x00\x00\x00") == dst_int + except (TypeError, ValueError): + return False + + +def ensure_obp_source_for_tg_store( + store: SubscriptionStore, + system_name: str, + bridge_key: str, + dst_id_b: bytes, + dst_int: int, + now: float, +) -> None: + """Ensure OBP has ACTIVE TS1 source row in main and #reflector tables.""" + for key in (bridge_key, "#" + bridge_key): + if not any(sub.table_key() == key for sub in store.snapshot()): + continue + channel_tgid = dst_int + patched = False + for sub in list(store.snapshot()): + if sub.system.value != system_name: + continue + if sub.table_key() != key: + continue + if int(sub.channel.slot) != 1: + continue + if not _tgid_match(bytes_3(int(sub.target_tgid.value)), dst_id_b, dst_int): + continue + if not sub.is_active(): + sub.state.phase = SubscriptionPhase.ACTIVE + store.upsert(sub) + patched = True + break + if not patched: + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(channel_tgid), slot=1), # type: ignore[arg-type] + system=SystemId(system_name), + target_tgid=TgId(dst_int), + role=SubscriptionRole.ECHO, + policy=ActivationPolicy.INBAND, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE, timer_expires_at=now), + bridge_key=key if key.startswith("#") else None, + timeout_seconds=None, + triggers=InbandTriggers(), + ) + ) diff --git a/src/adn_server/application/subscription/stat_trimmer_ops.py b/src/adn_server/application/subscription/stat_trimmer_ops.py new file mode 100644 index 0000000..54f6cb1 --- /dev/null +++ b/src/adn_server/application/subscription/stat_trimmer_ops.py @@ -0,0 +1,30 @@ +"""Store-native stat_trimmer_loop (P2-015).""" + +from __future__ import annotations + +import logging +from collections import defaultdict + +from adn_server.application.ports import SubscriptionStore +from adn_server.application.subscription.bridges_export import _legacy_to_type +from adn_server.domain.subscription import Subscription + +logger = logging.getLogger(__name__) + + +def apply_stat_trimmer_store(store: SubscriptionStore) -> None: + """Remove STAT-only bridge tables with no ON-active or OFF legs in use.""" + by_table: dict[str, list[Subscription]] = defaultdict(list) + for sub in store.snapshot(): + by_table[sub.table_key()].append(sub) + + for bridge_key, entries in list(by_table.items()): + has_stat = any(_legacy_to_type(sub) == "STAT" for sub in entries) + in_use = any( + (_legacy_to_type(sub) == "ON" and sub.is_active()) or _legacy_to_type(sub) == "OFF" + for sub in entries + ) + if has_stat and not in_use: + for sub in entries: + store.remove(sub.subscription_id) + logger.debug("(ROUTER) STAT bridge %s removed", bridge_key) diff --git a/tests/application/test_bridge_table_store_ops.py b/tests/application/test_bridge_table_store_ops.py new file mode 100644 index 0000000..54e46bd --- /dev/null +++ b/tests/application/test_bridge_table_store_ops.py @@ -0,0 +1,65 @@ +"""P2-015 slice 4: store-native bridge table and OBP source ops.""" + +from __future__ import annotations + +from adn_server.application.subscription.bridge_table_ops import ( + make_single_bridge_store, + make_static_tg_store, +) +from adn_server.application.subscription.obp_source_ops import ensure_obp_source_for_tg_store +from adn_server.application.subscription.store_sync import replace_store_from_bridges +from adn_server.domain import bytes_3 +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore +from tests.harness.deterministic import active_bridge, minimal_config + + +def _store(bridges: dict | None = None) -> InMemorySubscriptionStore: + store = InMemorySubscriptionStore() + if bridges: + replace_store_from_bridges(store, bridges) + return store + + +def test_make_static_tg_store_creates_active_off_leg() -> None: + config = minimal_config(("MASTER-A",)) + store = _store() + make_static_tg_store( + store, + 52090, + 2, + 10.0, + "MASTER-A", + config["SYSTEMS"], + now=1000.0, + single_mode=False, + ) + leg = next(s for s in store.snapshot() if s.system.value == "MASTER-A" and int(s.channel.slot) == 2) + assert leg.is_active() + + +def test_ensure_obp_source_activates_ts1_leg() -> None: + bridges = active_bridge(7305, (("OBP-CL", 1), ("MASTER-A", 2))) + for entry in bridges["7305"]: + if entry["SYSTEM"] == "OBP-CL": + entry["ACTIVE"] = False + store = _store(bridges) + ensure_obp_source_for_tg_store(store, "OBP-CL", "7305", bytes_3(7305), 7305, now=500.0) + obp = next(s for s in store.snapshot() if s.system.value == "OBP-CL") + assert obp.is_active() + + +def test_make_single_bridge_store_obp_leg() -> None: + config = minimal_config(("MASTER-A",)) + add_obp = { + "OBP-CL": { + "MODE": "OPENBRIDGE", + "ENABLED": True, + "IP": "127.0.0.1", + "PORT": 0, + } + } + config["SYSTEMS"].update(add_obp) + store = _store() + make_single_bridge_store(store, 7305, "MASTER-A", 2, 10.0, config["SYSTEMS"], now=1.0) + obp = next(s for s in store.snapshot() if s.system.value == "OBP-CL") + assert obp.is_active() diff --git a/tests/application/test_timer_store_ops_slice3.py b/tests/application/test_timer_store_ops_slice3.py new file mode 100644 index 0000000..95354d6 --- /dev/null +++ b/tests/application/test_timer_store_ops_slice3.py @@ -0,0 +1,93 @@ +"""P2-015: store-native stat_trimmer, bridge_debug, bridge_reset.""" + +from __future__ import annotations + +from adn_server.application.subscription.bridge_debug_ops import apply_bridge_debug_store +from adn_server.application.subscription.bridge_reset_ops import ( + deactivate_system_legs_store, + restore_prohibited_static_legs_store, +) +from adn_server.application.subscription.stat_trimmer_ops import apply_stat_trimmer_store +from adn_server.application.subscription.store_sync import replace_store_from_bridges +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import SubscriptionPhase +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore +from tests.harness.deterministic import active_bridge + + +def _store(bridges: dict) -> InMemorySubscriptionStore: + store = InMemorySubscriptionStore() + replace_store_from_bridges(store, bridges) + return store + + +def test_stat_trimmer_removes_unused_stat_table() -> None: + tg = bytes_3(12345) + bridges = { + "12345": [ + { + "SYSTEM": "OBP-CL", + "TS": 1, + "TGID": tg, + "ACTIVE": True, + "TIMEOUT": "", + "TO_TYPE": "STAT", + "ON": [], + "OFF": [], + "RESET": [], + "TIMER": 0, + }, + { + "SYSTEM": "MASTER-A", + "TS": 2, + "TGID": tg, + "ACTIVE": False, + "TIMEOUT": 600.0, + "TO_TYPE": "ON", + "ON": [tg], + "OFF": [], + "RESET": [], + "TIMER": 0, + }, + ] + } + store = _store(bridges) + apply_stat_trimmer_store(store) + assert store.snapshot() == () + + +def test_bridge_debug_removes_prohibited_numeric_keys() -> None: + bridges = active_bridge(52090, (("MASTER-A", 2),)) + bridges["5"] = list(bridges["52090"]) + store = _store(bridges) + apply_bridge_debug_store(store, {"MASTER-A": {"MODE": "MASTER"}}, now=1000.0) + assert all(s.table_key() != "5" for s in store.snapshot()) + assert any(s.table_key() == "52090" for s in store.snapshot()) + + +def test_deactivate_system_legs_store() -> None: + store = _store(active_bridge(52090, (("MASTER-A", 2), ("MASTER-B", 2)))) + deactivate_system_legs_store(store, "MASTER-A", now=5000.0) + master = next(s for s in store.snapshot() if s.system.value == "MASTER-A") + other = next(s for s in store.snapshot() if s.system.value == "MASTER-B") + assert master.state.phase == SubscriptionPhase.IDLE + assert other.is_active() + + +def test_restore_prohibited_static_leg() -> None: + bridges = active_bridge(9990, (("MASTER-A", 2),)) + for entry in bridges["9990"]: + if entry["SYSTEM"] == "MASTER-A": + entry["ACTIVE"] = False + entry["TO_TYPE"] = "ON" + store = _store(bridges) + sys_cfg = { + "MODE": "MASTER", + "ENABLED": True, + "TS2_STATIC": "9990", + "USE_ACL": False, + } + restore_prohibited_static_legs_store(store, "MASTER-A", sys_cfg, lambda *_a: True, now=100.0) + leg = next(s for s in store.snapshot() if s.system.value == "MASTER-A") + assert leg.is_active() + assert leg.role.value == "echo"