From 4e8a2d04b55203f9ed013d74b19dcd5fd4d7d24c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Wed, 27 May 2026 01:17:43 -0400 Subject: [PATCH] fix: parrot recording/playback end-of-call and re-key handling Add idle (5s) and max (180s) recording timers, synthetic VTERM when the hotspot ends without terminator, and non-blocking reactor playback. Handle hotspot re-key mid-PTT without a second VHEAD and renumber playback seq so bridge ingress does not drop frames. Restore prohibited static bridge legs (TG 9990 ECHO) after bridge reset and options refresh. --- .../application/bridge_use_cases.py | 65 +++ .../application/playback_use_cases.py | 385 ++++++++++++++++-- 2 files changed, 417 insertions(+), 33 deletions(-) diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index d217b79..89769af 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -770,9 +770,73 @@ class BridgeUseCases: del sys_cfg["_opt_key"] except KeyError: pass + try: + del sys_cfg["_options_static_apply_fp"] + except KeyError: + pass + self._restore_prohibited_static_bridge_legs(system_name) sys_cfg["_reset"] = False sys_cfg["_resetlog"] = False + def _restore_prohibited_static_bridge_legs(self, system_name: str) -> None: + """After BRIDGERESET / peer RPTO: restore static TGs in prohibited_tgs (parity with _make_echo_bridges).""" + prohibited_tgs = (0, 1, 2, 3, 4, 5, 9, 9990, 9991, 9992, 9993, 9994, 9995, 9996, 9997, 9998, 9999) + sys_cfg = self._config.get("SYSTEMS", {}).get(system_name, {}) + if sys_cfg.get("MODE") != "MASTER" or not sys_cfg.get("ENABLED", True): + return + bridges = self._router.get_bridges() + now = time.time() + for ts, static_key, acl_key in ( + (1, "TS1_STATIC", "TGID_TS1_ACL"), + (2, "TS2_STATIC", "TGID_TS2_ACL"), + ): + for tg_s in str(sys_cfg.get(static_key) or "").split(","): + tg_s = tg_s.strip() + if not tg_s: + continue + try: + tg = int(tg_s) + except ValueError: + continue + if tg not in prohibited_tgs: + continue + if sys_cfg.get("USE_ACL") and not self.acl_check( + bytes_3(tg), sys_cfg.get(acl_key, (True, [])) + ): + continue + bridge_key = str(tg) + if bridge_key not in bridges: + continue + timeout_sec = (1.0 / 6.0) * 60.0 + leg: dict[str, Any] = { + "SYSTEM": system_name, + "TS": ts, + "TGID": bytes_3(tg), + "ACTIVE": True, + "TIMEOUT": timeout_sec, + "TO_TYPE": "NONE", + "ON": [], + "OFF": [], + "RESET": [], + "TIMER": now + timeout_sec, + } + for i, entry in enumerate(bridges[bridge_key]): + if entry.get("SYSTEM") == system_name and entry.get("TS") == ts: + if entry.get("ACTIVE") and entry.get("TO_TYPE") == "NONE": + break + bridges[bridge_key][i] = leg + logger.info( + "(ROUTER) Restored service bridge leg: %s bridge %s TS %s", + system_name, bridge_key, ts, + ) + break + else: + bridges[bridge_key].append(leg) + logger.info( + "(ROUTER) Re-added service bridge leg: %s bridge %s TS %s", + system_name, bridge_key, ts, + ) + def _remove_bridge_system(self, system_name: str, bridges: dict[str, list[dict[str, Any]]]) -> None: """Remove all bridge entries for system (legacy remove_bridge_system).""" to_remove: list[str] = [] @@ -1224,6 +1288,7 @@ class BridgeUseCases: # this exact fingerprint after a previous RPTO in this process. _fp = f"{new_ts1}|{new_ts2}|{int(_tmout)}" if sys_cfg.get("_options_static_apply_fp") == _fp: + self._restore_prohibited_static_bridge_legs(system_name) return # Legacy: reset TGs that were removed (bridge_master.py 1736-1767) old_ts1 = str(sys_cfg.get("TS1_STATIC") or "").strip() diff --git a/src/adn_server/application/playback_use_cases.py b/src/adn_server/application/playback_use_cases.py index 7babd68..7c53fcb 100644 --- a/src/adn_server/application/playback_use_cases.py +++ b/src/adn_server/application/playback_use_cases.py @@ -21,28 +21,54 @@ # Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA ############################################################################### -"""Exact port of playback.py dmrd_received: record group voice, then play back with new stream ID.""" +"""Port of playback.py dmrd_received: record group voice, then play back with new stream ID.""" from __future__ import annotations import logging from random import randint -from time import sleep, time +from time import time from typing import Any -from ..domain import int_id, bytes_4, HBPF_DATA_SYNC, HBPF_SLT_VTERM +from twisted.internet import reactor +from twisted.internet.base import DelayedCall + +from ..domain import int_id, bytes_4, HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM logger = logging.getLogger(__name__) +# Legacy playback.py: sleep(2) before playback, sleep(0.06) between packets. +_PLAYBACK_DELAY_S = 2.0 +_PACKET_INTERVAL_S = 0.06 +# Match bridge stream_trimmer_loop RX idle (bridge_use_cases / legacy bridge_master). +_RECORD_IDLE_S = 5.0 +# HBP ingress source timeout (bridge_master.py / bridge_use_cases dmrd_received ~2183). +_SOURCE_MAX_S = 180.0 + class PlaybackUseCases: - """Exact port of legacy playback class. One instance per system.""" + """Legacy playback class behaviour; playback is scheduled on the reactor (non-blocking).""" def __init__(self, system_name: str, get_protocol: Any = None) -> None: self._system = system_name self._get_protocol = get_protocol self.STATUS: dict[str, Any] = {} self.CALL_DATA: list[bytes] = [] + self._rx_stream_by_slot: dict[int, bytes] = {} + self._record_stream: bytes = b"" + self._seen_record_streams: set[bytes] = set() + self._ignored_streams_logged: set[bytes] = set() + self._record_ctx: dict[str, Any] = {} + self._recording_active = False + self._last_record_time = 0.0 + self._playback_busy = False + self._playback_packets: list[bytes] = [] + self._playback_index = 0 + self._playback_stream_id = b"" + self._delay_call: DelayedCall | None = None + self._packet_call: DelayedCall | None = None + self._idle_call: DelayedCall | None = None + self._max_call: DelayedCall | None = None def dmrd_received( self, @@ -58,41 +84,334 @@ class PlaybackUseCases: stream_id: bytes, data: bytes, ) -> None: - """Exact port of playback.dmrd_received (playback.py lines 114-161).""" + """Port of playback.dmrd_received (playback.py lines 114-161).""" + if call_type != "group": + return + + # Legacy blocked the reactor during delay + playback; ignore new voice meanwhile. + if self._playback_busy: + return + pkt_time = time() + proto = self._get_protocol() if self._get_protocol else None - if call_type == "group": - proto = self._get_protocol() if self._get_protocol else None - slot_status = getattr(proto, "STATUS", {}).get(slot, {}) if proto else {} + if self._recording_active and self.CALL_DATA and self._recording_exceeded_max(pkt_time): + self._on_record_max_duration(proto) + return - if stream_id != slot_status.get("RX_STREAM_ID", b"\x00"): - self.STATUS["RX_START"] = pkt_time + slot_status = getattr(proto, "STATUS", {}).get(slot, {}) if proto else {} + record_stream = self._record_stream or self._rx_stream_by_slot.get( + slot, slot_status.get("RX_STREAM_ID", b"\x00"), + ) + is_vhead = frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD + idle_gap = pkt_time - self._last_record_time if self._last_record_time else 0.0 + + if stream_id != record_stream: + # Hotspot re-key mid-PTT: new VHEAD shortly after the last voice frame. + if ( + self._recording_active + and self.CALL_DATA + and is_vhead + and stream_id not in self._seen_record_streams + and idle_gap < _RECORD_IDLE_S + ): + prev = record_stream + self._record_stream = stream_id + self._rx_stream_by_slot[slot] = stream_id + self._seen_record_streams.add(stream_id) + logger.debug( + "(%s) Recording continues: stream %s -> %s (same PTT)", + self._system, int_id(prev), int_id(stream_id), + ) + # Same PTT re-key: track new stream but do not store a second VHEAD. + self._last_record_time = pkt_time + self._record_ctx = { + "slot": slot, + "rf_src": rf_src, + "peer_id": peer_id, + "dst_id": dst_id, + } + self._schedule_record_idle(proto) + return + # New PTT while prior recording never got VTERM (user released, hotspot cut, etc.). + if self._recording_active and self.CALL_DATA and is_vhead and idle_gap >= _RECORD_IDLE_S: logger.info( - "(%s) *START RECORDING* STREAM ID: %s SUB: %s REPEATER: %s TGID %s, TS %s", - self._system, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, + "(%s) New PTT after %.1fs idle; discarding prior recording (%d packets)", + self._system, idle_gap, len(self.CALL_DATA), ) - self.CALL_DATA.append(data) + self._reset_recording_state(slot) + self._begin_recording(slot, stream_id, pkt_time, rf_src, peer_id, dst_id, proto) + self._append_voice(data, pkt_time, proto, slot, rf_src, peer_id, dst_id) return + if self._recording_active and self.CALL_DATA: + if stream_id not in self._ignored_streams_logged: + self._ignored_streams_logged.add(stream_id) + logger.debug( + "(%s) Ignoring stream %s (current recording stream %s)", + self._system, int_id(stream_id), int_id(record_stream), + ) + return + if self.CALL_DATA: + logger.debug( + "(%s) Discarding stale recording (%d packets) before new stream %s", + self._system, len(self.CALL_DATA), int_id(stream_id), + ) + self._reset_recording_state(slot) + self._begin_recording(slot, stream_id, pkt_time, rf_src, peer_id, dst_id, proto) + self._append_voice(data, pkt_time, proto, slot, rf_src, peer_id, dst_id) + return - if (frame_type == HBPF_DATA_SYNC) and (dtype_vseq == HBPF_SLT_VTERM) and (slot_status.get("RX_TYPE") != HBPF_SLT_VTERM) and self.CALL_DATA: - call_duration = pkt_time - self.STATUS.get("RX_START", pkt_time) - self.CALL_DATA.append(data) - logger.info("(%s) *END RECORDING* STREAM ID: %s", self._system, int_id(stream_id)) - sleep(2) - _new_stream_id = bytes_4(randint(0x00, 0xFFFFFFFF)) - logger.info( - "(%s) *START PLAYBACK* STREAM ID: %s SUB: %s REPEATER: %s TGID %s, TS %s, Duration: %s", - self._system, int_id(_new_stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, call_duration, + if ( + frame_type == HBPF_DATA_SYNC + and dtype_vseq == HBPF_SLT_VTERM + and self._recording_active + and self.CALL_DATA + and stream_id == record_stream + ): + call_duration = pkt_time - self.STATUS.get("RX_START", pkt_time) + self.CALL_DATA.append(data) + logger.info("(%s) *END RECORDING* STREAM ID: %s", self._system, int_id(stream_id)) + self._commit_recording(proto, call_duration) + return + + if self._recording_active and self.CALL_DATA: + self._append_voice(data, pkt_time, proto, slot, rf_src, peer_id, dst_id) + + def _append_voice( + self, + data: bytes, + pkt_time: float, + proto: Any, + slot: int, + rf_src: bytes, + peer_id: bytes, + dst_id: bytes, + ) -> None: + self.CALL_DATA.append(data) + self._last_record_time = pkt_time + self._record_ctx = { + "slot": slot, + "rf_src": rf_src, + "peer_id": peer_id, + "dst_id": dst_id, + } + self._schedule_record_idle(proto) + + def _schedule_record_idle(self, proto: Any) -> None: + self._cancel_record_idle() + if not self._recording_active: + return + self._idle_call = reactor.callLater(_RECORD_IDLE_S, self._on_record_idle, proto) + + def _cancel_record_idle(self) -> None: + if self._idle_call is not None and self._idle_call.active(): + self._idle_call.cancel() + self._idle_call = None + + def _on_record_idle(self, proto: Any) -> None: + self._idle_call = None + if not self._recording_active or not self.CALL_DATA or self._playback_busy: + return + call_duration = time() - self.STATUS.get("RX_START", time()) + logger.info( + "(%s) *END RECORDING* idle timeout %.0fs (no VTERM) stream %s", + self._system, _RECORD_IDLE_S, int_id(self._record_stream), + ) + self._commit_recording(proto, call_duration) + + def _packet_is_vhead(self, data: bytes) -> bool: + if len(data) < 16: + return False + bits = data[15] + return ((bits & 0x30) >> 4) == HBPF_DATA_SYNC and (bits & 0xF) == HBPF_SLT_VHEAD + + def _packet_is_vterm(self, data: bytes) -> bool: + if len(data) < 16: + return False + bits = data[15] + return ((bits & 0x30) >> 4) == HBPF_DATA_SYNC and (bits & 0xF) == HBPF_SLT_VTERM + + def _make_vterm_packet(self, template: bytes, slot: int) -> bytes: + """Build a voice terminator from the last recorded frame (idle/max end has no VTERM).""" + if len(template) < 53: + return template + pkt = bytearray(template) + ts_bit = 0x80 if slot == 2 else 0x00 + pkt[15] = ts_bit | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM + pkt[4] = (int(pkt[4]) + 1) & 0xFF + return bytes(pkt) + + def _ensure_vterm(self, recorded: list[bytes], slot: int) -> list[bytes]: + if not recorded or self._packet_is_vterm(recorded[-1]): + return recorded + logger.debug( + "(%s) Appending synthetic VTERM for playback (%d packets, no terminator from source)", + self._system, len(recorded), + ) + return recorded + [self._make_vterm_packet(recorded[-1], slot)] + + def _commit_recording(self, proto: Any, call_duration: float) -> None: + ctx = self._record_ctx + slot = ctx.get("slot", 1) + recorded = self._ensure_vterm(list(self.CALL_DATA), slot) + self._reset_recording_state(slot) + if not recorded: + return + self._playback_busy = True + self._delay_call = reactor.callLater( + _PLAYBACK_DELAY_S, + self._start_playback, + proto, + recorded, + ctx.get("rf_src", b"\x00\x00\x00"), + ctx.get("peer_id", b"\x00\x00\x00\x00"), + ctx.get("dst_id", b"\x00\x00\x00"), + slot, + call_duration, + ) + + def _begin_recording( + self, + slot: int, + stream_id: bytes, + pkt_time: float, + rf_src: bytes, + peer_id: bytes, + dst_id: bytes, + proto: Any, + ) -> None: + self._recording_active = True + self._record_stream = stream_id + self._rx_stream_by_slot[slot] = stream_id + self._seen_record_streams = {stream_id} + self._ignored_streams_logged.clear() + self._last_record_time = pkt_time + self._record_ctx = { + "slot": slot, + "rf_src": rf_src, + "peer_id": peer_id, + "dst_id": dst_id, + } + self.STATUS["RX_START"] = pkt_time + logger.info( + "(%s) *START RECORDING* STREAM ID: %s SUB: %s REPEATER: %s TGID %s, TS %s", + self._system, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, + ) + self._schedule_record_max(proto) + self._schedule_record_idle(proto) + + def _recording_exceeded_max(self, pkt_time: float) -> bool: + rx_start = self.STATUS.get("RX_START", 0.0) + return bool(rx_start) and (pkt_time - rx_start) >= _SOURCE_MAX_S + + def _schedule_record_max(self, proto: Any) -> None: + self._cancel_record_max() + if not self._recording_active: + return + self._max_call = reactor.callLater(_SOURCE_MAX_S, self._on_record_max_duration, proto) + + def _cancel_record_max(self) -> None: + if self._max_call is not None and self._max_call.active(): + self._max_call.cancel() + self._max_call = None + + def _on_record_max_duration(self, proto: Any) -> None: + self._max_call = None + if not self._recording_active or not self.CALL_DATA or self._playback_busy: + return + call_duration = time() - self.STATUS.get("RX_START", time()) + logger.info( + "(%s) *END RECORDING* source timeout %.0fs (HBP max, no VTERM) stream %s", + self._system, _SOURCE_MAX_S, int_id(self._record_stream), + ) + self._commit_recording(proto, call_duration) + + def _reset_recording_state(self, slot: int) -> None: + self._cancel_record_idle() + self._cancel_record_max() + self.CALL_DATA = [] + self._recording_active = False + self._record_stream = b"" + self._seen_record_streams.clear() + self._ignored_streams_logged.clear() + self._last_record_time = 0.0 + self._record_ctx = {} + self._rx_stream_by_slot.pop(slot, None) + + def _start_playback( + self, + proto: Any, + recorded: list[bytes], + rf_src: bytes, + peer_id: bytes, + dst_id: bytes, + slot: int, + call_duration: float, + ) -> None: + self._delay_call = None + if not proto or not recorded: + self._finish_playback() + return + + self._playback_stream_id = bytes_4(randint(0x00, 0xFFFFFFFF)) + logger.info( + "(%s) *START PLAYBACK* STREAM ID: %s SUB: %s REPEATER: %s TGID %s, TS %s, Duration: %s", + self._system, + int_id(self._playback_stream_id), + int_id(rf_src), + int_id(peer_id), + int_id(dst_id), + slot, + call_duration, + ) + self._playback_packets = self._prepare_playback_packets(recorded) + self._playback_index = 0 + self._send_next_packet(proto) + + def _prepare_playback_packets(self, recorded: list[bytes]) -> list[bytes]: + """Rewrite stream ID, drop mid-call VHEADs (re-key), monotonic seq for bridge ingress.""" + out: list[bytes] = [] + seq = 1 + for i, pkt in enumerate(recorded): + if len(pkt) < 20: + continue + if i > 0 and self._packet_is_vhead(pkt): + logger.debug( + "(%s) Skipping mid-call VHEAD in playback (hotspot re-key)", + self._system, ) + continue + new_pkt = bytearray(pkt[:16] + self._playback_stream_id + pkt[20:]) + new_pkt[4] = seq & 0xFF + seq = (seq % 255) + 1 + out.append(bytes(new_pkt)) + return out + + def _send_next_packet(self, proto: Any) -> None: + self._packet_call = None + if self._playback_index >= len(self._playback_packets): + logger.info( + "(%s) *END PLAYBACK* STREAM ID: %s", + self._system, + int_id(self._playback_stream_id), + ) + self._finish_playback() + return + if proto: + proto.send_system(self._playback_packets[self._playback_index]) + self._playback_index += 1 + self._packet_call = reactor.callLater(_PACKET_INTERVAL_S, self._send_next_packet, proto) - if proto: - for i in self.CALL_DATA: - i = i[:16] + _new_stream_id + i[20:] - proto.send_system(i) - sleep(0.06) - self.CALL_DATA = [] - logger.info("(%s) *END PLAYBACK* STREAM ID: %s", self._system, int_id(_new_stream_id)) - - else: - if self.CALL_DATA: - self.CALL_DATA.append(data) + def _finish_playback(self) -> None: + if self._delay_call is not None and self._delay_call.active(): + self._delay_call.cancel() + if self._packet_call is not None and self._packet_call.active(): + self._packet_call.cancel() + self._cancel_record_idle() + self._cancel_record_max() + self._playback_busy = False + self._playback_packets = [] + self._playback_index = 0 + self._playback_stream_id = b"" + self._delay_call = None + self._packet_call = None