From 160f1ae8026e5f3ff4267ca23accce32d2b408ce Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 20 Mar 2026 21:13:42 +0000 Subject: [PATCH] fix: align OpenBridge and bridge routing with legacy server Match bridge_master/hblink behaviour: BRIDGES table iteration for to_target, OBP VTERM/_fin timing, stream seq/crc handling, and document related GLOBAL options. --- adn-server.example.yaml | 3 + .../application/bridge_use_cases.py | 629 ++++++++++++------ .../twisted_adapters/udp_hbp.py | 173 ++++- 3 files changed, 569 insertions(+), 236 deletions(-) diff --git a/adn-server.example.yaml b/adn-server.example.yaml index caa97f6..09fafce 100644 --- a/adn-server.example.yaml +++ b/adn-server.example.yaml @@ -17,6 +17,9 @@ GLOBAL: PASS_SECURITY: "" USERS_PASS: user_passwords.json HASH_ENCRYPT: encryption_key.secret + # Multi-OBP mesh hub: when true, inbound peer BCSQ does not block bridge TX from the mesh-designated source OBP + # (see _obp_mesh_first). Set false for strict legacy routerOBP behavior. + OBP_BRIDGE_MESH_OWNER_BYPASS_BCSQ: true REPORTS: REPORT: true diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index 5990b90..424a4e1 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -72,6 +72,13 @@ class BridgeUseCases: self._report_factory = report_factory self._on_bridge_deactivated = on_bridge_deactivated # (system_name: str) -> None; legacy disconnectedVoice self._send_bcsq = send_bcsq # (system_name, tgid, stream_id) -> None; legacy OBP send_bcsq from router + # (stream_id, tgid_bytes) -> {"owner": system_name, "last": monotonic}; cross-OBP first leg (bridge_master race) + self._obp_mesh_first: dict[tuple[bytes, bytes], dict[str, Any]] = {} + + def _obp_prune_mesh_first(self, now_m: float) -> None: + stale = [k for k, v in self._obp_mesh_first.items() if now_m - float(v.get("last", 0)) > 180.0] + for k in stale: + self._obp_mesh_first.pop(k, None) def get_bridges(self) -> dict[str, list[dict[str, Any]]]: """Return current BRIDGES.""" @@ -343,6 +350,15 @@ class BridgeUseCases: pass to_remove.append(stream_id) for stream_id in to_remove: + st_rm = obp_streams.get(stream_id) or {} + tgid_rm = st_rm.get("TGID", b"\x00\x00\x00") + self._obp_mesh_first.pop((stream_id, tgid_rm), None) + _syscfg = systems_cfg.get(system_name, {}) + _bmap = _syscfg.get("_bcsq") + if isinstance(_bmap, dict): + for _tgid_k, _sid in list(_bmap.items()): + if _sid == stream_id: + _bmap.pop(_tgid_k, None) obp_streams.pop(stream_id, None) continue for slot in (1, 2): @@ -438,6 +454,70 @@ class BridgeUseCases: """Check ID against ACL. Legacy acl_check.""" return self._router.acl_check(id_val, acl) + def _ensure_obp_source_for_tg( + self, + system_name: str, + bridge_key: str, + dst_id_b: bytes, + dst_int: int, + ) -> None: + """Ensure this OBP has an ACTIVE source row for TG (TS1) in main and #reflector bridges. + + remove_bridge_system / BRIDGERESET sets all rows for a system to ACTIVE False. Local MASTER + traffic still matches MASTER source rows; inbound OBP traffic needs these OBP rows re-enabled + or added (e.g. new OBP in config after bridge was built). + Same TG range as make_single_bridge OBP entries. + """ + systems_cfg = self._config.get("SYSTEMS", {}) + if systems_cfg.get(system_name, {}).get("MODE") != "OPENBRIDGE": + return + if not systems_cfg.get(system_name, {}).get("ENABLED", True): + return + if not (79 <= dst_int < 9990 or dst_int > 9999): + return + bridges = self._router.get_bridges() + now = time.time() + + def _tgid_match(entry: dict[str, Any]) -> bool: + tg = entry.get("TGID") + if tg == dst_id_b: + return True + try: + return int_id(tg or b"\x00\x00\x00") == dst_int + except (TypeError, ValueError): + return False + + def _patch(entries: list[dict[str, Any]]) -> None: + for e in entries: + if e.get("SYSTEM") != system_name: + continue + if e.get("TS") != 1: + continue + if not _tgid_match(e): + continue + if not e.get("ACTIVE"): + e["ACTIVE"] = True + return + entries.append( + { + "SYSTEM": system_name, + "TS": 1, + "TGID": dst_id_b, + "ACTIVE": True, + "TIMEOUT": "", + "TO_TYPE": "NONE", + "OFF": [], + "ON": [], + "RESET": [], + "TIMER": now, + } + ) + + for key in (bridge_key, "#" + bridge_key): + if key not in bridges: + continue + _patch(bridges[key]) + def make_single_bridge( self, _tgid: bytes | int, @@ -1080,8 +1160,20 @@ class BridgeUseCases: 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: - """Called by UDP MASTER when DMRD is received. Forward to other systems in same bridge (to_target).""" + """Called by UDP when DMRD is received. Forward to other systems in same bridge (to_target). + + 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 bridge_master 3080–3085: private call to ID 4000 only disconnects dynamics; do not route as PC. @@ -1119,28 +1211,29 @@ class BridgeUseCases: self.make_single_bridge(dst_id, system_name, slot, float(tmout)) self.apply_static_tg_to_bridge(dst_int) bridges = self._router.get_bridges() - # Legacy bridge_master 3198–3218: forward to current bridge (tgid) and reflector bridge (#tgid) - entries = list(bridges.get(bridge_key, [])) + list(bridges.get("#" + bridge_key, [])) - # Legacy: only forward if source rule matches (system, TGID, slot, ACTIVE True) — bridge.py 349, 664 + # Legacy bridge_master routerOBP ~2413-2418: scan every BRIDGES[_bridge] table; forward only within + # a table that contains a matching ACTIVE source row (not a flat merge of TG + #TG only). dst_id_b = dst_id if isinstance(dst_id, bytes) and len(dst_id) >= 3 else bytes_3(dst_int) - has_source = any( - entry.get("SYSTEM") == system_name - and entry.get("TS") == bridge_match_slot - and entry.get("ACTIVE") - and (entry.get("TGID") == dst_id_b or int_id(entry.get("TGID") or b"\x00\x00\x00") == dst_int) - for entry in entries - ) + + def _row_is_active_source(row: dict[str, Any]) -> bool: + return bool( + row.get("SYSTEM") == system_name + and row.get("TS") == bridge_match_slot + and row.get("ACTIVE") + and ( + row.get("TGID") == dst_id_b + or int_id(row.get("TGID") or b"\x00\x00\x00") == dst_int + ) + ) + + if source_is_obp: + self._ensure_obp_source_for_tg(system_name, bridge_key, dst_id_b, dst_int) + bridges = self._router.get_bridges() + has_source = any(_row_is_active_source(e) for elist in bridges.values() for e in elist) if not has_source and systems_cfg.get(system_name, {}).get("MODE") == "MASTER": self.options_config_for_system(system_name) bridges = self._router.get_bridges() - entries = list(bridges.get(bridge_key, [])) + list(bridges.get("#" + bridge_key, [])) - has_source = any( - entry.get("SYSTEM") == system_name - and entry.get("TS") == bridge_match_slot - and entry.get("ACTIVE") - and (entry.get("TGID") == dst_id_b or int_id(entry.get("TGID") or b"\x00\x00\x00") == dst_int) - for entry in entries - ) + has_source = any(_row_is_active_source(e) for elist in bridges.values() for e in elist) if not has_source: logger.debug( "(ROUTER) No matching source rule for TG %s from %s slot %s (ACTIVE), not forwarding", @@ -1180,37 +1273,29 @@ class BridgeUseCases: ) obp_streams[stream_id]["LOOPLOG"] = True return - # First OBP: only the OBP that received the stream first (min 1ST) forwards; others drop - hr_times = {} - if obp_streams is not None and stream_id in obp_streams and "1ST" in obp_streams[stream_id]: - hr_times[system_name] = obp_streams[stream_id]["1ST"] - dst_id_b = dst_id if isinstance(dst_id, bytes) and len(dst_id) >= 3 else bytes_3(int_id(dst_id)) - for other_name, proto in (protocols or {}).items(): - if other_name == system_name: - continue - if systems_cfg.get(other_name, {}).get("MODE") != "OPENBRIDGE": - continue - other_streams = getattr(proto, "_obp_streams", None) - if other_streams is None or stream_id not in other_streams: - continue - st = other_streams[stream_id] - if "1ST" not in st or (st.get("TGID") or b"") != dst_id_b: - continue - hr_times[other_name] = st["1ST"] - if hr_times: - first_system = min(hr_times, key=hr_times.get) - if first_system != system_name: - if obp_streams is not None and stream_id in obp_streams and not obp_streams[stream_id].get("LOOPLOG"): - logger.debug( - "(%s) OBP *LoopControl* FIRST OBP %s, STREAM ID: %s, TG %s, IGNORE THIS SOURCE", - system_name, first_system, int_id(stream_id), int_id(dst_id), - ) - obp_streams[stream_id]["LOOPLOG"] = True - if self._send_bcsq and systems_cfg.get(system_name, {}).get("ENHANCED_OBP"): - if obp_streams is not None and stream_id in obp_streams and "_bcsq" not in obp_streams[stream_id]: - self._send_bcsq(system_name, dst_id, stream_id) - obp_streams[stream_id]["_bcsq"] = True - return + # First OBP: single server-wide winner per (stream_id, TGID). Per-leg 1ST compare can both-forward on + # the first reactor tick before peer _obp_streams exists → duplicate TX → mesh BCSQ on every link. + _mesh_key = (stream_id, dst_id_b) + _now_m = time.monotonic() + self._obp_prune_mesh_first(_now_m) + _mf = self._obp_mesh_first.get(_mesh_key) + if _mf is None: + _mf = {"owner": system_name, "last": _now_m} + self._obp_mesh_first[_mesh_key] = _mf + elif _mf.get("owner") != system_name: + if obp_streams is not None and stream_id in obp_streams and not obp_streams[stream_id].get("LOOPLOG"): + logger.debug( + "(%s) OBP *LoopControl* FIRST OBP (mesh) %s, STREAM ID: %s, TG %s, IGNORE THIS SOURCE", + system_name, _mf.get("owner"), int_id(stream_id), int_id(dst_id), + ) + obp_streams[stream_id]["LOOPLOG"] = True + if self._send_bcsq and systems_cfg.get(system_name, {}).get("ENHANCED_OBP"): + if obp_streams is not None and stream_id in obp_streams and "_bcsq" not in obp_streams[stream_id]: + self._send_bcsq(system_name, dst_id, stream_id) + obp_streams[stream_id]["_bcsq"] = True + return + else: + _mf["last"] = _now_m # Legacy bridge.py: BRDG_EVENT for report/monitor (GROUP VOICE START/END RX/TX) if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): try: @@ -1246,13 +1331,36 @@ class BridgeUseCases: _bits = data[15] if len(data) > 15 else 0 protocols = self._get_protocols() if self._get_protocols else {} src_proto = protocols.get(system_name) if protocols else None - # Legacy routerHBP: _ber = _data[53:54]; _rssi = _data[54:55] - # Legacy routerOBP: _ber = b'\x00'; _rssi = b'\x00' (from dmrd_received defaults) - if source_is_obp: + # Legacy `bridge.py` `routerOBP.dmrd_received` forwards `_hops,_ber,_rssi,_source_server,_source_rptr` + # exactly as received from `hblink` (`bridge.py` ~486). Values are set in `hblink`: + # - DMRD v1 OBP: SERVER_ID, zeros rptr, empty hops (`hblink.py` ~338–345) + # - DMRE v5: packet fields + incremented hops (`hblink.py` ~592–596) + # Legacy `bridge_master` `routerHBP.dmrd_received`: ber/rssi from payload, SERVER_ID + peer as rptr (~2936–2940). + if source_is_obp and obp_use_parsed: + _hops = obp_hops + _ber = obp_ber + _rssi = obp_rssi + _sid = self._config.get("GLOBAL", {}).get("SERVER_ID") + if obp_source_server is not None: + _source_server = obp_source_server + elif 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 = obp_source_rptr + elif source_is_obp: _ber = b"\x00" _rssi = b"\x00" _hops = b"" - _source_server = b"\x00\x00\x00\x00" + _sid = self._config.get("GLOBAL", {}).get("SERVER_ID") + 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 = b"\x00\x00\x00\x00" else: _ber = data[53:54] if len(data) > 53 else b"\x00" @@ -1276,197 +1384,288 @@ class BridgeUseCases: source_lc = st.get(slot, {}).get("RX_LC") if not source_lc or len(source_lc) < 9: source_lc = b"\x00\x00\x20" + dst_id_b + rf_src - # Legacy bridge_master to_target: _sysIgnore — dedupe OBP (SYSTEM, TS) per packet across - # multiple to_target passes (current TG + reflector + other bridges). We merge those - # lists into one `entries` loop, so the same OpenBridge can appear twice (e.g. TG + #TG) - # and would get duplicate sends → echo/loops without this. + # Legacy bridge_master routerOBP: _sysIgnore accumulates across each to_target(BRIDGES[_bridge]) + # pass; dedupe (SYSTEM, TS) for OpenBridge targets so the same leg is not sent twice per packet. sys_ignore_obp: set[tuple[str, int]] = set() forwarded = [] - for entry in entries: - if entry.get("SYSTEM") == system_name: - continue - if not entry.get("ACTIVE", False): + _g = self._config.get("GLOBAL", {}) + _bypass_bcsq_mesh_owner = _g.get("OBP_BRIDGE_MESH_OWNER_BYPASS_BCSQ", True) + _obp_mesh_forward_owner = False + if source_is_obp and _bypass_bcsq_mesh_owner: + _om = self._obp_mesh_first.get((stream_id, dst_id_b)) + _obp_mesh_forward_owner = _om is not None and _om.get("owner") == system_name + for _bridge_table_name, _bridge_rows in list(bridges.items()): + if not any(_row_is_active_source(r) for r in _bridge_rows): continue - if not systems_cfg.get(entry["SYSTEM"], {}).get("ENABLED", True): - continue - _target_system = systems_cfg.get(entry["SYSTEM"], {}) - target_mode = _target_system.get("MODE") - tgt_proto = protocols.get(entry["SYSTEM"]) if protocols else None - _target_status = getattr(tgt_proto, "STATUS", None) if tgt_proto else None - - if target_mode == "OPENBRIDGE": - # ── Exact port of legacy bridge.py OBP target (lines 355-401 / 672-718) ── - entry_ts_obp = entry.get("TS") - if entry_ts_obp is None: - entry_ts_obp = 1 - _obp_key = (entry["SYSTEM"], int(entry_ts_obp)) - if _obp_key in sys_ignore_obp: + for entry in _bridge_rows: + if entry.get("SYSTEM") == system_name: continue - sys_ignore_obp.add(_obp_key) - target_tgid = entry.get("TGID") - if isinstance(target_tgid, int): - target_tgid = bytes_3(target_tgid) - if _target_status is not None: - if stream_id not in _target_status: - _target_status[stream_id] = { - "START": pkt_time, - "CONTENTION": False, - "RFS": rf_src, - "TGID": dst_id_b, - } - dst_lc = source_lc[0:3] + target_tgid + rf_src - _target_status[stream_id]["H_LC"] = bptc.encode_header_lc(dst_lc) - _target_status[stream_id]["T_LC"] = bptc.encode_terminator_lc(dst_lc) - _target_status[stream_id]["EMB_LC"] = bptc.encode_emblc(dst_lc) + if not entry.get("ACTIVE", False): + continue + if not systems_cfg.get(entry["SYSTEM"], {}).get("ENABLED", True): + continue + _target_system = systems_cfg.get(entry["SYSTEM"], {}) + target_mode = _target_system.get("MODE") + tgt_proto = protocols.get(entry["SYSTEM"]) if protocols else None + _target_status = getattr(tgt_proto, "STATUS", None) if tgt_proto else None + + if target_mode == "OPENBRIDGE": + # ── bridge_master.routerOBP.to_target OPENBRIDGE branch (~1850-1959) ── + entry_ts_obp = entry.get("TS") + if entry_ts_obp is None: + entry_ts_obp = 1 + _obp_key = (entry["SYSTEM"], int(entry_ts_obp)) + if _obp_key in sys_ignore_obp: + continue + sys_ignore_obp.add(_obp_key) + target_tgid = entry.get("TGID") + if isinstance(target_tgid, int): + target_tgid = bytes_3(target_tgid) + # If target has quenched us, don't send (~1856-1859). Mesh hub: peers often BCSQ every leg on + # duplicate paths; bypass for the server-designated forward OBP so TX continues to remotes. + _bcsq_map = _target_system.get("_bcsq") + if ( + not _obp_mesh_forward_owner + and isinstance(_bcsq_map, dict) + and dst_id_b in _bcsq_map + and _bcsq_map[dst_id_b] == stream_id + ): + logger.debug( + "(%s) OBP skip (BCSQ): target=%s TGID=%s stream=%s", + system_name, + entry["SYSTEM"], + int_id(dst_id_b), + int_id(stream_id), + ) + continue + # If target has missed keepalives (ENHANCED_OBP), don't send (~1861-1863) + if _target_system.get("ENHANCED_OBP") and ( + "_bcka" not in _target_system or _target_system["_bcka"] < pkt_time - 60 + ): + logger.debug( + "(%s) OBP skip (BCKA stale/missing): target=%s _bcka=%s pkt_time=%s", + system_name, + entry["SYSTEM"], + _target_system.get("_bcka"), + pkt_time, + ) + continue + # Talkgroup ACL (global + per-system TG1) (~1865-1873) + _global_cfg = self._config.get("GLOBAL", {}) + if _global_cfg.get("USE_ACL"): + if not self.acl_check(target_tgid, _global_cfg.get("TG1_ACL", (True, []))): + logger.debug( + "(%s) OBP skip (global ACL): target=%s TGID=%s", + system_name, + entry["SYSTEM"], + int_id(target_tgid), + ) + continue + if not self.acl_check(target_tgid, _target_system.get("TG1_ACL", (True, []))): + logger.debug( + "(%s) OBP skip (system ACL): target=%s TGID=%s", + system_name, + entry["SYSTEM"], + int_id(target_tgid), + ) + continue + if _target_status is not None: + if stream_id not in _target_status: + _target_status[stream_id] = { + "START": pkt_time, + "CONTENTION": False, + "RFS": rf_src, + "TGID": dst_id_b, + "RX_PEER": peer_id, + "EMB_LC": {1: b"\x00", 2: b"\x00", 3: b"\x00", 4: b"\x00"}, + "H_LC": b"\x00", + "T_LC": b"\x00", + } + try: + dst_lc = source_lc[0:3] + target_tgid + rf_src + except Exception: + logger.exception("(to_target) caught exception") + _target_status[stream_id]["LAST"] = pkt_time + return + _target_status[stream_id]["H_LC"] = bptc.encode_header_lc(dst_lc) + _target_status[stream_id]["T_LC"] = bptc.encode_terminator_lc(dst_lc) + _target_status[stream_id]["EMB_LC"] = bptc.encode_emblc(dst_lc) + logger.debug( + "(%s) Conference Bridge: %s, Call Bridged to OBP System: %s TS: %s, TGID: %s", + system_name, _bridge_table_name, entry["SYSTEM"], entry.get("TS", 1), int_id(target_tgid), + ) + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + try: + self._report_factory.send_bridge_event( + "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( + entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid) + ) + ) + except Exception: + pass + if "EMB_LC" not in _target_status[stream_id]: + try: + dst_lc = source_lc[0:3] + target_tgid + rf_src + _target_status[stream_id]["EMB_LC"] = bptc.encode_emblc(dst_lc) + except Exception: + logger.exception("(to_target) caught exception while creating EMB_LC") + return + if "H_LC" not in _target_status[stream_id]: + try: + dst_lc = source_lc[0:3] + target_tgid + rf_src + _target_status[stream_id]["H_LC"] = bptc.encode_header_lc(dst_lc) + except Exception: + logger.exception("(to_target) caught exception while creating H_LC") + return + if "T_LC" not in _target_status[stream_id]: + try: + dst_lc = source_lc[0:3] + target_tgid + rf_src + _target_status[stream_id]["T_LC"] = bptc.encode_terminator_lc(dst_lc) + except Exception: + logger.exception("(to_target) caught exception while creating T_LC") + return + _target_status[stream_id]["LAST"] = pkt_time + _tmp_bits = _bits & ~(1 << 7) + _tmp_data = b"".join([data[:8], target_tgid, data[11:15], bytes([_tmp_bits]), data[16:20]]) + dmrbits = bitarray(endian="big") + dmrbits.frombytes(dmrpkt) + if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: + dmrbits = _target_status[stream_id]["H_LC"][0:98] + dmrbits[98:166] + _target_status[stream_id]["H_LC"][98:197] + elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: + dmrbits = _target_status[stream_id]["T_LC"][0:98] + dmrbits[98:166] + _target_status[stream_id]["T_LC"][98:197] + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + try: + call_duration = pkt_time - _target_status[stream_id].get("START", pkt_time) + self._report_factory.send_bridge_event( + "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( + entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration + ) + ) + except Exception: + pass + elif dtype_vseq in (1, 2, 3, 4): + dmrbits = dmrbits[0:116] + _target_status[stream_id]["EMB_LC"][dtype_vseq] + dmrbits[148:264] + dmrpkt_out = dmrbits.tobytes() + _tmp_data = b"".join([_tmp_data, dmrpkt_out]) + else: + _tmp_bits = _bits & ~(1 << 7) + _tmp_data = b"".join([data[:8], target_tgid, data[11:15], bytes([_tmp_bits]), data[16:20]]) + _tmp_data = b"".join([_tmp_data, dmrpkt]) + try: + self._send_to_system( + entry["SYSTEM"], _tmp_data, + _hops=_hops, _ber=_ber, _rssi=_rssi, + _source_server=_source_server, _source_rptr=_source_rptr, + ) + forwarded.append(entry["SYSTEM"]) + except Exception as e: + logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), e) + + else: + # ── Exact port of legacy bridge.py HBP target (lines 403-486 / 720-803) ── + entry_ts = entry.get("TS") + if entry_ts is None: + continue + entry_tgid_b = entry.get("TGID") + if isinstance(entry_tgid_b, int): + entry_tgid_b = bytes_3(entry_tgid_b) + 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 (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))): + 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))): + 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): + 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): + 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: + _is_new_stream = (_ts_st.get("TX_STREAM_ID") != stream_id) + else: + src_status = getattr(src_proto, "STATUS", None) if src_proto else None + _is_new_stream = (src_status.get(slot, {}).get("RX_STREAM_ID") if src_status else b"") != stream_id + if _is_new_stream: + _ts_st["TX_START"] = pkt_time + _ts_st["TX_TGID"] = entry_tgid_b + _ts_st["TX_STREAM_ID"] = stream_id + _ts_st["TX_RFS"] = rf_src + _ts_st["TX_PEER"] = peer_id + dst_lc = source_lc[0:3] + entry_tgid_b + rf_src + _ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc) + _ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc) + _ts_st["TX_EMB_LC"] = bptc.encode_emblc(dst_lc) logger.info( - "(%s) Conference Bridge: %s, Call Bridged to OBP System: %s TS: %s, TGID: %s", - system_name, bridge_key, entry["SYSTEM"], entry.get("TS", 1), int_id(target_tgid), + "(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s", + system_name, _bridge_table_name, entry["SYSTEM"], entry_ts, int_id(entry_tgid_b), ) if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): try: self._report_factory.send_bridge_event( "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid) + entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b) ) ) except Exception: pass - _target_status[stream_id]["LAST"] = pkt_time - _tmp_bits = _bits & ~(1 << 7) - _tmp_data = b"".join([data[:8], target_tgid, data[11:15], bytes([_tmp_bits]), data[16:20]]) + _ts_st["TX_TIME"] = pkt_time + _ts_st["TX_TYPE"] = dtype_vseq + # Slot bit rewrite (legacy bridge.py 457-460 / 770-773) + _src_entry_ts = slot + if _src_entry_ts != entry_ts: + _tmp_bits = _bits ^ (1 << 7) + else: + _tmp_bits = _bits + _tmp_data = b"".join([data[:8], entry_tgid_b, data[11:15], bytes([_tmp_bits]), data[16:20]]) + # LC rewrite (exact port of legacy bridge.py 468-482 / 781-799) dmrbits = bitarray(endian="big") dmrbits.frombytes(dmrpkt) if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: - dmrbits = _target_status[stream_id]["H_LC"][0:98] + dmrbits[98:166] + _target_status[stream_id]["H_LC"][98:197] + dmrbits = _ts_st["TX_H_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_H_LC"][98:197] elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: - dmrbits = _target_status[stream_id]["T_LC"][0:98] + dmrbits[98:166] + _target_status[stream_id]["T_LC"][98:197] + dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197] if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): try: - call_duration = pkt_time - _target_status[stream_id].get("START", pkt_time) + call_duration = pkt_time - _ts_st.get("TX_START", pkt_time) self._report_factory.send_bridge_event( "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration + entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b), call_duration ) ) except Exception: pass elif dtype_vseq in (1, 2, 3, 4): - dmrbits = dmrbits[0:116] + _target_status[stream_id]["EMB_LC"][dtype_vseq] + dmrbits[148:264] + try: + dmrbits = dmrbits[0:116] + _ts_st["TX_EMB_LC"][dtype_vseq] + dmrbits[148:264] + except Exception: + pass dmrpkt_out = dmrbits.tobytes() + # bridge_master.routerOBP.to_target HBP branch: _tmp_data + dmrpkt only (~2041-2042); + # HBP source adds BER/RSSI from payload (bridge.py routerHBP ~800). if source_is_obp: _tmp_data = b"".join([_tmp_data, dmrpkt_out]) else: - _tmp_data = b"".join([_tmp_data, dmrpkt_out]) - else: - _tmp_bits = _bits & ~(1 << 7) - _tmp_data = b"".join([data[:8], target_tgid, data[11:15], bytes([_tmp_bits]), data[16:20]]) - _tmp_data = b"".join([_tmp_data, dmrpkt]) - try: - self._send_to_system( - entry["SYSTEM"], _tmp_data, - _hops=_hops, _ber=_ber, _rssi=_rssi, - _source_server=_source_server, _source_rptr=_source_rptr, - ) - forwarded.append(entry["SYSTEM"]) - except Exception as e: - logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), e) - - else: - # ── Exact port of legacy bridge.py HBP target (lines 403-486 / 720-803) ── - entry_ts = entry.get("TS") - if entry_ts is None: - continue - entry_tgid_b = entry.get("TGID") - if isinstance(entry_tgid_b, int): - entry_tgid_b = bytes_3(entry_tgid_b) - 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 (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))): - 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))): - 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): - 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): - 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: - _is_new_stream = (_ts_st.get("TX_STREAM_ID") != stream_id) - else: - src_status = getattr(src_proto, "STATUS", None) if src_proto else None - _is_new_stream = (src_status.get(slot, {}).get("RX_STREAM_ID") if src_status else b"") != stream_id - if _is_new_stream: - _ts_st["TX_START"] = pkt_time - _ts_st["TX_TGID"] = entry_tgid_b - _ts_st["TX_STREAM_ID"] = stream_id - _ts_st["TX_RFS"] = rf_src - _ts_st["TX_PEER"] = peer_id - dst_lc = source_lc[0:3] + entry_tgid_b + rf_src - _ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc) - _ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc) - _ts_st["TX_EMB_LC"] = bptc.encode_emblc(dst_lc) - logger.info( - "(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s", - system_name, bridge_key, entry["SYSTEM"], entry_ts, int_id(entry_tgid_b), - ) - if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): - try: - self._report_factory.send_bridge_event( - "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b) - ) - ) - except Exception: - pass - _ts_st["TX_TIME"] = pkt_time - _ts_st["TX_TYPE"] = dtype_vseq - # Slot bit rewrite (legacy bridge.py 457-460 / 770-773) - _src_entry_ts = slot - if _src_entry_ts != entry_ts: - _tmp_bits = _bits ^ (1 << 7) - else: - _tmp_bits = _bits - _tmp_data = b"".join([data[:8], entry_tgid_b, data[11:15], bytes([_tmp_bits]), data[16:20]]) - # LC rewrite (exact port of legacy bridge.py 468-482 / 781-799) - dmrbits = bitarray(endian="big") - dmrbits.frombytes(dmrpkt) - if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: - dmrbits = _ts_st["TX_H_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_H_LC"][98:197] - elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: - dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197] - if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): - try: - call_duration = pkt_time - _ts_st.get("TX_START", pkt_time) - self._report_factory.send_bridge_event( - "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b), call_duration - ) - ) - except Exception: - pass - elif dtype_vseq in (1, 2, 3, 4): + _tmp_data = b"".join([_tmp_data, dmrpkt_out, data[53:55]]) try: - dmrbits = dmrbits[0:116] + _ts_st["TX_EMB_LC"][dtype_vseq] + dmrbits[148:264] - except Exception: - pass - dmrpkt_out = dmrbits.tobytes() - # Legacy OBP source: _tmp_data + dmrpkt + b'\x00\x00' (bridge.py 483) - # Legacy HBP source: _tmp_data + dmrpkt + _data[53:55] (bridge.py 800) - if source_is_obp: - _tmp_data = b"".join([_tmp_data, dmrpkt_out, b"\x00\x00"]) - else: - _tmp_data = b"".join([_tmp_data, dmrpkt_out, data[53:55]]) - try: - self._send_to_system( - entry["SYSTEM"], _tmp_data, - _hops=_hops, _ber=_ber, _rssi=_rssi, - _source_server=_source_server, _source_rptr=_source_rptr, - ) - forwarded.append(entry["SYSTEM"]) - except Exception as e: - logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), e) + self._send_to_system( + entry["SYSTEM"], _tmp_data, + _hops=_hops, _ber=_ber, _rssi=_rssi, + _source_server=_source_server, _source_rptr=_source_rptr, + ) + forwarded.append(entry["SYSTEM"]) + except Exception as e: + logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), e) + # Legacy bridge_master routerOBP dmrd_received ~2420-2434: after to_target, VTERM resets lastSeq; + # _fin only when REPORT (blocks further packets on this stream_id from this source). + if source_is_obp and frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: + _src_p = protocols.get(system_name) if protocols else None + _obp_st = getattr(_src_p, "_obp_streams", None) if _src_p else None + if _obp_st is not None and stream_id in _obp_st: + _obp_st[stream_id]["lastSeq"] = False + if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): + _obp_st[stream_id]["_fin"] = True if forwarded: if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: logger.info( @@ -1475,8 +1674,8 @@ class BridgeUseCases: ) else: logger.debug( - "(ROUTER) No ACTIVE targets for TG %s from %s (bridge has %s entries)", - bridge_key, system_name, len(bridges.get(bridge_key, [])), + "(ROUTER) No ACTIVE targets forwarded for TG %s from %s (%s bridge tables)", + bridge_key, system_name, len(bridges), ) def _pvt_call_received( diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 6579e25..1feb927 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -181,6 +181,9 @@ class HBPProtocol(DatagramProtocol): self._config.get("TARGET_IP", ""), self._config.get("TARGET_PORT", ""), ) + # bridge_master.routerOBP.to_target skips ENHANCED targets when '_bcka' not in SYSTEMS[name]. + # Seed so cross-OBP forwarding works before the first inbound BCKA/DMR on *this* leg. + 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) @@ -228,7 +231,8 @@ class HBPProtocol(DatagramProtocol): if self._config.get("MODE") == "MASTER": self.send_peers(_packet, _hops, _ber, _rssi, _source_server, _source_rptr) elif self._config.get("MODE") == "OPENBRIDGE": - if "STUN" in self._CONFIG: + # Global STUN (config) or per-system BCST (hblink sets _config['_STUN'] on BCST RX) + if "STUN" in self._CONFIG or self._config.get("_STUN"): logger.info("(%s) Bridge STUNned, discarding", self._system) return if not _hops: @@ -236,7 +240,7 @@ class HBPProtocol(DatagramProtocol): if _packet[:3] == DMR and self._config.get("TARGET_IP"): _sid = self._CONFIG.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") _server_id = _sid if isinstance(_sid, bytes) and len(_sid) >= 4 else bytes_4(int(_sid) & 0xFFFFFFFF if isinstance(_sid, int) else 0) - _passphrase = self._config.get("PASSPHRASE", b"") + _passphrase = _get_passphrase_bytes(self._config) _target_addr = (self._config["TARGET_IP"], self._config["TARGET_PORT"]) if "VER" in self._config and self._config["VER"] > 4: _ver = VER.to_bytes(1, "big") @@ -948,8 +952,11 @@ class HBPProtocol(DatagramProtocol): _rf_src: bytes = b"\x00\x00\x00", _peer_id: bytes = b"\x00\x00\x00\x00", ) -> bool: - """Legacy bridge_master OBP: stream state, 1ST, _fin, duplicate/seq handling. Returns True if packet should be forwarded to bridge.""" + """Port of bridge_master.py dmrd_received: OBP STATUS[_stream_id] packet control (group/vcsbk).""" now = time.time() + _h = blake2b(digest_size=16) + _h.update(_data) + _pkt_crc = _h.digest() # Trim old streams (legacy: 180s) if _stream_id not in self._obp_streams: to_remove = [sid for sid, st in self._obp_streams.items() if st.get("LAST", 0) < now - 180] @@ -964,44 +971,90 @@ class HBPProtocol(DatagramProtocol): "lastSeq": False, "lastData": False, "packets": 0, + "loss": 0, + "crcs": set(), "_fin": False, "RFS": _rf_src, "RX_PEER": _peer_id, } + self._obp_streams[_stream_id]["crcs"].add(_pkt_crc) return True st = self._obp_streams[_stream_id] if st.get("_fin"): return False st["LAST"] = now - st["packets"] = st.get("packets", 0) + 1 - # Duplicate / out-of-order (legacy bridge_master 3256-3184) - if st.get("lastData") and st["lastData"] == _data and _seq > 1: + if "packets" in st: + st["packets"] = st["packets"] + 1 + # Duplicate handling# + # Handle inbound duplicates + # Duplicate complete packet + if st["lastData"] and st["lastData"] == _data and _seq > 1: + st["loss"] += 1 + logger.debug( + "(%s) *PacketControl* last packet is a complete duplicate of the previous one, disgarding. Stream ID:, %s TGID: %s, LOSS: %.2f%%", + self._system, + int_id(_stream_id), + int_id(_dst_id), + ((st["loss"] / st["packets"]) * 100), + ) + return False + # Duplicate SEQ number + if _seq and _seq == st["lastSeq"]: + st["loss"] += 1 logger.debug( - "(%s) *PacketControl* last packet is a complete duplicate, discarding. Stream ID: %s TGID: %s", - self._system, int_id(_stream_id), int_id(_dst_id), + "(%s) *PacketControl* Duplicate sequence number %s, disgarding. Stream ID:, %s TGID: %s, LOSS: %.2f%%", + self._system, + _seq, + int_id(_stream_id), + int_id(_dst_id), + ((st["loss"] / st["packets"]) * 100), ) return False - if _seq is not None and _seq == st.get("lastSeq"): + # Inbound out-of-order packets + if _seq and st["lastSeq"] and (_seq != 1) and (_seq < st["lastSeq"]): + st["loss"] += 1 logger.debug( - "(%s) *PacketControl* Duplicate sequence number %s, discarding. Stream ID: %s TGID: %s", - self._system, _seq, int_id(_stream_id), int_id(_dst_id), + "%s) *PacketControl* Out of order packet - last SEQ: %s, this SEQ: %s, disgarding. Stream ID:, %s TGID: %s, LOSS: %.2f%%", + self._system, + st["lastSeq"], + _seq, + int_id(_stream_id), + int_id(_dst_id), + ((st["loss"] / st["packets"]) * 100), ) return False - if _seq is not None and st.get("lastSeq") is not None and _seq != 1 and _seq < st["lastSeq"]: + # Duplicate DMR payload to previuos packet (by hash + if _seq > 0 and _pkt_crc in st["crcs"]: + st["loss"] += 1 logger.debug( - "(%s) *PacketControl* Out of order packet - last SEQ: %s, this SEQ: %s, discarding. Stream ID: %s TGID: %s", - self._system, st["lastSeq"], _seq, int_id(_stream_id), int_id(_dst_id), + "(%s) *PacketControl* DMR packet payload with hash: %s seen before in this stream, disgarding. Stream ID:, %s TGID: %s: SEQ:%s PACKETS: %s, LOSS: %.2f%% ", + self._system, + _pkt_crc, + int_id(_stream_id), + int_id(_dst_id), + _seq, + st["packets"], + ((st["loss"] / st["packets"]) * 100), ) return False - if _seq is not None and st.get("lastSeq") is not None and _seq > st["lastSeq"] + 1: + # Inbound missed packets + if _seq and st["lastSeq"] and _seq > (st["lastSeq"] + 1): + st["loss"] += 1 logger.debug( - "(%s) *PacketControl* Missed packet(s) - last SEQ: %s, this SEQ: %s. Stream ID: %s TGID: %s", - self._system, st["lastSeq"], _seq, int_id(_stream_id), int_id(_dst_id), + "(%s) *PacketControl* Missed packet(s) - last SEQ: %s, this SEQ: %s. Stream ID:, %s TGID: %s , LOSS: %.2f%%", + self._system, + st["lastSeq"], + _seq, + int_id(_stream_id), + int_id(_dst_id), + ((st["loss"] / st["packets"]) * 100), ) + # Save this sequence number st["lastSeq"] = _seq + # Save this packet st["lastData"] = _data - if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM: - st["_fin"] = True + st["crcs"].add(_pkt_crc) + # Legacy bridge_master routerOBP: _fin is set after to_target + REPORT END (~2430-2432), not here. return True def _obp_datagram_received(self, _packet: bytes, _sockaddr: tuple[str, int]) -> None: @@ -1118,7 +1171,33 @@ class HBPProtocol(DatagramProtocol): self.STATUS[_slot]["RX_RFS"] = _rf_src self.STATUS[_slot]["RX_TIME"] = time.time() 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 hblink DMRD v1: SERVER_ID + default rptr/hops/ber/rssi (`hblink.py` ~338–345, ~416) + _global = self._CONFIG.get("GLOBAL", {}) + _sid = _global.get("SERVER_ID", b"\x00\x00\x00\x00") + _obp_ss = ( + _sid + if isinstance(_sid, bytes) and len(_sid) >= 4 + else bytes_4(int(_sid) & 0xFFFFFFFF if isinstance(_sid, int) else 0) + ) + self._dmrd_received( + self._system, + _peer_id, + _rf_src, + _dst_id, + _seq, + _slot, + _call_type, + _frame_type, + _dtype_vseq, + _stream_id, + _data, + obp_use_parsed=True, + obp_hops=b"", + obp_source_server=_obp_ss, + obp_ber=b"\x00", + obp_rssi=b"\x00", + obp_source_rptr=b"\x00\x00\x00\x00", + ) self._config["_bcka"] = time.time() else: logger.warning("(%s) OpenBridge HMAC failed, packet discarded - OPCODE: %s SRC: %s", self._system, _packet[:4], _sockaddr) @@ -1127,6 +1206,9 @@ class HBPProtocol(DatagramProtocol): if len(_packet) < 69: return _data = _packet[:53] + # Legacy hblink OPENBRIDGE DMRE: BER/RSSI before version split (`hblink.py` ~432–433) + _ber = _packet[53:54] + _rssi = _packet[54:55] _embedded_version = _packet[55] if _embedded_version > 4: if len(_packet) < 89: @@ -1305,8 +1387,28 @@ class HBPProtocol(DatagramProtocol): self.STATUS[_slot]["RX_RFS"] = _rf_src self.STATUS[_slot]["RX_TIME"] = time.time() _data_dmrd = DMRD + _data[4:] + _hops_out = _inthops.to_bytes(1, "big") 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_dmrd) + # Legacy hblink DMRE: same fields passed to dmrd_received as after increment (`hblink.py` ~592–596) + self._dmrd_received( + self._system, + _peer_id, + _rf_src, + _dst_id, + _seq, + _slot, + _call_type, + _frame_type, + _dtype_vseq, + _stream_id, + _data_dmrd, + obp_use_parsed=True, + obp_hops=_hops_out, + obp_source_server=_source_server, + obp_ber=_ber, + obp_rssi=_rssi, + obp_source_rptr=_source_rptr, + ) self._config["_bcka"] = time.time() elif _packet[:4] == EOBP: logger.warning("(%s) *ProtoControl* KF7EEL EOBP protocol not supported", self._system) @@ -1324,6 +1426,35 @@ class HBPProtocol(DatagramProtocol): self._config.pop("_no_target_log_time", None) # reset so next "no target" logs once else: logger.info("(%s) *BridgeControl* BCKA invalid KeepAlive, packet discarded", self._system) + # Source quench — legacy hblink.py OPENBRIDGE ~629-639 (sets CONFIG['_bcsq'][tgid]=stream_id) + if _packet[:4] == BCSQ and len(_packet) >= 31: + _hash_bcsq = _packet[11:] + _tgid_bcsq = _packet[4:7] + _stream_bcsq = _packet[7:11] + _ckhs_bcsq = hmac_new(self._config["PASSPHRASE"], _packet[:11], sha1).digest() + if compare_digest(_hash_bcsq, _ckhs_bcsq): + if "_bcsq" not in self._config: + self._config["_bcsq"] = {} + self._config["_bcsq"][_tgid_bcsq] = _stream_bcsq + else: + logger.warning( + "(%s) *BridgeControl* BCSQ invalid Source Quench, packet discarded - SRC: %s", + self._system, + _sockaddr, + ) + # STUN — must match send_bcst: HMAC-SHA1 over opcode only (hblink.py ~282-285). RX used _packet[4:] in ~647 but that does not match TX. + if _packet[:4] == BCST and len(_packet) >= 24: + _hash_bcst = _packet[4:24] + _ckhs_bcst = hmac_new(self._config["PASSPHRASE"], _packet[:4], sha1).digest() + if compare_digest(_hash_bcst, _ckhs_bcst): + logger.trace("(%s) *BridgeControl* BCST STUN request received", self._system) + self._config["_STUN"] = True + else: + logger.warning( + "(%s) *BridgeControl* BCST invalid STUN, packet discarded - SRC: %s", + self._system, + _sockaddr, + ) if _packet[:4] == BCVE and len(_packet) >= 25: _ver = int.from_bytes(_packet[4:5], "big") _hash = _packet[5:25]