fix(obp): end TX legs on timeout/VTERM and demote max-hops spam

- Emit END,TX for OBP forward destinations when source stream times out in trimmer
- After VTERM END,RX, sweep remaining forward legs so monitor state stays consistent
- Log MAX HOPS drops at DEBUG to avoid per-packet warning floods
- Drop obsolete OBP_BRIDGE_MESH_OWNER_BYPASS_BCSQ comments from example yaml
pull/4/head
Rodrigo Pérez 6 months ago
parent 160f1ae802
commit 53d81140c6

@ -17,9 +17,6 @@ GLOBAL:
PASS_SECURITY: "<set-in-adn-server.yaml>" PASS_SECURITY: "<set-in-adn-server.yaml>"
USERS_PASS: user_passwords.json USERS_PASS: user_passwords.json
HASH_ENCRYPT: encryption_key.secret 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: REPORTS:
REPORT: true REPORT: true

@ -30,10 +30,14 @@ import logging
import re import re
import time import time
from collections import deque from collections import deque
from hashlib import blake2b
from time import perf_counter
from typing import Any from typing import Any
from bitarray import bitarray from bitarray import bitarray
from dmr_utils3 import bptc from dmr_utils3 import bptc
from dmr_utils3 import decode
from dmr_utils3.const import LC_OPT
from ..domain import int_id, bytes_3, bytes_4 from ..domain import int_id, bytes_3, bytes_4
from ..infrastructure.hbp_constants import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO from ..infrastructure.hbp_constants import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO
@ -42,6 +46,12 @@ from .ports import BridgeRouter
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _log_trace(msg: str, *args: Any) -> None:
"""Per-packet forwarding diagnostics (BCSQ/BCKA/ACL). Below DEBUG — enable TRACE in LOGGER config to see."""
if hasattr(logging, "TRACE"):
logger.log(logging.TRACE, msg, *args)
def _is_special_tg(bridge_key: str) -> bool: def _is_special_tg(bridge_key: str) -> bool:
"""True if bridge is special TGID 9990-9999 (excluded from infinite timer).""" """True if bridge is special TGID 9990-9999 (excluded from infinite timer)."""
if bridge_key and bridge_key[0:1] == "#": if bridge_key and bridge_key[0:1] == "#":
@ -72,13 +82,6 @@ class BridgeUseCases:
self._report_factory = report_factory self._report_factory = report_factory
self._on_bridge_deactivated = on_bridge_deactivated # (system_name: str) -> None; legacy disconnectedVoice 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 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]]]: def get_bridges(self) -> dict[str, list[dict[str, Any]]]:
"""Return current BRIDGES.""" """Return current BRIDGES."""
@ -314,6 +317,52 @@ class BridgeUseCases:
_system["TIMER"] = pkt_time _system["TIMER"] = pkt_time
logger.info("(%s) Bridge: %s set to ON with and \"OFF\" timer rule: timeout timer cancelled", system_name, _bridge) logger.info("(%s) Bridge: %s set to ON with and \"OFF\" timer rule: timeout timer cancelled", system_name, _bridge)
def _obp_emit_end_tx_for_forward_legs(self, stream_id: bytes, source_system: str, now: float) -> None:
"""Emit GROUP VOICE,END,TX for every OBP that still holds this stream as a to_target forward leg.
On idle timeout the trimmer sends END,RX for the source only. VTERM may never arrive for
forwarded legs, so the monitor would otherwise keep stale TX chips on destination rows.
Forward legs are identified by STATUS[stream_id] containing H_LC (see to_target OPENBRIDGE).
"""
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 {}
systems_cfg = self._config.get("SYSTEMS", {})
for tgt_name, tgt_proto in (protocols or {}).items():
if tgt_name == source_system:
continue
if systems_cfg.get(tgt_name, {}).get("MODE") != "OPENBRIDGE":
continue
tstatus = getattr(tgt_proto, "STATUS", None)
if not tstatus or stream_id not in tstatus:
continue
tst = tstatus[stream_id]
if not isinstance(tst, dict) or "H_LC" not in tst:
continue
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(
tgt_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: 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).""" """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") logger.debug("(ROUTER) Trimming inactive stream IDs from system lists")
@ -351,15 +400,15 @@ class BridgeUseCases:
to_remove.append(stream_id) to_remove.append(stream_id)
for stream_id in to_remove: for stream_id in to_remove:
st_rm = obp_streams.get(stream_id) or {} 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, {}) _syscfg = systems_cfg.get(system_name, {})
_bmap = _syscfg.get("_bcsq") _bmap = _syscfg.get("_bcsq")
if isinstance(_bmap, dict): if isinstance(_bmap, dict):
for _tgid_k, _sid in list(_bmap.items()): for _tgid_k, _sid in list(_bmap.items()):
if _sid == stream_id: if _sid == stream_id:
_bmap.pop(_tgid_k, None) _bmap.pop(_tgid_k, None)
self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now)
obp_streams.pop(stream_id, None) obp_streams.pop(stream_id, None)
getattr(protocol, "STATUS", {}).pop(stream_id, None)
continue continue
for slot in (1, 2): for slot in (1, 2):
_slot = protocol.STATUS.get(slot) _slot = protocol.STATUS.get(slot)
@ -1147,6 +1196,265 @@ class BridgeUseCases:
except Exception as e: except Exception as e:
logger.exception("(OPTIONS) caught exception: %s", e) logger.exception("(OPTIONS) caught exception: %s", e)
def _obp_wire_stream_dict(self, src_proto: Any, stream_id: bytes, st: dict[str, Any]) -> None:
"""Legacy routerOBP uses one STATUS[stream_id]; mirror into _obp_streams for trimmer/_fin."""
status = getattr(src_proto, "STATUS", None)
if status is not None:
status[stream_id] = st
obp = getattr(src_proto, "_obp_streams", None)
if obp is not None:
obp[stream_id] = st
def _obp_group_voice_router_obp(
self,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
seq: int,
slot: int,
call_type: str,
frame_type: int,
dtype_vseq: int,
stream_id: bytes,
data: bytes,
obp_hops: bytes,
) -> bool:
"""Port of bridge_master.routerOBP.dmrd_received group/vcsbk (~2269-2411). False = drop packet."""
pkt_time = time.time()
dmrpkt = data[20:53] if len(data) >= 53 else b""
_h = blake2b(digest_size=16)
_h.update(data)
_pkt_crc = _h.digest()
protocols = self._get_protocols() if self._get_protocols else {}
src_proto = protocols.get(system_name) if protocols else None
if not src_proto:
return True
systems_cfg = self._config.get("SYSTEMS", {})
_do_report = bool(self._config.get("REPORTS", {}).get("REPORT", True))
status = getattr(src_proto, "STATUS", None)
if status is None:
return True
if stream_id not in status:
st: dict[str, Any] = {
"START": pkt_time,
"CONTENTION": False,
"RFS": rf_src,
"TGID": dst_id,
"1ST": perf_counter(),
"lastSeq": False,
"lastData": False,
"RX_PEER": peer_id,
"packets": 0,
"loss": 0,
"crcs": set(),
}
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
try:
decoded = decode.voice_head_term(dmrpkt)
st["LC"] = decoded["LC"]
except Exception:
st["LC"] = LC_OPT + dst_id + rf_src
else:
st["LC"] = LC_OPT + dst_id + rf_src
self._obp_wire_stream_dict(src_proto, stream_id, st)
_inthops = int.from_bytes(obp_hops, "big") if obp_hops else 0
logger.info(
"(%s) *CALL START* STREAM ID: %s SUB: %s PEER: %s TGID %s TS %s HOPS %s",
system_name,
int_id(stream_id),
int_id(rf_src),
int_id(peer_id),
int_id(dst_id),
slot,
_inthops,
)
if _do_report and self._report_factory and hasattr(self._report_factory, "send_bridge_event"):
try:
self._report_factory.send_bridge_event(
"GROUP VOICE,START,RX,{},{},{},{},{},{}".format(
system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id)
)
)
except Exception:
pass
else:
st = status[stream_id]
if "packets" in st:
st["packets"] = st["packets"] + 1
if "_fin" in st:
if "_finlog" not in st:
logger.debug(
"(%s) OBP *LoopControl* STREAM ID: %s ALREADY FINISHED FROM THIS SOURCE, IGNORING",
system_name,
int_id(stream_id),
)
st["_finlog"] = True
return False
if st["START"] + 180 < pkt_time:
if "LOOPLOG" not in st or not st["LOOPLOG"]:
logger.info(
"(%s) OBP *TIMEOUT*, STREAM ID: %s, TG: %s, IGNORE THIS SOURCE",
system_name,
int_id(stream_id),
int_id(dst_id),
)
st["LOOPLOG"] = True
st["LAST"] = pkt_time
return False
hr_times: dict[str, float] = {}
_sysslot_last = 0
for other_name, proto in (protocols or {}).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:
_sysslot_last = _sysslot if isinstance(_sysslot, int) else _sysslot_last
slot_st = ostatus.get(_sysslot)
if not isinstance(slot_st, dict):
continue
if "RX_STREAM_ID" in slot_st and stream_id == slot_st.get("RX_STREAM_ID"):
if "LOOPLOG" not in st or not st["LOOPLOG"]:
logger.debug(
"(%s) OBP *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,
)
st["LOOPLOG"] = True
st["LAST"] = pkt_time
return False
else:
obp_status = getattr(proto, "STATUS", None)
if not obp_status:
continue
if (
stream_id in obp_status
and "1ST" in obp_status[stream_id]
and obp_status[stream_id].get("TGID") == dst_id
):
hr_times[other_name] = obp_status[stream_id]["1ST"]
fi = min(hr_times, key=hr_times.get, default=False)
hr_times.clear()
if not fi:
logger.warning(
"(%s) OBP *LoopControl* fi is empty for some reason : STREAM ID: %s, TG: %s, TS: %s",
system_name,
int_id(stream_id),
int_id(dst_id),
_sysslot_last,
)
return False
if system_name != fi:
if "LOOPLOG" not in st or not st["LOOPLOG"]:
call_duration = pkt_time - st["START"]
logger.debug(
"(%s) OBP *LoopControl* FIRST OBP %s, STREAM ID: %s, TG %s, IGNORE THIS SOURCE. PACKET RATE %0.2f/s",
system_name,
fi,
int_id(stream_id),
int_id(dst_id),
call_duration,
)
st["LOOPLOG"] = True
st["LAST"] = pkt_time
if systems_cfg.get(system_name, {}).get("ENHANCED_OBP") and "_bcsq" not in st:
if self._send_bcsq:
self._send_bcsq(system_name, dst_id, stream_id)
st["_bcsq"] = True
return False
if st["packets"] > 18 and (st["packets"] / st["START"] > 25):
logger.warning(
"(%s) *PacketControl* RATE DROP! Stream ID:, %s TGID: %s",
system_name,
int_id(stream_id),
int_id(dst_id),
)
pb = getattr(src_proto, "proxy_bad_peer", None)
if callable(pb):
pb()
return False
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%%",
system_name,
int_id(stream_id),
int_id(dst_id),
((st["loss"] / st["packets"]) * 100) if st.get("packets") else 0.0,
)
return False
if seq and seq == st["lastSeq"]:
st["loss"] += 1
logger.debug(
"(%s) *PacketControl* Duplicate sequence number %s, disgarding. Stream ID:, %s TGID: %s, LOSS: %.2f%%",
system_name,
seq,
int_id(stream_id),
int_id(dst_id),
((st["loss"] / st["packets"]) * 100) if st.get("packets") else 0.0,
)
return False
if seq and st["lastSeq"] and (seq != 1) and (seq < st["lastSeq"]):
st["loss"] += 1
logger.debug(
"(%s) *PacketControl* Out of order packet - last SEQ: %s, this SEQ: %s, disgarding. Stream ID:, %s TGID: %s, LOSS: %.2f%%",
system_name,
st["lastSeq"],
seq,
int_id(stream_id),
int_id(dst_id),
((st["loss"] / st["packets"]) * 100) if st.get("packets") else 0.0,
)
return False
if seq > 0 and _pkt_crc in st["crcs"]:
st["loss"] += 1
logger.debug(
"(%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%% ",
system_name,
_pkt_crc,
int_id(stream_id),
int_id(dst_id),
seq,
st["packets"],
((st["loss"] / st["packets"]) * 100) if st.get("packets") else 0.0,
)
return False
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 , LOSS: %.2f%%",
system_name,
st["lastSeq"],
seq,
int_id(stream_id),
int_id(dst_id),
((st["loss"] / st["packets"]) * 100) if st.get("packets") else 0.0,
)
st["lastSeq"] = seq
st["lastData"] = data
st = status[stream_id]
st["crcs"].add(_pkt_crc)
st["LAST"] = pkt_time
if self._config.get("GLOBAL", {}).get("GEN_STAT_BRIDGES"):
_di = int_id(dst_id)
_bk = str(_di)
if _di >= 5 and _di != 9 and _bk not in self._router.get_bridges():
logger.debug("(%s) Bridge for STAT TG %s does not exist. Creating", system_name, _di)
self.make_stat_bridge(dst_id)
return True
def dmrd_received( def dmrd_received(
self, self,
system_name: str, system_name: str,
@ -1229,6 +1537,22 @@ class BridgeUseCases:
if source_is_obp: if source_is_obp:
self._ensure_obp_source_for_tg(system_name, bridge_key, dst_id_b, dst_int) self._ensure_obp_source_for_tg(system_name, bridge_key, dst_id_b, dst_int)
bridges = self._router.get_bridges() bridges = self._router.get_bridges()
if source_is_obp and call_type in ("group", "vcsbk"):
if not self._obp_group_voice_router_obp(
system_name,
peer_id,
rf_src,
dst_id,
seq,
slot,
call_type,
frame_type,
dtype_vseq,
stream_id,
data,
obp_hops if obp_use_parsed else b"",
):
return
has_source = any(_row_is_active_source(e) for elist in bridges.values() for e in elist) 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": if not has_source and systems_cfg.get(system_name, {}).get("MODE") == "MASTER":
self.options_config_for_system(system_name) self.options_config_for_system(system_name)
@ -1241,88 +1565,35 @@ class BridgeUseCases:
) )
return return
pkt_time = time.time() pkt_time = time.time()
# Legacy bridge_master/bridge.py: Loop Control when source is OBP — ignore if we are echo (first HBP or first OBP wins) # Legacy bridge.py: BRDG_EVENT (OBP group/vcsbk START/END handled in _obp_group_voice_router_obp / post-forward VTERM)
if source_is_obp:
protocols = self._get_protocols() if self._get_protocols else {}
src_proto = protocols.get(system_name) if protocols else None
obp_streams = getattr(src_proto, "_obp_streams", None) if src_proto else None
if obp_streams is not None and stream_id in obp_streams:
if obp_streams[stream_id].get("_fin"):
if not obp_streams[stream_id].get("_finlog"):
logger.debug(
"(%s) OBP *LoopControl* STREAM ID: %s ALREADY FINISHED FROM THIS SOURCE, IGNORING",
system_name, int_id(stream_id),
)
obp_streams[stream_id]["_finlog"] = True
return
# First HBP: if any MASTER/PEER already has this stream_id on RX, we are the echo — drop (log once per stream, legacy LOOPLOG)
for other_name, proto in (protocols or {}).items():
if other_name == system_name:
continue
if systems_cfg.get(other_name, {}).get("MODE") == "OPENBRIDGE":
continue
status = getattr(proto, "STATUS", None)
if not status:
continue
for _slot in (1, 2):
if _slot in status and (status[_slot].get("RX_STREAM_ID") or b"") == stream_id:
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 HBP: %s, STREAM ID: %s, TG: %s, TS: %s, IGNORE THIS SOURCE",
system_name, other_name, int_id(stream_id), int_id(dst_id), _slot,
)
obp_streams[stream_id]["LOOPLOG"] = 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"): if self._report_factory and hasattr(self._report_factory, "send_bridge_event"):
try: try:
_obp_grp = source_is_obp and call_type in ("group", "vcsbk")
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
self._report_factory.send_bridge_event( if not _obp_grp:
"GROUP VOICE,START,RX,{},{},{},{},{},{}".format( self._report_factory.send_bridge_event(
system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id) "GROUP VOICE,START,RX,{},{},{},{},{},{}".format(
system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id)
)
) )
)
elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
duration = 0.0 if not _obp_grp:
protocols = self._get_protocols() if self._get_protocols else {} duration = 0.0
src_proto = protocols.get(system_name) if protocols else None protocols = self._get_protocols() if self._get_protocols else {}
if source_is_obp and src_proto and getattr(src_proto, "_obp_streams", None) and stream_id in src_proto._obp_streams: src_proto = protocols.get(system_name) if protocols else None
start = src_proto._obp_streams[stream_id].get("START") if src_proto and getattr(src_proto, "STATUS", None):
if start is not None: st = src_proto.STATUS
duration = pkt_time - start ent = st.get(stream_id)
elif src_proto and getattr(src_proto, "STATUS", None): start = ent.get("START") if isinstance(ent, dict) else None
st = src_proto.STATUS if start is None and slot in st:
start = st.get(stream_id, {}).get("START") or st.get(slot, {}).get("RX_START") start = st.get(slot, {}).get("RX_START")
if start is not None: if start is not None:
duration = pkt_time - start duration = pkt_time - start
self._report_factory.send_bridge_event( self._report_factory.send_bridge_event(
"GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( "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), duration system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), duration
)
) )
)
except Exception: except Exception:
pass pass
# ── Exact port of legacy bridge.py routerOBP/routerHBP forwarding to targets ── # ── Exact port of legacy bridge.py routerOBP/routerHBP forwarding to targets ──
@ -1379,7 +1650,8 @@ class BridgeUseCases:
if src_proto and getattr(src_proto, "STATUS", None): if src_proto and getattr(src_proto, "STATUS", None):
st = src_proto.STATUS st = src_proto.STATUS
if source_is_obp: if source_is_obp:
source_lc = st.get(stream_id, {}).get("LC") ent = st.get(stream_id)
source_lc = ent.get("LC") if isinstance(ent, dict) else None
else: else:
source_lc = st.get(slot, {}).get("RX_LC") source_lc = st.get(slot, {}).get("RX_LC")
if not source_lc or len(source_lc) < 9: if not source_lc or len(source_lc) < 9:
@ -1388,12 +1660,6 @@ class BridgeUseCases:
# pass; dedupe (SYSTEM, TS) for OpenBridge targets so the same leg is not sent twice per packet. # 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() sys_ignore_obp: set[tuple[str, int]] = set()
forwarded = [] forwarded = []
_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()): for _bridge_table_name, _bridge_rows in list(bridges.items()):
if not any(_row_is_active_source(r) for r in _bridge_rows): if not any(_row_is_active_source(r) for r in _bridge_rows):
continue continue
@ -1421,16 +1687,14 @@ class BridgeUseCases:
target_tgid = entry.get("TGID") target_tgid = entry.get("TGID")
if isinstance(target_tgid, int): if isinstance(target_tgid, int):
target_tgid = bytes_3(target_tgid) target_tgid = bytes_3(target_tgid)
# If target has quenched us, don't send (~1856-1859). Mesh hub: peers often BCSQ every leg on # If target has quenched us, don't send (~1856-1859).
# duplicate paths; bypass for the server-designated forward OBP so TX continues to remotes.
_bcsq_map = _target_system.get("_bcsq") _bcsq_map = _target_system.get("_bcsq")
if ( if (
not _obp_mesh_forward_owner isinstance(_bcsq_map, dict)
and isinstance(_bcsq_map, dict)
and dst_id_b in _bcsq_map and dst_id_b in _bcsq_map
and _bcsq_map[dst_id_b] == stream_id and _bcsq_map[dst_id_b] == stream_id
): ):
logger.debug( _log_trace(
"(%s) OBP skip (BCSQ): target=%s TGID=%s stream=%s", "(%s) OBP skip (BCSQ): target=%s TGID=%s stream=%s",
system_name, system_name,
entry["SYSTEM"], entry["SYSTEM"],
@ -1442,7 +1706,7 @@ class BridgeUseCases:
if _target_system.get("ENHANCED_OBP") and ( if _target_system.get("ENHANCED_OBP") and (
"_bcka" not in _target_system or _target_system["_bcka"] < pkt_time - 60 "_bcka" not in _target_system or _target_system["_bcka"] < pkt_time - 60
): ):
logger.debug( _log_trace(
"(%s) OBP skip (BCKA stale/missing): target=%s _bcka=%s pkt_time=%s", "(%s) OBP skip (BCKA stale/missing): target=%s _bcka=%s pkt_time=%s",
system_name, system_name,
entry["SYSTEM"], entry["SYSTEM"],
@ -1454,7 +1718,7 @@ class BridgeUseCases:
_global_cfg = self._config.get("GLOBAL", {}) _global_cfg = self._config.get("GLOBAL", {})
if _global_cfg.get("USE_ACL"): if _global_cfg.get("USE_ACL"):
if not self.acl_check(target_tgid, _global_cfg.get("TG1_ACL", (True, []))): if not self.acl_check(target_tgid, _global_cfg.get("TG1_ACL", (True, []))):
logger.debug( _log_trace(
"(%s) OBP skip (global ACL): target=%s TGID=%s", "(%s) OBP skip (global ACL): target=%s TGID=%s",
system_name, system_name,
entry["SYSTEM"], entry["SYSTEM"],
@ -1462,7 +1726,7 @@ class BridgeUseCases:
) )
continue continue
if not self.acl_check(target_tgid, _target_system.get("TG1_ACL", (True, []))): if not self.acl_check(target_tgid, _target_system.get("TG1_ACL", (True, []))):
logger.debug( _log_trace(
"(%s) OBP skip (system ACL): target=%s TGID=%s", "(%s) OBP skip (system ACL): target=%s TGID=%s",
system_name, system_name,
entry["SYSTEM"], entry["SYSTEM"],
@ -1657,15 +1921,53 @@ class BridgeUseCases:
forwarded.append(entry["SYSTEM"]) forwarded.append(entry["SYSTEM"])
except Exception as e: except Exception as e:
logger.warning("(ROUTER) send_to_system %s failed: %s", entry.get("SYSTEM"), 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; # Legacy bridge_master routerOBP ~2420-2434: after to_target, VTERM — CALL END log, END RX report, _fin, lastSeq
# _fin only when REPORT (blocks further packets on this stream_id from this source). if (
if source_is_obp and frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: source_is_obp
and call_type in ("group", "vcsbk")
and frame_type == HBPF_DATA_SYNC
and dtype_vseq == HBPF_SLT_VTERM
):
_src_p = protocols.get(system_name) if protocols else None _src_p = protocols.get(system_name) if protocols else None
_obp_st = getattr(_src_p, "_obp_streams", None) if _src_p 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: if _obp_st is not None and stream_id in _obp_st:
_obp_st[stream_id]["lastSeq"] = False ost = _obp_st[stream_id]
if self._report_factory and hasattr(self._report_factory, "send_bridge_event"): _end_t = time.time()
_obp_st[stream_id]["_fin"] = True call_duration = _end_t - ost.get("START", _end_t)
packet_rate = (ost.get("packets", 0) / call_duration) if call_duration else 0.0
loss_pct = ((ost.get("loss", 0) / ost["packets"]) * 100) if ost.get("packets") else 0.0
logger.info(
"(%s) *CALL END* STREAM ID: %s SUB: %s PEER: %s TGID %s, TS %s, Duration: %.2f, Packet rate: %.2f/s, Loss: %.2f%%",
system_name,
int_id(stream_id),
int_id(rf_src),
int_id(peer_id),
int_id(dst_id),
slot,
call_duration,
packet_rate,
loss_pct,
)
if self._config.get("REPORTS", {}).get("REPORT", True) 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),
call_duration,
)
)
ost["_fin"] = True
self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, _end_t)
except Exception:
pass
ost["lastSeq"] = False
if forwarded: if forwarded:
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
logger.info( logger.info(
@ -1673,7 +1975,7 @@ class BridgeUseCases:
bridge_key, system_name, ", ".join(forwarded), bridge_key, system_name, ", ".join(forwarded),
) )
else: else:
logger.debug( _log_trace(
"(ROUTER) No ACTIVE targets forwarded for TG %s from %s (%s bridge tables)", "(ROUTER) No ACTIVE targets forwarded for TG %s from %s (%s bridge tables)",
bridge_key, system_name, len(bridges), bridge_key, system_name, len(bridges),
) )

@ -940,122 +940,8 @@ class HBPProtocol(DatagramProtocol):
_packet = _packet + hmac_new(self._config["PASSPHRASE"], _packet, sha1).digest() _packet = _packet + hmac_new(self._config["PASSPHRASE"], _packet, sha1).digest()
self.transport.write(_packet, _addr) self.transport.write(_packet, _addr)
def _obp_stream_accept( def proxy_bad_peer(self) -> None:
self, """Legacy bridge_master routerOBP rate-drop hook; HBSYSTEM has peer handling — OBP noop."""
_stream_id: bytes,
_dst_id: bytes,
_seq: int,
_data: bytes,
_frame_type: int,
_dtype_vseq: int,
_dmrpkt: bytes,
_rf_src: bytes = b"\x00\x00\x00",
_peer_id: bytes = b"\x00\x00\x00\x00",
) -> bool:
"""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]
for sid in to_remove:
self._obp_streams.pop(sid, None)
if _stream_id not in self._obp_streams:
self._obp_streams[_stream_id] = {
"1ST": time.perf_counter(),
"TGID": _dst_id,
"START": now,
"LAST": now,
"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
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* 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
# Inbound out-of-order packets
if _seq and st["lastSeq"] and (_seq != 1) and (_seq < st["lastSeq"]):
st["loss"] += 1
logger.debug(
"%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
# Duplicate DMR payload to previuos packet (by hash
if _seq > 0 and _pkt_crc in st["crcs"]:
st["loss"] += 1
logger.debug(
"(%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
# 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 , 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
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: def _obp_datagram_received(self, _packet: bytes, _sockaddr: tuple[str, int]) -> None:
"""Port of hblink.py OPENBRIDGE.datagramReceived: DMRD v1 (53+HMAC), BCKA, BCVE.""" """Port of hblink.py OPENBRIDGE.datagramReceived: DMRD v1 (53+HMAC), BCKA, BCVE."""
@ -1138,38 +1024,7 @@ class HBPProtocol(DatagramProtocol):
"(%s) CALL RX (OBP) src %s -> TG %s slot %s", "(%s) CALL RX (OBP) src %s -> TG %s slot %s",
self._system, int_id(_rf_src), int_id(_dst_id), _slot, self._system, int_id(_rf_src), int_id(_dst_id), _slot,
) )
# Legacy bridge.py routerOBP: stream state with LC (STATUS[stream_id]) # Group/vcsbk stream state, LC, duplicates: bridge_use_cases._obp_group_voice_router_obp (legacy routerOBP.dmrd_received)
dmrpkt_obp = _data[20:53] if len(_data) >= 53 else b""
if _call_type == "group" or _call_type == "vcsbk":
if not self._obp_stream_accept(_stream_id, _dst_id, _seq, _data, _frame_type, _dtype_vseq, dmrpkt_obp, _rf_src, _peer_id):
self._config["_bcka"] = time.time()
return
# Legacy routerOBP: STATUS[stream_id] with LC (bridge.py 259-278)
if _stream_id not in self.STATUS:
self.STATUS[_stream_id] = {
"START": time.time(),
"CONTENTION": False,
"RFS": _rf_src,
"TGID": _dst_id,
"LAST": time.time(),
"lastSeq": False,
"lastData": False,
}
if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt_obp) >= 33:
try:
decoded_obp = decode.voice_head_term(dmrpkt_obp)
self.STATUS[_stream_id]["LC"] = decoded_obp["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
self.STATUS[_stream_id]["LAST"] = time.time()
# Also update slot STATUS for contention/RX tracking
if _slot in self.STATUS:
self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TIME"] = time.time()
if self._dmrd_received: if self._dmrd_received:
# Legacy hblink DMRD v1: SERVER_ID + default rptr/hops/ber/rssi (`hblink.py` ~338–345, ~416) # Legacy hblink DMRD v1: SERVER_ID + default rptr/hops/ber/rssi (`hblink.py` ~338–345, ~416)
_global = self._CONFIG.get("GLOBAL", {}) _global = self._CONFIG.get("GLOBAL", {})
@ -1298,7 +1153,13 @@ class HBPProtocol(DatagramProtocol):
return return
_inthops = (_hops if isinstance(_hops, int) else int.from_bytes(_hops, "big")) + 1 _inthops = (_hops if isinstance(_hops, int) else int.from_bytes(_hops, "big")) + 1
if _inthops > 10: if _inthops > 10:
logger.warning("(%s) MAX HOPS exceed, dropping. Hops: %s, DST: %s, SRC: %s", self._system, _inthops, _int_dst_id, _src_srv_int) logger.debug(
"(%s) MAX HOPS exceed, dropping. Hops: %s, DST: %s, SRC: %s",
self._system,
_inthops,
_int_dst_id,
_src_srv_int,
)
self._obp_send_bcsq(_dst_id, _stream_id) self._obp_send_bcsq(_dst_id, _stream_id)
return return
if _call_type != "unit": if _call_type != "unit":
@ -1356,36 +1217,6 @@ class HBPProtocol(DatagramProtocol):
return return
if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD: if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD:
logger.info("(%s) CALL RX (OBP DMRE) src %s -> TG %s slot %s", self._system, int_id(_rf_src), int_id(_dst_id), _slot) logger.info("(%s) CALL RX (OBP DMRE) src %s -> TG %s slot %s", self._system, int_id(_rf_src), int_id(_dst_id), _slot)
dmrpkt_obp = _data[20:53] if len(_data) >= 53 else b""
if _call_type == "group" or _call_type == "vcsbk":
if not self._obp_stream_accept(_stream_id, _dst_id, _seq, _data, _frame_type, _dtype_vseq, dmrpkt_obp, _rf_src, _peer_id):
self._config["_bcka"] = time.time()
return
# Legacy routerOBP: STATUS[stream_id] with LC (same as DMRD v1 path)
if _stream_id not in self.STATUS:
self.STATUS[_stream_id] = {
"START": time.time(),
"CONTENTION": False,
"RFS": _rf_src,
"TGID": _dst_id,
"LAST": time.time(),
"lastSeq": False,
"lastData": False,
}
if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt_obp) >= 33:
try:
decoded_obp = decode.voice_head_term(dmrpkt_obp)
self.STATUS[_stream_id]["LC"] = decoded_obp["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
self.STATUS[_stream_id]["LAST"] = time.time()
if _slot in self.STATUS:
self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TIME"] = time.time()
_data_dmrd = DMRD + _data[4:] _data_dmrd = DMRD + _data[4:]
_hops_out = _inthops.to_bytes(1, "big") _hops_out = _inthops.to_bytes(1, "big")
if self._dmrd_received: if self._dmrd_received:

Loading…
Cancel
Save

Powered by TurnKey Linux.