feat(plugins): voice_slot_for_tg and monitor TX for plugin voice

What a voice plugin needs from the core that it can't see itself:

- ServerContext.voice_slot_for_tg(tg): the MASTER slot to speak a TG on
  now, or None while every slot is busy. Same choice announcements make
  today (dynamic UA session or active bridge leg of the TG on that
  MASTER first, then TS2, then TS1), ported to PluginIngress so the
  announcement code can leave the core. Only for granted talkgroups.
- PluginIngress reports GROUP VOICE START/END,TX on the MASTER a plugin
  stream plays on, the line announcements send today: its hotspots hear
  the stream but no bridge leg reports it. A stream cut by a radio, or
  never terminated, still gets its END.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
pull/106/head
yo 4 days ago
parent c418b1a893
commit d76de25cbc

@ -35,3 +35,5 @@ class ServerContext:
call_later: Callable[..., Any]
# Present only for a plugin listed in PLUGINS.send (see domain/send.py).
send_dmrd: Callable[[bytes], bool] | None = None
# With group voice granted: the MASTER slot to speak a TG on now, None while all are busy.
voice_slot_for_tg: Callable[[int], int | None] | None = None

@ -29,7 +29,7 @@ from typing import Any, Callable
from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM
from ....domain.mesh_engine import server_id_bytes
from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd
from ...routing.helpers import slot_voice_held_by_other_stream
from ...routing.helpers import master_dynamic_tg_slots, slot_voice_held_by_other_stream
from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind
logger = logging.getLogger(__name__)
@ -54,17 +54,60 @@ class PluginIngress:
get_protocols: Callable[[], dict[str, Any]],
send_local: Callable[[str, bytes], None],
clock: Callable[[], float] = time.time,
send_routing_event: Callable[[str], None] | None = None,
) -> None:
self._routing = routing
self._config = config
self._get_protocols = get_protocols
self._send_local = send_local
self._clock = clock
self._send_routing_event = send_routing_event
# Plugin voice streams on air: stream_id -> (master, slot, tg, rf_src, start).
self._on_air: dict[bytes, tuple[str, int, int, int, float]] = {}
def set_routing(self, routing: Any) -> None:
"""Bootstrap builds the plugin manager before routing; routing is set before any load."""
self._routing = routing
def voice_slot_for_tg(self, tg: int) -> int | None:
"""The MASTER slot a plugin should speak ``tg`` on now, or None while every slot is busy.
Same choice as scheduled announcements: the slot where ``tg`` is a dynamic UA
session or has an active bridge leg on that MASTER first, then TS2, then TS1.
"""
master = announcement_ptt_system(self._config)
proto = self._get_protocols().get(master) if master else None
status = getattr(proto, "STATUS", None)
sys_cfg = self._config.get("SYSTEMS", {}).get(master or "", {})
if not status or sys_cfg.get("MODE") != "MASTER":
return None
dynamic = master_dynamic_tg_slots(sys_cfg, int(tg))
bridged = self._active_bridge_slots(int(tg), master)
now = self._clock()
for ts in dict.fromkeys([*sorted(dynamic, reverse=True), *sorted(bridged, reverse=True), 2, 1]):
slot = status.get(ts)
if slot and not self._slot_busy(slot, ts in dynamic or ts in bridged, now):
return ts
return None
@staticmethod
def _slot_busy(slot: dict[str, Any], tg_lives_here: bool, now: float) -> bool:
if slot_voice_held_by_other_stream(slot, b"", now):
return True # live voice on the slot, in or out
if slot.get("RX_TYPE") != HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM:
return True # an outside QSO still in its hang time
if tg_lives_here:
return False
return not (slot.get("RX_TYPE") == HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM)
def _active_bridge_slots(self, tg: int, master: str) -> set[int]:
table = self._routing.routing_table_for_report() if self._routing is not None else {}
return {
int(be["TS"])
for be in table.get(str(tg), [])
if isinstance(be, dict) and be.get("SYSTEM") == master and be.get("ACTIVE") and be.get("TS") is not None
}
def deliver(self, pkt: bytes, plugin: str) -> bool:
header = parse_dmrd_header(pkt)
kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq) if header else None
@ -84,10 +127,11 @@ class PluginIngress:
slot = getattr(proto, "STATUS", {}).get(header.slot) if proto is not None else None
if slot is None:
return False
if slot_voice_held_by_other_stream(slot, header.stream_id, now):
return False
if self._route(master, pkt, now, server_id, plugin) is not True:
if slot_voice_held_by_other_stream(slot, header.stream_id, now) or (
self._route(master, pkt, now, server_id, plugin) is not True
):
self._release(slot, header.stream_id)
self._off_air(header.stream_id, now)
return False
is_term = header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VTERM
slot["TX_TYPE"] = HBPF_SLT_VTERM if is_term else HBPF_SLT_VHEAD
@ -96,8 +140,39 @@ class PluginIngress:
slot["TX_TGID"] = header.dst_id
slot["TX_TIME"] = now
self._send_local(master, pkt)
if header.stream_id not in self._on_air:
self._on_air_start(master, header, now)
if is_term:
self._off_air(header.stream_id, now)
return True
def _on_air_start(self, master: str, header: Any, now: float) -> None:
"""Monitor TX on the MASTER itself: its hotspots hear the stream, which no bridge leg reports."""
if len(self._on_air) >= 64: # plugins that never sent a terminator
for stream_id in list(self._on_air)[:32]:
self._off_air(stream_id, now)
tg, rf_src = int_id(header.dst_id), int_id(header.rf_src)
self._on_air[header.stream_id] = (master, header.slot, tg, rf_src, now)
self._report("START", master, header.stream_id, header.slot, tg, rf_src)
def _off_air(self, stream_id: bytes, now: float) -> None:
entry = self._on_air.pop(stream_id, None)
if entry is not None:
master, slot, tg, rf_src, start = entry
self._report("END", master, stream_id, slot, tg, rf_src, now - start)
def _report(
self, action: str, master: str, stream_id: bytes, slot: int, tg: int, rf_src: int, duration: float | None = None
) -> None:
# Same line announcements send (VoiceUseCases._emit_announcement_voice_event).
if self._send_routing_event is None:
return
parts = ["GROUP VOICE", action, "TX", master, str(int_id(stream_id)), str(rf_src), str(rf_src), str(slot), str(tg)]
if duration is not None:
parts.append(f"{duration:.2f}")
parts.append("1")
self._send_routing_event(",".join(parts))
def _route(self, master: str, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool | None:
return inject_plugin_dmrd(self._routing, master, pkt, pkt_time=now, server_id=server_id, plugin=plugin)

@ -136,10 +136,18 @@ class PluginManager:
The sender re-reads the permission on every frame, so revoking it on SIGHUP
is immediate; granting it to an already loaded plugin needs that plugin reloaded.
"""
if self._sender_factory is None or send_permission(server_config, name) is None:
permission = send_permission(server_config, name) if self._sender_factory is not None else None
if permission is None:
return self._ctx
sender = self._sender_factory(name)
if permission.group_voice_tgs:
logger.info(
"(PLUGIN-MANAGER) %s may send unit data and group voice on TG %s (PLUGINS.send)",
name, sorted(permission.group_voice_tgs),
)
return replace(self._ctx, send_dmrd=sender, voice_slot_for_tg=getattr(sender, "voice_slot_for_tg", None))
logger.info("(PLUGIN-MANAGER) %s may send unit data (PLUGINS.send)", name)
return replace(self._ctx, send_dmrd=self._sender_factory(name))
return replace(self._ctx, send_dmrd=sender)
def shutdown_all(self) -> None:
for name in list(self._loaded):

@ -54,6 +54,7 @@ class PluginDmrdSender:
call_from_reactor: Callable[..., Any],
clock: Callable[[], float] = time.monotonic,
in_reactor_thread: Callable[[], bool] = lambda: False,
slot_for_tg: Callable[[int], int | None] | None = None,
) -> None:
self._plugin = plugin
self._config = server_config
@ -61,6 +62,7 @@ class PluginDmrdSender:
self._call_from_reactor = call_from_reactor
self._clock = clock
self._in_reactor_thread = in_reactor_thread
self._slot_for_tg = slot_for_tg
self._lock = threading.Lock()
self._tokens = math.inf # starts full: clamped to the rate on first use
self._refilled = clock()
@ -101,6 +103,13 @@ class PluginDmrdSender:
self._call_from_reactor(self._deliver, bytes(pkt), self._plugin)
return True
def voice_slot_for_tg(self, tg: int) -> int | None:
"""``ServerContext.voice_slot_for_tg``: only for granted talkgroups; reactor thread."""
permission = send_permission(self._config, self._plugin)
if permission is None or int(tg) not in permission.group_voice_tgs or self._slot_for_tg is None:
return None
return self._slot_for_tg(int(tg))
def _take_token(self, permission: SendPermission) -> bool:
with self._lock:
now = self._clock()

@ -455,14 +455,17 @@ def run_peer_server(
if callable(send_system):
send_system(pkt)
plugin_ingress = PluginIngress(None, config, lambda: protocols, _send_local) # routing set below
plugin_ingress = PluginIngress( # routing set below
None, config, lambda: protocols, _send_local, send_routing_event=reporting_use_cases.send_routing_event,
)
plugin_manager = PluginManager(
plugin_bus,
server_ctx,
project_root,
sender_factory=lambda name: PluginDmrdSender(
name, config, plugin_ingress.deliver, reactor.callFromThread, in_reactor_thread=isInIOThread,
name, config, plugin_ingress.deliver, reactor.callFromThread,
in_reactor_thread=isInIOThread, slot_for_tg=plugin_ingress.voice_slot_for_tg,
),
)
voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks)

@ -250,3 +250,12 @@ def test_a_rate_that_is_not_a_positive_number_grants_nothing(rate) -> None:
from adn_server.application.plugins.domain.send import send_permission
assert send_permission(_config(max_frames_per_s=rate), "d-aprs") is None
def test_voice_slot_query_only_for_granted_talkgroups() -> None:
sender = PluginDmrdSender(
"d-aprs", _config(group_voice_tgs=[213]), lambda *a: True, call_from_reactor=lambda *a: None,
slot_for_tg=lambda tg: 2,
)
assert sender.voice_slot_for_tg(213) == 2
assert sender.voice_slot_for_tg(214) is None

@ -187,3 +187,51 @@ def test_a_second_plugin_stream_cannot_talk_over_the_first() -> None:
sc, ingress, _ = _voice_scenario()
_play(sc, ingress, _beacon(stream=1)[:2])
assert _play(sc, ingress, _beacon(stream=2)[:1]) == [False]
def test_voice_slot_prefers_the_slot_where_the_tg_is_bridged_on_the_master() -> None:
config = minimal_config(("SYSTEM", "SYSTEM-B"))
config["SYSTEMS"]["SYSTEM"]["PEERS"] = {b"\x00\x00\x03\xe9": {"CALLSIGN": "HOTSPOT"}}
table = active_routing_table(TG, (("SYSTEM", 1), ("SYSTEM-B", 2)), timeout_minutes=10**6)
sc = DeterministicScenario(config=config, routing_table=table)
sc.routing.apply_startup_subscriptions()
for ts in (1, 2):
sc.protocols["SYSTEM"].STATUS[ts] = idle_hbp_slot() | {"RX_TYPE": HBPF_SLT_VTERM}
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None, clock=sc.clock.time)
assert ingress.voice_slot_for_tg(TG) == 1 # TG 213 lives on TS1 of this MASTER
assert ingress.voice_slot_for_tg(9999) == 2 # nowhere yet: TS2 first
def test_voice_slot_is_none_while_every_slot_is_busy() -> None:
sc, ingress, _ = _voice_scenario()
for ts in (1, 2):
sc.protocols["SYSTEM"].STATUS[ts] = {
"RX_TYPE": HBPF_SLT_VHEAD, "RX_TIME": sc.clock.time(), "TX_TYPE": HBPF_SLT_VTERM,
"RX_STREAM_ID": b"\x05\x05\x05\x05",
}
assert ingress.voice_slot_for_tg(TG) is None
def test_the_monitor_sees_the_beacon_on_the_master_it_plays_on() -> None:
sc, _, _ = _voice_scenario()
events: list[str] = []
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None,
clock=sc.clock.time, send_routing_event=events.append)
_play(sc, ingress, _beacon(stream=0x0C0C0C0C))
starts = [e for e in events if e.startswith("GROUP VOICE,START,TX,SYSTEM,")]
ends = [e for e in events if e.startswith("GROUP VOICE,END,TX,SYSTEM,")]
assert len(starts) == 1 and len(ends) == 1
assert starts[0].split(",")[4:9] == [str(0x0C0C0C0C), str(BEACON_ID), str(BEACON_ID), "2", str(TG)]
def test_a_beacon_cut_by_a_radio_still_ends_on_the_monitor() -> None:
sc, _, _ = _voice_scenario()
events: list[str] = []
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None,
clock=sc.clock.time, send_routing_event=events.append)
frames = _beacon()
_play(sc, ingress, frames[:3])
sc.protocols["SYSTEM"].STATUS[2].update(RX_TYPE=HBPF_SLT_VHEAD, RX_TIME=sc.clock.time() + 0.06, RX_STREAM_ID=b"\x09" * 4)
assert _play(sc, ingress, frames[3:4]) == [False]
assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1
assert sc.protocols["SYSTEM"].STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM

Loading…
Cancel
Save

Powered by TurnKey Linux.