Add DATABASE config, async DynamicTgStore, and restore on RPTC so SINGLE=0/1 dynamics survive hotspot disconnects and server restarts without blocking the DMRD voice path.pull/3/head
parent
e7e99a7903
commit
f7047620af
@ -0,0 +1,127 @@
|
||||
# ADN DMR Peer Server - application dynamic tg use cases
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Persist/restore per-peer dynamic TGs via ``DynamicTgStore`` port (voice path non-blocking)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from typing import Any
|
||||
|
||||
from adn_server.application.ports import DynamicTgStore
|
||||
from adn_server.application.routing.helpers import (
|
||||
_peer_ua_session_entry,
|
||||
peer_receives_group_tgid,
|
||||
peer_single_mode,
|
||||
purge_expired_peer_ua_sessions,
|
||||
restore_peer_ua_entries_to_memory,
|
||||
)
|
||||
from adn_server.domain import int_id
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class DynamicTgUseCases:
|
||||
def __init__(
|
||||
self,
|
||||
store: DynamicTgStore,
|
||||
*,
|
||||
on_restored: Callable[[bytes, str, dict[str, Any], list[DynamicTgEntry], float], None]
|
||||
| None = None,
|
||||
) -> None:
|
||||
self._store = store
|
||||
self._on_restored = on_restored
|
||||
|
||||
def persist_after_register(
|
||||
self,
|
||||
peer: dict[str, Any],
|
||||
peer_id: bytes,
|
||||
slot: int,
|
||||
tgid: int,
|
||||
sys_cfg: dict[str, Any],
|
||||
*,
|
||||
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)):
|
||||
return
|
||||
now = time.time()
|
||||
peer_int = int_id(peer_id)
|
||||
if peer_single_mode(peer, sys_cfg):
|
||||
session = _peer_ua_session_entry(sys_cfg, peer_id, int(slot))
|
||||
if not session:
|
||||
return
|
||||
self._store.replace_single_slot(
|
||||
DynamicTgEntry(
|
||||
int_id=peer_int,
|
||||
system_name=system_name,
|
||||
slot=int(slot),
|
||||
tgid=int(tgid),
|
||||
single_mode=True,
|
||||
expires_at=float(session["expires"]),
|
||||
updated_at=now,
|
||||
)
|
||||
)
|
||||
else:
|
||||
self._store.upsert(
|
||||
DynamicTgEntry(
|
||||
int_id=peer_int,
|
||||
system_name=system_name,
|
||||
slot=int(slot),
|
||||
tgid=int(tgid),
|
||||
single_mode=False,
|
||||
expires_at=None,
|
||||
updated_at=now,
|
||||
)
|
||||
)
|
||||
|
||||
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 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()
|
||||
|
||||
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
|
||||
]
|
||||
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)
|
||||
if tgids:
|
||||
logger.info(
|
||||
"(DYNAMIC_TG) Restored %s TG(s) for peer %s on %s: %s",
|
||||
len(tgids), int_id(peer_id), system_name, sorted(set(tgids)),
|
||||
)
|
||||
return tgids
|
||||
|
||||
return self._store.load_peer(int_id(peer_id), system_name).addCallback(_apply)
|
||||
|
||||
def purge_expired(self, config: dict[str, Any]) -> None:
|
||||
now = time.time()
|
||||
self._store.purge_expired(now)
|
||||
for sys_cfg in config.get("SYSTEMS", {}).values():
|
||||
if isinstance(sys_cfg, dict):
|
||||
purge_expired_peer_ua_sessions(sys_cfg, now=now)
|
||||
@ -0,0 +1,73 @@
|
||||
# ADN DMR Peer Server - application routing dynamic tg restore
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Bridge/subscription sync after dynamic TG rows are restored from persistence."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
from typing import Any
|
||||
|
||||
from adn_server.application.ports import SubscriptionStore
|
||||
from adn_server.application.subscription.routing_table_export import _legacy_to_type
|
||||
from adn_server.application.subscription.subscription_queries import store_has_table
|
||||
from adn_server.domain import bytes_3
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
from adn_server.domain.subscription import SubscriptionPhase
|
||||
|
||||
|
||||
def sync_restored_dynamic_bridges(
|
||||
entries: list[DynamicTgEntry],
|
||||
*,
|
||||
system_name: str,
|
||||
peer_id: bytes,
|
||||
sys_cfg: dict[str, Any],
|
||||
sub_store: SubscriptionStore,
|
||||
ensure_dynamic_relay: Callable[[bytes, str, int, float], None],
|
||||
ua_timer_minutes_for_peer: Callable[[str, bytes], float] | None,
|
||||
now: float,
|
||||
) -> None:
|
||||
"""Align bridge timers and create missing relay tables for restored dynamics."""
|
||||
for entry in entries:
|
||||
if not entry.single_mode or entry.expires_at is None:
|
||||
continue
|
||||
exp = float(entry.expires_at)
|
||||
if exp <= now:
|
||||
continue
|
||||
for sub in sub_store.snapshot():
|
||||
if (
|
||||
sub.table_key() == str(entry.tgid)
|
||||
and sub.system.value == system_name
|
||||
and int(sub.channel.slot) == int(entry.slot)
|
||||
):
|
||||
sub.state.timer_expires_at = exp
|
||||
if not sub.is_active() and _legacy_to_type(sub) == "ON":
|
||||
sub.state.phase = SubscriptionPhase.ACTIVE
|
||||
sub_store.upsert(sub)
|
||||
tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10))
|
||||
if ua_timer_minutes_for_peer is not None:
|
||||
tmout = float(ua_timer_minutes_for_peer(system_name, peer_id))
|
||||
seen: set[tuple[int, int]] = set()
|
||||
for entry in entries:
|
||||
key = (int(entry.tgid), int(entry.slot))
|
||||
if key in seen or store_has_table(sub_store, str(entry.tgid)):
|
||||
continue
|
||||
seen.add(key)
|
||||
ensure_dynamic_relay(bytes_3(entry.tgid), system_name, int(entry.slot), tmout)
|
||||
@ -0,0 +1,38 @@
|
||||
# ADN DMR Peer Server - domain dynamic tg
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Per-peer user-activated dynamic TG persistence."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class DynamicTgEntry:
|
||||
"""One persisted dynamic TG row for a hotspot peer."""
|
||||
|
||||
int_id: int
|
||||
system_name: str
|
||||
slot: int
|
||||
tgid: int
|
||||
single_mode: bool
|
||||
expires_at: float | None
|
||||
updated_at: float
|
||||
@ -0,0 +1,42 @@
|
||||
# ADN DMR Peer Server - infrastructure persistence database config
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""DATABASE block from config (MariaDB for dynamic TG persistence)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
|
||||
def _block(config: dict[str, Any], key: str) -> dict[str, Any]:
|
||||
raw = config.get(key) or {}
|
||||
return raw if isinstance(raw, dict) else {}
|
||||
|
||||
|
||||
def database_settings(config: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Resolved MariaDB settings from ``DATABASE``."""
|
||||
block = _block(config, "DATABASE")
|
||||
return {
|
||||
"db_server": str(block.get("DB_SERVER", "localhost")),
|
||||
"db_username": str(block.get("DB_USERNAME", "")),
|
||||
"db_password": str(block.get("DB_PASSWORD", "")),
|
||||
"db_name": str(block.get("DB_NAME", "")),
|
||||
"db_port": int(block.get("DB_PORT", 3306)),
|
||||
}
|
||||
@ -0,0 +1,109 @@
|
||||
# ADN DMR Peer Server - infrastructure persistence dynamic tg repository
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""MySQL ``peer_dynamic_tgs`` table adapter."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
from twisted.enterprise import adbapi
|
||||
from twisted.internet.defer import inlineCallbacks, returnValue
|
||||
|
||||
from adn_server.application.ports import DynamicTgStore
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _row_to_entry(row: tuple[Any, ...]) -> DynamicTgEntry:
|
||||
return DynamicTgEntry(
|
||||
int_id=int(row[0]),
|
||||
system_name=str(row[1]),
|
||||
slot=int(row[2]),
|
||||
tgid=int(row[3]),
|
||||
single_mode=bool(int(row[4])),
|
||||
expires_at=float(row[5]) if row[5] is not None else None,
|
||||
updated_at=float(row[6]),
|
||||
)
|
||||
|
||||
|
||||
class MysqlDynamicTgRepository(DynamicTgStore):
|
||||
"""``DynamicTgStore`` via Twisted adbapi."""
|
||||
|
||||
def __init__(self, pool: adbapi.ConnectionPool) -> None:
|
||||
self._pool = pool
|
||||
|
||||
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
|
||||
self._pool.runOperation(
|
||||
"""INSERT INTO peer_dynamic_tgs
|
||||
(int_id, system_name, slot, tgid, single_mode, expires_at, updated_at)
|
||||
VALUES (%s, %s, %s, %s, %s, %s, %s)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
single_mode=VALUES(single_mode),
|
||||
expires_at=VALUES(expires_at),
|
||||
updated_at=VALUES(updated_at)""",
|
||||
(
|
||||
entry.int_id,
|
||||
entry.system_name,
|
||||
entry.slot,
|
||||
entry.tgid,
|
||||
int(entry.single_mode),
|
||||
expires,
|
||||
updated,
|
||||
),
|
||||
).addErrback(lambda f: logger.error("(DYNAMIC_TG) upsert: %s", f.getTraceback()))
|
||||
|
||||
def replace_single_slot(self, entry: DynamicTgEntry) -> None:
|
||||
self._pool.runOperation(
|
||||
"DELETE FROM peer_dynamic_tgs WHERE int_id=%s AND system_name=%s AND slot=%s AND single_mode=1",
|
||||
(entry.int_id, entry.system_name, entry.slot),
|
||||
).addCallback(lambda _: self.upsert(entry)).addErrback(
|
||||
lambda f: logger.error("(DYNAMIC_TG) replace_single_slot: %s", f.getTraceback())
|
||||
)
|
||||
|
||||
def delete_peer_slot(self, int_id: int, system_name: str, slot: int) -> None:
|
||||
self._pool.runOperation(
|
||||
"DELETE FROM peer_dynamic_tgs WHERE int_id=%s AND system_name=%s AND slot=%s",
|
||||
(int_id, system_name, slot),
|
||||
).addErrback(
|
||||
lambda f: logger.error("(DYNAMIC_TG) delete_peer_slot: %s", f.getTraceback())
|
||||
)
|
||||
|
||||
@inlineCallbacks
|
||||
def load_peer(self, int_id: int, system_name: str) -> Any:
|
||||
rows = yield self._pool.runQuery(
|
||||
"""SELECT int_id, system_name, slot, tgid, single_mode, expires_at, updated_at
|
||||
FROM peer_dynamic_tgs
|
||||
WHERE int_id=%s AND system_name=%s""",
|
||||
(int_id, system_name),
|
||||
)
|
||||
returnValue([_row_to_entry(row) for row in (rows or [])])
|
||||
|
||||
def purge_expired(self, now: float) -> None:
|
||||
cutoff = int(now)
|
||||
self._pool.runOperation(
|
||||
"DELETE FROM peer_dynamic_tgs WHERE single_mode=1 AND expires_at IS NOT NULL AND expires_at <= %s",
|
||||
(cutoff,),
|
||||
).addErrback(lambda f: logger.error("(DYNAMIC_TG) purge_expired: %s", f.getTraceback()))
|
||||
@ -0,0 +1,88 @@
|
||||
# ADN DMR Peer Server - infrastructure persistence mysql pool helpers
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Shared Twisted adbapi MySQL pool helpers."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from twisted.enterprise import adbapi
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def create_mysql_pool(
|
||||
host: str,
|
||||
user: str,
|
||||
password: str,
|
||||
db_name: str,
|
||||
port: int,
|
||||
) -> adbapi.ConnectionPool:
|
||||
"""Create Twisted adbapi pool using ``MySQLdb`` (``mysqlclient`` package)."""
|
||||
return adbapi.ConnectionPool(
|
||||
"MySQLdb",
|
||||
host=host,
|
||||
user=user,
|
||||
passwd=password,
|
||||
db=db_name,
|
||||
port=port,
|
||||
charset="utf8mb4",
|
||||
)
|
||||
|
||||
|
||||
def verify_database_sync(
|
||||
host: str,
|
||||
user: str,
|
||||
password: str,
|
||||
db_name: str,
|
||||
port: int,
|
||||
) -> bool:
|
||||
"""Blocking startup check: MariaDB reachable and ``peer_dynamic_tgs`` exists."""
|
||||
try:
|
||||
import MySQLdb
|
||||
except ImportError as err:
|
||||
logger.critical(
|
||||
"(DATABASE) mysqlclient required for dynamic TG persistence: %s",
|
||||
err,
|
||||
)
|
||||
return False
|
||||
try:
|
||||
conn = MySQLdb.connect(
|
||||
host=host,
|
||||
user=user,
|
||||
passwd=password,
|
||||
db=db_name,
|
||||
port=port,
|
||||
charset="utf8mb4",
|
||||
)
|
||||
cur = conn.cursor()
|
||||
cur.execute("SELECT 1 FROM peer_dynamic_tgs LIMIT 1")
|
||||
cur.close()
|
||||
conn.close()
|
||||
logger.info("(DATABASE) peer_dynamic_tgs table: OK")
|
||||
return True
|
||||
except Exception as err:
|
||||
logger.critical(
|
||||
"(DATABASE) startup check failed: %s "
|
||||
"(configure DATABASE in adn-server.yaml and run adn-monitor db_bootstrap --update)",
|
||||
err,
|
||||
)
|
||||
return False
|
||||
@ -0,0 +1,65 @@
|
||||
# ADN DMR Peer Server - tests application dynamic tg persist
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Persist SINGLE=1 absolute expiry to MariaDB."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases
|
||||
from adn_server.application.routing.helpers import register_peer_ua_session
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
from tests.application.test_peer_single_downlink import _peer_id, _sys_cfg
|
||||
|
||||
|
||||
class _CaptureStore:
|
||||
def __init__(self) -> None:
|
||||
self.replaced: list[DynamicTgEntry] = []
|
||||
|
||||
def upsert(self, entry: DynamicTgEntry) -> None:
|
||||
pass
|
||||
|
||||
def replace_single_slot(self, entry: DynamicTgEntry) -> None:
|
||||
self.replaced.append(entry)
|
||||
|
||||
def delete_peer_slot(self, int_id: int, system_name: str, slot: int) -> None:
|
||||
pass
|
||||
|
||||
def load_peer(self, int_id: int, system_name: str):
|
||||
return []
|
||||
|
||||
def purge_expired(self, now: float) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def test_persist_after_register_stores_absolute_expires_from_memory() -> None:
|
||||
peer_id = _peer_id()
|
||||
sys_cfg = _sys_cfg()
|
||||
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=60;"}
|
||||
now = 1_700_000_000.0
|
||||
register_peer_ua_session(peer, peer_id, 2, 7304, sys_cfg, now=now)
|
||||
store = _CaptureStore()
|
||||
uc = DynamicTgUseCases(store)
|
||||
uc.persist_after_register(
|
||||
peer, peer_id, 2, 7304, sys_cfg, system_name="MASTER-A",
|
||||
)
|
||||
assert len(store.replaced) == 1
|
||||
entry = store.replaced[0]
|
||||
assert entry.expires_at == now + 60.0 * 60.0
|
||||
assert entry.single_mode is True
|
||||
@ -0,0 +1,151 @@
|
||||
# ADN DMR Peer Server - tests application dynamic tg restore
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""Restore persisted dynamic TG rows into in-memory UA stores."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from adn_server.application.routing.helpers import (
|
||||
peer_should_receive_group_voice,
|
||||
purge_expired_peer_ua_sessions,
|
||||
restore_peer_ua_entries_to_memory,
|
||||
)
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
from tests.application.test_peer_single_downlink import _peer_id, _sys_cfg
|
||||
|
||||
|
||||
def test_restore_single_mode_respects_expiry() -> None:
|
||||
peer_id = _peer_id()
|
||||
sys_cfg = _sys_cfg()
|
||||
now = 1_000_000.0
|
||||
entries = [
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7305,
|
||||
single_mode=True,
|
||||
expires_at=now + 300.0,
|
||||
updated_at=now,
|
||||
),
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=730,
|
||||
single_mode=True,
|
||||
expires_at=now - 1.0,
|
||||
updated_at=now - 60.0,
|
||||
),
|
||||
]
|
||||
restored = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, entries, now=now + 10)
|
||||
assert restored == [7305]
|
||||
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
|
||||
assert peer_should_receive_group_voice(
|
||||
peer, 2, 7305, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 20
|
||||
)
|
||||
assert not peer_should_receive_group_voice(
|
||||
peer, 2, 730, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 20
|
||||
)
|
||||
|
||||
|
||||
def test_restore_single_zero_multi_dynamic() -> None:
|
||||
peer_id = _peer_id()
|
||||
sys_cfg = _sys_cfg()
|
||||
now = 1_000_000.0
|
||||
entries = [
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7304,
|
||||
single_mode=False,
|
||||
expires_at=None,
|
||||
updated_at=now,
|
||||
),
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7306,
|
||||
single_mode=False,
|
||||
expires_at=None,
|
||||
updated_at=now,
|
||||
),
|
||||
]
|
||||
restored = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, entries, now=now)
|
||||
assert sorted(restored) == [7304, 7306]
|
||||
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
|
||||
assert peer_should_receive_group_voice(
|
||||
peer, 2, 7304, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 10
|
||||
)
|
||||
assert peer_should_receive_group_voice(
|
||||
peer, 2, 7306, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 10
|
||||
)
|
||||
|
||||
|
||||
def test_restore_single_preserves_absolute_expires_timestamp() -> None:
|
||||
"""SINGLE=1: reconnect restores the same wall-clock expiry, not a fresh TIMER."""
|
||||
peer_id = _peer_id()
|
||||
sys_cfg: dict = {}
|
||||
activated_at = 1_700_000_000.0
|
||||
expires_at = activated_at + 60.0 * 60.0 # TIMER=60 at activation
|
||||
reconnect_at = activated_at + 34.0 * 60.0 # 26 minutes remaining
|
||||
entries = [
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7304,
|
||||
single_mode=True,
|
||||
expires_at=expires_at,
|
||||
updated_at=activated_at,
|
||||
),
|
||||
]
|
||||
restore_peer_ua_entries_to_memory(sys_cfg, peer_id, entries, now=reconnect_at)
|
||||
pk_sessions = sys_cfg["_PEER_UA_SESSIONS"]
|
||||
per_peer = next(iter(pk_sessions.values()))
|
||||
assert per_peer[2]["expires"] == expires_at
|
||||
assert per_peer[2]["expires"] - reconnect_at == pytest.approx(26.0 * 60.0, rel=0.01)
|
||||
|
||||
|
||||
def test_purge_expired_peer_ua_sessions() -> None:
|
||||
peer_id = _peer_id()
|
||||
sys_cfg = _sys_cfg()
|
||||
now = 1_000_000.0
|
||||
entries = [
|
||||
DynamicTgEntry(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7304,
|
||||
single_mode=True,
|
||||
expires_at=now + 60.0,
|
||||
updated_at=now,
|
||||
),
|
||||
]
|
||||
restore_peer_ua_entries_to_memory(sys_cfg, peer_id, entries, now=now)
|
||||
purge_expired_peer_ua_sessions(sys_cfg, now=now + 120)
|
||||
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
|
||||
assert not peer_should_receive_group_voice(
|
||||
peer, 2, 7304, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 130
|
||||
)
|
||||
@ -0,0 +1,96 @@
|
||||
# ADN DMR Peer Server - tests infrastructure dynamic tg repository
|
||||
#
|
||||
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
|
||||
#
|
||||
###############################################################################
|
||||
# This program is free software; you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation; either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program; if not, write to the Free Software Foundation,
|
||||
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
|
||||
###############################################################################
|
||||
|
||||
"""DynamicTgStore SQL operations (mock pool)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pytest
|
||||
from twisted.internet.defer import succeed
|
||||
|
||||
from adn_server.domain.dynamic_tg import DynamicTgEntry
|
||||
from adn_server.infrastructure.persistence.dynamic_tg_repository import MysqlDynamicTgRepository
|
||||
|
||||
|
||||
def _entry(**kwargs) -> DynamicTgEntry:
|
||||
defaults = dict(
|
||||
int_id=730039101,
|
||||
system_name="MASTER-A",
|
||||
slot=2,
|
||||
tgid=7305,
|
||||
single_mode=True,
|
||||
expires_at=9_999_999.0,
|
||||
updated_at=1_000_000.0,
|
||||
)
|
||||
defaults.update(kwargs)
|
||||
return DynamicTgEntry(**defaults)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def repo() -> MysqlDynamicTgRepository:
|
||||
pool = MagicMock()
|
||||
pool.runOperation.return_value = succeed(None)
|
||||
pool.runQuery.return_value = succeed([])
|
||||
return MysqlDynamicTgRepository(pool)
|
||||
|
||||
|
||||
def test_upsert_issues_insert(repo: MysqlDynamicTgRepository) -> None:
|
||||
repo.upsert(_entry())
|
||||
sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001
|
||||
assert "INSERT INTO peer_dynamic_tgs" in sql
|
||||
assert args[0] == 730039101
|
||||
assert args[3] == 7305
|
||||
|
||||
|
||||
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
|
||||
assert "DELETE FROM peer_dynamic_tgs" in delete_sql
|
||||
|
||||
|
||||
def test_delete_peer_slot(repo: MysqlDynamicTgRepository) -> None:
|
||||
repo.delete_peer_slot(730039101, "MASTER-A", 2)
|
||||
sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001
|
||||
assert "DELETE FROM peer_dynamic_tgs" in sql
|
||||
assert args == (730039101, "MASTER-A", 2)
|
||||
|
||||
|
||||
def test_purge_expired(repo: MysqlDynamicTgRepository) -> None:
|
||||
repo.purge_expired(1_000_000.0)
|
||||
sql, args = repo._pool.runOperation.call_args[0] # noqa: SLF001
|
||||
assert "expires_at <= %s" in sql
|
||||
assert args == (1_000_000,)
|
||||
|
||||
|
||||
def test_load_peer_maps_rows(repo: MysqlDynamicTgRepository) -> None:
|
||||
repo._pool.runQuery.return_value = succeed( # noqa: SLF001
|
||||
[(730039101, "MASTER-A", 2, 7304, 0, None, 1_000_000.0)]
|
||||
)
|
||||
result: list = []
|
||||
d = repo.load_peer(730039101, "MASTER-A")
|
||||
d.addCallback(result.append)
|
||||
from twisted.internet import reactor
|
||||
|
||||
reactor.runUntilCurrent()
|
||||
assert len(result) == 1
|
||||
assert result[0][0].tgid == 7304
|
||||
assert result[0][0].single_mode is False
|
||||
Loading…
Reference in new issue