perf(plugins): build only the events some plugin wants; batch deferred ones

As soon as any plugin was loaded, every voice frame built a full
CallLegContext and event (~13 us) and scheduled its own callLater
(~2-7 us), even for plugins that never look at voice.

- A plugin may declare `events` (the event classes it handles). The bus
  keeps the union and `wants()` / `wants_any()`; the voice and data
  bridges check it before building anything, and emit() only calls
  plugins that take that event. Undeclared plugins and internal
  handlers still get everything.
- emit_deferred batches: one call_later per reactor tick, same order.
- The group context's per-stream constant part (IDs, mode, proxy flag,
  aliases) is resolved once per stream; each event still gets its own
  `extra` dict.
- Fix: _plugin_started / _plugin_ended grew by one entry per call
  forever; END now drops the START key and ended keys are capped.

Per voice frame, group voice to 6 OBP + a MASTER (48 us with no plugin):
any plugin 66.0 -> 59.4 us; a data-only plugin 49.4 us.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
pull/105/head
yo 4 days ago
parent 044ba970df
commit d74ef04bfc

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

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

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

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

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

@ -30,6 +30,10 @@ class ServerPlugin(Protocol):
"""Drop-in plugin loaded from plugins/<name>/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."""

@ -0,0 +1,194 @@
# ADN DMR Peer Server - tests plugin bus event interest and batching
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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 == {}
Loading…
Cancel
Save

Powered by TurnKey Linux.