fix: complete TG 4000 reset for monitor and dynamic TG persistence

Emit INGRESS BRDG_EVENT so SINGLE=0 UA chips clear without stuck TX;
wipe all peer dynamic rows from memory and MariaDB on reset. Never store
TG 4000 as a UA session. Require DATABASE only for full peer-server configs.
pull/3/head
Rodrigo Pérez 3 months ago
parent ddf5a91262
commit 5d57a548b2

@ -30,6 +30,7 @@ from typing import Any
from adn_server.application.ports import DynamicTgStore
from adn_server.application.routing.helpers import (
_peer_ua_session_entry,
is_ua_session_tgid,
peer_receives_group_tgid,
peer_single_mode,
purge_expired_peer_ua_sessions,
@ -63,7 +64,7 @@ class DynamicTgUseCases:
system_name: str,
) -> None:
"""Enqueue DB write (async); memory already updated by register_peer_ua_session."""
if int(tgid) <= 0 or int(tgid) == 4000 or peer_receives_group_tgid(peer, slot, int(tgid)):
if not is_ua_session_tgid(int(tgid)) or peer_receives_group_tgid(peer, slot, int(tgid)):
return
now = time.time()
peer_int = int_id(peer_id)
@ -98,6 +99,9 @@ class DynamicTgUseCases:
def delete_peer_slot(self, peer_id: bytes, system_name: str, slot: int) -> None:
self._store.delete_peer_slot(int_id(peer_id), system_name, int(slot))
def delete_peer(self, peer_id: bytes, system_name: str) -> None:
self._store.delete_peer(int_id(peer_id), system_name)
def restore_peer(self, peer_id: bytes, system_name: str, sys_cfg: dict[str, Any]) -> Any:
"""Load from persistence on reconnect (RPTC); returns async handle from port."""
now = time.time()

@ -382,6 +382,10 @@ class DynamicTgStore(ABC):
def delete_peer_slot(self, int_id: int, system_name: str, slot: int) -> None:
...
@abstractmethod
def delete_peer(self, int_id: int, system_name: str) -> None:
...
@abstractmethod
def load_peer(self, int_id: int, system_name: str) -> Any:
"""Returns Deferred firing with ``list[DynamicTgEntry]``."""

@ -103,6 +103,11 @@ def tg4000_reset_on_vhead(int_dst_id: int, frame_type: int, dtype_vseq: int) ->
)
def is_ua_session_tgid(tgid: int) -> bool:
"""True when a keyed TG may be stored as a user-activated dynamic session."""
return int(tgid) > 0 and int(tgid) != 4000
def obp_target_bcsq_quenches_stream(
systems_cfg: dict[str, Any], target_name: str, dst_id_b: bytes, stream_id: bytes
) -> bool:
@ -309,7 +314,7 @@ def register_peer_ua_multi_tg(
if peer_single_mode(peer, sys_cfg):
return
tgid_i = int(tgid)
if tgid_i <= 0 or tgid_i == 4000:
if not is_ua_session_tgid(tgid_i):
return
if peer_receives_group_tgid(peer, slot, tgid_i):
return
@ -355,6 +360,8 @@ def register_peer_ua_session(
now: float | None = None,
) -> None:
"""Track UA TG for this hotspot (SINGLE=1 exclusive; SINGLE=0 multi-dynamic set)."""
if not is_ua_session_tgid(tgid):
return
if not peer_single_mode(peer, sys_cfg):
register_peer_ua_multi_tg(peer, peer_id, slot, tgid, sys_cfg)
return
@ -391,7 +398,7 @@ def seed_peer_ua_session_from_status(
if bytes_4(int_id(peer_id)) != bytes_4(int_id(rx_peer)):
return
rx_tgid = int_id(status_slot.get("RX_TGID", b"\x00\x00\x00"))
if rx_tgid <= 0:
if not is_ua_session_tgid(rx_tgid):
return
connected_at = float(peer.get("CONNECTED", 0) or 0)
rx_time = float(status_slot.get("RX_TIME", 0) or 0)
@ -443,7 +450,7 @@ def export_peer_ua_sessions(
continue
exp = float(entry.get("expires", 0) or 0)
tgid = int(entry.get("tgid", 0) or 0)
if tgid > 0 and exp > pkt_time:
if is_ua_session_tgid(tgid) and exp > pkt_time:
out[str(slot)] = {"tgid": tgid, "expires_at": exp}
return out
@ -461,6 +468,8 @@ def restore_peer_ua_entries_to_memory(
restored: list[int] = []
for entry in entries:
tgid = int(entry.tgid)
if not is_ua_session_tgid(tgid):
continue
slot = int(entry.slot)
if entry.single_mode:
expires = entry.expires_at

@ -262,6 +262,23 @@ def _validate_proxy(proxy_cfg: dict[str, Any] | None, systems: dict[str, Any], e
)
def _config_requires_database(config: dict[str, Any]) -> bool:
"""True for ``run_peer_server`` configs (proxy/master); not echo-only PEER fleets."""
proxy = config.get("PROXY")
if isinstance(proxy, dict) and proxy:
return True
systems = config.get("SYSTEMS")
if not isinstance(systems, dict):
return False
for sys_cfg in systems.values():
if not isinstance(sys_cfg, dict):
continue
mode = str(sys_cfg.get("MODE", "")).upper()
if mode in ("MASTER", "OPENBRIDGE"):
return True
return False
def _validate_database(db_cfg: Any, errors: list[str]) -> None:
if db_cfg is None:
errors.append("DATABASE: required block missing in adn-server.yaml")
@ -335,7 +352,8 @@ def validate_config(config: dict[str, Any], *, config_path: str | None = None) -
proxy_cfg = config.get("PROXY")
_validate_proxy(proxy_cfg if isinstance(proxy_cfg, dict) else None, systems if isinstance(systems, dict) else {}, errors)
_validate_database(config.get("DATABASE"), errors)
if _config_requires_database(config):
_validate_database(config.get("DATABASE"), errors)
if errors:
header = f"Configuration error in {config_path}:" if config_path else "Configuration error:"

@ -91,6 +91,14 @@ class MysqlDynamicTgRepository(DynamicTgStore):
lambda f: logger.error("(DYNAMIC_TG) delete_peer_slot: %s", f.getTraceback())
)
def delete_peer(self, int_id: int, system_name: str) -> None:
self._pool.runOperation(
"DELETE FROM peer_dynamic_tgs WHERE int_id=%s AND system_name=%s",
(int_id, system_name),
).addErrback(
lambda f: logger.error("(DYNAMIC_TG) delete_peer: %s", f.getTraceback())
)
@inlineCallbacks
def load_peer(self, int_id: int, system_name: str) -> Any:
rows = yield self._pool.runQuery(

@ -50,6 +50,7 @@ from ...application.routing.helpers import (
peer_should_receive_group_voice,
peer_single_exclusive_tgid,
register_peer_ua_session,
resolve_voice_peer_id,
seed_peer_ua_session_from_status,
tg4000_reset_on_vhead,
)
@ -399,13 +400,50 @@ class HBPProtocol(DatagramProtocol):
session_codec=self._mesh_session_codec(),
)
def _apply_tg4000_reset(self, peer_id: bytes, slot: int, call_type: str) -> None:
def _emit_tg4000_routing_event(
self,
peer_id: bytes,
rf_src: bytes,
stream_id: bytes,
slot: int,
) -> None:
"""BRDG_EVENT for monitor SINGLE=0 clear (skipped when dmrd_received returns early).
INGRESS (not START): monitor clears UA sessions on dest 4000 without lighting TRX chips.
"""
report = self._report
if report is None or not self._CONFIG.get("REPORTS", {}).get("REPORT", True):
return
if not hasattr(report, "send_routing_event"):
return
systems_cfg = self._CONFIG.get("SYSTEMS", {})
report_peer = resolve_voice_peer_id(peer_id, rf_src, self._system, systems_cfg)
report.send_routing_event(
"GROUP VOICE,INGRESS,RX,{},{},{},{},{},4000".format(
self._system,
int_id(stream_id),
int_id(report_peer),
int_id(rf_src),
slot,
)
)
def _apply_tg4000_reset(
self,
peer_id: bytes,
slot: int,
call_type: str,
*,
rf_src: bytes,
stream_id: bytes,
) -> None:
"""Clear per-peer UA dynamics and stale STATUS; deactivate bridges on 4000."""
peer = self._peers.get(peer_id, {})
clear_peer_ua_sessions(peer, self._config, peer_id, slot=slot)
clear_peer_rx_status_slots(self.STATUS, peer_id, slot=slot)
clear_peer_ua_sessions(peer, self._config, peer_id)
clear_peer_rx_status_slots(self.STATUS, peer_id)
if self._dynamic_tg_uc is not None:
self._dynamic_tg_uc.delete_peer_slot(peer_id, self._system, slot)
self._dynamic_tg_uc.delete_peer(peer_id, self._system)
self._emit_tg4000_routing_event(peer_id, rf_src, stream_id, slot)
if self._on_in_band_signalling:
self._on_in_band_signalling(self._system, slot, bytes_3(4000), time.time())
_kind = "Private call to ID" if call_type == "unit" else "Group call to TG"
@ -432,12 +470,17 @@ class HBPProtocol(DatagramProtocol):
call_type: str,
frame_type: int,
dtype_vseq: int,
*,
rf_src: bytes,
stream_id: bytes,
) -> bool:
"""Legacy early return for TG/ID 4000; reset once per PTT on voice header only."""
if int_dst_id != 4000:
return False
if tg4000_reset_on_vhead(int_dst_id, frame_type, dtype_vseq):
self._apply_tg4000_reset(peer_id, slot, call_type)
self._apply_tg4000_reset(
peer_id, slot, call_type, rf_src=rf_src, stream_id=stream_id,
)
return True
def _peer_should_receive_dmrd(self, peer_id: bytes, packet: bytes) -> bool:
@ -975,6 +1018,7 @@ class HBPProtocol(DatagramProtocol):
# TG 4000: reset after REPEAT so peers see the packet (legacy order)
if self._handle_tg4000_packet(
_peer_id, _slot, _int_dst_id, _call_type, _frame_type, _dtype_vseq,
rf_src=_rf_src, stream_id=_stream_id,
):
return
if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD:
@ -1375,6 +1419,7 @@ class HBPProtocol(DatagramProtocol):
# TG 4000: reset after ACL/SUB_MAP (legacy order — routerHBP.dmrd_received)
if self._handle_tg4000_packet(
_peer_id, _slot, _int_dst_id, _call_type, _frame_type, _dtype_vseq,
rf_src=_rf_src, stream_id=_stream_id,
):
return
if _call_type == "group" and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD:

@ -41,6 +41,9 @@ class _CaptureStore:
def delete_peer_slot(self, int_id: int, system_name: str, slot: int) -> None:
pass
def delete_peer(self, int_id: int, system_name: str) -> None:
pass
def load_peer(self, int_id: int, system_name: str):
return []

@ -34,6 +34,7 @@ from adn_server.application.routing.helpers import (
peer_should_receive_group_voice,
peer_single_blocks_group_voice,
peer_single_blocks_uplink,
peer_single_exclusive_tgid,
register_peer_ua_multi_tg,
register_peer_ua_session,
seed_peer_ua_session_from_status,
@ -372,6 +373,21 @@ def test_single_zero_tg4000_clears_multi_dynamic() -> None:
)
def test_register_peer_ua_session_ignores_4000() -> None:
"""TG 4000 resets dynamics; it must never be stored as a UA session."""
peer_id = _peer_id()
sys_cfg = _sys_cfg()
now = 1_000_000.0
peer_multi = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
register_peer_ua_session(peer_multi, peer_id, 2, 4000, sys_cfg, now=now)
pk = peer_id
assert sys_cfg.get("_PEER_UA_MULTI_TGS", {}).get(pk, {}).get(2, set()) == set()
peer_single = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
register_peer_ua_session(peer_single, peer_id, 2, 4000, sys_cfg, now=now)
assert peer_single_exclusive_tgid(peer_single, 2, sys_cfg, peer_id=peer_id, now=now) is None
def test_new_tx_replaces_single_session_tg() -> None:
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
sys_cfg = _sys_cfg()

@ -25,11 +25,23 @@ from __future__ import annotations
from typing import Any
def minimal_database_config() -> dict[str, Any]:
"""Minimal DATABASE block for tests (full peer server configs)."""
return {
"DB_SERVER": "localhost",
"DB_USERNAME": "hbmon",
"DB_PASSWORD": "test",
"DB_NAME": "hbmon",
"DB_PORT": 3306,
}
def minimal_valid_config(**overrides: Any) -> dict[str, Any]:
"""Minimal config dict that passes validate_config (integrated proxy required)."""
config: dict[str, Any] = {
"GLOBAL": {"SERVER_ID": 1},
"REPORTS": {"REPORT": False},
"DATABASE": minimal_database_config(),
"SYSTEMS": {
"HOTSPOT": {
"MODE": "MASTER",

@ -74,6 +74,14 @@ def test_delete_peer_slot(repo: MysqlDynamicTgRepository) -> None:
assert args == (730039101, "MASTER-A", 2)
def test_delete_peer(repo: MysqlDynamicTgRepository) -> None:
repo.delete_peer(730039101, "MASTER-A")
sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001
assert "DELETE FROM peer_dynamic_tgs" in sql
assert "slot" not in sql
assert args == (730039101, "MASTER-A")
def test_purge_expired(repo: MysqlDynamicTgRepository) -> None:
repo.purge_expired(1_000_000.0)
sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001

Loading…
Cancel
Save

Powered by TurnKey Linux.