diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 03deb3a..6972147 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -45,6 +45,30 @@ from .ports import BridgeRouter logger = logging.getLogger(__name__) +# While loop-control loser, re-send BCSQ periodically so peers stop forwarding if first UDP was lost (legacy sends once). +_BCSQ_LOSER_RESEND_SEC = 2.0 + + +def _obp_target_bcsq_quenches_stream( + systems_cfg: dict[str, Any], target_name: str, dst_id_b: bytes, stream_id: bytes +) -> bool: + """True if target OBP config has _bcsq[tgid]==stream_id (bytes key or same int TG).""" + m = systems_cfg.get(target_name, {}).get("_bcsq") + if not isinstance(m, dict) or not m: + return False + tid = dst_id_b[:3] if isinstance(dst_id_b, bytes) and len(dst_id_b) >= 3 else bytes_3(int_id(dst_id_b)) + if m.get(tid) == stream_id: + return True + for k, v in m.items(): + if v != stream_id: + continue + try: + if isinstance(k, bytes) and len(k) >= 3 and int_id(k) == int_id(tid): + return True + except Exception: + continue + return False + def _log_trace(msg: str, *args: Any) -> None: """Per-packet forwarding diagnostics (BCSQ/BCKA/ACL). Below DEBUG — enable TRACE in LOGGER config to see.""" @@ -363,6 +387,49 @@ class BridgeUseCases: pass tstatus.pop(stream_id, None) + def on_obp_bcsq_received(self, system_name: str, tgid: bytes, stream_id: bytes) -> None: + """After valid BCSQ on this OBP leg: emit END,TX and clear forward STATUS if present (monitor parity; no VTERM to peer).""" + if not bool(self._config.get("REPORTS", {}).get("REPORT", True)): + return + report = self._report_factory + if not report or not hasattr(report, "send_bridge_event"): + return + protocols = self._get_protocols() if self._get_protocols else {} + tgt_proto = protocols.get(system_name) + if not tgt_proto: + return + if self._config.get("SYSTEMS", {}).get(system_name, {}).get("MODE") != "OPENBRIDGE": + return + tstatus = getattr(tgt_proto, "STATUS", None) + if not tstatus or stream_id not in tstatus: + return + tst = tstatus[stream_id] + if not isinstance(tst, dict) or "H_LC" not in tst: + return + if tst.get("TGID", b"\x00\x00\x00") != tgid: + return + now = time.time() + rfs = tst.get("RFS", b"\x00\x00\x00") + peer = tst.get("RX_PEER", b"\x00\x00\x00\x00") + tgid_b = tst.get("TGID", b"\x00\x00\x00") + start = tst.get("START", now) + duration = max(0.0, now - start) + try: + report.send_bridge_event( + "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( + system_name, + int_id(stream_id), + int_id(peer), + int_id(rfs), + 1, + int_id(tgid_b), + duration, + ) + ) + except Exception: + pass + tstatus.pop(stream_id, None) + def stream_trimmer_loop(self) -> None: """Trim old stream state (legacy stream_trimmer_loop, 5s). RX/TX timeout per system/slot; OBP streams (legacy bridge.py 181-240).""" logger.debug("(ROUTER) Trimming inactive stream IDs from system lists") @@ -1364,11 +1431,29 @@ class BridgeUseCases: call_duration, ) st["LOOPLOG"] = True + if _do_report and self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + try: + self._report_factory.send_bridge_event( + "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( + system_name, + int_id(stream_id), + int_id(peer_id), + int_id(rf_src), + slot, + int_id(dst_id), + max(0.0, pkt_time - st.get("START", pkt_time)), + ) + ) + except Exception: + pass st["LAST"] = pkt_time - if systems_cfg.get(system_name, {}).get("ENHANCED_OBP") and "_bcsq" not in st: - if self._send_bcsq: + if systems_cfg.get(system_name, {}).get("ENHANCED_OBP") and self._send_bcsq: + now_sq = time.time() + last_sq = float(st.get("_bcsq_last", 0.0)) + if "_bcsq" not in st or (now_sq - last_sq >= _BCSQ_LOSER_RESEND_SEC): self._send_bcsq(system_name, dst_id, stream_id) - st["_bcsq"] = True + st["_bcsq_last"] = now_sq + st["_bcsq"] = True return False if st["packets"] > 18 and (st["packets"] / st["START"] > 25): @@ -1688,12 +1773,7 @@ class BridgeUseCases: if isinstance(target_tgid, int): target_tgid = bytes_3(target_tgid) # If target has quenched us, don't send (~1856-1859). - _bcsq_map = _target_system.get("_bcsq") - if ( - isinstance(_bcsq_map, dict) - and dst_id_b in _bcsq_map - and _bcsq_map[dst_id_b] == stream_id - ): + if _obp_target_bcsq_quenches_stream(systems_cfg, entry["SYSTEM"], dst_id_b, stream_id): _log_trace( "(%s) OBP skip (BCSQ): target=%s TGID=%s stream=%s", system_name, diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index eac3821..764a040 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -143,6 +143,7 @@ class HBPProtocol(DatagramProtocol): on_in_band_signalling: Callable[[str, int, bytes, float], None] | None = None, on_options_received: Callable[[str], None] | None = None, on_deactivate_dynamic_bridges: Callable[[str], None] | None = None, + on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None, ) -> None: self._CONFIG = config self._system = system_name @@ -155,12 +156,14 @@ class HBPProtocol(DatagramProtocol): self._on_in_band_signalling = on_in_band_signalling self._on_options_received = on_options_received self._on_deactivate_dynamic_bridges = on_deactivate_dynamic_bridges + self._on_obp_bcsq_received = on_obp_bcsq_received self._config = config.get("SYSTEMS", {}).get(system_name, {}) if self._config.get("MODE") == "OPENBRIDGE": self._laststrid = deque([], 20) self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} # Legacy bridge_master: OBP stream state by stream_id for loop control (1ST) and duplicate handling self._obp_streams = {} + self._bcsq_log_once: set[tuple[bytes, bytes]] = set() else: self._laststrid = {1: b"", 2: b""} self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} @@ -932,13 +935,47 @@ class HBPProtocol(DatagramProtocol): else: logger.debug("(%s) *BridgeControl* not sending BCVE, TARGET not currently known", self._system) + def _obp_sync_target_sock_from_peer(self, _sockaddr: tuple[str, int]) -> None: + """If RELAX_CHECKS accepted traffic from a different IP:port than TARGET_SOCK, sync (same idea as BCKA). + Ensures BCSQ and outbound DMR go to the peer address we actually receive from.""" + if self._config.get("MODE") != "OPENBRIDGE" or not self._config.get("RELAX_CHECKS"): + return + if not _sockaddr or not _sockaddr[0]: + return + cur = self._config.get("TARGET_SOCK") + if cur == _sockaddr: + return + h, p = _sockaddr[0], int(_sockaddr[1]) + logger.info( + "(%s) *BridgeControl* OBP peer address sync to %s:%s (RELAX_CHECKS; was %s:%s)", + self._system, + h, + p, + (cur[0] if cur and cur[0] else "?"), + (cur[1] if cur and len(cur) > 1 else "?"), + ) + self._config["TARGET_IP"] = h + self._config["TARGET_PORT"] = p + self._config["TARGET_SOCK"] = (h, p) + def _obp_send_bcsq(self, _tgid: bytes, _stream_id: bytes) -> None: """Legacy send_bcsq: BCSQ + tgid + stream_id + HMAC-SHA1. Uses TARGET_SOCK (IP only).""" _addr = self._config.get("TARGET_SOCK") + if not _addr or not _addr[0]: + tip = self._config.get("TARGET_IP") + tport = int(self._config.get("TARGET_PORT", 62044)) + if tip: + _addr = (tip, tport) + self._config["TARGET_SOCK"] = _addr if _addr and _addr[0]: _packet = BCSQ + _tgid + _stream_id _packet = _packet + hmac_new(self._config["PASSPHRASE"], _packet, sha1).digest() self.transport.write(_packet, _addr) + else: + logger.warning( + "(%s) *BridgeControl* BCSQ not sent: no TARGET_SOCK/TARGET_IP — peer cannot be quenched", + self._system, + ) def proxy_bad_peer(self) -> None: """Legacy bridge_master routerOBP rate-drop hook; HBSYSTEM has peer handling — OBP noop.""" @@ -957,6 +994,7 @@ class HBPProtocol(DatagramProtocol): _hash = _packet[53:73] _ckhs = hmac_new(self._config["PASSPHRASE"], _data, sha1).digest() if compare_digest(_hash, _ckhs) and (_sockaddr == self._config.get("TARGET_SOCK") or self._config.get("RELAX_CHECKS")): + self._obp_sync_target_sock_from_peer(_sockaddr) _peer_id = _data[11:15] if self._config.get("NETWORK_ID") != _peer_id: if _stream_id not in self._laststrid: @@ -1092,6 +1130,7 @@ class HBPProtocol(DatagramProtocol): if not (compare_digest(_hash, _ckhs) and (_sockaddr == self._config.get("TARGET_SOCK") or self._config.get("RELAX_CHECKS"))): logger.warning("(%s) OpenBridge DMRE BLAKE2b failed, packet discarded - SRC: %s", self._system, _sockaddr) return + self._obp_sync_target_sock_from_peer(_sockaddr) _peer_id = _data[11:15] if self._config.get("NETWORK_ID") != _peer_id: if _stream_id not in self._laststrid: @@ -1267,6 +1306,23 @@ class HBPProtocol(DatagramProtocol): if "_bcsq" not in self._config: self._config["_bcsq"] = {} self._config["_bcsq"][_tgid_bcsq] = _stream_bcsq + if self._config.get("MODE") == "OPENBRIDGE": + _key = (_stream_bcsq, _tgid_bcsq) + _once = getattr(self, "_bcsq_log_once", None) + if isinstance(_once, set) and _key not in _once: + _once.add(_key) + logger.info( + "(%s) *BridgeControl* BCSQ accepted: stream_id=%s TGID=%s (peer quenched; forwarding on this OBP stops for this stream/TG)", + self._system, + int_id(_stream_bcsq), + int_id(_tgid_bcsq), + ) + cb = self._on_obp_bcsq_received + if cb is not None: + try: + cb(self._system, _tgid_bcsq, _stream_bcsq) + except Exception: + logger.exception("(%s) on_obp_bcsq_received failed", self._system) else: logger.warning( "(%s) *BridgeControl* BCSQ invalid Source Quench, packet discarded - SRC: %s", @@ -1316,6 +1372,7 @@ def HBPProtocolFactory( on_in_band_signalling: Callable[[str, int, bytes, float], None] | None = None, on_options_received: Callable[[str], None] | None = None, on_deactivate_dynamic_bridges: Callable[[str], None] | None = None, + on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None, ) -> HBPProtocol: """Create HBP protocol instance (legacy: one HBSYSTEM per system).""" return HBPProtocol( @@ -1330,4 +1387,5 @@ def HBPProtocolFactory( on_in_band_signalling=on_in_band_signalling, on_options_received=on_options_received, on_deactivate_dynamic_bridges=on_deactivate_dynamic_bridges, + on_obp_bcsq_received=on_obp_bcsq_received, ) diff --git a/src/adn_server/main.py b/src/adn_server/main.py index 170b0c4..2f005fb 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -455,6 +455,7 @@ def main() -> None: on_in_band_signalling=bridge_use_cases.apply_in_band_signalling, on_options_received=bridge_use_cases.options_config_for_system, on_deactivate_dynamic_bridges=bridge_use_cases.deactivate_all_dynamic_bridges, + on_obp_bcsq_received=bridge_use_cases.on_obp_bcsq_received, ) protocols[system_name] = protocol reactor.listenUDP(udp_port, protocol, interface=ip or "0.0.0.0")