fix(obp): improve BCSQ delivery and loop-control feedback

- Sync TARGET_SOCK from actual peer when RELAX_CHECKS accepts DMRD/DMRE
- Fallback TARGET_IP/PORT for BCSQ send; warn if no destination
- Resend BCSQ every 2s while loop-control loser (UDP loss)
- Emit END,RX on first loser transition for monitor
- Robust _bcsq TG key matching in to_target
pull/4/head
Rodrigo Pérez 6 months ago
parent 53d81140c6
commit 820150b904

@ -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,

@ -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,
)

@ -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")

Loading…
Cancel
Save

Powered by TurnKey Linux.