diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index 919d2ad..e457a3e 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -27,7 +27,10 @@ from __future__ import annotations from abc import ABC, abstractmethod from collections.abc import Sequence -from typing import Any, Protocol +from typing import TYPE_CHECKING, Any, Protocol + +if TYPE_CHECKING: + from adn_server.domain.subscription import AudioChannel, Subscription, SubscriptionId, SubscriptionPhase, SystemId class ConfigLoader(ABC): @@ -256,3 +259,47 @@ class TalkerAliasEmblcEncoder(Protocol): def encode_blocks(self, blocks: dict[int, bytes]) -> tuple[list[dict[int, Any]], int]: ... + + +class SubscriptionStore(ABC): + """Authoritative in-memory subscription registry (Phase 2; replaces BRIDGES dict over time).""" + + @abstractmethod + def get(self, sub_id: "SubscriptionId") -> "Subscription | None": + ... + + @abstractmethod + def upsert(self, subscription: "Subscription") -> None: + ... + + @abstractmethod + def remove(self, sub_id: "SubscriptionId") -> bool: + ... + + @abstractmethod + def clear(self) -> None: + ... + + @abstractmethod + def replace_all(self, subscriptions: Sequence["Subscription"]) -> None: + ... + + @abstractmethod + def snapshot(self) -> tuple["Subscription", ...]: + ... + + @abstractmethod + def list_by_channel(self, channel: "AudioChannel") -> tuple["Subscription", ...]: + ... + + @abstractmethod + def list_by_system(self, system: "SystemId") -> tuple["Subscription", ...]: + ... + + @abstractmethod + def list_active(self) -> tuple["Subscription", ...]: + ... + + @abstractmethod + def list_by_phase(self, phase: "SubscriptionPhase") -> tuple["Subscription", ...]: + ... diff --git a/src/adn_server/domain/__init__.py b/src/adn_server/domain/__init__.py index 02d873c..5cd07ea 100644 --- a/src/adn_server/domain/__init__.py +++ b/src/adn_server/domain/__init__.py @@ -22,6 +22,16 @@ ############################################################################### from .entities import BridgeEntry, StreamState, SystemConfig +from .subscription import ( + ActivationPolicy, + AudioChannel, + Subscription, + SubscriptionId, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, +) from .errors import ACLError, ConfigError, DomainError from .hbp_protocol import ( HBPF_DATA_SYNC, @@ -40,6 +50,14 @@ __all__ = [ "BridgeEntry", "StreamState", "SystemConfig", + "ActivationPolicy", + "AudioChannel", + "Subscription", + "SubscriptionId", + "SubscriptionPhase", + "SubscriptionRole", + "SubscriptionState", + "SystemId", "DmrId", "TgId", "Slot", diff --git a/src/adn_server/domain/subscription.py b/src/adn_server/domain/subscription.py new file mode 100644 index 0000000..451ffd5 --- /dev/null +++ b/src/adn_server/domain/subscription.py @@ -0,0 +1,94 @@ +"""Subscription model: logical TG channels and per-system participation (Phase 2).""" + +from __future__ import annotations + +from dataclasses import dataclass +from enum import Enum + +from .value_objects import Slot, TgId + + +class SubscriptionRole(str, Enum): + """How a system participates in an audio channel.""" + + SINK = "sink" + SOURCE = "source" + PASSIVE_STAT = "passive_stat" + ECHO = "echo" + + +class ActivationPolicy(str, Enum): + """What may activate this subscription (separate from session state).""" + + USER_ACTIVATED = "user_activated" + STATIC = "static" + INBAND = "inband" + OPENBRIDGE_STAT = "openbridge_stat" + + +class SubscriptionPhase(str, Enum): + """Runtime session phase for a subscription leg.""" + + IDLE = "idle" + ACTIVE = "active" + HANGTIME = "hangtime" + + +@dataclass(frozen=True, slots=True) +class SystemId: + """Enabled system name (legacy BRIDGES row ``SYSTEM``).""" + + value: str + + def __str__(self) -> str: + return self.value + + +@dataclass(frozen=True, slots=True) +class AudioChannel: + """Logical talkgroup channel: ``(tgid, slot)``.""" + + tgid: TgId + slot: Slot + + +@dataclass(frozen=True, slots=True) +class SubscriptionId: + """Stable key for one system leg on a channel.""" + + channel: AudioChannel + system: SystemId + + +@dataclass(slots=True) +class SubscriptionState: + """Mutable session state (legacy ACTIVE / TIMER / hangtime).""" + + phase: SubscriptionPhase = SubscriptionPhase.IDLE + timer_expires_at: float | None = None + + +@dataclass(slots=True) +class Subscription: + """One system participating in a channel with LC rewrite and activation rules.""" + + channel: AudioChannel + system: SystemId + target_tgid: TgId + role: SubscriptionRole + policy: ActivationPolicy + state: SubscriptionState + bridge_key: str | None = None + timeout_seconds: float | None = None + + @property + def subscription_id(self) -> SubscriptionId: + return SubscriptionId(channel=self.channel, system=self.system) + + def is_active(self) -> bool: + return self.state.phase == SubscriptionPhase.ACTIVE + + def table_key(self) -> str: + if self.bridge_key is not None: + return self.bridge_key + return str(self.channel.tgid.value) diff --git a/src/adn_server/infrastructure/subscription_store.py b/src/adn_server/infrastructure/subscription_store.py new file mode 100644 index 0000000..44bfb4b --- /dev/null +++ b/src/adn_server/infrastructure/subscription_store.py @@ -0,0 +1,51 @@ +"""In-memory subscription store (Phase 2 authority; not wired to BRIDGES yet).""" + +from __future__ import annotations + +from collections.abc import Sequence + +from adn_server.application.ports import SubscriptionStore +from adn_server.domain.subscription import ( + AudioChannel, + Subscription, + SubscriptionId, + SubscriptionPhase, + SystemId, +) + + +class InMemorySubscriptionStore(SubscriptionStore): + """Hold subscriptions keyed by ``SubscriptionId``; no Twisted or YAML.""" + + def __init__(self) -> None: + self._items: dict[SubscriptionId, Subscription] = {} + + def get(self, sub_id: SubscriptionId) -> Subscription | None: + return self._items.get(sub_id) + + def upsert(self, subscription: Subscription) -> None: + self._items[subscription.subscription_id] = subscription + + def remove(self, sub_id: SubscriptionId) -> bool: + return self._items.pop(sub_id, None) is not None + + def clear(self) -> None: + self._items.clear() + + def replace_all(self, subscriptions: Sequence[Subscription]) -> None: + self._items = {sub.subscription_id: sub for sub in subscriptions} + + def snapshot(self) -> tuple[Subscription, ...]: + return tuple(self._items.values()) + + def list_by_channel(self, channel: AudioChannel) -> tuple[Subscription, ...]: + return tuple(sub for sub in self._items.values() if sub.channel == channel) + + def list_by_system(self, system: SystemId) -> tuple[Subscription, ...]: + return tuple(sub for sub in self._items.values() if sub.system == system) + + def list_active(self) -> tuple[Subscription, ...]: + return tuple(sub for sub in self._items.values() if sub.is_active()) + + def list_by_phase(self, phase: SubscriptionPhase) -> tuple[Subscription, ...]: + return tuple(sub for sub in self._items.values() if sub.state.phase == phase) diff --git a/tests/domain/test_subscription.py b/tests/domain/test_subscription.py new file mode 100644 index 0000000..80c7734 --- /dev/null +++ b/tests/domain/test_subscription.py @@ -0,0 +1,40 @@ +"""Domain subscription entities.""" + +from __future__ import annotations + +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + Subscription, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) + + +def _sub(*, phase: SubscriptionPhase = SubscriptionPhase.IDLE, active: bool = False) -> Subscription: + state = SubscriptionState( + phase=SubscriptionPhase.ACTIVE if active else phase, + timer_expires_at=100.0 if active else None, + ) + return Subscription( + channel=AudioChannel(tgid=TgId(730444), slot=2), + system=SystemId("OBP-CL"), + target_tgid=TgId(730444), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.INBAND, + state=state, + ) + + +def test_subscription_id_stable(): + sub = _sub() + assert sub.subscription_id.channel == sub.channel + assert sub.subscription_id.system == sub.system + + +def test_is_active(): + assert not _sub().is_active() + assert _sub(active=True).is_active() diff --git a/tests/infrastructure/test_subscription_store.py b/tests/infrastructure/test_subscription_store.py new file mode 100644 index 0000000..99dba82 --- /dev/null +++ b/tests/infrastructure/test_subscription_store.py @@ -0,0 +1,84 @@ +"""In-memory subscription store.""" + +from __future__ import annotations + +from adn_server.domain.subscription import ( + ActivationPolicy, + AudioChannel, + Subscription, + SubscriptionId, + SubscriptionPhase, + SubscriptionRole, + SubscriptionState, + SystemId, + TgId, +) +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + +def _leg( + *, + tgid: int, + slot: int, + system: str, + target: int | None = None, + phase: SubscriptionPhase = SubscriptionPhase.IDLE, +) -> Subscription: + return Subscription( + channel=AudioChannel(tgid=TgId(tgid), slot=slot), # type: ignore[arg-type] + system=SystemId(system), + target_tgid=TgId(target if target is not None else tgid), + role=SubscriptionRole.SINK, + policy=ActivationPolicy.STATIC, + state=SubscriptionState(phase=phase), + ) + + +def test_upsert_get_remove(): + store = InMemorySubscriptionStore() + sub = _leg(tgid=730444, slot=2, system="MASTER-A") + store.upsert(sub) + sub_id = SubscriptionId(channel=sub.channel, system=sub.system) + assert store.get(sub_id) is sub + assert store.remove(sub_id) is True + assert store.get(sub_id) is None + assert store.remove(sub_id) is False + + +def test_replace_all_and_snapshot(): + store = InMemorySubscriptionStore() + a = _leg(tgid=1, slot=1, system="A") + b = _leg(tgid=2, slot=2, system="B") + store.replace_all([a, b]) + snap = store.snapshot() + assert len(snap) == 2 + assert a in snap and b in snap + store.replace_all([a]) + assert store.snapshot() == (a,) + + +def test_list_by_channel_system_and_active(): + store = InMemorySubscriptionStore() + channel = AudioChannel(tgid=TgId(730444), slot=2) + idle = _leg(tgid=730444, slot=2, system="OBP-CL", phase=SubscriptionPhase.IDLE) + active = _leg(tgid=730444, slot=2, system="ECHO", phase=SubscriptionPhase.ACTIVE) + other = _leg(tgid=9990, slot=1, system="ECHO") + store.replace_all([idle, active, other]) + + by_channel = store.list_by_channel(channel) + assert len(by_channel) == 2 + assert idle in by_channel and active in by_channel + + by_system = store.list_by_system(SystemId("ECHO")) + assert len(by_system) == 2 + assert active in by_system and other in by_system + + assert store.list_active() == (active,) + assert store.list_by_phase(SubscriptionPhase.IDLE) == (idle, other) + + +def test_clear(): + store = InMemorySubscriptionStore() + store.upsert(_leg(tgid=1, slot=1, system="X")) + store.clear() + assert store.snapshot() == ()