From f446527e578c8bf548b0016c7b216a1c819d754a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 22 Jun 2026 13:22:59 -0400 Subject: [PATCH] feat: dynamic TG restore, OBP cross-slot downlink, and DB migration 005 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Export ua_multi_tgs in topology for monitor chip restore after server restart - Defer OBP→HBP slot contention on inject-only proxies; clear other-slot SINGLE session on new TX - Fix infinite UA timer persistence (expires_at NULL) and restore from DB (expires_at=0 sentinel) - Push CONFIG to monitor after async dynamic-TG restore on RPTC - Migration 005_peer_dynamic_tgs_expires_null (shared with adn-monitor) --- docs/en/monitor/configuration.md | 1 + docs/en/server/user-guide/configuration.md | 2 +- docs/es/monitor/configuration.md | 2 +- docs/es/server/user-guide/configuration.md | 2 +- schemas/examples/topology.json | 1 + schemas/report-v2.json | 8 ++ .../application/dynamic_tg_use_cases.py | 16 ++- src/adn_server/application/report/payloads.py | 7 +- src/adn_server/application/routing/helpers.py | 82 +++++++++++- .../application/routing_use_cases.py | 23 +++- .../persistence/dynamic_tg_repository.py | 20 ++- .../infrastructure/persistence/mysql_pool.py | 45 ++++++- .../twisted_adapters/udp_hbp.py | 18 ++- tests/application/test_dynamic_tg_restore.py | 26 ++++ .../test_obp_hbp_cross_slot_downlink.py | 121 ++++++++++++++++++ .../application/test_peer_single_downlink.py | 19 +++ .../test_dynamic_tg_repository.py | 8 ++ .../infrastructure/test_mysql_pool_ensure.py | 6 +- 18 files changed, 378 insertions(+), 29 deletions(-) create mode 100644 tests/application/test_obp_hbp_cross_slot_downlink.py diff --git a/docs/en/monitor/configuration.md b/docs/en/monitor/configuration.md index b170b73..fcf4401 100644 --- a/docs/en/monitor/configuration.md +++ b/docs/en/monitor/configuration.md @@ -109,6 +109,7 @@ Obsolete **`WEBSOCKET_SERVER`** YAML (Twisted on a separate port) is ignored; us | Migration | Table / change | |-----------|----------------| | **`004_peer_dynamic_tgs`** | **`peer_dynamic_tgs`** — per-peer dynamic TG rows written by **adn-server 2.0.0-rc.3+** (shared schema). | +| **`005_peer_dynamic_tgs_expires_null`** | Data fix: `expires_at = 0` → `NULL` for `single_mode = 1` (infinite UA timer rows). | **adn-server** also ensures **`peer_dynamic_tgs`** exists on startup (idempotent). Either path is sufficient; both can run against the same **`hbmon`** database. diff --git a/docs/en/server/user-guide/configuration.md b/docs/en/server/user-guide/configuration.md index 728b7c8..747a712 100644 --- a/docs/en/server/user-guide/configuration.md +++ b/docs/en/server/user-guide/configuration.md @@ -211,7 +211,7 @@ Do **not** run standalone **`adn-proxy`** on the same **`LISTEN_PORT`** when the **Uses one shared connection pool** for: -- **Dynamic TG persistence** — table **`peer_dynamic_tgs`** (per-peer user-activated TGs across hotspot reconnects). The server **creates the table on startup** if missing (migration id **`004_peer_dynamic_tgs`**, same schema as adn-monitor). +- **Dynamic TG persistence** — table **`peer_dynamic_tgs`** (per-peer user-activated TGs across hotspot reconnects). The server **creates the table on startup** if missing (migration id **`004_peer_dynamic_tgs`**, same schema as adn-monitor). Migration **`005_peer_dynamic_tgs_expires_null`** normalizes legacy infinite-timer rows (`expires_at = 0` → `NULL` for `single_mode = 1`). - **Integrated self-service** — table **`Clients`** when **`SELF_SERVICE.USE_SELFSERVICE: true`**. Startup aborts with a clear log if MariaDB is unreachable or **`DATABASE`** is incomplete. Install **`mysqlclient`** (`pip install -e ".[selfservice]"` includes it). diff --git a/docs/es/monitor/configuration.md b/docs/es/monitor/configuration.md index 55906ed..6232a78 100644 --- a/docs/es/monitor/configuration.md +++ b/docs/es/monitor/configuration.md @@ -120,7 +120,7 @@ No hay Alembic: el monitor usa **`schema_migrations`** y comprobaciones en **`in - **Replace:** carga en `{tabla}_import` con **commit cada 10 000 filas** (la tabla live sigue legible); swap atómico `RENAME TABLE` (bloqueo metadata breve). - **Merge** (ficheros locales): `INSERT IGNORE` con **commit cada 2 000 filas**. -Migraciones: `001_clients_callsign`, `002_clients_options_width`, `003_alias_pk_only`, **`004_peer_dynamic_tgs`** (tabla compartida con **adn-server 2.0.0-rc.3+**). +Migraciones: `001_clients_callsign`, `002_clients_options_width`, `003_alias_pk_only`, **`004_peer_dynamic_tgs`** (tabla compartida con **adn-server 2.0.0-rc.3+**), **`005_peer_dynamic_tgs_expires_null`** (datos: `expires_at=0` → `NULL` en filas `single_mode=1`). **adn-server** también asegura **`peer_dynamic_tgs`** al arrancar (idempotente). Cualquiera de los dos caminos basta; ambos pueden usar la misma base **`hbmon`**. diff --git a/docs/es/server/user-guide/configuration.md b/docs/es/server/user-guide/configuration.md index c1041df..ab03486 100644 --- a/docs/es/server/user-guide/configuration.md +++ b/docs/es/server/user-guide/configuration.md @@ -211,7 +211,7 @@ Se arranca siempre que exista un bloque **`PROXY`** (ver `adn-server.example.yam **Un solo pool** compartido para: -- **Persistencia de TG dinámicos** — tabla **`peer_dynamic_tgs`** (TG activados por usuario por hotspot entre reconexiones). El servidor **crea la tabla al arranque** si falta (migración **`004_peer_dynamic_tgs`**, mismo esquema que adn-monitor). +- **Persistencia de TG dinámicos** — tabla **`peer_dynamic_tgs`** (TG activados por usuario por hotspot entre reconexiones). El servidor **crea la tabla al arranque** si falta (migración **`004_peer_dynamic_tgs`**, mismo esquema que adn-monitor). La migración **`005_peer_dynamic_tgs_expires_null`** corrige filas legacy de timer infinito (`expires_at = 0` → `NULL` cuando `single_mode = 1`). - **Self-service integrado** — tabla **`Clients`** con **`SELF_SERVICE.USE_SELFSERVICE: true`**. El arranque aborta con log claro si MariaDB no responde o **`DATABASE`** está incompleto. Instala **`mysqlclient`** (`pip install -e ".[selfservice]"` lo incluye). diff --git a/schemas/examples/topology.json b/schemas/examples/topology.json index 2afc531..731a2b0 100644 --- a/schemas/examples/topology.json +++ b/schemas/examples/topology.json @@ -19,6 +19,7 @@ "single_mode": false, "ua_timer_min": 10, "ua_sessions": {}, + "ua_multi_tgs": {}, "rf_mode": "duplex" } ] diff --git a/schemas/report-v2.json b/schemas/report-v2.json index 0716d24..cdfc360 100644 --- a/schemas/report-v2.json +++ b/schemas/report-v2.json @@ -144,6 +144,14 @@ } } }, + "ua_multi_tgs": { + "type": "object", + "description": "Active SINGLE=0 multi-dynamic TG ids per slot (empty when none).", + "additionalProperties": { + "type": "array", + "items": { "$ref": "#/$defs/dmr_id" } + } + }, "rf_mode": { "type": "string", "enum": ["simplex", "duplex"], diff --git a/src/adn_server/application/dynamic_tg_use_cases.py b/src/adn_server/application/dynamic_tg_use_cases.py index 62213f4..fe44de4 100644 --- a/src/adn_server/application/dynamic_tg_use_cases.py +++ b/src/adn_server/application/dynamic_tg_use_cases.py @@ -38,10 +38,21 @@ from adn_server.application.routing.helpers import ( ) from adn_server.domain import int_id from adn_server.domain.dynamic_tg import DynamicTgEntry +from adn_server.domain.ua_timer import ua_session_never_expires logger = logging.getLogger(__name__) +def _dynamic_tg_entry_active(entry: DynamicTgEntry, now: float) -> bool: + if not entry.single_mode: + return True + if entry.expires_at is None: + return True + if ua_session_never_expires(float(entry.expires_at)): + return True + return float(entry.expires_at) > now + + class DynamicTgUseCases: def __init__( self, @@ -107,10 +118,7 @@ class DynamicTgUseCases: now = time.time() def _apply(entries: list[DynamicTgEntry]) -> list[int]: - active = [ - e for e in entries - if not e.single_mode or e.expires_at is None or float(e.expires_at) > now - ] + active = [e for e in entries if _dynamic_tg_entry_active(e, now)] tgids = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, active, now=now) if self._on_restored is not None: self._on_restored(peer_id, system_name, sys_cfg, active, now=now) diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index 75b9b69..5188430 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -28,7 +28,11 @@ from typing import Any from adn_server.domain.ua_timer import normalize_ua_timer_minutes, ua_timer_is_infinite -from adn_server.application.routing.helpers import export_peer_ua_sessions, peer_rf_mode +from adn_server.application.routing.helpers import ( + export_peer_ua_multi_tgs, + export_peer_ua_sessions, + peer_rf_mode, +) from adn_server.domain import int_id REPORT_PROTOCOL = 2 @@ -375,6 +379,7 @@ def _topology_peer_row( if not ua_timer_is_infinite(timer): row["ua_timer_min"] = timer row["ua_sessions"] = export_peer_ua_sessions(yaml_cfg, peer_key) + row["ua_multi_tgs"] = export_peer_ua_multi_tgs(yaml_cfg, peer_key) row["rf_mode"] = peer_rf_mode(peer) return row diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index fbcafb4..1cc40f3 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -186,6 +186,29 @@ def slot_status_peer_owner(slot_st: dict[str, Any]) -> bytes | None: return None +def inject_only_defer_obp_hbp_slot_contention( + config: dict[str, Any], + target_system: str, + target_system_cfg: dict[str, Any], + *, + source_is_obp: bool, +) -> bool: + """Whether OBP→MASTER ``to_target`` should skip global slot STATUS contention. + + Inject-only proxies defer slot checks to ``send_peer`` (same as REPEAT): + per-peer ``hbp_slot_blocks_group_voice_for_peer`` + OPTIONS/UA slot remap. + Global STATUS on the bridge wire TS would block cross-slot downlink while + another peer is active on that TS even though recipients listen on the other TS. + """ + if not source_is_obp: + return False + if target_system_cfg.get("MODE") != "MASTER": + return False + from ..proxy.deployment import is_proxy_inject_only + + return is_proxy_inject_only(config, target_system) + + def hbp_slot_blocks_group_voice_for_peer( slot_st: dict[str, Any], peer_id: bytes, @@ -579,6 +602,9 @@ def register_peer_ua_session( expires_at = UA_SESSION_NEVER_EXPIRES_AT else: expires_at = pkt_time + float(timer_min) * 60.0 + # One exclusive dynamic TG per hotspot (either RF slot); new local TX replaces all others. + other_slot = 2 if int(slot) == 1 else 1 + clear_peer_ua_sessions(peer, sys_cfg, peer_id, slot=other_slot) _write_peer_ua_session( peer, peer_id, @@ -668,6 +694,48 @@ def export_peer_ua_sessions( return out +def export_peer_ua_multi_tgs( + sys_cfg: dict[str, Any], + peer_id: bytes | int, +) -> dict[str, list[int]]: + """Active SINGLE=0 multi-dynamic TG sets for monitor snapshot.""" + pk = bytes_4(int_id(peer_id)) + store = sys_cfg.get("_PEER_UA_MULTI_TGS") + if not isinstance(store, dict): + return {} + per_peer = store.get(pk) + if not isinstance(per_peer, dict): + return {} + out: dict[str, list[int]] = {} + for slot in (1, 2): + tg_set = per_peer.get(slot) + if isinstance(tg_set, set) and tg_set: + tgids = sorted( + int(t) for t in tg_set if is_ua_session_tgid(int(t)) + ) + if tgids: + out[str(slot)] = tgids + return out + + +def sync_peer_ua_memory_from_store( + peer: dict[str, Any], + peer_id: bytes, + sys_cfg: dict[str, Any], +) -> None: + """Copy persisted UA rows from ``sys_cfg`` onto the live peer dict after DB restore.""" + pk = bytes_4(int_id(peer_id)) + store = sys_cfg.get("_PEER_UA_SESSIONS") + if isinstance(store, dict): + per_peer = store.get(pk) + if isinstance(per_peer, dict) and per_peer: + ua = peer.setdefault("_UA_SESSION", {}) + ua.clear() + for slot, entry in per_peer.items(): + if isinstance(entry, dict): + ua[slot] = dict(entry) + + def restore_peer_ua_entries_to_memory( sys_cfg: dict[str, Any], peer_id: bytes, @@ -685,13 +753,23 @@ def restore_peer_ua_entries_to_memory( continue slot = int(entry.slot) if entry.single_mode: + from adn_server.domain.ua_timer import UA_SESSION_NEVER_EXPIRES_AT, ua_session_never_expires + expires = entry.expires_at - if expires is not None and float(expires) <= pkt_time: + if ( + expires is not None + and not ua_session_never_expires(float(expires)) + and float(expires) <= pkt_time + ): continue per_peer = sys_cfg.setdefault("_PEER_UA_SESSIONS", {}).setdefault(pk, {}) + if expires is None or ua_session_never_expires(float(expires)): + exp_mem = UA_SESSION_NEVER_EXPIRES_AT + else: + exp_mem = float(expires) per_peer[slot] = { "tgid": tgid, - "expires": float(expires) if expires is not None else 0.0, + "expires": exp_mem, } restored.append(tgid) else: diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 223eaca..ace70dc 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -44,6 +44,7 @@ from .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAli from .talker_alias_use_cases import TalkerAliasUseCases from .routing.helpers import ( hbp_slot_blocks_group_voice, + inject_only_defer_obp_hbp_slot_contention, slot_has_active_voice, is_private_subscriber_dst, is_unit_data_ingress, @@ -616,12 +617,22 @@ class RoutingUseCases( _ts_st["TX_TYPE"] = HBPF_SLT_VTERM # Slot contention: active QSO blocks any other stream; post-VTERM uses GROUP_HANGTIME. _group_hangtime = float(_target_system.get("GROUP_HANGTIME", 0) or 0) - if not _closing_bridge_leg and hbp_slot_blocks_group_voice( - _ts_st, - entry_tgid_b, - stream_id, - pkt_time, - _group_hangtime, + _defer_slot_contention = inject_only_defer_obp_hbp_slot_contention( + self._config, + entry["SYSTEM"], + _target_system, + source_is_obp=source_is_obp, + ) + if ( + not _closing_bridge_leg + and not _defer_slot_contention + and hbp_slot_blocks_group_voice( + _ts_st, + entry_tgid_b, + stream_id, + pkt_time, + _group_hangtime, + ) ): if not _src_stream_st.get("CONTENTION"): _src_stream_st["CONTENTION"] = True diff --git a/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py b/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py index 3e9b58f..4677317 100644 --- a/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py +++ b/src/adn_server/infrastructure/persistence/dynamic_tg_repository.py @@ -35,6 +35,22 @@ from adn_server.domain.dynamic_tg import DynamicTgEntry logger = logging.getLogger(__name__) +def _expires_at_db_value(expires_at: float | None) -> int | None: + """Map runtime expiry to signed INT column (NULL = no wall-clock expiry).""" + from adn_server.domain.ua_timer import ua_session_never_expires + + if expires_at is None: + return None + if ua_session_never_expires(float(expires_at)): + return None + exp = int(expires_at) + if exp <= 0: + return None + if exp > 2147483647: + return 2147483647 + return exp + + def _row_to_entry(row: tuple[Any, ...]) -> DynamicTgEntry: return DynamicTgEntry( int_id=int(row[0]), @@ -55,7 +71,9 @@ class MysqlDynamicTgRepository(DynamicTgStore): def upsert(self, entry: DynamicTgEntry) -> None: updated = int(entry.updated_at or time.time()) - expires = int(entry.expires_at) if entry.expires_at is not None else None + expires = _expires_at_db_value( + float(entry.expires_at) if entry.expires_at is not None else None + ) self._pool.runOperation( """INSERT INTO peer_dynamic_tgs (int_id, system_name, slot, tgid, single_mode, expires_at, updated_at) diff --git a/src/adn_server/infrastructure/persistence/mysql_pool.py b/src/adn_server/infrastructure/persistence/mysql_pool.py index e75e2a9..032910f 100644 --- a/src/adn_server/infrastructure/persistence/mysql_pool.py +++ b/src/adn_server/infrastructure/persistence/mysql_pool.py @@ -33,6 +33,7 @@ from adn_server.infrastructure.persistence.database_config import validate_datab logger = logging.getLogger(__name__) _PEER_DYNAMIC_TGS_MIGRATION = "004_peer_dynamic_tgs" +_PEER_DYNAMIC_TGS_EXPIRES_NULL_MIGRATION = "005_peer_dynamic_tgs_expires_null" _MYSQL_HINTS: dict[int, str] = { 1045: "check DATABASE.DB_USERNAME and DB_PASSWORD in adn-server.yaml", @@ -61,6 +62,10 @@ _CREATE_PEER_DYNAMIC_TGS = """CREATE TABLE IF NOT EXISTS peer_dynamic_tgs ( KEY idx_expires (expires_at) ) DEFAULT CHARSET=utf8mb4""" +_UPDATE_PEER_DYNAMIC_EXPIRES_NULL = """UPDATE peer_dynamic_tgs +SET expires_at = NULL +WHERE single_mode = 1 AND expires_at = 0""" + def describe_mysql_error(err: Exception) -> str: """Turn a MySQL/MariaDB exception into an actionable startup message.""" @@ -114,14 +119,40 @@ def _mark_migration(cursor: Any, migration_id: str) -> None: ) -def _ensure_peer_dynamic_tgs_on_cursor(cursor: Any) -> None: - """Server-owned table; same migration id/DLL as adn-monitor ``004_peer_dynamic_tgs``.""" +def _table_exists(cursor: Any, table: str) -> bool: + cursor.execute( + "SELECT 1 FROM information_schema.TABLES " + "WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = %s", + (table,), + ) + return cursor.fetchone() is not None + + +def apply_migrations_on_cursor(cursor: Any) -> None: + """Incremental migrations shared with adn-monitor ``schema_migrations`` ids.""" cursor.execute(_CREATE_SCHEMA_MIGRATIONS) - if _migration_applied(cursor, _PEER_DYNAMIC_TGS_MIGRATION): - return - cursor.execute(_CREATE_PEER_DYNAMIC_TGS) - _mark_migration(cursor, _PEER_DYNAMIC_TGS_MIGRATION) - logger.info("(DATABASE) applied migration %s (peer_dynamic_tgs)", _PEER_DYNAMIC_TGS_MIGRATION) + + if not _migration_applied(cursor, _PEER_DYNAMIC_TGS_MIGRATION): + cursor.execute(_CREATE_PEER_DYNAMIC_TGS) + _mark_migration(cursor, _PEER_DYNAMIC_TGS_MIGRATION) + logger.info( + "(DATABASE) applied migration %s (peer_dynamic_tgs)", + _PEER_DYNAMIC_TGS_MIGRATION, + ) + + if not _migration_applied(cursor, _PEER_DYNAMIC_TGS_EXPIRES_NULL_MIGRATION): + if _table_exists(cursor, "peer_dynamic_tgs"): + cursor.execute(_UPDATE_PEER_DYNAMIC_EXPIRES_NULL) + _mark_migration(cursor, _PEER_DYNAMIC_TGS_EXPIRES_NULL_MIGRATION) + logger.info( + "(DATABASE) applied migration %s (peer_dynamic_tgs expires_at cleanup)", + _PEER_DYNAMIC_TGS_EXPIRES_NULL_MIGRATION, + ) + + +def _ensure_peer_dynamic_tgs_on_cursor(cursor: Any) -> None: + """Server-owned table; same migration ids/DLL as adn-monitor.""" + apply_migrations_on_cursor(cursor) def ensure_database_sync( diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index dee18fb..c9349e0 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -53,6 +53,7 @@ from ...application.routing.helpers import ( peer_should_receive_group_voice, peer_single_exclusive_tgid, register_peer_ua_session, + sync_peer_ua_memory_from_store, remap_dmrd_to_peer_static_slot, resolve_voice_peer_id, repeat_downlink_report_slot, @@ -1418,11 +1419,22 @@ class HBPProtocol(DatagramProtocol): self._push_config_to_monitor() if self._dynamic_tg_uc is not None: _sys_cfg = self._CONFIG.get("SYSTEMS", {}).get(self._system, {}) + _peer_for_restore = _peer_id + + def _after_dynamic_tg_restore(tgids: list[int]) -> list[int]: + if tgids: + peer = self._peers.get(_peer_for_restore) + if peer is not None: + sync_peer_ua_memory_from_store( + peer, _peer_for_restore, _sys_cfg, + ) + self._mark_downlink_index_dirty() + self._push_config_to_monitor() + return tgids + self._dynamic_tg_uc.restore_peer( _peer_id, self._system, _sys_cfg, - ).addCallback( - lambda tgids: self._mark_downlink_index_dirty() if tgids else tgids - ) + ).addCallback(_after_dynamic_tg_restore) else: self.transport.write(b"".join([MSTNAK, _peer_id]), _sockaddr) logger.info("(%s) Peer info from Radio ID that has not logged in: %s", self._system, int_id(_peer_id)) diff --git a/tests/application/test_dynamic_tg_restore.py b/tests/application/test_dynamic_tg_restore.py index eafa663..51a49b9 100644 --- a/tests/application/test_dynamic_tg_restore.py +++ b/tests/application/test_dynamic_tg_restore.py @@ -128,6 +128,32 @@ def test_restore_single_preserves_absolute_expires_timestamp() -> None: assert per_peer[2]["expires"] - reconnect_at == pytest.approx(26.0 * 60.0, rel=0.01) +def test_restore_single_infinite_expires_at_zero() -> None: + """TIMER=0 / infinite SINGLE sessions use expires_at=0 and must restore from DB.""" + from adn_server.domain.ua_timer import UA_SESSION_NEVER_EXPIRES_AT + + peer_id = _peer_id() + sys_cfg = _sys_cfg() + now = 1_000_000.0 + entries = [ + DynamicTgEntry( + int_id=730039101, + system_name="MASTER-A", + slot=1, + tgid=7144, + single_mode=True, + expires_at=UA_SESSION_NEVER_EXPIRES_AT, + updated_at=now, + ), + ] + restored = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, entries, now=now + 10) + assert restored == [7144] + peer = {"OPTIONS": b"TS2=714,71442;SINGLE=1;TIMER=0;"} + assert peer_should_receive_group_voice( + peer, 1, 7144, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 20, + ) + + def test_purge_expired_peer_ua_sessions() -> None: peer_id = _peer_id() sys_cfg = _sys_cfg() diff --git a/tests/application/test_obp_hbp_cross_slot_downlink.py b/tests/application/test_obp_hbp_cross_slot_downlink.py new file mode 100644 index 0000000..8054eec --- /dev/null +++ b/tests/application/test_obp_hbp_cross_slot_downlink.py @@ -0,0 +1,121 @@ +# ADN DMR Peer Server - OBP → HBP cross-slot downlink +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY + +from __future__ import annotations + +from adn_server.application.routing.helpers import inject_only_defer_obp_hbp_slot_contention +from adn_server.domain import HBPF_SLT_VHEAD, bytes_3, bytes_4 +from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time +from tests.harness.scenarios import obp_bridge_scenario + + +def _two_peer_master(scenario: DeterministicScenario) -> None: + master = scenario.config["SYSTEMS"]["MASTER-A"] + master["TS2_STATIC"] = "730444,7144" + master["PEERS"] = { + bytes_4(730001): {"CONNECTION": "YES", "OPTIONS": b"TS2=7144;"}, + bytes_4(730002): {"CONNECTION": "YES", "OPTIONS": b"TS1=730444;"}, + } + + +def test_defer_helper_requires_inject_only_obp_to_master() -> None: + cfg = {"PROXY": {"TARGET_SYSTEM": "MASTER-A"}, "SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}} + sys_cfg = {"MODE": "MASTER", "PEERS": {bytes_4(1): {"CONNECTION": "YES"}}} + assert inject_only_defer_obp_hbp_slot_contention( + cfg, "MASTER-A", sys_cfg, source_is_obp=True, + ) + assert not inject_only_defer_obp_hbp_slot_contention( + cfg, "MASTER-A", sys_cfg, source_is_obp=False, + ) + assert not inject_only_defer_obp_hbp_slot_contention( + {}, "MASTER-A", sys_cfg, source_is_obp=True, + ) + + +def _seed_busy_slot_2(scenario: DeterministicScenario, peer_id: int) -> None: + proto = scenario.protocols["MASTER-A"] + t = scenario.clock.time() + proto.STATUS[2] = { + "RX_STREAM_ID": bytes_4(0xAAAAAAAA), + "RX_PEER": bytes_4(peer_id), + "RX_TGID": bytes_3(7144), + "RX_TIME": t, + "RX_TYPE": HBPF_SLT_VHEAD, + "TX_STREAM_ID": bytes_4(0xAAAAAAAA), + "TX_TGID": bytes_3(7144), + "TX_TIME": t, + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_PEER": bytes_4(peer_id), + } + + +def test_obp_hbp_forwards_when_bridge_ts_busy_but_peer_listens_other_ts() -> None: + """OBP bridge leg TS2 must not block downlink to a TS1-static peer (inject-only).""" + scenario = obp_bridge_scenario("OBP-CL", tg=730444) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + _two_peer_master(scenario) + _seed_busy_slot_2(scenario, 730001) + base = PacketSpec( + peer_id=73010, + rf_src=3340062, + dst_id=730444, + slot=1, + stream_id=0xBBBBBBBB, + ) + with patch_routing_wall_time(scenario.clock): + scenario.inject_obp("OBP-CL", DeterministicScenario.voice_head_spec(base)) + assert len(scenario.capture.for_system("MASTER-A")) > 0 + + +def test_obp_hbp_forwards_dynamic_ua_on_slot_1_when_bridge_ts2_busy() -> None: + """Dynamic TG keyed on TS1 must receive OBP downlink even when TS2 is occupied.""" + scenario = obp_bridge_scenario("OBP-CL", tg=730444) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + master = scenario.config["SYSTEMS"]["MASTER-A"] + master["TS2_STATIC"] = "730444,7144" + p_busy = bytes_4(730001) + p_dyn = bytes_4(730002) + master["PEERS"] = { + p_busy: {"CONNECTION": "YES", "OPTIONS": b"TS2=7144;"}, + p_dyn: { + "CONNECTION": "YES", + "OPTIONS": b"TS2=730;SINGLE=1;TIMER=300;", + "_UA_SESSION": { + 1: {"tgid": 730444, "expires": scenario.clock.time() + 300.0}, + }, + }, + } + _seed_busy_slot_2(scenario, 730001) + base = PacketSpec( + peer_id=73010, + rf_src=3340062, + dst_id=730444, + slot=1, + stream_id=0xDDDDDDDD, + ) + with patch_routing_wall_time(scenario.clock): + scenario.inject_obp("OBP-CL", DeterministicScenario.voice_head_spec(base)) + assert len(scenario.capture.for_system("MASTER-A")) > 0 + + +def test_obp_hbp_single_peer_inject_only_defers_like_repeat() -> None: + """Sole inject-only hotspot: OBP uses per-peer slot checks (REPEAT parity).""" + scenario = obp_bridge_scenario("OBP-CL", tg=730444) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + master = scenario.config["SYSTEMS"]["MASTER-A"] + master["TS2_STATIC"] = "730444,7144" + master["PEERS"] = { + bytes_4(730002): {"CONNECTION": "YES", "OPTIONS": b"TS1=730444;"}, + } + _seed_busy_slot_2(scenario, 730002) + base = PacketSpec( + peer_id=73010, + rf_src=3340062, + dst_id=730444, + slot=1, + stream_id=0xCCCCCCCC, + ) + with patch_routing_wall_time(scenario.clock): + scenario.inject_obp("OBP-CL", DeterministicScenario.voice_head_spec(base)) + assert len(scenario.capture.for_system("MASTER-A")) > 0 diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index 423c0ea..d212305 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -498,3 +498,22 @@ def test_new_tx_replaces_single_session_tg() -> None: assert not peer_should_receive_group_voice( peer, 2, 7305, peer_id=peer_id, connected_count=8, sys_cfg=sys_cfg, now=now + 121 ) + + +def test_single_tx_on_other_slot_replaces_indigo_and_unblocks_downlink() -> None: + """Dual-slot hotspot: TX 7144 TS1 must clear stale 730444 session on TS2 (SINGLE=1).""" + peer = {"OPTIONS": b"TS2=714,71442;SINGLE=1;TIMER=300;"} + sys_cfg = _sys_cfg() + peer_id = _peer_id() + now = 1_000_000.0 + register_peer_ua_session(peer, peer_id, 2, 730444, sys_cfg, now=now) + assert not peer_should_receive_group_voice( + peer, 1, 7144, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 10, + ) + register_peer_ua_session(peer, peer_id, 1, 7144, sys_cfg, now=now + 20) + assert peer_should_receive_group_voice( + peer, 1, 7144, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 30, + ) + assert not peer_should_receive_group_voice( + peer, 2, 730444, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 30, + ) diff --git a/tests/infrastructure/test_dynamic_tg_repository.py b/tests/infrastructure/test_dynamic_tg_repository.py index 98c4766..d4aafb4 100644 --- a/tests/infrastructure/test_dynamic_tg_repository.py +++ b/tests/infrastructure/test_dynamic_tg_repository.py @@ -61,6 +61,14 @@ def test_upsert_issues_insert(repo: MysqlDynamicTgRepository) -> None: assert args[3] == 7305 +def test_upsert_infinite_expires_stored_as_null(repo: MysqlDynamicTgRepository) -> None: + from adn_server.domain.ua_timer import UA_SESSION_NEVER_EXPIRES_AT + + repo.upsert(_entry(expires_at=UA_SESSION_NEVER_EXPIRES_AT)) + _sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001 + assert args[5] is None + + def test_replace_single_slot_deletes_then_upserts(repo: MysqlDynamicTgRepository) -> None: repo.replace_single_slot(_entry()) delete_sql = repo._pool.runOperation.call_args_list[0][0][0] # noqa: SLF001 diff --git a/tests/infrastructure/test_mysql_pool_ensure.py b/tests/infrastructure/test_mysql_pool_ensure.py index 106b26a..5df0e5c 100644 --- a/tests/infrastructure/test_mysql_pool_ensure.py +++ b/tests/infrastructure/test_mysql_pool_ensure.py @@ -44,18 +44,20 @@ def test_describe_mysql_error_unknown_database() -> None: def test_ensure_peer_dynamic_tgs_applies_migration_once() -> None: cur = MagicMock() - cur.fetchone.side_effect = [None, (1,)] # migration missing, then exists on re-check path + cur.fetchone.side_effect = [None, None, (1,)] # 004 missing, 005 missing, table exists _ensure_peer_dynamic_tgs_on_cursor(cur) executed = [call[0][0] for call in cur.execute.call_args_list] assert any("schema_migrations" in sql for sql in executed) assert any("peer_dynamic_tgs" in sql for sql in executed) + assert any("UPDATE peer_dynamic_tgs" in sql for sql in executed) assert any("INSERT IGNORE INTO schema_migrations" in sql for sql in executed) -def test_ensure_peer_dynamic_tgs_skips_when_migration_marked() -> None: +def test_ensure_peer_dynamic_tgs_skips_when_migrations_marked() -> None: cur = MagicMock() cur.fetchone.return_value = (1,) _ensure_peer_dynamic_tgs_on_cursor(cur) executed = [call[0][0] for call in cur.execute.call_args_list] assert any("schema_migrations" in sql for sql in executed) assert not any("CREATE TABLE IF NOT EXISTS peer_dynamic_tgs" in sql for sql in executed) + assert not any("UPDATE peer_dynamic_tgs" in sql for sql in executed)