Merge pull request #107 from Amateur-Digital-Network/perf/talker-alias-embed

perf(talker-alias): stop re-decoding a stream's alias once it is complete
pull/108/head
ce5rpy 4 days ago committed by GitHub
commit 9bef975eea
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

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

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

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

@ -0,0 +1,83 @@
# ADN DMR Peer Server - tests talker alias embedded LC fast path
#
# 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
###############################################################################
"""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
Loading…
Cancel
Save

Powered by TurnKey Linux.