diff --git a/docs/en/server/user-guide/plugins.md b/docs/en/server/user-guide/plugins.md index baa4f0f..899cf2e 100644 --- a/docs/en/server/user-guide/plugins.md +++ b/docs/en/server/user-guide/plugins.md @@ -68,6 +68,7 @@ Defined in `src/adn_server/application/plugins/domain/protocol.py`: | Method | Thread | Role | |--------|--------|------| | `name: str` | — | Plugin identifier (directory name) | +| `events` *(optional)* | — | Event classes the plugin handles, e.g. `(VoiceCallFrame, VoiceCallEnd)`. The server then skips building the others — see [Declaring the events you need](#declaring-the-events-you-need). Without it, every event is delivered | | `on_load(bus, config, server_ctx)` | reactor | Read config; subscribe to bus if needed | | `on_event(event)` | **reactor** | Handle bus events — **O(1), no blocking I/O** | | `on_reload(config)` | reactor | Optional hot-reload after `config.yaml` change | @@ -144,6 +145,28 @@ Pacing (about 60 ms per burst) is the plugin's job, e.g. with `call_later`. - Events are emitted **after voice/data has been forwarded** (`emit_deferred` — next reactor tick). - Uncaught exceptions increment a per-plugin trip counter; after repeated failures the plugin is **disabled** and `on_shutdown()` is called (circuit breaker). - `bus.subscribe(handler)` is available for internal handlers; plugins normally implement `on_event` only. +- Deferred events are batched: all those of one reactor tick are delivered by a single `call_later`. + +### Declaring the events you need + +Building an event costs the server work on every frame (a voice call is ~17 frames per second, per stream), whether or not a plugin looks at it. A plugin that lists the events it handles lets the server skip the rest: + +```python +from adn_server.application.plugins.domain.events import UnitDataFrame + +class DAprsPlugin: + name = "d-aprs" + events = (UnitDataFrame,) # no voice events are built for this plugin +``` + +Measured on a Raspberry Pi 5, group voice routed to 6 OpenBridges and a MASTER (48 µs per frame with no plugin): + +| Plugins loaded | Cost per voice frame | +|---|---| +| One that declares only unit data events | +1.4 µs | +| One that declares voice events, or declares nothing | +11 µs | + +`events` is read when the plugin is registered (after `on_load`). Internal handlers added with `bus.subscribe` receive every event. --- diff --git a/docs/es/server/user-guide/plugins.md b/docs/es/server/user-guide/plugins.md index 8cf4120..677674c 100644 --- a/docs/es/server/user-guide/plugins.md +++ b/docs/es/server/user-guide/plugins.md @@ -68,6 +68,7 @@ Definido en `src/adn_server/application/plugins/domain/protocol.py`: | Método | Hilo | Función | |--------|------|---------| | `name: str` | — | Identificador (nombre del directorio) | +| `events` *(opcional)* | — | Clases de evento que maneja el plugin, p. ej. `(VoiceCallFrame, VoiceCallEnd)`. El servidor deja entonces de construir las demás — ver [Declarar los eventos que necesitas](#declarar-los-eventos-que-necesitas). Sin él, se entregan todos | | `on_load(bus, config, server_ctx)` | reactor | Leer config; suscribirse al bus si hace falta | | `on_event(event)` | **reactor** | Manejar eventos — **O(1), sin I/O bloqueante** | | `on_reload(config)` | reactor | Hot-reload opcional tras cambio de `config.yaml` | @@ -144,6 +145,28 @@ El ritmo (unos 60 ms por ráfaga) lo marca el plugin, por ejemplo con `call_late - Los eventos se emiten **después del forward** de voz/datos (`emit_deferred` — siguiente tick del reactor). - Excepciones no capturadas incrementan un contador; tras repetir fallos el plugin se **deshabilita** y se llama `on_shutdown()` (circuit breaker). - `bus.subscribe(handler)` existe para handlers internos; los plugins suelen usar solo `on_event`. +- Los eventos diferidos se agrupan: todos los de un ciclo del reactor se entregan con un solo `call_later`. + +### Declarar los eventos que necesitas + +Construir un evento le cuesta trabajo al servidor en cada trama (una llamada de voz son ~17 tramas por segundo y por stream), la mire o no un plugin. Un plugin que declara los eventos que maneja le ahorra al servidor el resto: + +```python +from adn_server.application.plugins.domain.events import UnitDataFrame + +class DAprsPlugin: + name = "d-aprs" + events = (UnitDataFrame,) # para este plugin no se construye ningún evento de voz +``` + +Medido en una Raspberry Pi 5, voz de grupo enrutada a 6 OpenBridge y un MASTER (48 µs por trama sin plugins): + +| Plugins cargados | Coste por trama de voz | +|---|---| +| Uno que declara solo eventos de datos | +1,4 µs | +| Uno que declara eventos de voz, o no declara nada | +11 µs | + +`events` se lee al registrar el plugin (después de `on_load`). Los handlers internos añadidos con `bus.subscribe` reciben todos los eventos. --- diff --git a/src/adn_server/application/plugins/application/bus.py b/src/adn_server/application/plugins/application/bus.py index f535952..60b350b 100644 --- a/src/adn_server/application/plugins/application/bus.py +++ b/src/adn_server/application/plugins/application/bus.py @@ -47,6 +47,13 @@ class PluginBus: self._trip_counts: dict[str, int] = {} self._disabled: set[str] = set() self._call_later: Callable[..., Any] | None = None + # Event types some subscriber wants: bridges skip building the rest. A plugin + # without an ``events`` attribute, or any internal handler, wants everything. + self._plugin_events: dict[int, frozenset[type] | None] = {} + self._wants_all = False + self._wanted: frozenset[type] = frozenset() + # Deferred events of one reactor tick, delivered by a single call_later. + self._pending: list[object] = [] def set_call_later(self, call_later: Callable[..., Any]) -> None: self._call_later = call_later @@ -57,16 +64,34 @@ class PluginBus: if plugin in self._plugins: return self._plugins.append(plugin) + events = getattr(plugin, "events", None) + self._plugin_events[id(plugin)] = None if events is None else frozenset(events) + self._update_interest() def unregister_plugin(self, plugin: ServerPlugin) -> None: self._plugins = [p for p in self._plugins if p is not plugin] + self._plugin_events.pop(id(plugin), None) + self._update_interest() def subscribe(self, handler: Callable[[object], None]) -> None: self._handlers.append(handler) + self._update_interest() def has_subscribers(self) -> bool: return bool(self._plugins or self._handlers) + def wants(self, event_type: type) -> bool: + """True when some subscriber takes ``event_type``: build the event only then.""" + return self._wants_all or event_type in self._wanted + + def wants_any(self, *event_types: type) -> bool: + return self._wants_all or not self._wanted.isdisjoint(event_types) + + def _update_interest(self) -> None: + declared = [self._plugin_events.get(id(p)) for p in self._plugins] + self._wants_all = bool(self._handlers) or any(events is None for events in declared) + self._wanted = frozenset().union(*(events for events in declared if events is not None)) + def emit(self, event: object) -> None: if not self.has_subscribers(): return @@ -75,9 +100,13 @@ class PluginBus: handler(event) except Exception: logger.exception("(PLUGIN-BUS) handler failed") + event_type = type(event) for plugin in list(self._plugins): if plugin.name in self._disabled: continue + events = self._plugin_events.get(id(plugin)) + if events is not None and event_type not in events: + continue try: plugin.on_event(event) except Exception: @@ -85,12 +114,19 @@ class PluginBus: self._trip(plugin) def emit_deferred(self, event: object) -> None: - """Schedule emit on next reactor tick (after forward).""" - if not self.has_subscribers(): + """Emit on the next reactor tick (after forward), batched: one call_later per tick.""" + if not self.wants(type(event)): + return + if self._call_later is None: + self.emit(event) return - if self._call_later is not None: - self._call_later(0, self.emit, event) - else: + self._pending.append(event) + if len(self._pending) == 1: + self._call_later(0, self._flush) + + def _flush(self) -> None: + pending, self._pending = self._pending, [] + for event in pending: self.emit(event) def _trip(self, plugin: ServerPlugin) -> None: diff --git a/src/adn_server/application/plugins/application/data_bridge.py b/src/adn_server/application/plugins/application/data_bridge.py index 9f716dd..5a58851 100644 --- a/src/adn_server/application/plugins/application/data_bridge.py +++ b/src/adn_server/application/plugins/application/data_bridge.py @@ -31,6 +31,8 @@ from ..application.bus import PluginBus from ..domain.events import CallLegContext, UnitDataEnd, UnitDataFrame, UnitDataStart from .bridge_common import alias_extra, int_byte, unit_data_label +_DATA_EVENTS = (UnitDataStart, UnitDataFrame, UnitDataEnd) + class DataPluginBridge: """Emits unit-data events to PluginBus after routing forward.""" @@ -69,13 +71,8 @@ class DataPluginBridge: synthetic: bool = False, ) -> None: """``synthetic``: the frame was sent by a plugin (``ServerContext.send_dmrd``).""" - if not self._bus.has_subscribers(): + if not self._bus.wants_any(*_DATA_EVENTS): return - systems_cfg = self._config.get("SYSTEMS", {}) - mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER") - server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0 - bits = data[15] if len(data) > 15 else 0 - dmrpkt = data[20:53] if len(data) >= 53 else b"" sid = int_id(stream_id) stream_key = (system_name, slot) prev = self._active_streams.get(stream_key) @@ -84,6 +81,14 @@ class DataPluginBridge: is_new_stream = prev is None or prev[0] != sid if is_new_stream: self._active_streams[stream_key] = (sid, pkt_time, synthetic) + starts = is_new_stream or dtype_vseq == 6 + if not (starts and self._bus.wants(UnitDataStart)) and not self._bus.wants(UnitDataFrame): + return # stream tracked for UnitDataEnd; nothing else to build + systems_cfg = self._config.get("SYSTEMS", {}) + mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER") + server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0 + bits = data[15] if len(data) > 15 else 0 + dmrpkt = data[20:53] if len(data) >= 53 else b"" ctx = CallLegContext( call_family="DATA", direction="RX", @@ -112,7 +117,7 @@ class DataPluginBridge: }, ) label = unit_data_label(dtype_vseq) - if is_new_stream or dtype_vseq == 6: + if starts: self._bus.emit( UnitDataStart( context=ctx, diff --git a/src/adn_server/application/plugins/application/voice_bridge.py b/src/adn_server/application/plugins/application/voice_bridge.py index c90ed43..825867c 100644 --- a/src/adn_server/application/plugins/application/voice_bridge.py +++ b/src/adn_server/application/plugins/application/voice_bridge.py @@ -33,6 +33,10 @@ from ..application.bus import PluginBus from ..domain.events import CallLegContext, VoiceCallEnd, VoiceCallFrame, VoiceCallStart from .bridge_common import alias_extra, int_byte, stream_talker_alias +_VOICE_EVENTS = (VoiceCallStart, VoiceCallFrame, VoiceCallEnd) +_ENDED_KEEP = 4096 +_STREAM_CTX_KEEP = 512 + class VoicePluginBridge: def __init__( @@ -45,14 +49,23 @@ class VoicePluginBridge: self._bus = bus self._config = config self._get_dmra_blocks = get_dmra_blocks + # Streams whose START / END went out, so each is emitted once. END drops the + # START key; ended keys are kept, oldest first, only up to _ENDED_KEEP. self._plugin_started: set[tuple[str, int]] = set() - self._plugin_ended: set[tuple[str, int]] = set() + self._plugin_ended: dict[tuple[str, int], None] = {} + # Per-stream constant part of the event context, oldest first, bounded. + self._stream_ctx: dict[tuple, dict[str, Any]] = {} def update_config(self, config: dict[str, Any]) -> None: self._config = config + def _forget_stream(self, system_name: str, stream_id: bytes) -> None: + for key in [k for k in self._stream_ctx if k[0] == system_name and k[1] == stream_id]: + del self._stream_ctx[key] + def has_subscribers(self) -> bool: - return self._bus.has_subscribers() + """True when some subscriber takes voice events at all.""" + return self._bus.wants_any(*_VOICE_EVENTS) def _end_extra( self, @@ -94,23 +107,31 @@ class VoicePluginBridge: forwarded: tuple[str, ...] | list[str], extra: dict[str, Any] | None = None, ) -> dict[str, Any]: - systems_cfg = self._config.get("SYSTEMS", {}) - sys_cfg = systems_cfg.get(system_name, {}) - mode = sys_cfg.get("MODE", "MASTER") - server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0 + # What stays the same for every frame of a stream is resolved once (IDs, mode, + # proxy flag, aliases); only per-frame fields are recomputed. + key = (system_name, stream_id, peer_id, rf_src, dst_id, slot, bool(synthetic_announcement)) + fixed = self._stream_ctx.get(key) + if fixed is None: + if len(self._stream_ctx) >= _STREAM_CTX_KEEP: + del self._stream_ctx[next(iter(self._stream_ctx))] + sys_cfg = self._config.get("SYSTEMS", {}).get(system_name, {}) + fixed = self._stream_ctx[key] = dict( + call_family="GROUP", + direction="RX", + origin_system=system_name, + system_mode=str(sys_cfg.get("MODE", "MASTER")), + peer_id=int_id(peer_id), + src_id=int_id(rf_src), + dst_id=int_id(dst_id), + slot=slot, + stream_id=int_id(stream_id), + server_id=int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0, + is_synthetic=bool(synthetic_announcement), + is_proxy_ingress=is_proxy_inject_only(self._config, system_name), + extra=alias_extra(self._config, peer_id, rf_src, dst_id), + ) return dict( - call_family="GROUP", - direction="RX", - origin_system=system_name, - system_mode=str(mode), - peer_id=int_id(peer_id), - src_id=int_id(rf_src), - dst_id=int_id(dst_id), - slot=slot, - stream_id=int_id(stream_id), - server_id=server_id, - is_synthetic=bool(synthetic_announcement), - is_proxy_ingress=is_proxy_inject_only(self._config, system_name), + fixed, pkt_time=pkt_time, obp_source_server_id=int_byte(obp_source_server) if source_is_obp else None, obp_hops=int_byte(obp_hops) if source_is_obp else None, @@ -118,7 +139,7 @@ class VoicePluginBridge: ber=int_byte(obp_ber), rssi=int_byte(obp_rssi), forwarded_systems=tuple(forwarded), - extra=extra if extra is not None else alias_extra(self._config, peer_id, rf_src, dst_id), + extra=dict(extra if extra is not None else fixed["extra"]), # each event owns its dict ) def emit_group_voice_start( @@ -141,7 +162,7 @@ class VoicePluginBridge: forwarded: tuple[str, ...] | list[str] = (), voice_phase: str | None = None, ) -> bool: - if not self._bus.has_subscribers(): + if not self.has_subscribers(): return False ctx_base = self._group_ctx_base( system_name=system_name, @@ -189,7 +210,7 @@ class VoicePluginBridge: synthetic_announcement: bool = False, forwarded: tuple[str, ...] | list[str] = (), ) -> bool: - if not self._bus.has_subscribers(): + if not self.has_subscribers(): return False ctx_base = self._group_ctx_base( system_name=system_name, @@ -211,7 +232,11 @@ class VoicePluginBridge: key = (ctx_base["origin_system"], ctx_base["stream_id"]) if key in self._plugin_ended: return False - self._plugin_ended.add(key) + self._plugin_ended[key] = None + self._plugin_started.discard(key) + self._forget_stream(system_name, stream_id) + if len(self._plugin_ended) > _ENDED_KEEP: + del self._plugin_ended[next(iter(self._plugin_ended))] end_extra = self._end_extra( origin_system=ctx_base["origin_system"], stream_id=ctx_base["stream_id"], @@ -248,9 +273,8 @@ class VoicePluginBridge: synthetic_announcement: bool, forwarded: list[str], ) -> None: - if not self._bus.has_subscribers(): - return - if frame_type not in (HBPF_VOICE, HBPF_VOICE_SYNC): + # Once per voice frame: skip building the event unless someone takes it. + if frame_type not in (HBPF_VOICE, HBPF_VOICE_SYNC) or not self._bus.wants(VoiceCallFrame): return ctx_base = self._group_ctx_base( system_name=system_name, @@ -296,7 +320,8 @@ class VoicePluginBridge: phase: str, duration_s: float = 0.0, ) -> None: - if not self._bus.has_subscribers(): + wanted = self._bus.wants(VoiceCallFrame) if phase == "FRAME" else self.has_subscribers() + if not wanted: return systems_cfg = self._config.get("SYSTEMS", {}) mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER") diff --git a/src/adn_server/application/plugins/domain/protocol.py b/src/adn_server/application/plugins/domain/protocol.py index 6f29012..c24d7a4 100644 --- a/src/adn_server/application/plugins/domain/protocol.py +++ b/src/adn_server/application/plugins/domain/protocol.py @@ -30,6 +30,10 @@ class ServerPlugin(Protocol): """Drop-in plugin loaded from plugins//plugin/.""" name: str + # Optional: the event classes this plugin handles, e.g. (VoiceCallFrame, VoiceCallEnd). + # Declaring them lets the server skip building every other event; without it the + # plugin receives all of them. Read when the plugin is registered, after on_load. + # events: tuple[type, ...] def on_load(self, bus: Any, config: dict[str, Any], server_ctx: Any) -> None: """Subscribe to bus; read plugin config.""" diff --git a/tests/application/test_plugin_bus_interest.py b/tests/application/test_plugin_bus_interest.py new file mode 100644 index 0000000..a813541 --- /dev/null +++ b/tests/application/test_plugin_bus_interest.py @@ -0,0 +1,194 @@ +# ADN DMR Peer Server - tests plugin bus event interest and batching +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Plugins pay only for the events they take; deferred events cost one call_later per tick.""" + +from __future__ import annotations + +import pytest + +from adn_server.application.plugins.application import voice_bridge as vb +from adn_server.application.plugins.application.bus import PluginBus +from adn_server.application.plugins.application.data_bridge import DataPluginBridge +from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge +from adn_server.application.plugins.domain.events import ( + UnitDataEnd, + UnitDataFrame, + UnitDataStart, + VoiceCallEnd, + VoiceCallFrame, + VoiceCallStart, +) +from adn_server.domain import HBPF_DATA_SYNC, HBPF_VOICE + + +class _Plugin: + def __init__(self, name: str, events=None) -> None: + self.name = name + if events is not None: + self.events = events + self.got: list[object] = [] + + def on_event(self, event: object) -> None: + self.got.append(event) + + def on_shutdown(self) -> None: + pass + + +CONFIG = {"GLOBAL": {"SERVER_ID": 2131}, "SYSTEMS": {"SYSTEM": {"MODE": "MASTER"}}} + + +def _voice_frame(bridge: VoicePluginBridge, stream: int = 7) -> None: + bridge.notify_group_after_forward( + system_name="SYSTEM", peer_id=b"\0\0\0\1", rf_src=b"\0\0\1", dst_id=b"\0\0\xd5", slot=2, + frame_type=HBPF_VOICE, dtype_vseq=1, stream_id=stream.to_bytes(4, "big"), data=b"DMRD" + bytes(51), + pkt_time=1.0, source_is_obp=False, obp_hops=b"", obp_source_server=None, obp_ber=b"\0", + obp_rssi=b"\0", obp_source_rptr=b"\0\0\0\0", synthetic_announcement=False, forwarded=[], + ) + + +def _data_frame(bridge: DataPluginBridge, dtype: int = 7) -> None: + bridge.notify_after_forward( + route="unit", system_name="SYSTEM", peer_id=b"\0\0\0\1", rf_src=b"\0\0\1", dst_id=b"\x0d\xbb\xa7", + seq=1, slot=2, frame_type=HBPF_DATA_SYNC, dtype_vseq=dtype, stream_id=b"\0\0\0\x09", + data=b"DMRD" + bytes(51), pkt_time=1.0, source_is_obp=False, obp_hops=b"", obp_source_server=None, + obp_ber=b"\0", obp_rssi=b"\0", obp_source_rptr=b"\0\0\0\0", forwarded=[], + ) + + +def test_a_plugin_without_events_still_gets_everything() -> None: + bus = PluginBus() + legacy = _Plugin("legacy") + bus.register_plugin(legacy) + assert bus.wants(VoiceCallFrame) and bus.wants(UnitDataEnd) + _voice_frame(VoicePluginBridge(bus, CONFIG)) + assert [type(e) for e in legacy.got] == [VoiceCallFrame] + + +def test_declared_events_are_the_only_ones_delivered() -> None: + bus = PluginBus() + data = _Plugin("d-aprs", events=(UnitDataFrame,)) + voice = _Plugin("igate", events=(VoiceCallFrame, VoiceCallEnd)) + bus.register_plugin(data) + bus.register_plugin(voice) + _voice_frame(VoicePluginBridge(bus, CONFIG)) + _data_frame(DataPluginBridge(bus, CONFIG)) + assert [type(e) for e in voice.got] == [VoiceCallFrame] + assert [type(e) for e in data.got] == [UnitDataFrame] + + +def test_voice_frames_are_not_even_built_when_nobody_takes_them(monkeypatch) -> None: + bus = PluginBus() + bus.register_plugin(_Plugin("d-aprs", events=(UnitDataStart, UnitDataFrame, UnitDataEnd))) + bridge = VoicePluginBridge(bus, CONFIG) + + def boom(*a, **k): + raise AssertionError("voice event built for nobody") + + monkeypatch.setattr(bridge, "_group_ctx_base", boom) + monkeypatch.setattr(vb, "CallLegContext", boom) + _voice_frame(bridge) + assert not bridge.has_subscribers() + + +def test_unregistering_the_last_interested_plugin_stops_the_work() -> None: + bus = PluginBus() + igate = _Plugin("igate", events=(VoiceCallFrame,)) + bus.register_plugin(igate) + assert bus.wants(VoiceCallFrame) + bus.unregister_plugin(igate) + assert not bus.wants(VoiceCallFrame) and not bus.wants_any(VoiceCallStart, UnitDataFrame) + + +def test_an_internal_handler_takes_everything() -> None: + bus = PluginBus() + bus.subscribe(lambda e: None) + assert bus.wants(VoiceCallFrame) and bus.wants(UnitDataStart) + + +def test_deferred_events_of_a_tick_share_one_call_later_and_keep_their_order() -> None: + scheduled: list = [] + bus = PluginBus() + bus.set_call_later(lambda delay, fn, *a: scheduled.append((fn, a))) + plugin = _Plugin("p") + bus.register_plugin(plugin) + events = [object(), object(), object()] + for e in events: + bus.emit_deferred(e) + assert len(scheduled) == 1 and plugin.got == [] + fn, args = scheduled.pop() + fn(*args) + assert plugin.got == events + bus.emit_deferred(events[0]) # next tick schedules again + assert len(scheduled) == 1 + + +def test_a_deferred_event_nobody_takes_is_not_scheduled() -> None: + scheduled: list = [] + bus = PluginBus() + bus.set_call_later(lambda delay, fn, *a: scheduled.append(fn)) + bus.register_plugin(_Plugin("d-aprs", events=(UnitDataFrame,))) + bus.emit_deferred(VoiceCallFrame(context=None, dmrpkt=b"", frame_type=0, dtype_vseq=0)) + assert scheduled == [] + + +@pytest.mark.parametrize("calls", [3, vb._ENDED_KEEP + 10]) +def test_stream_bookkeeping_does_not_grow_without_bound(calls: int) -> None: + bus = PluginBus() + bus.register_plugin(_Plugin("p")) + bridge = VoicePluginBridge(bus, CONFIG) + kw = dict(system_name="SYSTEM", peer_id=b"\0\0\0\1", rf_src=b"\0\0\1", dst_id=b"\0\0\xd5", slot=2, + pkt_time=1.0, source_is_obp=False) + for n in range(calls): + sid = n.to_bytes(4, "big") + assert bridge.emit_group_voice_start(stream_id=sid, **kw) + assert bridge.emit_group_voice_end(stream_id=sid, duration_s=1.0, **kw) + assert not bridge.emit_group_voice_end(stream_id=sid, duration_s=1.0, **kw) # still once + assert bridge._plugin_started == set() + assert len(bridge._plugin_ended) == min(calls, vb._ENDED_KEEP) + + +def test_a_stream_resolves_its_aliases_once_and_each_event_owns_its_extra(monkeypatch) -> None: + calls: list[int] = [] + real = vb.alias_extra + monkeypatch.setattr(vb, "alias_extra", lambda *a: calls.append(1) or real(*a)) + bus = PluginBus() + plugin = _Plugin("igate", events=(VoiceCallFrame,)) + bus.register_plugin(plugin) + bridge = VoicePluginBridge(bus, {**CONFIG, "_SUB_IDS": {1: "C31AG"}}) + for _ in range(6): + _voice_frame(bridge) + assert len(calls) == 1 + assert plugin.got[0].context.extra == {"src_callsign": "C31AG"} + plugin.got[0].context.extra["mutated"] = True + assert "mutated" not in plugin.got[1].context.extra + + +def test_the_cached_context_is_dropped_when_the_stream_ends() -> None: + bus = PluginBus() + bus.register_plugin(_Plugin("p")) + bridge = VoicePluginBridge(bus, CONFIG) + _voice_frame(bridge, stream=7) + assert bridge._stream_ctx + bridge.emit_group_voice_end(system_name="SYSTEM", peer_id=b"\0\0\0\1", rf_src=b"\0\0\1", dst_id=b"\0\0\xd5", + slot=2, stream_id=(7).to_bytes(4, "big"), pkt_time=2.0, duration_s=1.0, + source_is_obp=False) + assert bridge._stream_ctx == {}