diff --git a/src/adn_server/application/dynamic_tg_use_cases.py b/src/adn_server/application/dynamic_tg_use_cases.py index f3023dd..324aac4 100644 --- a/src/adn_server/application/dynamic_tg_use_cases.py +++ b/src/adn_server/application/dynamic_tg_use_cases.py @@ -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() diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index 59ba14f..416a82d 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -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]``.""" diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 3c6d320..0f87294 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -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 diff --git a/src/adn_server/infrastructure/config_validator.py b/src/adn_server/infrastructure/config_validator.py index f2280f3..4c2fb7d 100644 --- a/src/adn_server/infrastructure/config_validator.py +++ b/src/adn_server/infrastructure/config_validator.py @@ -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:" diff --git a/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py b/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py index 0475bf9..3e9b58f 100644 --- a/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py +++ b/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py @@ -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( diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 57f982b..4b2753f 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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: diff --git a/tests/application/test_dynamic_tg_persist.py b/tests/application/test_dynamic_tg_persist.py index 3ba0c85..09e191e 100644 --- a/tests/application/test_dynamic_tg_persist.py +++ b/tests/application/test_dynamic_tg_persist.py @@ -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 [] diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index 0764a16..7de9885 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -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() diff --git a/tests/conftest.py b/tests/conftest.py index 2a63812..bbd59d0 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -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", diff --git a/tests/infrastructure/test_dynamic_tg_repository.py b/tests/infrastructure/test_dynamic_tg_repository.py index 708bbc8..98c4766 100644 --- a/tests/infrastructure/test_dynamic_tg_repository.py +++ b/tests/infrastructure/test_dynamic_tg_repository.py @@ -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