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.pull/1/head
parent
6f8c7548d7
commit
9e0b44128b
@ -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=()),
|
||||
)
|
||||
)
|
||||
@ -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,
|
||||
)
|
||||
@ -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
|
||||
@ -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(),
|
||||
)
|
||||
)
|
||||
@ -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)
|
||||
@ -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()
|
||||
@ -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"
|
||||
Loading…
Reference in new issue