feat(subscription): VoiceIngress router and forward legs (V2-P2-003/004)

Add immutable VoiceIngress/ForwardLeg domain types and a pure
SubscriptionRouter that resolves active forward targets from the store.
pull/1/head
Rodrigo Pérez 4 months ago
parent d0abe7c09e
commit a46bec7d1e

@ -1,5 +1,6 @@
"""Subscription application helpers."""
from .bridges_export import export_bridges, subscription_to_legacy_row
from .router import SubscriptionRouter
__all__ = ["export_bridges", "subscription_to_legacy_row"]
__all__ = ["SubscriptionRouter", "export_bridges", "subscription_to_legacy_row"]

@ -0,0 +1,47 @@
"""Resolve voice ingress to forward legs using the subscription store."""
from __future__ import annotations
from adn_server.application.ports import SubscriptionStore
from adn_server.domain.subscription import AudioChannel, SubscriptionPhase
from adn_server.domain.voice_routing import ForwardLeg, VoiceIngress
class SubscriptionRouter:
"""Pure router: no Twisted, no BRIDGES dict mutation (legacy parity for forward targets)."""
def __init__(self, store: SubscriptionStore) -> None:
self._store = store
def resolve(self, ingress: VoiceIngress) -> tuple[ForwardLeg, ...]:
"""Return active forward legs when the source subscription is ACTIVE on the dst channel."""
table_key = str(ingress.dst_tgid.value)
match_slot: int = 1 if ingress.source_is_obp else int(ingress.slot)
source_channel = AudioChannel(tgid=ingress.dst_tgid, slot=match_slot) # type: ignore[arg-type]
if not self._source_is_active(ingress.source_system, source_channel):
return ()
legs: list[ForwardLeg] = []
for sub in self._store.snapshot():
if sub.table_key() != table_key:
continue
if sub.system.value == ingress.source_system:
continue
if not sub.is_active():
continue
legs.append(
ForwardLeg(
target_system=sub.system.value,
slot=sub.channel.slot,
target_tgid=sub.target_tgid,
)
)
return tuple(legs)
def _source_is_active(self, source_system: str, channel: AudioChannel) -> bool:
for sub in self._store.list_by_channel(channel):
if sub.system.value != source_system:
continue
return sub.state.phase == SubscriptionPhase.ACTIVE
return False

@ -32,6 +32,7 @@ from .subscription import (
SubscriptionState,
SystemId,
)
from .voice_routing import ForwardLeg, VoiceIngress
from .errors import ACLError, ConfigError, DomainError
from .hbp_protocol import (
HBPF_DATA_SYNC,
@ -58,6 +59,8 @@ __all__ = [
"SubscriptionRole",
"SubscriptionState",
"SystemId",
"VoiceIngress",
"ForwardLeg",
"DmrId",
"TgId",
"Slot",

@ -0,0 +1,29 @@
"""Immutable voice routing messages (Phase 2)."""
from __future__ import annotations
from dataclasses import dataclass
from .value_objects import DmrId, Slot, TgId
@dataclass(frozen=True, slots=True)
class VoiceIngress:
"""Normalized group/voice packet entering the subscription router."""
source_system: str
slot: Slot
dst_tgid: TgId
source_is_obp: bool = False
stream_id: int | None = None
peer_id: DmrId | None = None
src_id: DmrId | None = None
@dataclass(frozen=True, slots=True)
class ForwardLeg:
"""One outbound leg: forward ingress audio to ``target_system`` on ``slot`` with LC rewrite."""
target_system: str
slot: Slot
target_tgid: TgId

@ -0,0 +1,87 @@
"""Resolve forward legs from subscription store."""
from __future__ import annotations
from adn_server.application.subscription.router import SubscriptionRouter
from adn_server.domain.subscription import (
ActivationPolicy,
AudioChannel,
Subscription,
SubscriptionPhase,
SubscriptionRole,
SubscriptionState,
SystemId,
TgId,
)
from adn_server.domain.voice_routing import VoiceIngress
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _sub(
*,
tgid: int,
slot: int,
system: str,
active: bool,
target: int | None = None,
) -> Subscription:
phase = SubscriptionPhase.ACTIVE if active else SubscriptionPhase.IDLE
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.INBAND,
state=SubscriptionState(phase=phase),
)
def test_resolve_returns_empty_when_source_inactive():
store = InMemorySubscriptionStore()
store.replace_all(
[
_sub(tgid=730444, slot=1, system="MASTER-A", active=False),
_sub(tgid=730444, slot=1, system="OBP-CL", active=True),
]
)
router = SubscriptionRouter(store)
ingress = VoiceIngress(source_system="MASTER-A", slot=1, dst_tgid=TgId(730444))
assert router.resolve(ingress) == ()
def test_resolve_forwards_to_other_active_legs():
store = InMemorySubscriptionStore()
store.replace_all(
[
_sub(tgid=730444, slot=1, system="MASTER-A", active=True),
_sub(tgid=730444, slot=1, system="OBP-CL", active=True),
_sub(tgid=730444, slot=1, system="MASTER-B", active=False),
]
)
router = SubscriptionRouter(store)
ingress = VoiceIngress(source_system="MASTER-A", slot=1, dst_tgid=TgId(730444))
legs = router.resolve(ingress)
assert len(legs) == 1
assert legs[0].target_system == "OBP-CL"
assert legs[0].slot == 1
assert int(legs[0].target_tgid) == 730444
def test_obp_ingress_uses_ts1_for_source_match():
store = InMemorySubscriptionStore()
store.replace_all(
[
_sub(tgid=730444, slot=1, system="OBP-CL", active=True),
_sub(tgid=730444, slot=1, system="MASTER-A", active=True),
]
)
router = SubscriptionRouter(store)
ingress = VoiceIngress(
source_system="OBP-CL",
slot=2,
dst_tgid=TgId(730444),
source_is_obp=True,
)
legs = router.resolve(ingress)
assert len(legs) == 1
assert legs[0].target_system == "MASTER-A"
Loading…
Cancel
Save

Powered by TurnKey Linux.