feat(subscription): VoiceIngress builder from DMRD args (V2-P2-004)

Add build_voice_ingress for dmrd_received parity, bridge_match_slot on
VoiceIngress, and focused builder tests without runtime wiring.
pull/1/head
Rodrigo Pérez 4 months ago
parent 12a21d389a
commit fe71492d11

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

@ -0,0 +1,52 @@
"""Build immutable ``VoiceIngress`` from DMRD receive parameters (legacy ``dmrd_received``)."""
from __future__ import annotations
from adn_server.domain import int_id
from adn_server.domain.voice_routing import VoiceIngress
from adn_server.domain.value_objects import DmrId, TgId
_BRIDGE_CALL_TYPES = frozenset({"group", "vcsbk"})
def build_voice_ingress(
*,
source_system: str,
system_mode: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
call_type: str,
stream_id: bytes = b"",
) -> VoiceIngress | None:
"""Map ``BridgeUseCases.dmrd_received`` args to a routable ingress, or ``None`` if not bridged."""
if call_type not in _BRIDGE_CALL_TYPES:
return None
match_slot = 1 if int(slot) == 1 else 2
return VoiceIngress(
source_system=source_system,
slot=match_slot, # type: ignore[arg-type]
dst_tgid=TgId(int_id(dst_id)),
source_is_obp=system_mode == "OPENBRIDGE",
call_type=call_type,
stream_id=_stream_id(stream_id),
peer_id=_optional_dmr_id(peer_id),
src_id=_optional_dmr_id(rf_src),
)
def _optional_dmr_id(raw: bytes) -> DmrId | None:
if not raw:
return None
value = int_id(raw)
if value == 0:
return None
return DmrId(value)
def _stream_id(raw: bytes) -> int | None:
if not raw:
return None
value = int_id(raw)
return value if value != 0 else None

@ -14,7 +14,7 @@ class SubscriptionRouter:
def resolve(self, ingress: VoiceIngress) -> tuple[ForwardLeg, ...]:
"""Return active forward legs when the source has an ACTIVE row on the dst TG (legacy to_target)."""
match_slot: int = 1 if ingress.source_is_obp else int(ingress.slot)
match_slot = ingress.bridge_match_slot
dst_tgid = ingress.dst_tgid.value
tables = self.bridge_tables_with_active_source(
ingress.source_system,

@ -15,10 +15,16 @@ class VoiceIngress:
slot: Slot
dst_tgid: TgId
source_is_obp: bool = False
call_type: str = "group"
stream_id: int | None = None
peer_id: DmrId | None = None
src_id: DmrId | None = None
@property
def bridge_match_slot(self) -> int:
"""Slot used for BRIDGES source lookup (legacy: OBP always TS1)."""
return 1 if self.source_is_obp else int(self.slot)
@dataclass(frozen=True, slots=True)
class ForwardLeg:

@ -0,0 +1,109 @@
"""Build VoiceIngress from DMRD parameters (P2-004)."""
from __future__ import annotations
from adn_server.application.subscription.ingress import build_voice_ingress
from adn_server.application.subscription.router import SubscriptionRouter
from adn_server.domain import bytes_3, bytes_4
from adn_server.domain.subscription import (
ActivationPolicy,
AudioChannel,
Subscription,
SubscriptionPhase,
SubscriptionRole,
SubscriptionState,
SystemId,
TgId,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def test_build_voice_ingress_hbp_group():
ingress = build_voice_ingress(
source_system="MASTER-A",
system_mode="MASTER",
peer_id=bytes_4(1234),
rf_src=bytes_3(73010),
dst_id=bytes_3(730444),
slot=2,
call_type="group",
stream_id=bytes_4(99),
)
assert ingress is not None
assert ingress.source_system == "MASTER-A"
assert ingress.slot == 2
assert ingress.source_is_obp is False
assert ingress.bridge_match_slot == 2
assert int(ingress.dst_tgid) == 730444
assert ingress.call_type == "group"
assert ingress.stream_id == 99
assert ingress.peer_id is not None and int(ingress.peer_id) == 1234
assert ingress.src_id is not None and int(ingress.src_id) == 73010
def test_build_voice_ingress_obp_preserves_packet_slot_for_metadata():
ingress = build_voice_ingress(
source_system="OBP-CL",
system_mode="OPENBRIDGE",
peer_id=bytes_4(1),
rf_src=bytes_3(73010),
dst_id=bytes_3(730444),
slot=2,
call_type="group",
)
assert ingress is not None
assert ingress.slot == 2
assert ingress.source_is_obp is True
assert ingress.bridge_match_slot == 1
def test_build_voice_ingress_none_for_unit_call():
assert (
build_voice_ingress(
source_system="MASTER-A",
system_mode="MASTER",
peer_id=bytes_4(1),
rf_src=bytes_3(73010),
dst_id=bytes_3(730444),
slot=1,
call_type="unit",
)
is None
)
def test_built_ingress_resolves_same_as_manual_obp_ingress():
store = InMemorySubscriptionStore()
store.replace_all(
[
Subscription(
channel=AudioChannel(tgid=TgId(730444), slot=1),
system=SystemId("OBP-CL"),
target_tgid=TgId(730444),
role=SubscriptionRole.PASSIVE_STAT,
policy=ActivationPolicy.INBAND,
state=SubscriptionState(phase=SubscriptionPhase.ACTIVE),
),
Subscription(
channel=AudioChannel(tgid=TgId(730444), slot=1),
system=SystemId("MASTER-A"),
target_tgid=TgId(730444),
role=SubscriptionRole.SINK,
policy=ActivationPolicy.INBAND,
state=SubscriptionState(phase=SubscriptionPhase.ACTIVE),
),
]
)
ingress = build_voice_ingress(
source_system="OBP-CL",
system_mode="OPENBRIDGE",
peer_id=bytes_4(1),
rf_src=bytes_3(73010),
dst_id=bytes_3(730444),
slot=2,
call_type="group",
)
assert ingress is not None
legs = SubscriptionRouter(store).resolve(ingress)
assert len(legs) == 1
assert legs[0].target_system == "MASTER-A"
Loading…
Cancel
Save

Powered by TurnKey Linux.