diff --git a/docs/en/server/user-guide/configuration.md b/docs/en/server/user-guide/configuration.md index 5726b47..f2aeb7b 100644 --- a/docs/en/server/user-guide/configuration.md +++ b/docs/en/server/user-guide/configuration.md @@ -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. diff --git a/docs/en/server/user-guide/talker-alias.md b/docs/en/server/user-guide/talker-alias.md index 8b8c8d0..45686bb 100644 --- a/docs/en/server/user-guide/talker-alias.md +++ b/docs/en/server/user-guide/talker-alias.md @@ -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. diff --git a/docs/es/server/user-guide/configuration.md b/docs/es/server/user-guide/configuration.md index 7e6b8f2..7911d5c 100644 --- a/docs/es/server/user-guide/configuration.md +++ b/docs/es/server/user-guide/configuration.md @@ -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. diff --git a/pyproject.toml b/pyproject.toml index 32b924e..c437104 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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" diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 8492991..2ec557d 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -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") diff --git a/src/adn_server/application/playback_use_cases.py b/src/adn_server/application/playback_use_cases.py index af07732..e6761b0 100644 --- a/src/adn_server/application/playback_use_cases.py +++ b/src/adn_server/application/playback_use_cases.py @@ -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 diff --git a/src/adn_server/application/talker_alias_use_cases.py b/src/adn_server/application/talker_alias_use_cases.py index 0a8c925..cd8769c 100644 --- a/src/adn_server/application/talker_alias_use_cases.py +++ b/src/adn_server/application/talker_alias_use_cases.py @@ -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 diff --git a/src/adn_server/domain/talker_alias.py b/src/adn_server/domain/talker_alias.py index 5044840..0d12fe0 100644 --- a/src/adn_server/domain/talker_alias.py +++ b/src/adn_server/domain/talker_alias.py @@ -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 diff --git a/src/adn_server/infrastructure/__init__.py b/src/adn_server/infrastructure/__init__.py index 1d67bd1..439606a 100644 --- a/src/adn_server/infrastructure/__init__.py +++ b/src/adn_server/infrastructure/__init__.py @@ -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"] diff --git a/src/adn_server/infrastructure/config_reload.py b/src/adn_server/infrastructure/config_reload.py index ff5b055..69acc51 100644 --- a/src/adn_server/infrastructure/config_reload.py +++ b/src/adn_server/infrastructure/config_reload.py @@ -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", {}) diff --git a/src/adn_server/infrastructure/logging_config.py b/src/adn_server/infrastructure/logging_config.py index ec79e45..9a26557 100644 --- a/src/adn_server/infrastructure/logging_config.py +++ b/src/adn_server/infrastructure/logging_config.py @@ -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") diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 49100be..a7c27f3 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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, ) diff --git a/src/adn_server/main.py b/src/adn_server/main.py index 4d45b83..1fcbafb 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -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: