feat(subscription): domain model and in-memory store (V2-P2-001)

Add AudioChannel, Subscription entities, SubscriptionStore port, and
InMemorySubscriptionStore with unit tests (no runtime wiring yet).
pull/1/head
Rodrigo Pérez 4 months ago
parent 471e1795bf
commit c734b1f0b0

@ -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", ...]:
...

@ -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",

@ -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)

@ -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)

@ -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()

@ -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() == ()
Loading…
Cancel
Save

Powered by TurnKey Linux.