diff --git a/src/adn_server/application/subscription/__init__.py b/src/adn_server/application/subscription/__init__.py new file mode 100644 index 0000000..1bab8e0 --- /dev/null +++ b/src/adn_server/application/subscription/__init__.py @@ -0,0 +1,5 @@ +"""Subscription application helpers.""" + +from .bridges_export import export_bridges, subscription_to_legacy_row + +__all__ = ["export_bridges", "subscription_to_legacy_row"] diff --git a/src/adn_server/application/subscription/bridges_export.py b/src/adn_server/application/subscription/bridges_export.py new file mode 100644 index 0000000..458f8db --- /dev/null +++ b/src/adn_server/application/subscription/bridges_export.py @@ -0,0 +1,89 @@ +"""One-way export: SubscriptionStore → legacy ``BRIDGES`` dict (D-08).""" + +from __future__ import annotations + +import time +from typing import Any + +from adn_server.application.ports import SubscriptionStore +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import ( + ActivationPolicy, + Subscription, + SubscriptionPhase, + SubscriptionRole, +) + +_DEFAULT_TIMEOUT_SEC = 600.0 + + +def subscription_to_legacy_row(sub: Subscription, *, now: float | None = None) -> dict[str, Any]: + """Map one subscription to a legacy ``BRIDGES[table][i]`` row.""" + epoch = time.time() if now is None else now + channel_b = bytes_3(int(sub.channel.tgid)) + target_b = bytes_3(int(sub.target_tgid)) + to_type = _legacy_to_type(sub) + timeout = _legacy_timeout(sub, to_type) + on_list, off_list = _legacy_on_off(sub, channel_b, to_type) + active = sub.state.phase == SubscriptionPhase.ACTIVE + timer = sub.state.timer_expires_at + if timer is None: + timer = epoch + timeout if isinstance(timeout, (int, float)) else epoch + + return { + "SYSTEM": sub.system.value, + "TS": sub.channel.slot, + "TGID": target_b, + "ACTIVE": active, + "TIMEOUT": timeout, + "TO_TYPE": to_type, + "ON": on_list, + "OFF": off_list, + "RESET": [], + "TIMER": timer, + } + + +def export_bridges( + store: SubscriptionStore, + *, + now: float | None = None, +) -> dict[str, list[dict[str, Any]]]: + """Build a legacy ``BRIDGES`` snapshot from the subscription store (export only).""" + epoch = time.time() if now is None else now + bridges: dict[str, list[dict[str, Any]]] = {} + for sub in store.snapshot(): + key = sub.table_key() + bridges.setdefault(key, []).append(subscription_to_legacy_row(sub, now=epoch)) + return bridges + + +def _legacy_to_type(sub: Subscription) -> str: + if sub.role == SubscriptionRole.ECHO: + return "NONE" + if sub.role == SubscriptionRole.PASSIVE_STAT: + return "STAT" + if sub.policy == ActivationPolicy.STATIC and sub.state.phase == SubscriptionPhase.ACTIVE: + return "OFF" + return "ON" + + +def _legacy_timeout(sub: Subscription, to_type: str) -> float | str: + if to_type == "STAT": + return "" + if sub.timeout_seconds is not None: + return sub.timeout_seconds + return _DEFAULT_TIMEOUT_SEC + + +def _legacy_on_off( + sub: Subscription, + channel_b: bytes, + to_type: str, +) -> tuple[list[bytes], list[bytes]]: + del sub + if to_type == "STAT": + return [], [] + if to_type == "NONE": + return [], [] + return [channel_b], [] diff --git a/tests/application/test_bridges_export.py b/tests/application/test_bridges_export.py new file mode 100644 index 0000000..ac81755 --- /dev/null +++ b/tests/application/test_bridges_export.py @@ -0,0 +1,107 @@ +"""One-way BRIDGES export from subscriptions.""" + +from __future__ import annotations + +from adn_server.application.subscription.bridges_export import export_bridges, subscription_to_legacy_row +from adn_server.domain import bytes_3 +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + Subscription, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + +def test_subscription_to_legacy_row_echo(): + now = 1_700_000_000.0 + sub = Subscription( + channel=AudioChannel(tgid=TgId(9990), slot=2), + system=SystemId("ECHO"), + target_tgid=TgId(9990), + role=SubscriptionRole.ECHO, + policy=ActivationPolicy.INBAND, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE, timer_expires_at=now + 120), + timeout_seconds=120.0, + ) + row = subscription_to_legacy_row(sub, now=now) + assert row["SYSTEM"] == "ECHO" + assert row["TS"] == 2 + assert row["TGID"] == bytes_3(9990) + assert row["ACTIVE"] is True + assert row["TO_TYPE"] == "NONE" + assert row["TIMEOUT"] == 120.0 + assert row["TIMER"] == now + 120 + + +def test_subscription_to_legacy_row_stat_obp(): + sub = Subscription( + channel=AudioChannel(tgid=TgId(730444), slot=1), + system=SystemId("OBP-CL"), + target_tgid=TgId(730444), + role=SubscriptionRole.PASSIVE_STAT, + policy=ActivationPolicy.OPENBRIDGE_STAT, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE), + ) + row = subscription_to_legacy_row(sub) + assert row["TO_TYPE"] == "STAT" + assert row["TIMEOUT"] == "" + assert row["ON"] == [] + + +def test_export_bridges_groups_by_table_key(): + store = InMemorySubscriptionStore() + channel = AudioChannel(tgid=TgId(730444), slot=1) + store.upsert( + Subscription( + channel=channel, + system=SystemId("MASTER-A"), + target_tgid=TgId(730444), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.INBAND, + state=SubscriptionState(phase=SubscriptionPhase.IDLE), + ) + ) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(730444), slot=2), + system=SystemId("MASTER-A"), + target_tgid=TgId(730444), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.INBAND, + state=SubscriptionState(phase=SubscriptionPhase.IDLE), + ) + ) + store.upsert( + Subscription( + channel=AudioChannel(tgid=TgId(9990), slot=1), + system=SystemId("OBP-CL"), + target_tgid=TgId(9990), + role=SubscriptionRole.PASSIVE_STAT, + policy=ActivationPolicy.OPENBRIDGE_STAT, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE), + bridge_key="730444", + ) + ) + bridges = export_bridges(store) + assert set(bridges.keys()) == {"730444"} + assert len(bridges["730444"]) == 3 + + +def test_static_active_uses_off_to_type(): + sub = Subscription( + channel=AudioChannel(tgid=TgId(12345), slot=1), + system=SystemId("MASTER-A"), + target_tgid=TgId(12345), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.STATIC, + state=SubscriptionState(phase=SubscriptionPhase.ACTIVE), + ) + row = subscription_to_legacy_row(sub) + assert row["TO_TYPE"] == "OFF" + assert row["ACTIVE"] is True + assert row["ON"] == [bytes_3(12345)]