fix(talker-alias): MMDVMHost parity, both-mode fallback, and OBP inject

- both: passthrough when source sends TA (DMRA or embedded voice), inject
  template otherwise; OBP sources inject immediately at VHEAD
- Always rewrite destination group embedded LC on bridge forward (legacy
  parity; fixes OBP TG mismatch packet loss introduced by TA branch)
- Fix dmr_utils3 encode_emblc index bug for lossless injected TA text
- Log decoded TA from voice; strict MMDVM DMRA wire layout
- chore: ruff per-file E402 ignore for main/parrot entrypoints
pull/4/head
Rodrigo Pérez 4 months ago
parent 07695fffd5
commit 1f902aefa4

@ -33,7 +33,7 @@ With **systemd**, add to your unit file:
ExecReload=/bin/kill -HUP $MAINPID
```
**Reload applies:** `GLOBAL`, `REPORTS`, `ALIASES`, per-system settings, **new/removed SYSTEMS** (including `GENERATOR` expansion and new OpenBridge legs), and updated bind addresses (listener restart for that system only).
**Reload applies:** `GLOBAL`, `REPORTS`, `ALIASES`, **`LOGGER.LOG_LEVEL`** (without process restart), per-system settings, **new/removed SYSTEMS** (including `GENERATOR` expansion and new OpenBridge legs), and updated bind addresses (listener restart for that system only).
**Not reloaded:** `adn-voice.yaml` (separate 15 s loop), Python code, subscriber alias files (separate periodic reload). **BRIDGES** table is not rebuilt on reload — restart if bridge rules changed in a way that requires a full reset.

@ -14,13 +14,18 @@ When enabled on a **MASTER** system, the server can:
| Mode | Behaviour |
|------|-----------|
| **`both`** (default) | Pass through TA from the source hotspot/radio when all four `DMRA` blocks were received; otherwise inject from `subscriber_ids` + template. |
| **`passthrough`** | Relay buffered `DMRA` only. |
| **`inject`** | Always build TA from the configured template and alias data. |
| **`both`** (default) | Prefer the source's own TA, fall back to inject. The server briefly waits (≈2 s after `VHEAD`) for the source's TA, decoded from either a valid MMDVM `DMRA` buffer or the embedded LC inside voice (FLCO 4–7). If found, it is **passed through unchanged**. If the source sent **no** TA within that window (e.g. a radio that does not support Talker Alias), ADN **injects** the configured template (embedded LC + `DMRA`). When the **source is an OpenBridge** system (which never carries TA), there is nothing to wait for, so ADN injects immediately at `VHEAD`. |
| **`passthrough`** | Relay the source's TA only (buffered `DMRA` and the source embedded LC, both untouched). Never inject the template. |
| **`inject`** | Always build TA from the configured template and alias data, overwriting the embedded LC and sending `DMRA` packets. |
On **bridge forward** at voice header (`VHEAD`), the server sends four `DMRA` packets to each HBP target (**MASTER** peers or **PEER** upstream) once per stream, then forwards `DMRD` as usual.
**MMDVMHost / DMRGateway (Pi-Star, WPSD):** stock MMDVMHost does **not** consume standalone downlink `DMRA` UDP; it decodes Talker Alias from **embedded LC inside `DMRD` voice** (FLCO 4–7). When TA is enabled, ADN injects TA into the embedded LC of voice bursts **B–E** (dtype 1–4) on **bridge forward** to HBP targets, in addition to optional standalone `DMRA` packets for clients that support them.
**MMDVMHost / DMRGateway (Pi-Star, WPSD):** stock MMDVMHost does **not** consume standalone downlink `DMRA` UDP; it decodes Talker Alias from **embedded LC inside `DMRD` voice** (FLCO 4–7).
The embedded LC of voice bursts **B–E** (dtype 1–4) is always rewritten with the **destination** group LC (re-encoded for the forwarded TGID, as legacy `bridge.py` does) — this is required for the receiving MMDVM to accept the voice; a stale embedded LC from the source TG causes packet loss. The Talker Alias is then **overlaid on alternate superframes**:
- **`inject`** and the **`both`** no-source-TA fallback use the configured template.
- **`passthrough`** and **`both`** with a source TA re-encode the source's own Talker Alias (decoded from its `DMRA`/embedded-voice blocks) and overlay that, so the radio's alias reaches the far side while the group LC stays correct for the new TG.
On the **same MASTER**, when **`REPEAT`** copies group voice to other logged-in hotspots, the server also sends those four `DMRA` packets on `VHEAD` (excluding the transmitting peer). Bridge forwarding to the same system shares the same once-per-stream dedupe, so TA is not sent twice.

@ -33,7 +33,7 @@ Con **systemd**, en la unidad:
ExecReload=/bin/kill -HUP $MAINPID
```
**Se recarga:** `GLOBAL`, `REPORTS`, `ALIASES`, parámetros por system, **systems nuevos/eliminados** (incluida expansión `GENERATOR` y OBP nuevos), y cambios de IP/puerto (solo reinicia el listener de ese system).
**Se recarga:** `GLOBAL`, `REPORTS`, `ALIASES`, **`LOGGER.LOG_LEVEL`** (sin reiniciar el proceso), parámetros por system, **systems nuevos/eliminados** (incluida expansión `GENERATOR` y OBP nuevos), y cambios de IP/puerto (solo reinicia el listener de ese system).
**No se recarga:** `adn-voice.yaml` (loop aparte cada 15 s), código Python, ficheros de alias (recarga periódica). La tabla **BRIDGES** no se reconstruye — reinicia si cambiaste reglas de bridge que exijan reset completo.

@ -27,5 +27,10 @@ docs = ["mkdocs>=1.6", "mkdocs-material>=9.5", "pymdown-extensions>=10.3"]
[tool.setuptools.packages.find]
where = ["src"]
[tool.ruff.lint.per-file-ignores]
# Entrypoints adjust sys.path before importing the package (intentional E402).
"src/adn_server/main.py" = ["E402"]
"src/adn_server/parrot_main.py" = ["E402"]
[project.scripts]
adn-server = "adn_server.main:main"

@ -41,7 +41,7 @@ from dmr_utils3.const import LC_OPT
from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4, int_id
from ..domain.talker_alias import DMRA_BLOCK_COUNT
from .ports import BridgeRouter
from .talker_alias_use_cases import TalkerAliasUseCases
from .talker_alias_use_cases import TalkerAliasUseCases, passthrough_complete, talker_alias_settings
logger = logging.getLogger(__name__)
@ -95,6 +95,7 @@ class BridgeUseCases:
send_bcsq: Any = None,
send_dmra_to_system: Any = None,
get_dmra_blocks: Any = None,
call_later: Any = None,
) -> None:
self._router = bridge_router
self._config = config
@ -105,7 +106,137 @@ class BridgeUseCases:
self._send_bcsq = send_bcsq # (system_name, tgid, stream_id) -> None; legacy OBP send_bcsq from router
self._send_dmra_to_system = send_dmra_to_system
self._get_dmra_blocks = get_dmra_blocks
self._call_later = call_later
self._talker_alias = TalkerAliasUseCases(config)
# (source_system, stream_id) -> {rf_src, peer, targets, timer}
self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {}
def _get_stream_dmra_blocks(self, source_system: str, stream_id: bytes) -> dict[int, bytes] | None:
if not self._get_dmra_blocks:
return None
return self._get_dmra_blocks(source_system, stream_id)
def _both_ta_key(self, source_system: str, stream_id: bytes) -> tuple[str, bytes]:
return (source_system, stream_id)
def _source_cannot_carry_ta(self, source_system: str) -> bool:
"""OpenBridge carries no Talker Alias (no DMRA UDP nor embedded LC TA).
In `both` mode there is nothing to wait for, so inject the template right away
instead of deferring (OBP streams are often short and would expire first).
"""
return self._config.get("SYSTEMS", {}).get(source_system, {}).get("MODE") == "OPENBRIDGE"
def _cancel_both_ta_wait(self, source_system: str, stream_id: bytes) -> None:
wait = self._both_ta_wait.pop(self._both_ta_key(source_system, stream_id), None)
if wait and wait.get("timer") and getattr(wait["timer"], "cancel", None):
try:
wait["timer"].cancel()
except Exception:
pass
def _register_both_ta_wait(
self,
source_system: str,
target_system: str,
rf_src: bytes,
stream_id: bytes,
source_peer: bytes,
) -> None:
"""Defer DMRA UDP + embed inject until MMDVM fragments arrive (both mode)."""
if not self._call_later:
return
key = self._both_ta_key(source_system, stream_id)
wait = self._both_ta_wait.setdefault(
key,
{"rf_src": rf_src, "peer": source_peer, "targets": set()},
)
wait["rf_src"] = rf_src
wait["peer"] = source_peer
wait["targets"].add(target_system)
if wait.get("timer"):
return
# Wait long enough to detect a slow source TA (MMDVM emits TA ~1 s in) before
# falling back to inject; passthrough is applied earlier via on_dmra_fragment_stored.
wait["timer"] = self._call_later(2.0, self._finalize_both_ta, source_system, stream_id)
def _finalize_both_ta(self, source_system: str, stream_id: bytes) -> None:
key = self._both_ta_key(source_system, stream_id)
wait = self._both_ta_wait.pop(key, None)
if not wait:
return
rf_src = wait["rf_src"]
peer = wait["peer"]
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
# No source TA within the window: fall back to inject (template).
fallback_inject = not (blocks and passthrough_complete(blocks))
for target_system in wait["targets"]:
if not self._talker_alias.should_send_on_vhead(target_system, stream_id):
continue
self._send_talker_alias_to_target(
source_system, target_system, rf_src, stream_id, peer,
force=True, fallback_inject=fallback_inject,
)
self._apply_both_ta_embed(source_system, rf_src, stream_id, force_inject=fallback_inject)
def _apply_both_ta_embed(
self,
source_system: str,
rf_src: bytes,
stream_id: bytes,
*,
force_inject: bool = False,
) -> None:
"""Install embed LC state once the TA decision is known (passthrough or inject)."""
if not self._get_protocols:
return
for proto in self._get_protocols().values():
status = getattr(proto, "STATUS", None)
if not isinstance(status, dict):
continue
for st in status.values():
if not isinstance(st, dict):
continue
if st.get("TX_STREAM_ID") != stream_id or st.get("TX_RFS") != rf_src:
continue
if st.get("TX_TA_ON"):
continue
target = st.get("_ta_target_system", source_system)
self._init_talker_alias_embed(
st, source_system, target, rf_src, stream_id, force_inject=force_inject,
)
def on_dmra_fragment_stored(
self,
source_system: str,
peer_id: bytes,
rf_src: bytes,
stream_id: bytes,
) -> None:
"""Source TA may complete after VHEAD (DMRA UDP or decoded from voice).
Once the buffer is complete, overlay the source TA on the outgoing embedded LC
(passthrough/both). For `both` with a pending wait, also relay DMRA now and cancel
the inject fallback.
"""
if talker_alias_settings(self._config, source_system)["mode"] not in ("both", "passthrough"):
return
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
if not blocks or not passthrough_complete(blocks):
return
key = self._both_ta_key(source_system, stream_id)
wait = self._both_ta_wait.get(key)
if not wait:
self._apply_both_ta_embed(source_system, rf_src, stream_id)
return
wait["peer"] = peer_id
self._cancel_both_ta_wait(source_system, stream_id)
for target_system in wait["targets"]:
if self._talker_alias.should_send_on_vhead(target_system, stream_id):
self._send_talker_alias_to_target(
source_system, target_system, rf_src, stream_id, peer_id, force=True,
)
self._apply_both_ta_embed(source_system, rf_src, stream_id)
def _send_talker_alias_to_target(
self,
@ -114,6 +245,9 @@ class BridgeUseCases:
rf_src: bytes,
stream_id: bytes,
source_peer: bytes,
*,
force: bool = False,
fallback_inject: bool = False,
) -> None:
"""Emit DMRA to an HBP target on VHEAD (once per target stream)."""
if not self._send_dmra_to_system:
@ -123,12 +257,28 @@ class BridgeUseCases:
return
if not self._talker_alias.should_send_on_vhead(target_system, stream_id):
return
if (
not force
and talker_alias_settings(self._config, source_system)["mode"] == "both"
and not (self._get_stream_dmra_blocks(source_system, stream_id) and passthrough_complete(
self._get_stream_dmra_blocks(source_system, stream_id) or {}
))
):
if self._source_cannot_carry_ta(source_system):
# OBP source can never supply TA: inject the template immediately.
fallback_inject = True
else:
self._register_both_ta_wait(
source_system, target_system, rf_src, stream_id, source_peer,
)
return
packets = self._talker_alias.packets_for_stream(
source_system,
rf_src,
stream_id,
self._get_dmra_blocks,
target_system=target_system,
fallback_inject=fallback_inject,
)
if not packets:
return
@ -138,6 +288,7 @@ class BridgeUseCases:
except Exception as e:
logger.warning("(ROUTER) send_dmra_to_system %s failed: %s", target_system, e)
return
self._talker_alias.mark_dmra_sent(target_system, stream_id)
sid = int_id(stream_id)
if peer_count:
logger.debug(
@ -164,10 +315,13 @@ class BridgeUseCases:
def clear_talker_alias_stream(self, system_name: str, stream_id: bytes) -> None:
"""Release per-stream TA dedupe state after VTERM."""
self._cancel_both_ta_wait(system_name, stream_id)
self._talker_alias.clear_stream(system_name, stream_id)
if not self._get_protocols:
return
proto = self._get_protocols().get(system_name)
if proto is not None and hasattr(proto, "clear_ta_stream_buffer"):
proto.clear_ta_stream_buffer(stream_id)
if proto is None:
return
status = getattr(proto, "STATUS", None)
@ -185,8 +339,31 @@ class BridgeUseCases:
target_system: str,
rf_src: bytes,
stream_id: bytes,
*,
force_inject: bool = False,
) -> None:
"""Prepare per-stream embedded TA state for DMRD voice injection."""
"""Prepare per-stream embedded TA state for DMRD voice injection.
The group LC is always rewritten for the destination TG by ``_rewrite_embed_lc``;
here we only set ``TX_TA_EMB`` when a Talker Alias should be overlaid. The source
TA (passthrough/both) is re-encoded from its decoded DMRA/voice blocks once they
arrive; the template is used for ``inject`` and the ``both`` fallback. If no TA is
available yet (e.g. at VHEAD, before the source TA has been decoded), TX_TA_EMB is
left unset and only the destination group LC is emitted until it becomes available.
"""
st["_ta_source_system"] = source_system
st["_ta_target_system"] = target_system
settings = talker_alias_settings(self._config, source_system)
if not settings["enabled"]:
return
target_mode = self._config.get("SYSTEMS", {}).get(target_system, {}).get("MODE")
if (
settings["mode"] == "both"
and self._source_cannot_carry_ta(source_system)
and target_mode in ("MASTER", "PEER")
):
# OBP source can never supply TA: inject the template immediately.
force_inject = True
st.pop("TX_TA_EMB", None)
st.pop("TX_TA_PHASE", None)
st.pop("TX_TA_ON", None)
@ -196,6 +373,7 @@ class BridgeUseCases:
stream_id,
self._get_dmra_blocks,
target_system=target_system,
fallback_inject=force_inject,
)
if emblcs:
st["TX_TA_EMB"], st["TX_TA_BLOCK_COUNT"] = emblcs
@ -203,8 +381,21 @@ class BridgeUseCases:
# First B1–B4 cycle carries group LC; TA on the next cycle.
st["TX_TA_ON"] = False
def _rewrite_embed_lc(self, dmrbits: bitarray, st: dict[str, Any], dtype_vseq: int, emb_key: str) -> None:
"""Replace embedded LC on voice bursts B–E; alternate group LC and TA blocks."""
def _rewrite_embed_lc(
self,
dmrbits: bitarray,
st: dict[str, Any],
dtype_vseq: int,
emb_key: str,
) -> None:
"""Replace embedded LC on voice bursts B–E (legacy bridge.py parity).
Group-call superframes always carry the **destination** group LC (``emb_key``),
re-encoded for the rewritten TGID — this is required for the voice to be accepted
by the receiving MMDVM (a mismatched embedded LC causes packet loss). When a
Talker Alias is available (``TX_TA_EMB``: injected template or the source TA
re-encoded from its DMRA/voice blocks) it is overlaid on alternate superframes.
"""
if dtype_vseq not in (1, 2, 3, 4):
return
ta_emb = st.get("TX_TA_EMB")

@ -65,6 +65,7 @@ class PlaybackUseCases:
self._playback_packets: list[bytes] = []
self._playback_index = 0
self._playback_stream_id = b""
self._playback_ta_from_stream = b""
self._delay_call: DelayedCall | None = None
self._packet_call: DelayedCall | None = None
self._idle_call: DelayedCall | None = None
@ -197,6 +198,14 @@ class PlaybackUseCases:
dst_id: bytes,
) -> None:
self.CALL_DATA.append(data)
if proto is not None and hasattr(proto, "store_ta_from_voice_burst") and len(data) >= 53:
bits = data[15]
frame_type = (bits & 0x30) >> 4
vseq = bits & 0xF
if frame_type != HBPF_DATA_SYNC and vseq in (1, 2, 3, 4):
proto.store_ta_from_voice_burst(
peer_id, rf_src, self._record_stream, vseq, data[20:53],
)
self._last_record_time = pkt_time
self._record_ctx = {
"slot": slot,
@ -263,6 +272,7 @@ class PlaybackUseCases:
ctx = self._record_ctx
slot = ctx.get("slot", 1)
recorded = self._ensure_vterm(list(self.CALL_DATA), slot)
self._playback_ta_from_stream = self._record_stream
self._reset_recording_state(slot)
if not recorded:
return
@ -363,6 +373,9 @@ class PlaybackUseCases:
return
self._playback_stream_id = bytes_4(randint(0x00, 0xFFFFFFFF))
if proto is not None and self._playback_ta_from_stream and hasattr(proto, "copy_ta_stream_buffer"):
proto.copy_ta_stream_buffer(self._playback_ta_from_stream, self._playback_stream_id)
self._playback_ta_from_stream = b""
logger.info(
"(%s) *START PLAYBACK* STREAM ID: %s SUB: %s REPEATER: %s TGID %s, TS %s, Duration: %.2f",
self._system,
@ -424,5 +437,6 @@ class PlaybackUseCases:
self._playback_packets = []
self._playback_index = 0
self._playback_stream_id = b""
self._playback_ta_from_stream = b""
self._delay_call = None
self._packet_call = None

@ -13,12 +13,16 @@ from ..domain.talker_alias import (
DMRA_BLOCK_COUNT,
DMRA_PAYLOAD_LEN,
buffer_from_blocks,
buffer_from_wire_blocks,
build_dmra_packet,
build_dmra_packets,
decode_7bit,
decode_ta_from_blocks,
encode_talker_alias_emblc,
encode_talker_alias_emblc_from_blocks,
required_ta_block_count,
talker_alias_decode_complete,
truncate_talker_alias,
is_ta_header_byte,
)
logger = logging.getLogger(__name__)
@ -77,18 +81,40 @@ def format_talker_alias_text(config: dict[str, Any], rf_src: bytes) -> str:
return truncate_talker_alias(text)
def _passthrough_buf_ready(blocks: dict[int, bytes], buf: bytes) -> bool:
need = required_ta_block_count(buf)
if need < 1 or need > DMRA_BLOCK_COUNT:
return False
if not all(i in blocks and len(blocks[i]) >= DMRA_PAYLOAD_LEN for i in range(need)):
return False
return talker_alias_decode_complete(buf)
def passthrough_complete(blocks: dict[int, bytes]) -> bool:
"""True when all four TA blocks were received."""
return all(i in blocks and len(blocks[i]) >= DMRA_PAYLOAD_LEN for i in range(DMRA_BLOCK_COUNT))
"""True when required TA blocks (1–4) are present and fully decode."""
if not blocks:
return False
buf = buffer_from_wire_blocks(blocks)
if _passthrough_buf_ready(blocks, buf):
return True
if blocks.get(0) and is_ta_header_byte(blocks[0][0]):
return _passthrough_buf_ready(blocks, buffer_from_blocks(blocks))
return False
def passthrough_packets_from_blocks(rf_src: bytes, blocks: dict[int, bytes]) -> list[bytes]:
"""Rebuild four DMRA packets from buffered block payloads."""
packets: list[bytes] = []
for block_id in range(DMRA_BLOCK_COUNT):
payload = blocks.get(block_id, b"\x00" * DMRA_PAYLOAD_LEN)
packets.append(build_dmra_packet(rf_src, block_id, payload))
return packets
"""Rebuild DMRA packets from buffered wire payloads (only blocks needed for TA)."""
if passthrough_complete(blocks):
buf = buffer_from_wire_blocks(blocks)
if not talker_alias_decode_complete(buf) and blocks.get(0) and is_ta_header_byte(blocks[0][0]):
buf = buffer_from_blocks(blocks)
count = required_ta_block_count(buf)
return [
build_dmra_packet(rf_src, block_id, blocks[block_id])
for block_id in range(count)
if block_id in blocks
]
return []
class TalkerAliasUseCases:
@ -104,11 +130,10 @@ class TalkerAliasUseCases:
self._embed_logged.discard((system_name, stream_id))
def should_send_on_vhead(self, target_system: str, stream_id: bytes) -> bool:
key = (target_system, stream_id)
if key in self._sent_streams:
return False
self._sent_streams.add(key)
return True
return (target_system, stream_id) not in self._sent_streams
def mark_dmra_sent(self, target_system: str, stream_id: bytes) -> None:
self._sent_streams.add((target_system, stream_id))
def packets_for_stream(
self,
@ -118,6 +143,7 @@ class TalkerAliasUseCases:
get_passthrough_blocks: Any,
*,
target_system: str | None = None,
fallback_inject: bool = False,
) -> list[bytes] | None:
"""Return DMRA packets to send on VHEAD, or None if TA disabled / nothing to send."""
settings = talker_alias_settings(self._config, source_system)
@ -131,7 +157,7 @@ class TalkerAliasUseCases:
if mode == "passthrough":
if not have_passthrough:
return None
text = decode_7bit(buffer_from_blocks(blocks))
text = decode_ta_from_blocks(blocks)
logger.debug(
"(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
@ -144,20 +170,28 @@ class TalkerAliasUseCases:
source_system, text, via, target, int_id(stream_id),
)
return build_dmra_packets(rf_src, text)
# both
# both: prefer the source's own TA. If a valid MMDVM DMRA buffer arrived,
# relay it. Otherwise the source's embedded LC (e.g. MMDVM voice) is passed
# through unchanged in the DMRD voice, so do NOT inject a template here.
if have_passthrough:
text = decode_7bit(buffer_from_blocks(blocks))
text = decode_ta_from_blocks(blocks)
logger.debug(
"(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
)
return passthrough_packets_from_blocks(rf_src, blocks)
text = format_talker_alias_text(self._config, rf_src)
if fallback_inject:
text = format_talker_alias_text(self._config, rf_src)
logger.debug(
"(%s) *TALKER ALIAS* inject '%s' (no source TA) via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
)
return build_dmra_packets(rf_src, text)
logger.debug(
"(%s) *TALKER ALIAS* inject '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
"(%s) *TALKER ALIAS* passthrough (source embedded TA) via %s -> %s stream %s",
source_system, via, target, int_id(stream_id),
)
return build_dmra_packets(rf_src, text)
return None
def embedded_emblc_for_stream(
self,
@ -167,43 +201,41 @@ class TalkerAliasUseCases:
get_passthrough_blocks: Any,
*,
target_system: str | None = None,
fallback_inject: bool = False,
) -> tuple[list[dict[int, Any]], int] | None:
"""Return (encode_emblc dicts, block count 1–4) for embedded TA in DMRD, or None."""
"""Return (encode_emblc dicts, block count 1–4) for the embedded TA to overlay, or None.
The group LC is rewritten separately for the destination TG; this only supplies the
Talker Alias overlaid on alternate superframes:
- **passthrough / both with source TA:** re-encode the source's decoded TA blocks
(round-trips losslessly via the fixed ``encode_emblc``).
- **inject / both fallback:** the configured template.
- otherwise (no TA available yet): ``None`` — only the destination group LC is sent.
"""
settings = talker_alias_settings(self._config, source_system)
if not settings["enabled"]:
return None
mode = settings["mode"]
target = target_system or source_system
via = "repeat" if source_system == target else "bridge"
mode = settings["mode"]
log_key = (source_system, stream_id)
log_embed = log_key not in self._embed_logged
def _log_embed(text: str, passthrough: bool) -> None:
if not log_embed:
def _log(text: str, suffix: str) -> None:
if log_key in self._embed_logged:
return
self._embed_logged.add(log_key)
kind = "passthrough" if passthrough else "inject"
logger.debug(
"(%s) *TALKER ALIAS* embed %s '%s' via %s -> %s stream %s",
source_system, kind, text, via, target, int_id(stream_id),
"(%s) *TALKER ALIAS* embed inject '%s'%s via %s -> %s stream %s",
source_system, text, suffix, via, target, int_id(stream_id),
)
blocks = get_passthrough_blocks(source_system, stream_id) if get_passthrough_blocks else None
have_passthrough = bool(blocks and passthrough_complete(blocks))
if mode == "passthrough":
if not have_passthrough:
return None
text = decode_7bit(buffer_from_blocks(blocks))
_log_embed(text, True)
return encode_talker_alias_emblc_from_blocks(blocks)
if mode == "inject":
text = format_talker_alias_text(self._config, rf_src)
_log_embed(text, False)
return encode_talker_alias_emblc(text)
if have_passthrough:
text = decode_7bit(buffer_from_blocks(blocks))
_log_embed(text, True)
if mode in ("passthrough", "both") and blocks and passthrough_complete(blocks):
_log(decode_ta_from_blocks(blocks), " (source TA)")
return encode_talker_alias_emblc_from_blocks(blocks)
text = format_talker_alias_text(self._config, rf_src)
_log_embed(text, False)
return encode_talker_alias_emblc(text)
if mode == "inject" or fallback_inject:
suffix = "" if mode == "inject" else " (no source TA)"
_log(format_talker_alias_text(self._config, rf_src), suffix)
return encode_talker_alias_emblc(format_talker_alias_text(self._config, rf_src))
return None

@ -121,7 +121,7 @@ def required_ta_block_count(buf: bytes) -> int:
def buffer_from_blocks(blocks: dict[int, bytes]) -> bytes:
"""Merge up to four block payloads into a 28-byte buffer."""
"""Merge up to four block payloads into a 28-byte buffer (ADN inject layout)."""
buf = bytearray(DMRA_BUF_LEN)
for block_id, payload in blocks.items():
if 0 <= block_id < DMRA_BLOCK_COUNT and payload:
@ -130,6 +130,46 @@ def buffer_from_blocks(blocks: dict[int, bytes]) -> bytes:
return bytes(buf)
def buffer_from_wire_blocks(blocks: dict[int, bytes]) -> bytes:
"""Merge MMDVMHost DMRA wire payloads into the 28-byte TA buffer.
``writeTalkerAlias`` copies ``data + 2`` (7 bytes) from the 9-byte embedded LC;
``CDMRTA::add(blockId, data + 2, 7)`` stores those at ``m_buf[blockId * 7]``.
The UDP payload is therefore the same layout as ``buffer_from_blocks``.
"""
return buffer_from_blocks(blocks)
def is_ta_header_byte(byte0: int) -> bool:
"""True if byte looks like ETSI TA header (UTF-8/ISO formats; not 7-bit)."""
if (byte0 & 1) != 0:
return False
fmt = (byte0 >> 6) & 0x03
size = (byte0 >> 1) & 0x1F
return fmt in (1, 2) and 1 <= size <= TALKER_ALIAS_MAX_LEN
def talker_alias_decode_complete(buf: bytes) -> bool:
"""True when TA buffer decodes with expected size (MMDVM DMRTA parity)."""
if len(buf) < 1 or not is_ta_header_byte(buf[0]):
return False
ta_size = (buf[0] >> 1) & 0x1F
text = decode_ta(buf).rstrip("\x00").strip()
return len(text) >= ta_size
def decode_ta_from_blocks(blocks: dict[int, bytes]) -> str:
"""Decode TA from buffered DMRA blocks (MMDVM wire layout, then ADN layout)."""
buf = buffer_from_wire_blocks(blocks)
if talker_alias_decode_complete(buf):
return decode_ta(buf).rstrip("\x00").strip()
if blocks.get(0) and is_ta_header_byte(blocks[0][0]):
buf = buffer_from_blocks(blocks)
if talker_alias_decode_complete(buf):
return decode_ta(buf).rstrip("\x00").strip()
return ""
def build_dmra_packets(rf_src: bytes, text: str) -> list[bytes]:
"""Build HBP DMRA packets (1–4) for server injection."""
rf = rf_src[:3] if len(rf_src) >= 3 else rf_src.ljust(3, b"\x00")[:3]
@ -151,13 +191,81 @@ def build_dmra_packet(rf_src: bytes, block_id: int, payload: bytes) -> bytes:
def parse_dmra_packet(data: bytes) -> tuple[bytes, int, bytes] | None:
"""Parse DMRA packet; returns (rf_src, block_id, payload) or None."""
"""Parse MMDVMHost HBP DMRA (15 bytes), per ``DMRNetwork::writeTalkerAlias``.
Layout: ``DMRA`` | src_id[24-bit BE @4–6] | type @7 (0–3) | memcpy(@8, data+2, 7).
``data`` is the 9-byte embedded LC (``CDMREmbeddedData::getRawData``); wire bytes equal
``CDMRTA::add(blockId, data+2, 7)`` → ``m_buf[blockId*7 : blockId*7+7]``.
"""
if len(data) < DMRA_PACKET_LEN or data[:4] != DMRA_OPCODE:
return None
rf_src = data[4:7]
block_id = data[7]
payload = data[8:15]
return rf_src, block_id, payload
if block_id > 3:
return None
return data[4:7], block_id, data[8:15]
def store_ta_block(blocks: dict[int, bytes], block_id: int, payload: bytes) -> bool:
"""Store one wire fragment (7 bytes at ``m_buf[block_id*7]``)."""
if block_id < 0 or block_id > 3:
return False
blocks[block_id] = payload[:DMRA_PAYLOAD_LEN]
return True
def store_ta_from_embed_lc(blocks: dict[int, bytes], block_id: int, lc9: bytes) -> bool:
"""Store from 9-byte embedded LC (MMDVM ``data`` in ``DMRSlot`` TA FLCO 4–7 path)."""
if block_id < 0 or block_id > 3 or len(lc9) < 9:
return False
return store_ta_block(blocks, block_id, lc9[2:9])
def talker_alias_block_id_from_lc(lc: bytes) -> int | None:
"""FLCO 4–7 (MMDVM TALKER_ALIAS_HEADER..BLOCK3) → block index 0–3, else None."""
if len(lc) < 9:
return None
flco = lc[0]
if flco < FLCO_TALKER_ALIAS_HEADER or flco > FLCO_TALKER_ALIAS_BLOCK3:
return None
return flco - FLCO_TALKER_ALIAS_HEADER
def try_buffer_ta_from_voice_fragments(
acc: dict[int, "bitarray"],
vseq: int,
dmrpkt: bytes,
blocks: dict[int, bytes],
) -> bool:
"""Reassemble embedded LC across voice bursts B–E (vseq 1–4) and store if it is a TA.
``dmr_utils3.bptc.decode_emblc`` is correct on properly FEC-encoded input (the encode
helper is the buggy one), so a TA sent by a real radio decodes losslessly here.
Returns True when a TA block was decoded and stored.
"""
from dmr_utils3 import bptc, decode
if vseq not in (1, 2, 3, 4) or len(dmrpkt) < 33:
return False
try:
embed = decode.voice(dmrpkt)["EMBED"]
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)):
return False
try:
lc = bptc.decode_emblc(acc[1] + acc[2] + acc[3] + acc[4])
except Exception:
acc.clear()
return False
acc.clear()
block_id = talker_alias_block_id_from_lc(lc)
if block_id is None:
return False
return store_ta_from_embed_lc(blocks, block_id, lc)
def talker_alias_lc_bytes(block_id: int, payload7: bytes) -> bytes:
@ -167,23 +275,90 @@ def talker_alias_lc_bytes(block_id: int, payload7: bytes) -> bytes:
return bytes([FLCO_TALKER_ALIAS_HEADER + block, 0x00]) + payload
def encode_emblc(_lc: bytes) -> dict[int, bitarray]:
"""Embedded-LC BPTC encode with the ``dmr_utils3`` bug fixed.
Upstream ``dmr_utils3.bptc.encode_emblc`` builds segment D's second row from
``_binlc[24]`` instead of ``_binlc[25]`` (duplicating bit 24, dropping bit 25), which
corrupts one bit of the embedded LC (e.g. injected ``Rodrigo`` decodes as ``Rodrigg``).
This is a faithful copy with that single index corrected, so injected Talker Alias
round-trips losslessly through ``decode_emblc`` and on real radios.
"""
from bitarray import bitarray
from dmr_utils3 import crc, hamming
_csum = crc.csum5(_lc)
_binlc = bitarray(endian="big")
_binlc.frombytes(_lc)
_binlc.insert(32, _csum[0])
_binlc.insert(43, _csum[1])
_binlc.insert(54, _csum[2])
_binlc.insert(65, _csum[3])
_binlc.insert(76, _csum[4])
for index in range(0, 112, 16):
for hindex, hbit in zip(
range(index + 11, index + 16), hamming.enc_16114(_binlc[index:index + 11])
):
_binlc.insert(hindex, hbit)
for index in range(0, 16):
_binlc.insert(
index + 112,
_binlc[index] ^ _binlc[index + 16] ^ _binlc[index + 32] ^ _binlc[index + 48]
^ _binlc[index + 64] ^ _binlc[index + 80] ^ _binlc[index + 96],
)
def _seg(rows: tuple[tuple[int, ...], ...]) -> bitarray:
out = bitarray(endian="big")
for row in rows:
out.extend([_binlc[i] for i in row])
return out
emblc_b = _seg((
(0, 16, 32, 48, 64, 80, 96, 112),
(1, 17, 33, 49, 65, 81, 97, 113),
(2, 18, 34, 50, 66, 82, 98, 114),
(3, 19, 35, 51, 67, 83, 99, 115),
))
emblc_c = _seg((
(4, 20, 36, 52, 68, 84, 100, 116),
(5, 21, 37, 53, 69, 85, 101, 117),
(6, 22, 38, 54, 70, 86, 102, 118),
(7, 23, 39, 55, 71, 87, 103, 119),
))
emblc_d = _seg((
(8, 24, 40, 56, 72, 88, 104, 120),
(9, 25, 41, 57, 73, 89, 105, 121), # fixed: bit 25 (dmr_utils3 wrongly uses 24)
(10, 26, 42, 58, 74, 90, 106, 122),
(11, 27, 43, 59, 75, 91, 107, 123),
))
emblc_e = _seg((
(12, 28, 44, 60, 76, 92, 108, 124),
(13, 29, 45, 61, 77, 93, 109, 125),
(14, 30, 46, 62, 78, 94, 110, 126),
(15, 31, 47, 63, 79, 95, 111, 127),
))
return {1: emblc_b, 2: emblc_c, 3: emblc_d, 4: emblc_e}
def encode_talker_alias_emblc(text: str) -> tuple[list[dict[int, bitarray]], int]:
"""Embedded-LC dicts for TA blocks 0..N-1 and block count N (1–4)."""
from dmr_utils3 import bptc
encoded = encode_utf8(text)
blocks = blocks_from_buffer(encoded)
count = required_ta_block_count(encoded)
emblcs = [bptc.encode_emblc(talker_alias_lc_bytes(i, blocks[i])) for i in range(count)]
emblcs = [encode_emblc(talker_alias_lc_bytes(i, blocks[i])) for i in range(count)]
return emblcs, count
def encode_talker_alias_emblc_from_blocks(blocks: dict[int, bytes]) -> tuple[list[dict[int, bitarray]], int]:
"""Build embedded TA LC dicts from buffered DMRA payloads."""
from dmr_utils3 import bptc
buf = buffer_from_blocks(blocks)
buf_blocks = blocks_from_buffer(buf)
"""Build embedded TA LC dicts from buffered DMRA wire payloads."""
buf = buffer_from_wire_blocks(blocks)
count = required_ta_block_count(buf)
emblcs = [bptc.encode_emblc(talker_alias_lc_bytes(i, buf_blocks[i])) for i in range(count)]
if not talker_alias_decode_complete(buf) and blocks.get(0) and is_ta_header_byte(blocks[0][0]):
buf = buffer_from_blocks(blocks)
count = required_ta_block_count(buf)
emblcs = [
encode_emblc(talker_alias_lc_bytes(i, blocks[i]))
for i in range(count)
if i in blocks
]
return emblcs, count

@ -22,6 +22,6 @@
###############################################################################
from .config_loader import YamlConfigLoader
from .logging_config import reopen_file_handlers, setup_logging
from .logging_config import reapply_log_level, reopen_file_handlers, setup_logging
__all__ = ["YamlConfigLoader", "reopen_file_handlers", "setup_logging"]
__all__ = ["YamlConfigLoader", "reapply_log_level", "reopen_file_handlers", "setup_logging"]

@ -12,6 +12,7 @@ from typing import Any, Callable
from ..domain.errors import ConfigError
from .config_loader import YamlConfigLoader
from .logging_config import reapply_log_level
from .config_normalizer import (
apply_talker_alias_defaults,
ensure_system_runtime_config,
@ -146,6 +147,9 @@ def reload_server_config(
raise
merge_top_level_config(config, incoming)
if "LOGGER" in incoming:
level_name = reapply_log_level(config.get("LOGGER", {}))
log.info("(CONFIG-RELOAD) LOG_LEVEL applied: %s", level_name)
old_systems = dict(config.get("SYSTEMS", {}))
new_systems = incoming.get("SYSTEMS", {})

@ -72,6 +72,27 @@ def reopen_file_handlers(logger: logging.Logger | None = None) -> int:
return count
def reapply_log_level(log_config: dict[str, Any]) -> str:
"""Apply LOGGER.LOG_LEVEL after SIGHUP reload (handlers unchanged).
Returns the level name applied (e.g. ``DEBUG``).
"""
if not logging_enabled(log_config):
level = logging.CRITICAL
else:
level = getattr(logging, (log_config.get("LOG_LEVEL", "INFO")).upper(), logging.INFO)
log_name = log_config.get("LOG_NAME", "ADN")
root = logging.getLogger()
root.setLevel(level)
app_logger = logging.getLogger(log_name)
app_logger.setLevel(level)
for handler in root.handlers:
handler.setLevel(level)
for handler in app_logger.handlers:
handler.setLevel(level)
return logging.getLevelName(level)
def setup_logging(log_config: dict[str, Any]) -> logging.Logger:
"""Configure logging from CONFIG['LOGGER']. Returns application logger."""
log_name = log_config.get("LOG_NAME", "ADN")

@ -43,7 +43,13 @@ from twisted.internet import reactor, task
from twisted.internet.protocol import DatagramProtocol
from ...domain import bytes_4, int_id
from ...domain.talker_alias import DMRA_PACKET_LEN, parse_dmra_packet
from ...domain.talker_alias import (
DMRA_PACKET_LEN,
decode_ta_from_blocks,
parse_dmra_packet,
store_ta_block,
try_buffer_ta_from_voice_fragments,
)
from ..hbp_constants import (
BC,
BCKA,
@ -146,6 +152,7 @@ class HBPProtocol(DatagramProtocol):
on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None,
on_talker_alias_local_repeat: Callable[[str, bytes, bytes, bytes], None] | None = None,
on_talker_alias_stream_end: Callable[[str, bytes], None] | None = None,
on_dmra_fragment_stored: Callable[[str, bytes, bytes, bytes], None] | None = None,
) -> None:
self._CONFIG = config
self._system = system_name
@ -161,6 +168,7 @@ class HBPProtocol(DatagramProtocol):
self._on_obp_bcsq_received = on_obp_bcsq_received
self._on_talker_alias_local_repeat = on_talker_alias_local_repeat
self._on_talker_alias_stream_end = on_talker_alias_stream_end
self._on_dmra_fragment_stored = on_dmra_fragment_stored
self._config = config.get("SYSTEMS", {}).get(system_name, {})
if self._config.get("MODE") == "OPENBRIDGE":
self._laststrid = deque([], 20)
@ -177,10 +185,14 @@ class HBPProtocol(DatagramProtocol):
self._peers = self._config.setdefault("PEERS", {})
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()
else:
self._peers = {}
self._dmra_by_stream = {}
self._dmra_rf_stream = {}
self._ta_voice_acc = {}
self._ta_decoded_logged = set()
if self._config.get("MODE") == "PEER":
self._dmra_downlink: dict[bytes, dict[str, Any]] = {}
if self._config.get("MODE") == "PEER":
@ -250,8 +262,6 @@ class HBPProtocol(DatagramProtocol):
if not parsed:
return
rf_src, block_id, payload = parsed
if block_id > 3:
return
stream_id = self._dmra_rf_stream.get((peer_id, rf_src))
if not stream_id:
return
@ -260,9 +270,62 @@ class HBPProtocol(DatagramProtocol):
stream_id,
{"blocks": {}, "rf_src": rf_src, "peer": peer_id, "last": now},
)
entry["blocks"][block_id] = payload
if not store_ta_block(entry["blocks"], block_id, payload):
return
entry["last"] = now
entry["rf_src"] = rf_src
if self._on_dmra_fragment_stored:
self._on_dmra_fragment_stored(self._system, peer_id, rf_src, stream_id)
def store_ta_from_voice_burst(
self,
peer_id: bytes,
rf_src: bytes,
stream_id: bytes,
vseq: int,
dmrpkt: bytes,
) -> None:
"""Buffer TA from embedded LC in voice bursts B–E (MMDVM DMRSlot path)."""
if self._config.get("MODE") != "MASTER" or not stream_id:
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()},
)
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:
text = decode_ta_from_blocks(entry["blocks"])
if text:
self._ta_decoded_logged.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),
)
if self._on_dmra_fragment_stored:
self._on_dmra_fragment_stored(self._system, peer_id, rf_src, stream_id)
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)
def copy_ta_stream_buffer(self, from_stream: bytes, to_stream: bytes) -> None:
"""Carry decoded TA blocks from recording stream to parrot playback stream."""
entry = self._dmra_by_stream.get(from_stream)
if not entry or not entry.get("blocks"):
return
self._dmra_by_stream[to_stream] = {
"blocks": dict(entry["blocks"]),
"rf_src": entry.get("rf_src", b""),
"peer": entry.get("peer", b""),
"last": time.time(),
}
def get_dmra_blocks(self, stream_id: bytes) -> dict[int, bytes] | None:
"""Return buffered DMRA block payloads for a stream, if any."""
@ -281,6 +344,8 @@ class HBPProtocol(DatagramProtocol):
for stream_id in list(self._dmra_by_stream):
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)
def send_dmra_to_peers(self, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
"""Send DMRA packets to logged-in peers (MASTER downlink). Returns peer count."""
@ -567,6 +632,15 @@ class HBPProtocol(DatagramProtocol):
if sub_map is not None:
sub_map[_rf_src] = (self._system, _slot, pkt_time)
self.note_dmrd_stream(_peer_id, _rf_src, _stream_id)
if (
_call_type in ("group", "vcsbk")
and _frame_type != HBPF_DATA_SYNC
and _dtype_vseq in (1, 2, 3, 4)
and len(_data) >= 53
):
self.store_ta_from_voice_burst(
_peer_id, _rf_src, _stream_id, _dtype_vseq, _data[20:53],
)
if self._config.get("REPEAT", True):
pkt = [_data[:11], b"", _data[15:]]
for _peer in self._peers:
@ -838,8 +912,16 @@ class HBPProtocol(DatagramProtocol):
_peer_id = pid
break
if _peer_id is not None and len(_data) >= DMRA_PACKET_LEN:
logger.debug("(%s) Peer has sent Talker Alias packet %s", self._system, _data)
self.store_dmra_packet(_peer_id, _data)
if parse_dmra_packet(_data):
logger.debug("(%s) Peer has sent Talker Alias packet %s", self._system, _data)
self.store_dmra_packet(_peer_id, _data)
else:
logger.debug(
"(%s) Peer DMRA ignored (MMDVM expects byte7=0-3, got %s); raw %s",
self._system,
_data[7],
_data,
)
elif _command == PRIN:
logger.info("(%s) *ProxyInfo* Connection from IP:Port: %s", self._system, _data.decode("utf8", errors="replace")[4:])
@ -1575,6 +1657,7 @@ def HBPProtocolFactory(
on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None,
on_talker_alias_local_repeat: Callable[[str, bytes, bytes, bytes], None] | None = None,
on_talker_alias_stream_end: Callable[[str, bytes], None] | None = None,
on_dmra_fragment_stored: Callable[[str, bytes, bytes, bytes], None] | None = None,
) -> HBPProtocol:
"""Create HBP protocol instance (legacy: one HBSYSTEM per system)."""
return HBPProtocol(
@ -1592,4 +1675,5 @@ def HBPProtocolFactory(
on_obp_bcsq_received=on_obp_bcsq_received,
on_talker_alias_local_repeat=on_talker_alias_local_repeat,
on_talker_alias_stream_end=on_talker_alias_stream_end,
on_dmra_fragment_stored=on_dmra_fragment_stored,
)

@ -322,6 +322,7 @@ def main() -> None:
send_bcsq=send_bcsq,
send_dmra_to_system=send_dmra_to_system,
get_dmra_blocks=get_dmra_blocks,
call_later=reactor.callLater,
)
bridge_use_cases.apply_startup_bridges()
report_factory.set_bridges(bridge_router.get_bridges())
@ -481,6 +482,7 @@ def main() -> None:
on_obp_bcsq_received=bridge_use_cases.on_obp_bcsq_received,
on_talker_alias_local_repeat=bridge_use_cases.send_talker_alias_local_repeat,
on_talker_alias_stream_end=bridge_use_cases.clear_talker_alias_stream,
on_dmra_fragment_stored=bridge_use_cases.on_dmra_fragment_stored,
)
def _listen_system(_name: str, bind: BindSpec, protocol: Any) -> Any:

Loading…
Cancel
Save

Powered by TurnKey Linux.