fix: legacy parity — HBP protocol, bridge routing, voice and main wiring

HBP protocol (udp_hbp.py):
- Add dereg(), proxy_IPBlackList/proxy_bad_peer, PRBL import
- Add errback on all LoopingCalls (MASTER/PEER/OBP)
- Fix TG 4000 deactivation order (after REPEAT/ACL, with return)
- Gate recording on dmrd_received acceptance
- Fix playFileOnRequest trigger (VTERM+unit, not VHEAD)
- Add PEER stream state + RX_LC setup before dmrd_received
- Add proxy blacklist on failed login/callsign mismatch

Bridge routing (bridge_use_cases.py):
- Port full HBP ingress controls (collision, counters, dedup, loss)
- Add _reset gate, contention flag, OBP unit-data loop control
- Fix _send_data_to_obp _source_server for DMRE unit data
- Add OBP two-stage stream trimmer (5s _to / 180s remove)
- Add private voice START,TX reports and 7-digit gate

Voice (voice_use_cases.py):
- Fix play_file_on_request peer field to bytes_4(9)
- Add missing sleep(1) after pkt_gen

Config (config_normalizer.py):
- Add TG2_ACL=PERMIT:ALL default for OPENBRIDGE

Main (main.py):
- Add graceful deregister on shutdown (MSTCL/RPTCL)
- Add reactor.suggestThreadPoolSize(100)
- Fix looping errback to stop reactor on error
- Add D-APRS/SC alias defaults
- Gate bridgeDebug loop on DEBUG_BRIDGES config
pull/4/head
Rodrigo Pérez 6 months ago
parent bb771a744e
commit 8e6f03d680

@ -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 <FROM> SUB: %s PEER: %s <TO> 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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.