diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 80a4905..f1db340 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -466,31 +466,41 @@ class BridgeUseCases: for system_name, protocol in protocols.items(): if not getattr(protocol, "STATUS", None): continue - # OBP: legacy bridge.py 181-202 — send GROUP VOICE,END,RX on timeout so report/monitor clears TG + # OBP: legacy bridge_master stream_trimmer_loop ~631-682 — two-stage lifecycle: + # Stage 1 (5s idle, no _to, no _fin): set _to=True, emit END,RX, continue. + # Stage 2 (180s idle): remove stream entry. if systems_cfg.get(system_name, {}).get("MODE") == "OPENBRIDGE": obp_streams = getattr(protocol, "_obp_streams", None) - if obp_streams and report and hasattr(report, "send_bridge_event"): + if obp_streams: to_remove = [] for stream_id, st in list(obp_streams.items()): last = st.get("LAST", 0) + # Stage 2: finished streams older than 180s → remove if st.get("_fin") and last < now - 180: to_remove.append(stream_id) continue - if last < now - 5: + # Stage 2: timed-out streams older than 180s → remove + if st.get("_to") and last < now - 180: + to_remove.append(stream_id) + continue + # Stage 1: 5s idle, not yet timed out → mark _to, emit END + if "_to" not in st and "_fin" not in st and last < now - 5: try: rfs = st.get("RFS", b"\x00\x00\x00") peer = st.get("RX_PEER", b"\x00\x00\x00\x00") tgid = st.get("TGID", b"\x00\x00\x00") start = st.get("START", now) duration = max(0.0, last - start) - report.send_bridge_event( - "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( - system_name, int_id(stream_id), int_id(peer), int_id(rfs), 1, int_id(tgid), duration + if report and hasattr(report, "send_bridge_event"): + report.send_bridge_event( + "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( + system_name, int_id(stream_id), int_id(peer), int_id(rfs), 1, int_id(tgid), duration + ) ) - ) except Exception: - pass - to_remove.append(stream_id) + st["LAST"] = now + st["_to"] = True + continue for stream_id in to_remove: st_rm = obp_streams.get(stream_id) or {} _syscfg = systems_cfg.get(system_name, {}) @@ -554,7 +564,7 @@ class BridgeUseCases: sys_cfg = systems_cfg.get(system_name, {}) if sys_cfg.get("_reset"): logger.info("(BRIDGERESET) Bridge reset for %s - no peers", system_name) - self._remove_bridge_system(system_name, bridges) + self.remove_bridge_system(system_name) try: del sys_cfg["_opt_key"] except KeyError: @@ -1415,6 +1425,39 @@ class BridgeUseCases: if obp is not None: obp[stream_id] = st + def _is_stream_known(self, system_name: str, stream_id: bytes, slot: int = 0) -> bool: + """Return True if *stream_id* is already tracked. + + Legacy routerHBP (bridge_master.py ~3054): ``_stream_id != STATUS[_slot]['RX_STREAM_ID']`` + Legacy routerOBP (bridge_master.py ~2193): ``_stream_id not in self.STATUS`` + + HBP protocols pre-seed ``STATUS[stream_id]`` before dmrd_received, so checking + ``stream_id in STATUS`` would always be True. Instead we compare against the + per-slot ``RX_STREAM_ID`` for HBP (MASTER/PEER) and ``stream_id in STATUS`` for OBP. + """ + systems_cfg = self._config.get("SYSTEMS", {}) + sys_mode = systems_cfg.get(system_name, {}).get("MODE", "") + if sys_mode == "OPENBRIDGE": + protocols = self._get_protocols() if self._get_protocols else {} + proto = protocols.get(system_name) + if proto is None: + return False + status = getattr(proto, "STATUS", None) + if status is None: + return False + return stream_id in status + protocols = self._get_protocols() if self._get_protocols else {} + proto = protocols.get(system_name) + if proto is None: + return False + status = getattr(proto, "STATUS", None) + if status is None or not isinstance(status, dict): + return False + slot_st = status.get(slot) + if not isinstance(slot_st, dict): + return False + return slot_st.get("RX_STREAM_ID") == stream_id + def _obp_group_voice_router_obp( self, system_name: str, @@ -1721,21 +1764,38 @@ class BridgeUseCases: obp_ber: bytes = b"\x00", obp_rssi: bytes = b"\x00", obp_source_rptr: bytes = b"\x00\x00\x00\x00", - ) -> None: + ) -> bool: """Called by UDP when DMRD is received. Forward to other systems in same bridge (to_target). + Returns True if the packet was accepted (passed ingress controls), False/None if dropped. Legacy `hblink.dmrd_received` passes `_hash,_hops,_source_server,_ber,_rssi,_source_rptr` after parsing OPENBRIDGE DMRD v1 / DMRE (`hblink.py` ~309–416, ~592–596). When `obp_use_parsed` is True, the OBP path uses those values (1:1 with `bridge.py` `routerOBP.dmrd_received` → `send_system`). """ if not self._send_to_system: return + # Legacy routerHBP.dmrd_received ~3015-3020: log once on first packet after _reset. + # Only applies to HBP sources (routerOBP has no _reset gate). + sys_cfg = self._config.get("SYSTEMS", {}).get(system_name, {}) + if sys_cfg.get("MODE") != "OPENBRIDGE" and sys_cfg.get("_reset") and not sys_cfg.get("_resetlog"): + logger.info("(%s) disallow transmission until reset cycle is complete", system_name) + sys_cfg["_resetlog"] = True + return # Legacy bridge_master 3080–3085: private call to ID 4000 only disconnects dynamics; do not route as PC. if call_type == "unit" and int_id(dst_id) == 4000: return if call_type == "unit": - self._pvt_call_received(system_name, peer_id, rf_src, dst_id, seq, slot, frame_type, dtype_vseq, stream_id, data) - return + if dtype_vseq in (6, 7, 8) or (dtype_vseq == 3 and not self._is_stream_known(system_name, stream_id, slot)): + self._unit_data_received( + system_name, peer_id, rf_src, dst_id, seq, slot, + frame_type, dtype_vseq, stream_id, data, + obp_use_parsed=obp_use_parsed, obp_hops=obp_hops, + obp_source_server=obp_source_server, + obp_ber=obp_ber, obp_rssi=obp_rssi, obp_source_rptr=obp_source_rptr, + ) + elif len(str(int_id(dst_id))) == 7: + self._pvt_call_received(system_name, peer_id, rf_src, dst_id, seq, slot, frame_type, dtype_vseq, stream_id, data) + return True bridge_key = str(int_id(dst_id)) bridges = self._router.get_bridges() dst_int = int_id(dst_id) @@ -1814,8 +1874,131 @@ class BridgeUseCases: "(ROUTER) No matching source rule for TG %s from %s slot %s (ACTIVE), not forwarding", bridge_key, system_name, bridge_match_slot, ) - return + return True pkt_time = time.time() + + # HBP group ingress controls (legacy routerHBP.dmrd_received ~3270-3399) + if not source_is_obp and call_type in ("group", "vcsbk"): + protocols = self._get_protocols() if self._get_protocols else {} + src_proto = protocols.get(system_name) + if src_proto: + _slot_st = getattr(src_proto, "STATUS", {}).get(slot, {}) + # New stream detection (legacy ~3270-3316) + _is_new_stream = stream_id != _slot_st.get("RX_STREAM_ID") + if _is_new_stream: + _slot_st["packets"] = 0 + _slot_st["loss"] = 0 + _slot_st["crcs"] = set() + _slot_st["LOOPLOG"] = False + _slot_st.pop("_bcsq", None) + _slot_st["lastSeq"] = False + _slot_st["lastData"] = False + # Collision check (legacy ~3276-3278) + if ( + _slot_st.get("RX_TYPE") != HBPF_SLT_VTERM + and pkt_time < (_slot_st.get("RX_TIME", 0) + STREAM_TO) + and rf_src != _slot_st.get("RX_RFS", b"\x00") + ): + logger.warning( + "(%s) Packet received with STREAM ID: %s SUB: %s PEER: %s TGID %s, SLOT %s collided with existing call", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, + ) + return + # Increment packet counter (legacy ~3317) + _slot_st["packets"] = _slot_st.get("packets", 0) + 1 + _pkts = _slot_st["packets"] + _rx_start = _slot_st.get("RX_START", pkt_time) + # Rate limit (legacy ~3331-3336): >18 packets and rate >25/s + if _pkts > 18 and _rx_start < pkt_time: + _rate = _pkts / (pkt_time - _rx_start) + if _rate > 25: + logger.warning( + "(%s) *PacketControl* RATE DROP! Stream ID: %s TGID: %s", + system_name, int_id(stream_id), int_id(dst_id), + ) + _slot_st["LAST"] = pkt_time + return + # 180s timeout (legacy ~3338-3344) + if _rx_start + 180 < pkt_time: + if not _slot_st.get("LOOPLOG"): + logger.info( + "(%s) HBP *SOURCE TIMEOUT* STREAM ID: %s, TG: %s, TS: %s, IGNORE THIS SOURCE", + system_name, int_id(stream_id), int_id(dst_id), slot, + ) + _slot_st["LOOPLOG"] = True + _slot_st["LAST"] = pkt_time + return + # HBP + OBP loop control (legacy ~3346-3368) + for other_name, proto in protocols.items(): + if other_name == system_name: + continue + omode = systems_cfg.get(other_name, {}).get("MODE") + ostatus = getattr(proto, "STATUS", None) + if not ostatus: + continue + if omode != "OPENBRIDGE": + for _sysslot in ostatus: + ss = ostatus.get(_sysslot) + if isinstance(ss, dict) and stream_id == ss.get("RX_STREAM_ID"): + if not _slot_st.get("LOOPLOG"): + logger.debug( + "(%s) HBP *LoopControl* FIRST HBP: %s, STREAM ID: %s, TG: %s, TS: %s, IGNORE THIS SOURCE", + system_name, other_name, int_id(stream_id), int_id(dst_id), _sysslot, + ) + _slot_st["LOOPLOG"] = True + _slot_st["LAST"] = pkt_time + return + else: + if ( + stream_id in ostatus + and "1ST" in ostatus[stream_id] + and ostatus[stream_id].get("TGID") == dst_id + ): + if not _slot_st.get("LOOPLOG"): + logger.debug( + "(%s) HBP *LoopControl* FIRST OBP %s, STREAM ID: %s, TG %s, IGNORE THIS SOURCE", + system_name, other_name, int_id(stream_id), int_id(dst_id), + ) + _slot_st["LOOPLOG"] = True + _slot_st["LAST"] = pkt_time + if ( + systems_cfg.get(system_name, {}).get("ENHANCED_OBP") + and "_bcsq" not in _slot_st + ): + if src_proto and hasattr(src_proto, "_obp_send_bcsq"): + src_proto._obp_send_bcsq(dst_id, stream_id) + _slot_st["_bcsq"] = True + return + # Duplicate handling (legacy ~3370-3399) + if _slot_st.get("lastData") and _slot_st["lastData"] == data and seq > 1: + _slot_st["loss"] = _slot_st.get("loss", 0) + 1 + logger.debug( + "(%s) *PacketControl* last packet is a complete duplicate, discarding. Stream ID: %s TGID: %s", + system_name, int_id(stream_id), int_id(dst_id), + ) + return + if seq and seq == _slot_st.get("lastSeq"): + _slot_st["loss"] = _slot_st.get("loss", 0) + 1 + return + if seq and _slot_st.get("lastSeq") and seq != 1 and seq < _slot_st.get("lastSeq", 0): + _slot_st["loss"] = _slot_st.get("loss", 0) + 1 + return + _pkt_crc = None + if seq > 0: + _h = blake2b(digest_size=16) + _h.update(data) + _pkt_crc = _h.digest() + if "crcs" in _slot_st and _pkt_crc in _slot_st["crcs"]: + _slot_st["loss"] = _slot_st.get("loss", 0) + 1 + return + # Missed packets (legacy ~3392-3394): just increment loss, don't drop + if seq and _slot_st.get("lastSeq") and seq > (_slot_st.get("lastSeq", 0) + 1): + _slot_st["loss"] = _slot_st.get("loss", 0) + 1 + _slot_st["lastSeq"] = seq + _slot_st["lastData"] = data + if _pkt_crc and "crcs" in _slot_st: + _slot_st["crcs"].add(_pkt_crc) + # Legacy bridge.py: BRDG_EVENT (OBP group/vcsbk START/END handled in _obp_group_voice_router_obp / post-forward VTERM) if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): try: @@ -2059,14 +2242,30 @@ class BridgeUseCases: if _target_status is None or entry_ts not in _target_status: continue _ts_st = _target_status[entry_ts] - # Contention handling (exact port of legacy bridge.py 413-432 / 730-745) + if source_is_obp: + _src_stream_st = getattr(src_proto, "STATUS", {}).get(stream_id, {}) if src_proto else {} + else: + _src_stream_st = getattr(src_proto, "STATUS", {}).get(slot, {}) if src_proto else {} + # Contention handling (exact port of legacy bridge_master.py ~2056-2075) if (entry_tgid_b != _ts_st.get("RX_TGID", b"\x00\x00\x00")) and ((pkt_time - _ts_st.get("RX_TIME", 0)) < float(_target_system.get("GROUP_HANGTIME", 0))): + if not _src_stream_st.get("CONTENTION"): + _src_stream_st["CONTENTION"] = True + logger.info("(%s) Call not routed to TGID %s, target active or in group hangtime: HBSystem: %s, TS: %s, TGID: %s", system_name, int_id(entry_tgid_b), entry.get("SYSTEM"), entry_ts, int_id(_ts_st.get("RX_TGID", b""))) continue if (entry_tgid_b != _ts_st.get("TX_TGID", b"\x00\x00\x00")) and ((pkt_time - _ts_st.get("TX_TIME", 0)) < float(_target_system.get("GROUP_HANGTIME", 0))): + if not _src_stream_st.get("CONTENTION"): + _src_stream_st["CONTENTION"] = True + logger.info("(%s) Call not routed to TGID %s, target in group hangtime: HBSystem: %s, TS: %s, TGID: %s", system_name, int_id(entry_tgid_b), entry.get("SYSTEM"), entry_ts, int_id(_ts_st.get("TX_TGID", b""))) continue if (entry_tgid_b == _ts_st.get("RX_TGID", b"\x00\x00\x00")) and ((pkt_time - _ts_st.get("RX_TIME", 0)) < STREAM_TO): + if not _src_stream_st.get("CONTENTION"): + _src_stream_st["CONTENTION"] = True + logger.info("(%s) Call not routed to TGID %s, matching call already active on target: HBSystem: %s, TS: %s, TGID: %s", system_name, int_id(entry_tgid_b), entry.get("SYSTEM"), entry_ts, int_id(_ts_st.get("RX_TGID", b""))) continue if (entry_tgid_b == _ts_st.get("TX_TGID", b"\x00\x00\x00")) and (rf_src != _ts_st.get("TX_RFS", b"")) and ((pkt_time - _ts_st.get("TX_TIME", 0)) < STREAM_TO): + if not _src_stream_st.get("CONTENTION"): + _src_stream_st["CONTENTION"] = True + logger.info("(%s) Call not routed for subscriber %s, call route in progress on target: HBSystem: %s, TS: %s, TGID: %s, SUB: %s", system_name, int_id(rf_src), entry.get("SYSTEM"), entry_ts, int_id(_ts_st.get("TX_TGID", b"")), int_id(_ts_st.get("TX_RFS", b""))) continue # New stream detection — legacy OBP uses _target_status[TS]['TX_STREAM_ID'], legacy HBP uses self.STATUS[_slot]['RX_STREAM_ID'] if source_is_obp: @@ -2197,6 +2396,336 @@ class BridgeUseCases: "(ROUTER) Bridged TG %s from %s -> %s", bridge_key, system_name, ", ".join(forwarded), ) + return True + + # ── Unit DATA path (SMS, GPS, CSBK) — legacy routerOBP/routerHBP unit data branch ── + + def _send_data_to_obp( + self, + source_system: str, + target: str, + data: bytes, + dmrpkt: bytes, + pkt_time: float, + stream_id: bytes, + dst_id: bytes, + peer_id: bytes, + rf_src: bytes, + bits: int, + slot: int, + hops: bytes = b"", + ber: bytes = b"\x00", + rssi: bytes = b"\x00", + source_server: bytes = b"\x00\x00\x00\x00", + source_rptr: bytes = b"\x00\x00\x00\x00", + ) -> None: + """Legacy sendDataToOBP: forward a unit-data packet to an OPENBRIDGE target.""" + systems_cfg = self._config.get("SYSTEMS", {}) + _target_system = systems_cfg.get(target, {}) + if _target_system.get("ENHANCED_OBP") and "_bcka" in _target_system and _target_system["_bcka"] < pkt_time - 60: + return + protocols = self._get_protocols() if self._get_protocols else {} + target_proto = protocols.get(target) + if not target_proto: + return + _target_status = getattr(target_proto, "STATUS", {}) + if stream_id not in _target_status: + _target_status[stream_id] = { + "START": pkt_time, + "CONTENTION": False, + "RFS": rf_src, + "TGID": dst_id, + "RX_PEER": peer_id, + "packets": 0, + } + _target_status[stream_id]["LAST"] = pkt_time + _tmp_bits = bits ^ (1 << 7) if slot == 2 else bits + _tmp_data = b"".join([data[:15], _tmp_bits.to_bytes(1, "big"), data[16:20], dmrpkt]) + try: + self._send_to_system(target, _tmp_data, _hops=hops, _ber=ber, _rssi=rssi, _source_server=source_server, _source_rptr=source_rptr) + except Exception as exc: + logger.warning("(%s) send_data_to_obp %s failed: %s", source_system, target, exc) + return + logger.debug("(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s", source_system, target, int_id(dst_id)) + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + try: + self._report_factory.send_bridge_event( + "UNIT DATA,DATA,TX,{},{},{},{},{},{}".format( + target, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id), + ) + ) + except Exception: + pass + + def _send_data_to_hbp( + self, + source_system: str, + d_system: str, + d_slot: int, + dst_id: bytes, + tmp_bits: int, + data: bytes, + dmrpkt: bytes, + rf_src: bytes, + stream_id: bytes, + peer_id: bytes, + ) -> None: + """Legacy sendDataToHBP: forward a unit-data packet to an HBP (MASTER/PEER) target.""" + _tmp_data = b"".join([data[:15], tmp_bits.to_bytes(1, "big"), data[16:20], dmrpkt]) + try: + self._send_to_system(d_system, _tmp_data) + except Exception as exc: + logger.warning("(%s) send_data_to_hbp %s failed: %s", source_system, d_system, exc) + return + logger.debug("(%s) UNIT Data Bridged to HBP System: %s DST_ID: %s", source_system, d_system, int_id(dst_id)) + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + try: + self._report_factory.send_bridge_event( + "UNIT DATA,DATA,TX,{},{},{},{},{},{}".format( + d_system, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id), + ) + ) + except Exception: + pass + + def _unit_data_received( + self, + system_name: str, + peer_id: bytes, + rf_src: bytes, + dst_id: bytes, + seq: int, + slot: int, + frame_type: int, + dtype_vseq: int, + stream_id: bytes, + data: bytes, + *, + obp_use_parsed: bool = False, + obp_hops: bytes = b"", + obp_source_server: bytes | None = None, + obp_ber: bytes = b"\x00", + obp_rssi: bytes = b"\x00", + obp_source_rptr: bytes = b"\x00\x00\x00\x00", + ) -> None: + """Legacy routerOBP/routerHBP unit data branch: DATA-GATEWAY + OBP fan-out + SUB_MAP/hotspot.""" + pkt_time = time.time() + dmrpkt = data[20:53] if len(data) >= 53 else b"" + _bits = data[15] if len(data) > 15 else 0 + _int_dst_id = int_id(dst_id) + systems_cfg = self._config.get("SYSTEMS", {}) + source_is_obp = systems_cfg.get(system_name, {}).get("MODE") == "OPENBRIDGE" + global_cfg = self._config.get("GLOBAL", {}) + if source_is_obp and obp_source_server is not None: + _source_server = obp_source_server + else: + _sid = global_cfg.get("SERVER_ID", b"\x00\x00\x00\x00") + if isinstance(_sid, bytes) and len(_sid) >= 4: + _source_server = _sid + elif _sid is not None: + _source_server = bytes_4(int(_sid) & 0xFFFFFFFF) + else: + _source_server = b"\x00\x00\x00\x00" + _source_rptr = peer_id if not source_is_obp else obp_source_rptr + + if source_is_obp: + _hops = obp_hops + _ber = obp_ber + _rssi = obp_rssi + else: + _hops = b"" + _ber = data[53:54] if len(data) > 53 else b"\x00" + _rssi = data[54:55] if len(data) > 54 else b"\x00" + + # OBP unit-data loop control (legacy routerOBP ~2201-2255) + if source_is_obp: + protocols = self._get_protocols() if self._get_protocols else {} + src_proto = protocols.get(system_name) + if src_proto: + status = getattr(src_proto, "STATUS", None) + if status is not None: + if stream_id not in status: + status[stream_id] = { + "START": pkt_time, "CONTENTION": False, "RFS": rf_src, + "TGID": dst_id, "1ST": perf_counter(), "lastSeq": False, + "lastData": False, "RX_PEER": peer_id, "packets": 0, "crcs": set(), + } + status[stream_id]["LAST"] = pkt_time + status[stream_id]["packets"] = status[stream_id].get("packets", 0) + 1 + # HBP loop check: if any HBP already has this RX_STREAM_ID, drop + for other_name, proto in protocols.items(): + omode = systems_cfg.get(other_name, {}).get("MODE") + if other_name != system_name and omode != "OPENBRIDGE": + ostatus = getattr(proto, "STATUS", None) + if not ostatus: + continue + for _sysslot in ostatus: + slot_st = ostatus.get(_sysslot) + if isinstance(slot_st, dict) and stream_id == slot_st.get("RX_STREAM_ID"): + if not status[stream_id].get("LOOPLOG"): + logger.debug( + "(%s) OBP UNIT *LoopControl* FIRST HBP: %s, STREAM ID: %s, TG: %s, TS: %s, IGNORE THIS SOURCE", + system_name, other_name, int_id(stream_id), int_id(dst_id), _sysslot, + ) + status[stream_id]["LOOPLOG"] = True + return + # OBP earliest-wins: compare 1ST across OBP systems + hr_times: dict[str, float] = {} + for other_name, proto in protocols.items(): + omode = systems_cfg.get(other_name, {}).get("MODE") + if omode == "OPENBRIDGE": + obp_st = getattr(proto, "STATUS", None) + if obp_st and stream_id in obp_st and "1ST" in obp_st[stream_id]: + if obp_st[stream_id].get("TGID") == dst_id: + hr_times[other_name] = obp_st[stream_id]["1ST"] + fi = min(hr_times, key=hr_times.get, default=False) + if not fi: + logger.warning( + "(%s) OBP UNIT *LoopControl* fi is empty: STREAM ID: %s, TG: %s", + system_name, int_id(stream_id), int_id(dst_id), + ) + return + if system_name != fi: + if not status[stream_id].get("LOOPLOG"): + logger.debug( + "(%s) OBP UNIT *LoopControl* FIRST OBP %s, STREAM ID: %s, TG %s, IGNORE THIS SOURCE", + system_name, fi, int_id(stream_id), int_id(dst_id), + ) + status[stream_id]["LOOPLOG"] = True + return + + # Dtype-specific logging (legacy routerOBP ~2257-2280, routerHBP ~3052-3082) + if dtype_vseq == 3: + logger.info( + "(%s) *UNIT CSBK* STREAM ID: %s SUB: %s PEER: %s DST_ID %s TS %s", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, + ) + elif dtype_vseq == 6: + logger.info( + "(%s) *UNIT DATA HEADER* STREAM ID: %s SUB: %s PEER: %s DST_ID %s TS %s", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, + ) + elif dtype_vseq == 7: + logger.info( + "(%s) *UNIT VCSBK 1/2 DATA BLOCK* STREAM ID: %s SUB: %s PEER: %s TGID %s TS %s", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, + ) + elif dtype_vseq == 8: + logger.info( + "(%s) *UNIT VCSBK 3/4 DATA BLOCK* STREAM ID: %s SUB: %s PEER: %s TGID %s TS %s", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, + ) + else: + logger.info( + "(%s) *UNKNOWN DATA TYPE* STREAM ID: %s SUB: %s PEER: %s TGID %s TS %s", + system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, + ) + + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + _dtype_labels = {3: "UNIT CSBK", 6: "UNIT DATA HEADER", 7: "UNIT VCSBK 1/2 DATA BLOCK", 8: "UNIT VCSBK 3/4 DATA BLOCK"} + _label = _dtype_labels.get(dtype_vseq, "UNIT DATA") + try: + self._report_factory.send_bridge_event( + "{},DATA,RX,{},{},{},{},{},{}".format( + _label, system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, _int_dst_id, + ) + ) + except Exception: + pass + + # DATA-GATEWAY forwarding (legacy ~2281-2284 / ~3083-3087) + if global_cfg.get("DATA_GATEWAY"): + dg_cfg = systems_cfg.get("DATA-GATEWAY", {}) + if dg_cfg.get("MODE") == "OPENBRIDGE" and dg_cfg.get("ENABLED"): + logger.debug("(%s) DATA packet sent to DATA-GATEWAY", system_name) + self._send_data_to_obp( + system_name, "DATA-GATEWAY", data, dmrpkt, pkt_time, stream_id, + dst_id, peer_id, rf_src, _bits, slot, + hops=_hops, ber=_ber, rssi=_rssi, + source_server=_source_server, source_rptr=_source_rptr, + ) + + # Fan-out to all OBP systems with VER > 1 and dst_id >= 1000000 (legacy ~2286-2295 / ~3088-3097) + protocols = self._get_protocols() if self._get_protocols else {} + for sys_name, sys_cfg in systems_cfg.items(): + if sys_name == system_name: + continue + if sys_name == "DATA-GATEWAY": + continue + if sys_cfg.get("MODE") == "OPENBRIDGE" and sys_cfg.get("VER", 1) > 1 and _int_dst_id >= 1000000: + self._send_data_to_obp( + system_name, sys_name, data, dmrpkt, pkt_time, stream_id, + dst_id, peer_id, rf_src, _bits, slot, + hops=_hops, ber=_ber, rssi=_rssi, + source_server=_source_server, source_rptr=_source_rptr, + ) + + # SUB_MAP lookup (legacy ~2297-2312 / ~3099-3114) + sub_map = self._config.get("_SUB_MAP", {}) + if dst_id in sub_map: + _d_system, _d_slot, _d_time = sub_map[dst_id] + _d_proto = protocols.get(_d_system) + if _d_proto: + _dst_slot = getattr(_d_proto, "STATUS", {}).get(_d_slot, {}) + logger.info("(%s) SUB_MAP matched, System: %s Slot: %s, Time: %s", system_name, _d_system, _d_slot, _d_time) + _d_sys_cfg = systems_cfg.get(_d_system, {}) + if ( + _dst_slot.get("RX_TYPE") == HBPF_SLT_VTERM + and _dst_slot.get("TX_TYPE") == HBPF_SLT_VTERM + and (pkt_time - _dst_slot.get("TX_TIME", 0) > _d_sys_cfg.get("GROUP_HANGTIME", 5)) + ): + _tmp_bits = _bits ^ (1 << 7) if slot != _d_slot else _bits + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + else: + logger.debug("(%s) UNIT Data not bridged to HBP - target busy: %s DST_ID: %s", system_name, _d_system, _int_dst_id) + else: + # Hotspot 6/7-digit peer ID match (legacy ~3131-3168 / ~2314-2345) + for _d_system, _d_sys_cfg in systems_cfg.items(): + if _d_sys_cfg.get("MODE") != "MASTER": + continue + _d_proto = protocols.get(_d_system) + if not _d_proto: + continue + _peers = _d_sys_cfg.get("PEERS", {}) + _matched = False + for _to_peer in _peers: + _int_to_peer = int_id(_to_peer) + _dst_str = str(_int_dst_id) + _to_str = str(_int_to_peer) + if len(_dst_str) == 6: + if _to_str[:6] == _dst_str: + _d_slot = 2 + _dst_slot = getattr(_d_proto, "STATUS", {}).get(_d_slot, {}) + logger.info("(%s) User Peer Hotspot ID (6-digit) matched, System: %s Slot: %s", system_name, _d_system, _d_slot) + if ( + _dst_slot.get("RX_TYPE") == HBPF_SLT_VTERM + and _dst_slot.get("TX_TYPE") == HBPF_SLT_VTERM + and (pkt_time - _dst_slot.get("TX_TIME", 0) > _d_sys_cfg.get("GROUP_HANGTIME", 5)) + ): + _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + else: + logger.debug("(%s) UNIT Data not bridged to HBP on slot %s - target busy: %s DST_ID: %s", system_name, _d_slot, _d_system, _int_dst_id) + _matched = True + break + elif len(_dst_str) >= 7: + if _to_str[:7] == _dst_str[:7]: + _d_slot = 2 + _dst_slot = getattr(_d_proto, "STATUS", {}).get(_d_slot, {}) + logger.info("(%s) User Peer Hotspot ID (7-digit) matched, System: %s Slot: %s", system_name, _d_system, _d_slot) + if ( + _dst_slot.get("RX_TYPE") == HBPF_SLT_VTERM + and _dst_slot.get("TX_TYPE") == HBPF_SLT_VTERM + and (pkt_time - _dst_slot.get("TX_TIME", 0) > _d_sys_cfg.get("GROUP_HANGTIME", 5)) + ): + _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + else: + logger.debug("(%s) UNIT Data not bridged to HBP on slot %s - target busy: %s DST_ID: %s", system_name, _d_slot, _d_system, _int_dst_id) + _matched = True + break + if _matched: + break def _pvt_call_received( self, @@ -2255,6 +2784,8 @@ class BridgeUseCases: _target_status = getattr(target_proto, "STATUS", {}) _target_system = systems_cfg.get(_target, {}) if _target_system.get("MODE") == "OPENBRIDGE": + if _target_system.get("ENHANCED_OBP") and "_bcka" in _target_system and _target_system["_bcka"] < pkt_time - 60: + continue if stream_id not in _target_status: _target_status[stream_id] = { "START": pkt_time, @@ -2268,6 +2799,13 @@ class BridgeUseCases: "(%s) PRIVATE call bridged to OBP System: %s TS: %s, UNIT: %s", system_name, _target, slot if _target_system.get("BOTH_SLOTS") else 1, int_id(dst_id), ) + report = self._report_factory + if report and hasattr(report, "send_bridge_event"): + report.send_bridge_event( + "PRIVATE VOICE,START,TX,{},{},{},{},{},{}".format( + _target, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), + ).encode("utf-8", "ignore") + ) _target_status[stream_id]["LAST"] = pkt_time if _target_system.get("BOTH_SLOTS"): _tmp_bits = _bits @@ -2300,6 +2838,13 @@ class BridgeUseCases: ts_st["TX_RFS"] = rf_src ts_st["TX_PEER"] = peer_id logger.info("(%s) PRIVATE call bridged to HBP System: %s TS: %s, DST: %s", system_name, _target, slot, int_id(dst_id)) + report = self._report_factory + if report and hasattr(report, "send_bridge_event"): + report.send_bridge_event( + "PRIVATE VOICE,START,TX,{},{},{},{},{},{}".format( + _target, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), + ).encode("utf-8", "ignore") + ) ts_st["TX_TIME"] = pkt_time ts_st["TX_TYPE"] = dtype_vseq send_data = data diff --git a/src/adn_server/application/voice_use_cases.py b/src/adn_server/application/voice_use_cases.py index 435b478..fc8f047 100644 --- a/src/adn_server/application/voice_use_cases.py +++ b/src/adn_server/application/voice_use_cases.py @@ -32,7 +32,7 @@ import time from datetime import datetime from typing import Any, Callable -from ..domain import bytes_3, HBPF_SLT_VHEAD, HBPF_SLT_VTERM +from ..domain import bytes_3, bytes_4, HBPF_SLT_VHEAD, HBPF_SLT_VTERM from .ports import VoiceProvider logger = logging.getLogger(__name__) @@ -654,10 +654,8 @@ class VoiceUseCases: logger.info("(%s) Playing on-demand AMBE file: %s (ID: %s)", system, file_number, file_number) time.sleep(1) _say = [pairs] - server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") - if not isinstance(server_id, bytes): - server_id = bytes_3(int(server_id)) - speech = self.pkt_gen(bytes_3(5000), bytes_3(9), server_id, 1, _say) + speech = self.pkt_gen(bytes_3(5000), bytes_3(9), bytes_4(9), 1, _say) + time.sleep(1) _slot = protocol.STATUS.get(2) if not _slot: return @@ -705,10 +703,7 @@ class VoiceUseCases: else: _say.append(words.get("notlinked") or silence) _say.append(silence) - server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") - if not isinstance(server_id, bytes): - server_id = bytes_3(int(server_id)) - speech = self.pkt_gen(bytes_3(5000), bytes_3(9), server_id, 1, _say) + speech = self.pkt_gen(bytes_3(5000), bytes_3(9), bytes_4(9), 1, _say) time.sleep(1) _slot = protocol.STATUS.get(2) if not _slot: diff --git a/src/adn_server/infrastructure/config_normalizer.py b/src/adn_server/infrastructure/config_normalizer.py index 56c941f..e673eaa 100644 --- a/src/adn_server/infrastructure/config_normalizer.py +++ b/src/adn_server/infrastructure/config_normalizer.py @@ -162,3 +162,4 @@ def normalize_obp_config(config: dict) -> None: sys_cfg.setdefault("ENHANCED_OBP", True) if "TG1_ACL" not in sys_cfg and "TGID_ACL" in sys_cfg: sys_cfg["TG1_ACL"] = sys_cfg["TGID_ACL"] + sys_cfg.setdefault("TG2_ACL", "PERMIT:ALL") diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 764a040..7fdea5b 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -64,6 +64,7 @@ from ..hbp_constants import ( MSTN, MSTP, MSTC, + PRBL, PRIN, RPTACK, RPTA, @@ -189,17 +190,21 @@ class HBPProtocol(DatagramProtocol): self._config["_bcka"] = time.time() if self._config.get("ENHANCED_OBP"): self._bcka_loop = task.LoopingCall(self._obp_send_bcka) - self._bcka_loop.start(10) + _bcka_d = self._bcka_loop.start(10) + _bcka_d.addErrback(self._looping_err_handle) self._bcve_loop = task.LoopingCall(self._obp_send_bcve) - self._bcve_loop.start(60) + _bcve_d = self._bcve_loop.start(60) + _bcve_d.addErrback(self._looping_err_handle) elif self._config.get("MODE") == "MASTER": ping_time = self._CONFIG.get("GLOBAL", {}).get("PING_TIME", 10) self._maintenance_loop = task.LoopingCall(self._master_maintenance_loop) - self._maintenance_loop.start(ping_time) + _maint_d = self._maintenance_loop.start(ping_time) + _maint_d.addErrback(self._looping_err_handle) elif self._config.get("MODE") == "PEER": ping_time = self._CONFIG.get("GLOBAL", {}).get("PING_TIME", 10) self._maintenance_loop = task.LoopingCall(self._peer_maintenance_loop) - self._maintenance_loop.start(ping_time) + _maint_d = self._maintenance_loop.start(ping_time) + _maint_d.addErrback(self._looping_err_handle) # ── Exact port of hblink.py send_peers / send_peer / send_master / send_system ── @@ -297,6 +302,19 @@ class HBPProtocol(DatagramProtocol): _slot["TX_TIME"] = _pkt_time self.send_system(pkt) + def dereg(self) -> None: + """Graceful de-registration (legacy hblink.py dereg / master_dereg / peer_dereg).""" + mode = self._config.get("MODE") + if mode == "MASTER": + for _peer in self._peers: + self.send_peer(_peer, b"".join([MSTCL, _peer])) + logger.info("(%s) De-Registration sent to Peer: %s (%s)", self._system, self._peers[_peer].get("CALLSIGN", b""), self._peers[_peer].get("RADIO_ID", b"")) + elif mode == "PEER": + self.send_master(b"".join([RPTCL, self._config.get("RADIO_ID", b"\x00\x00\x00\x00")])) + logger.info("(%s) De-Registration sent to Master: %s:%s", self._system, self._config.get("MASTER_SOCKADDR", ("?", "?"))[0], self._config.get("MASTER_SOCKADDR", ("?", "?"))[1]) + else: + logger.info("(%s) is mode %s. No De-Registration required, continuing shutdown", self._system, mode) + def validate_id(self, _peer_id: bytes): """Legacy validate_id (hblink.py lines 862-883). Returns True or callsign string or False.""" if "ALLOW_UNREG_ID" not in self._config: @@ -316,6 +334,17 @@ class HBPProtocol(DatagramProtocol): return _peer_ids[_int_peer_id] return False + def proxy_IPBlackList(self, peer_id: bytes, sockaddr: tuple[str, int]) -> None: + """Legacy hblink.py proxy_IPBlackList: send PRBL to proxy to blacklist a peer's IP for 5 min.""" + _bltime = str(time.time() + 300) + _prpacket = b"".join([PRBL, peer_id, _bltime.encode("UTF-8")]) + self.transport.write(_prpacket, sockaddr) + + def proxy_bad_peer(self) -> None: + """Legacy hblink.py proxy_BadPeer: blacklist all current peer IPs (called on rate-limit violations).""" + for _pi in getattr(self, "_peers", {}): + self.proxy_IPBlackList(_pi, self._peers[_pi]["SOCKADDR"]) + def datagramReceived(self, data: bytes, addr: tuple[str, int]) -> None: if len(data) < 4: return @@ -399,17 +428,6 @@ class HBPProtocol(DatagramProtocol): return pkt_time = time.time() _int_dst_id = int_id(_dst_id) - # TG / private ID 4000: deactivate dynamic bridges (legacy bridge_master 3080–3177). - # Must run before ACL: private 4000 is rarely in TG allow lists. - if _int_dst_id == 4000 and self._on_deactivate_dynamic_bridges: - _kind = "Private call to ID" if _call_type == "unit" else "Group call to TG" - logger.info( - "(%s) %s 4000 received on TS %s — deactivating all dynamic bridges", - self._system, - _kind, - _slot, - ) - self._on_deactivate_dynamic_bridges(self._system) # ACL (legacy order and _laststrid) if self._router and _global.get("USE_ACL"): if not self._router.acl_check(_rf_src, _global.get("SUB_ACL", (True, []))): @@ -447,20 +465,21 @@ class HBPProtocol(DatagramProtocol): sub_map = self._CONFIG.get("_SUB_MAP") if sub_map is not None: sub_map[_rf_src] = (self._system, _slot, pkt_time) - # TG 9991-9999: information services (legacy peer_datagramReceived, also needed in MASTER) - if 9991 <= _int_dst_id <= 9999 and self._on_play_file_request: - if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: - reactor.callInThread(self._on_play_file_request, self._system, str(_int_dst_id)) if self._config.get("REPEAT", True): pkt = [_data[:11], b"", _data[15:]] for _peer in self._peers: if _peer != _peer_id: pkt[1] = _peer self.transport.write(b"".join(pkt), self._peers[_peer]["SOCKADDR"]) - _voice = self._CONFIG.get("VOICE", {}) - if self._on_handle_recording and _voice.get("RECORDING_ENABLED") and int_id(_dst_id) == _voice.get("RECORDING_TG") and _slot == _voice.get("RECORDING_TIMESLOT", 2): - dmrpkt = _data[20:] if len(_data) > 20 else _data - self._on_handle_recording(dmrpkt, _frame_type, _dtype_vseq, _stream_id, pkt_time, _rf_src, _int_dst_id, _slot) + # TG 4000: deactivate after REPEAT so peers see the packet (legacy order) + if _int_dst_id == 4000 and self._on_deactivate_dynamic_bridges: + _kind = "Private call to ID" if _call_type == "unit" else "Group call to TG" + logger.info( + "(%s) %s 4000 received on TS %s — deactivating all dynamic bridges", + self._system, _kind, _slot, + ) + self._on_deactivate_dynamic_bridges(self._system) + return if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: logger.info( "(%s) CALL RX peer %s src %s -> TG %s slot %s", @@ -498,11 +517,24 @@ class HBPProtocol(DatagramProtocol): self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src else: self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src + _accepted = False if self._dmrd_received: - self._dmrd_received( + _accepted = self._dmrd_received( self._system, _peer_id, _rf_src, _dst_id, _seq, _slot, _call_type, _frame_type, _dtype_vseq, _stream_id, _data, ) + _voice = self._CONFIG.get("VOICE", {}) + if _accepted and self._on_handle_recording and _voice.get("RECORDING_ENABLED") and int_id(_dst_id) == _voice.get("RECORDING_TG") and _slot == _voice.get("RECORDING_TIMESLOT", 2): + dmrpkt = _data[20:53] if len(_data) >= 53 else _data[20:] + self._on_handle_recording(dmrpkt, _frame_type, _dtype_vseq, _stream_id, pkt_time, _rf_src, _int_dst_id, _slot) + if ( + _call_type == "unit" + and _frame_type == HBPF_DATA_SYNC + and _dtype_vseq == HBPF_SLT_VTERM + and 9991 <= _int_dst_id <= 9999 + and self._on_play_file_request + ): + reactor.callInThread(self._on_play_file_request, str(_int_dst_id), self._system) if ( _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM @@ -564,6 +596,8 @@ class HBPProtocol(DatagramProtocol): logger.info("(%s) Sent Challenge Response to %s for login: %s", self._system, int_id(_peer_id), self._peers[_peer_id]["SALT"]) else: self.transport.write(b"".join([MSTNAK, _peer_id]), _sockaddr) + if self._config.get("PROXY_CONTROL"): + self.proxy_IPBlackList(_peer_id, _sockaddr) logger.warning("(%s) Invalid Login from %s Radio ID: %s Denied by Registation ACL or not registered ID", self._system, _sockaddr[0], int_id(_peer_id)) if self._CONFIG.get("SYSTEMS", {}).get(self._system): self._CONFIG["SYSTEMS"][self._system]["_reset"] = True @@ -648,6 +682,8 @@ class HBPProtocol(DatagramProtocol): _this_peer["PACKAGE_ID"] = _data[262:302] if ("ALLOW_UNREG_ID" in self._config and not self._config["ALLOW_UNREG_ID"]) and _this_peer["CALLSIGN"].decode("utf8", errors="replace").rstrip() != self.validate_id(_peer_id): del self._peers[_peer_id] + if self._config.get("PROXY_CONTROL"): + self.proxy_IPBlackList(_peer_id, _sockaddr) self.transport.write(b"".join([MSTNAK, _peer_id]), _sockaddr) self._CONFIG.setdefault("SYSTEMS", {}).setdefault(self._system, {})["_reset"] = True logger.info("(%s) Callsign does not match subscriber database: ID: %s, Sent Call: %s, DB call %s", self._system, int_id(_peer_id), _this_peer["CALLSIGN"].decode("utf8", errors="replace").rstrip(), self.validate_id(_peer_id)) @@ -757,15 +793,6 @@ class HBPProtocol(DatagramProtocol): return pkt_time = time.time() _int_dst_id = int_id(_dst_id) - if _int_dst_id == 4000 and self._on_deactivate_dynamic_bridges: - _kind = "Private call to ID" if _call_type == "unit" else "Group call to TG" - logger.info( - "(%s) %s 4000 received on TS %s — deactivating all dynamic bridges", - self._system, - _kind, - _slot, - ) - self._on_deactivate_dynamic_bridges(self._system) _global = self._CONFIG.get("GLOBAL", {}) if self._router and _global.get("USE_ACL"): if not self._router.acl_check(_rf_src, _global.get("SUB_ACL", (True, []))): @@ -803,27 +830,71 @@ class HBPProtocol(DatagramProtocol): sub_map = self._CONFIG.get("_SUB_MAP") if sub_map is not None: sub_map[_rf_src] = (self._system, _slot, pkt_time) - if self._on_play_file_request and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: - if 9991 <= _int_dst_id <= 9999: - reactor.callInThread(self._on_play_file_request, str(_int_dst_id), self._system) + # TG 4000: deactivate after ACL/SUB_MAP (legacy order — routerHBP.dmrd_received) + if _int_dst_id == 4000 and self._on_deactivate_dynamic_bridges: + _kind = "Private call to ID" if _call_type == "unit" else "Group call to TG" + logger.info( + "(%s) %s 4000 received on TS %s — deactivating all dynamic bridges", + self._system, _kind, _slot, + ) + self._on_deactivate_dynamic_bridges(self._system) + return if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: logger.info( "(%s) CALL RX (from master) src %s -> TG %s slot %s", self._system, int_id(_rf_src), int_id(_dst_id), _slot, ) + # Stream state + LC setup (same as MASTER, needed by dmrd_received for LC rewrite) + if _stream_id not in self.STATUS: + self.STATUS[_stream_id] = { + "START": pkt_time, "CONTENTION": False, "RFS": _rf_src, + "TGID": _dst_id, "LAST": pkt_time, "lastSeq": False, "lastData": False, + } + dmrpkt = _data[20:53] if len(_data) >= 53 else b"" + if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt) >= 33: + try: + decoded = decode.voice_head_term(dmrpkt) + self.STATUS[_stream_id]["LC"] = decoded["LC"] + except Exception: + self.STATUS[_stream_id]["LC"] = LC_OPT + _dst_id + _rf_src + else: + self.STATUS[_stream_id]["LC"] = LC_OPT + _dst_id + _rf_src + if _slot in self.STATUS: + if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): + self.STATUS[_slot]["RX_START"] = pkt_time + if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt) >= 33: + try: + decoded_slot = decode.voice_head_term(dmrpkt) + self.STATUS[_slot]["RX_LC"] = decoded_slot["LC"] + except Exception: + self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src + else: + self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src + _accepted = False if self._dmrd_received: - self._dmrd_received(self._system, _peer_id, _rf_src, _dst_id, _seq, _slot, _call_type, _frame_type, _dtype_vseq, _stream_id, _data) - # Legacy in-band signalling on voice terminator (bridge.py 807-866): run before updating STATUS + _accepted = self._dmrd_received(self._system, _peer_id, _rf_src, _dst_id, _seq, _slot, _call_type, _frame_type, _dtype_vseq, _stream_id, _data) + _voice = self._CONFIG.get("VOICE", {}) + if _accepted and self._on_handle_recording and _voice.get("RECORDING_ENABLED") and int_id(_dst_id) == _voice.get("RECORDING_TG") and _slot == _voice.get("RECORDING_TIMESLOT", 2): + dmrpkt = _data[20:53] if len(_data) >= 53 else _data[20:] + self._on_handle_recording(dmrpkt, _frame_type, _dtype_vseq, _stream_id, pkt_time, _rf_src, _int_dst_id, _slot) + if ( + _call_type == "unit" + and _frame_type == HBPF_DATA_SYNC + and _dtype_vseq == HBPF_SLT_VTERM + and 9991 <= _int_dst_id <= 9999 + and self._on_play_file_request + ): + reactor.callInThread(self._on_play_file_request, str(_int_dst_id), self._system) if ( _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM - and self.STATUS - and self.STATUS.get(_slot, {}).get("RX_TYPE") != HBPF_SLT_VTERM + and _slot in self.STATUS + and self.STATUS[_slot].get("RX_TYPE") != HBPF_SLT_VTERM and self._on_in_band_signalling ): self._on_in_band_signalling(self._system, _slot, _dst_id, pkt_time) - if self.STATUS: - if _stream_id != self.STATUS[_slot]["RX_STREAM_ID"]: + if _slot in self.STATUS: + if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): self.STATUS[_slot]["RX_START"] = pkt_time self.STATUS[_slot]["RX_PEER"] = _peer_id self.STATUS[_slot]["RX_SEQ"] = _seq @@ -916,6 +987,10 @@ class HBPProtocol(DatagramProtocol): else: logger.error("(%s) Received an invalid command in packet: %s", self._system, _data.hex() if hasattr(_data, "hex") else _data) + def _looping_err_handle(self, failure) -> None: + """Legacy loopingErrHandle: log-only (hblink.py ~202/~720).""" + logger.error("(%s) Unhandled error in timed loop.\n %s", self._system, failure) + def _obp_send_bcka(self) -> None: """Legacy send_bcka: BCKA + HMAC-SHA1 to TARGET. Uses TARGET_SOCK (IP only; hostnames resolved at startup or on first peer packet).""" _addr = self._config.get("TARGET_SOCK") @@ -977,9 +1052,6 @@ class HBPProtocol(DatagramProtocol): self._system, ) - def proxy_bad_peer(self) -> None: - """Legacy bridge_master routerOBP rate-drop hook; HBSYSTEM has peer handling — OBP noop.""" - def _obp_datagram_received(self, _packet: bytes, _sockaddr: tuple[str, int]) -> None: """Port of hblink.py OPENBRIDGE.datagramReceived: DMRD v1 (53+HMAC), BCKA, BCVE.""" if _packet[:3] == DMR and _packet[:4] == DMRD and len(_packet) >= 73: diff --git a/src/adn_server/main.py b/src/adn_server/main.py index 2f005fb..22e5eb1 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -134,8 +134,10 @@ def _make_echo_bridges(config: dict) -> dict: def _looping_errback(logger: logging.Logger, failure): - """Errback for LoopingCalls (legacy loopingErrHandle).""" - logger.error("(GLOBAL) Unhandled error in timed loop: %s", failure.getTraceback()) + """Errback for LoopingCalls (legacy loopingErrHandle). Stops reactor to avoid memory leaks.""" + logger.error("(GLOBAL) STOPPING REACTOR TO AVOID MEMORY LEAK: Unhandled error in timed loop.\n %s", failure) + from twisted.internet import reactor as _reactor + _reactor.stop() def main() -> None: @@ -174,6 +176,8 @@ def main() -> None: peer_ids, subscriber_ids, talkgroup_ids, local_subscriber_ids, server_ids, checksums = ( alias_loader.load_aliases(config) ) + subscriber_ids[900999] = "D-APRS" + subscriber_ids[4294967295] = "SC" config["_SUB_IDS"] = subscriber_ids config["_PEER_IDS"] = peer_ids config["_TG_IDS"] = talkgroup_ids @@ -323,9 +327,10 @@ def main() -> None: task.LoopingCall(lambda: threads.deferToThread(reporting_use_cases.ka_reporting_loop)).start(60).addErrback(_looping_errback, logger) # bridgeDebug (legacy 66s) — remove invalid bridges, fix >1 active dial per MASTER - task.LoopingCall( - lambda: (bridge_use_cases.bridge_debug_loop(), report_factory.set_bridges(bridge_router.get_bridges()), report_factory.send_bridge()) - ).start(66).addErrback(_looping_errback, logger) + if config.get("GLOBAL", {}).get("DEBUG_BRIDGES"): + task.LoopingCall( + lambda: (bridge_use_cases.bridge_debug_loop(), report_factory.set_bridges(bridge_router.get_bridges()), report_factory.send_bridge()) + ).start(66).addErrback(_looping_errback, logger) # Alias reload (STALE_DAYS -> seconds) alias_interval = (aliases_cfg.get("STALE_DAYS") or 1) * 86400 @@ -401,6 +406,12 @@ def main() -> None: def sig_handler(sig, frame): logger.info("(GLOBAL) SHUTDOWN: CONFBRIDGE IS TERMINATING WITH SIGNAL %s", sig) + for _sys_name in protocols: + logger.info("(GLOBAL) SHUTDOWN: DE-REGISTER SYSTEM: %s", _sys_name) + try: + protocols[_sys_name].dereg() + except Exception: + pass config["GLOBAL"]["_KILL_SERVER"] = True if reactor.running: reactor.stop() @@ -463,6 +474,7 @@ def main() -> None: logger.info("(GLOBAL) UDP %s listening on %s:%s", system_name, ip or "*", udp_port) logger.info("(GLOBAL) ADN DMR Peer Server started. Use adn-dmr-server as reference.") + reactor.suggestThreadPoolSize(100) reactor.run()