Export SubscriptionStore snapshots to legacy BRIDGES dict rows for future hot-path parity without writing subscriptions from BRIDGES input.pull/1/head
parent
c734b1f0b0
commit
d0abe7c09e
@ -0,0 +1,5 @@
|
||||
"""Subscription application helpers."""
|
||||
|
||||
from .bridges_export import export_bridges, subscription_to_legacy_row
|
||||
|
||||
__all__ = ["export_bridges", "subscription_to_legacy_row"]
|
||||
@ -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], []
|
||||
@ -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)]
|
||||
Loading…
Reference in new issue