From 703b013a62f5acd03d897150f33cd267193bf9df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 25 Sep 2026 00:34:01 -0300 Subject: [PATCH] perf(talker-alias): stop re-decoding a stream's alias once it is complete - Extract a voice burst's embedded LC from the 5 bytes that hold it instead of running the full voice() decode. - Once a stream's alias has decoded, later bursts only refresh its expiry instead of reassembling it again. Co-Authored-By: Claude Opus 5.5 --- src/adn_server/domain/dmr/decode.py | 7 ++ src/adn_server/domain/talker_alias.py | 8 +- .../twisted_adapters/udp_hbp.py | 25 +++--- tests/talker_alias/test_embed_fast_path.py | 83 +++++++++++++++++++ 4 files changed, 109 insertions(+), 14 deletions(-) create mode 100644 tests/talker_alias/test_embed_fast_path.py diff --git a/src/adn_server/domain/dmr/decode.py b/src/adn_server/domain/dmr/decode.py index c68d987..1fc71df 100644 --- a/src/adn_server/domain/dmr/decode.py +++ b/src/adn_server/domain/dmr/decode.py @@ -94,6 +94,13 @@ def voice(_string): return {'AMBE': ambe, 'CC': cc, 'LCSS': lcss, 'EMBED': embed} +def voice_embed(_string): + """Bits 116-148 of a voice burst (``voice()['EMBED']``), from the 5 bytes that hold them.""" + bits = bitarray(endian='big') + bits.frombytes(_string[14:19]) + return bits[4:36] + + def to_bytes(_bits): add_bits = 8 - (len(_bits) % 8) if add_bits < 8: diff --git a/src/adn_server/domain/talker_alias.py b/src/adn_server/domain/talker_alias.py index b8481ef..ef13486 100644 --- a/src/adn_server/domain/talker_alias.py +++ b/src/adn_server/domain/talker_alias.py @@ -24,6 +24,8 @@ from __future__ import annotations from typing import TYPE_CHECKING, Any +from .dmr import bptc, decode + if TYPE_CHECKING: from bitarray import bitarray @@ -292,18 +294,16 @@ def try_buffer_ta_from_voice_fragments( decodes losslessly here. Returns True when a TA block was decoded and stored. """ - from .dmr import bptc, decode - if vseq not in (1, 2, 3, 4) or len(dmrpkt) < 33: return False try: - embed = decode.voice(dmrpkt)["EMBED"] + embed = decode.voice_embed(dmrpkt) except Exception: return False if vseq == 1: acc.clear() acc[vseq] = embed - if not all(i in acc for i in (1, 2, 3, 4)): + if len(acc) < 4: # keys are vseq 1-4 only, and vseq 1 clears it return False try: lc = bptc.decode_emblc(acc[1] + acc[2] + acc[3] + acc[4]) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 32cf760..7d9deb0 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -298,12 +298,12 @@ class HBPProtocol(DatagramProtocol): self._dmra_by_stream: dict[bytes, dict[str, Any]] = {} self._dmra_rf_stream: dict[tuple[bytes, bytes], bytes] = {} self._ta_voice_acc: dict[bytes, dict[int, Any]] = {} - self._ta_decoded_logged: set[bytes] = set() + self._ta_complete: set[bytes] = set() else: self._dmra_by_stream = {} self._dmra_rf_stream = {} self._ta_voice_acc = {} - self._ta_decoded_logged = set() + self._ta_complete = set() if self._config.get("MODE") == "PEER": self._dmra_downlink: dict[bytes, dict[str, Any]] = {} if self._config.get("MODE") == "PEER": @@ -947,19 +947,24 @@ class HBPProtocol(DatagramProtocol): return if not self._CONFIG.get("GLOBAL", {}).get("TALKER_ALIAS", False): return - entry = self._dmra_by_stream.setdefault( - stream_id, - {"blocks": {}, "rf_src": rf_src, "peer": peer_id, "last": time.time()}, - ) + entry = self._dmra_by_stream.get(stream_id) + if entry is None: + entry = self._dmra_by_stream[stream_id] = { + "blocks": {}, "rf_src": rf_src, "peer": peer_id, "last": time.time() + } + elif stream_id in self._ta_complete: + # A stream's alias does not change once decoded; only keep it from expiring. + entry["last"] = time.time() + return acc = self._ta_voice_acc.setdefault(stream_id, {}) if try_buffer_ta_from_voice_fragments(acc, vseq, dmrpkt, entry["blocks"]): entry["last"] = time.time() entry["rf_src"] = rf_src entry["peer"] = peer_id - if stream_id not in self._ta_decoded_logged: + if stream_id not in self._ta_complete: text = decode_ta_from_blocks(entry["blocks"]) if text: - self._ta_decoded_logged.add(stream_id) + self._ta_complete.add(stream_id) logger.debug( "(%s) *TALKER ALIAS* decoded '%s' from embedded voice (src %s stream %s)", self._system, text, int_id(rf_src), int_id(stream_id), @@ -970,7 +975,7 @@ class HBPProtocol(DatagramProtocol): def clear_ta_stream_buffer(self, stream_id: bytes) -> None: self._dmra_by_stream.pop(stream_id, None) self._ta_voice_acc.pop(stream_id, None) - self._ta_decoded_logged.discard(stream_id) + self._ta_complete.discard(stream_id) def copy_ta_stream_buffer(self, from_stream: bytes, to_stream: bytes) -> None: """Carry decoded TA blocks from recording stream to echo playback stream.""" @@ -1002,7 +1007,7 @@ class HBPProtocol(DatagramProtocol): if self._dmra_by_stream[stream_id].get("last", 0) < cutoff: del self._dmra_by_stream[stream_id] self._ta_voice_acc.pop(stream_id, None) - self._ta_decoded_logged.discard(stream_id) + self._ta_complete.discard(stream_id) def send_dmra_to_peers( self, diff --git a/tests/talker_alias/test_embed_fast_path.py b/tests/talker_alias/test_embed_fast_path.py new file mode 100644 index 0000000..103e32b --- /dev/null +++ b/tests/talker_alias/test_embed_fast_path.py @@ -0,0 +1,83 @@ +# ADN DMR Peer Server - tests talker alias embedded LC fast path +# +# 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 +############################################################################### + +"""Embedded talker alias runs on 4 of every 6 voice frames, so it must stay cheap.""" + +from __future__ import annotations + +import os +import time + +from tests.harness.obp_ingress import build_config +from tests.talker_alias.test_mmdvm_wire import _voice_pkt_with_embed + +import adn_server.infrastructure.twisted_adapters.udp_hbp as udp_hbp +from adn_server.domain.dmr import decode +from adn_server.domain.dmr.bptc import encode_emblc +from adn_server.domain.talker_alias import encode_utf8, talker_alias_lc_bytes +from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol + +STREAM = b"\x12\x34\x56\x78" +SRC = b"\x21\x70\x03" +PEER = b"\x00\x00\x08\x24" + + +def test_voice_embed_matches_the_full_decode() -> None: + for _ in range(5000): + pkt = os.urandom(33) + assert decode.voice_embed(pkt) == decode.voice(pkt)["EMBED"] + + +def _protocol() -> HBPProtocol: + config = build_config({"kind": "v5"}) + config["GLOBAL"]["TALKER_ALIAS"] = True + return HBPProtocol("OBP-FR", config, router=None, dmrd_received=lambda *a, **k: None) + + +def _send_superframe(protocol: HBPProtocol) -> None: + frags = encode_emblc(talker_alias_lc_bytes(0, encode_utf8("CE5RPY")[0:7])) + for vseq in (1, 2, 3, 4): + protocol.store_ta_from_voice_burst(PEER, SRC, STREAM, vseq, _voice_pkt_with_embed(frags[vseq])) + + +def test_a_complete_alias_is_not_decoded_again(monkeypatch) -> None: + protocol = _protocol() + _send_superframe(protocol) + assert STREAM in protocol._ta_complete + + calls = [] + original = udp_hbp.try_buffer_ta_from_voice_fragments + monkeypatch.setattr( + udp_hbp, "try_buffer_ta_from_voice_fragments", lambda *a: (calls.append(a), original(*a))[1] + ) + _send_superframe(protocol) + assert calls == [] + + +def test_a_complete_alias_does_not_expire_during_a_long_call() -> None: + """Skipping the decode must still refresh ``last``, or the trimmer drops the alias of + any call longer than its 180 s window while the call is still going.""" + protocol = _protocol() + _send_superframe(protocol) + protocol._dmra_by_stream[STREAM]["last"] = time.time() - 200 # already past the window + + _send_superframe(protocol) + protocol.trim_dmra_streams(max_age=180.0) + assert STREAM in protocol._dmra_by_stream