Merge pull request #41 from ce5rpy/develop

Release 2.3.0
pull/43/head
ce5rpy 3 months ago committed by GitHub
commit 556626e8e8
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -37,7 +37,8 @@ from adn_server.application.routing.helpers import (
restore_peer_ua_entries_to_memory,
)
from adn_server.domain import int_id
from adn_server.domain.dynamic_tg import DynamicTgEntry
from adn_server.domain.dynamic_tg import DynamicTgEntry, is_persisted_dynamic_row
from adn_server.domain.value_objects import bytes_4
logger = logging.getLogger(__name__)
@ -109,7 +110,8 @@ class DynamicTgUseCases:
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
if is_persisted_dynamic_row(e)
and (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:
@ -129,3 +131,30 @@ class DynamicTgUseCases:
for sys_cfg in config.get("SYSTEMS", {}).values():
if isinstance(sys_cfg, dict):
purge_expired_peer_ua_sessions(sys_cfg, now=now)
def process_reload_queue(
self,
*,
try_purge: Callable[[int, str, bytes], bool],
) -> Any:
"""Poll ``need_reload`` rows and apply TG-4000-equivalent purge when peer is online."""
def _on_rows(rows: list[tuple[int, str]] | None) -> None:
for peer_int, system_name in rows or []:
peer_id = bytes_4(peer_int)
try:
if try_purge(peer_int, system_name, peer_id):
logger.info(
"(DYNAMIC_TG) Applied monitor reload for peer %s on %s",
peer_int,
system_name,
)
except Exception as err:
logger.warning(
"(DYNAMIC_TG) reload for peer %s on %s failed: %s",
peer_int,
system_name,
err,
)
return self._store.select_need_reload().addCallback(_on_rows)

@ -407,6 +407,10 @@ class DynamicTgStore(ABC):
def load_peer(self, int_id: int, system_name: str) -> Any:
"""Returns Deferred firing with ``list[DynamicTgEntry]``."""
@abstractmethod
def select_need_reload(self) -> Any:
"""Returns Deferred firing with ``list[tuple[int, str]]`` (int_id, system_name)."""
@abstractmethod
def purge_expired(self, now: float) -> None:
...

@ -36,3 +36,13 @@ class DynamicTgEntry:
single_mode: bool
expires_at: float | None
updated_at: float
need_reload: bool = False
def is_persisted_dynamic_row(entry: DynamicTgEntry) -> bool:
"""True for real dynamic TG rows (exclude monitor reload control rows)."""
if entry.need_reload:
return False
if int(entry.slot) == 0 and int(entry.tgid) == 0:
return False
return True

@ -590,7 +590,14 @@ def run_peer_server(
apply_proxy_config_reload(proxy_state, config, logger=logger)
_wire_proxy_report_slots(report_factory, proxy_state)
elif proxy_enabled and proxy_target_system(config):
proxy_state = start_proxy_service(config, protocols, logger=logger)
proxy_state = start_proxy_service(
config,
protocols,
logger=logger,
mysql_pool=mysql_pool,
dynamic_tg_uc=dynamic_tg_uc,
purge_peer_dynamic=_purge_peer_dynamic_tgs,
)
_wire_proxy_report_slots(report_factory, proxy_state)
else:
_wire_proxy_report_slots(report_factory, None)
@ -708,10 +715,33 @@ def run_peer_server(
_wire_monitor_downlink_ctx(report_factory, protocols)
def _purge_peer_dynamic_tgs(peer_id: bytes, system_name: str) -> bool:
proto = protocols.get(system_name)
if proto is None:
return False
purge = getattr(proto, "purge_peer_dynamic_tgs", None)
if not callable(purge):
return False
return bool(purge(peer_id))
def _try_purge_dynamic_reload(peer_int: int, system_name: str, peer_id: bytes) -> bool:
proto = protocols.get(system_name)
if proto is None:
return False
peers = getattr(proto, "_peers", None)
if not isinstance(peers, dict) or peer_id not in peers:
return False
return _purge_peer_dynamic_tgs(peer_id, system_name)
if proxy_enabled:
try:
proxy_state = start_proxy_service(
config, protocols, logger=logger, mysql_pool=mysql_pool,
config,
protocols,
logger=logger,
mysql_pool=mysql_pool,
dynamic_tg_uc=dynamic_tg_uc,
purge_peer_dynamic=_purge_peer_dynamic_tgs,
)
except Exception as exc:
logger.error("(PROXY) failed to start integrated proxy: %s", exc)
@ -724,6 +754,13 @@ def run_peer_server(
reactor.addSystemEventTrigger("before", "shutdown", _stop_proxy)
_wire_proxy_report_slots(report_factory, proxy_state)
if proxy_state is None or proxy_state.self_service is None:
def dynamic_reload_loop() -> None:
dynamic_tg_uc.process_reload_queue(try_purge=_try_purge_dynamic_reload)
task.LoopingCall(dynamic_reload_loop).start(10).addErrback(_looping_errback, logger)
logger.info("(DYNAMIC_TG) need_reload poll loop started (every 10 seconds)")
logger.info("(GLOBAL) ADN DMR Peer Server started. Use adn-dmr-server as reference.")
reactor.suggestThreadPoolSize(100)
reactor.run()

@ -34,6 +34,10 @@ from adn_server.domain.dynamic_tg import DynamicTgEntry
logger = logging.getLogger(__name__)
_SELECT_COLS = (
"int_id, system_name, slot, tgid, single_mode, expires_at, updated_at, need_reload"
)
def _row_to_entry(row: tuple[Any, ...]) -> DynamicTgEntry:
return DynamicTgEntry(
@ -44,6 +48,7 @@ def _row_to_entry(row: tuple[Any, ...]) -> DynamicTgEntry:
single_mode=bool(int(row[4])),
expires_at=float(row[5]) if row[5] is not None else None,
updated_at=float(row[6]),
need_reload=bool(int(row[7])) if len(row) > 7 else False,
)
@ -58,12 +63,13 @@ class MysqlDynamicTgRepository(DynamicTgStore):
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)
(int_id, system_name, slot, tgid, single_mode, expires_at, updated_at, need_reload)
VALUES (%s, %s, %s, %s, %s, %s, %s, 0)
ON DUPLICATE KEY UPDATE
single_mode=VALUES(single_mode),
expires_at=VALUES(expires_at),
updated_at=VALUES(updated_at)""",
updated_at=VALUES(updated_at),
need_reload=0""",
(
entry.int_id,
entry.system_name,
@ -102,13 +108,20 @@ class MysqlDynamicTgRepository(DynamicTgStore):
@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
f"""SELECT {_SELECT_COLS}
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 [])])
@inlineCallbacks
def select_need_reload(self) -> Any:
rows = yield self._pool.runQuery(
"SELECT DISTINCT int_id, system_name FROM peer_dynamic_tgs WHERE need_reload=1"
)
returnValue([(int(row[0]), str(row[1])) for row in (rows or [])])
def purge_expired(self, now: float) -> None:
cutoff = int(now)
self._pool.runOperation(

@ -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_NEED_RELOAD_MIGRATION = "006_peer_dynamic_tgs_need_reload"
_MYSQL_HINTS: dict[int, str] = {
1045: "check DATABASE.DB_USERNAME and DB_PASSWORD in adn-server.yaml",
@ -114,14 +115,57 @@ def _mark_migration(cursor: Any, migration_id: str) -> None:
)
def _column_exists(cursor: Any, table: str, column: str) -> bool:
cursor.execute(
"SELECT 1 FROM information_schema.COLUMNS "
"WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = %s AND COLUMN_NAME = %s",
(table, column),
)
return cursor.fetchone() is not None
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 _index_exists(cursor: Any, table: str, index_name: str) -> bool:
cursor.execute(
"SELECT 1 FROM information_schema.STATISTICS "
"WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = %s AND INDEX_NAME = %s",
(table, index_name),
)
return cursor.fetchone() is not 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``."""
"""Server-owned table; same migration ids/DLL as adn-monitor."""
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_NEED_RELOAD_MIGRATION):
if _table_exists(cursor, "peer_dynamic_tgs"):
if not _column_exists(cursor, "peer_dynamic_tgs", "need_reload"):
cursor.execute(
"ALTER TABLE peer_dynamic_tgs "
"ADD COLUMN need_reload TINYINT(1) NOT NULL DEFAULT 0"
)
if not _index_exists(cursor, "peer_dynamic_tgs", "idx_need_reload_peer"):
cursor.execute(
"ALTER TABLE peer_dynamic_tgs "
"ADD INDEX idx_need_reload_peer (need_reload, int_id, system_name)"
)
_mark_migration(cursor, _PEER_DYNAMIC_TGS_NEED_RELOAD_MIGRATION)
logger.info(
"(DATABASE) applied migration %s (peer_dynamic_tgs.need_reload)",
_PEER_DYNAMIC_TGS_NEED_RELOAD_MIGRATION,
)
def ensure_database_sync(

@ -24,7 +24,7 @@ from __future__ import annotations
import logging
from dataclasses import dataclass, field
from typing import Any
from typing import Any, Callable
from twisted.internet import reactor
from twisted.internet.interfaces import IDelayedCall
@ -139,6 +139,8 @@ def _build_self_service(
*,
logger: logging.Logger,
mysql_pool: Any | None = None,
dynamic_tg_uc: Any | None = None,
purge_peer_dynamic: Callable[[bytes, str], bool] | None = None,
) -> ProxySelfServiceBridge | None:
ss = self_service_settings(config)
if not ss["enabled"]:
@ -170,6 +172,8 @@ def _build_self_service(
pbkdf2_salt=ss["pbkdf2_salt"],
pbkdf2_iterations=ss["pbkdf2_iterations"],
logger=logger,
dynamic_tg_uc=dynamic_tg_uc,
purge_peer_dynamic=purge_peer_dynamic,
)
return bridge
@ -180,6 +184,8 @@ def start_proxy_service(
*,
logger: logging.Logger,
mysql_pool: Any | None = None,
dynamic_tg_uc: Any | None = None,
purge_peer_dynamic: Callable[[bytes, str], bool] | None = None,
) -> ProxyServiceState:
"""Start LISTEN_PORT fan-in and inject into ``PROXY.TARGET_SYSTEM`` MASTER."""
runtime = _proxy_runtime_snapshot(config)
@ -297,6 +303,8 @@ def start_proxy_service(
state.client_sender,
logger=logger,
mysql_pool=mysql_pool,
dynamic_tg_uc=dynamic_tg_uc,
purge_peer_dynamic=purge_peer_dynamic,
)
if self_service_bridge is not None:
fanin._self_service = self_service_bridge # noqa: SLF001

@ -25,7 +25,7 @@ from __future__ import annotations
import logging
import struct
from hashlib import pbkdf2_hmac
from typing import Any
from typing import Any, Callable
from twisted.internet import reactor
from twisted.internet.defer import inlineCallbacks
@ -67,6 +67,8 @@ class ProxySelfServiceBridge:
pbkdf2_salt: str = "ADN",
pbkdf2_iterations: int = 2000,
logger: logging.Logger | None = None,
dynamic_tg_uc: Any = None,
purge_peer_dynamic: Callable[[bytes, str], bool] | None = None,
) -> None:
self._store = store
self._use_cases = use_cases
@ -77,6 +79,8 @@ class ProxySelfServiceBridge:
self._log = logger or logging.getLogger(__name__)
self._opt_timers: dict[bytes, IDelayedCall] = {}
self._loop_calls: list[LoopingCall] = []
self._dynamic_tg_uc = dynamic_tg_uc
self._purge_peer_dynamic = purge_peer_dynamic
# Peers that sent PASS= this session — only they get MySQL OPTIONS push.
self._mysql_option_peers: set[bytes] = set()
@ -93,6 +97,10 @@ class ProxySelfServiceBridge:
"(SELF_SERVICE) DB options on PASS= (immediate), send_opts every 10s, "
"lst_seen + reconcile_logged_in every 2min"
)
if self._dynamic_tg_uc is not None and self._purge_peer_dynamic is not None:
self._log.info(
"(SELF_SERVICE) peer_dynamic_tgs.need_reload polled on send_opts (TG 4000 parity)"
)
def stop_loops(self) -> None:
for call in self._loop_calls:
@ -344,9 +352,22 @@ class ProxySelfServiceBridge:
self._store.updt_tbl("rst_mod", pid)
self._inject_rpto(pid, options)
self._log.info("(SELF_SERVICE) Options update sent for: %s", int_id(pid))
yield self._process_dynamic_reload()
except Exception as err:
self._log.warning("(SELF_SERVICE) send_opts error: %s", err)
@inlineCallbacks
def _process_dynamic_reload(self) -> None:
if self._dynamic_tg_uc is None or self._purge_peer_dynamic is None:
return
def _try_purge(peer_int: int, system_name: str, peer_id: bytes) -> bool:
if self._use_cases.resolve_client(peer_id) is None:
return False
return bool(self._purge_peer_dynamic(peer_id, system_name))
yield self._dynamic_tg_uc.process_reload_queue(try_purge=_try_purge)
def lst_seen(self) -> None:
slots = self._use_cases.list_slots()
dmrid_list = [(slot.peer_id,) for slot in slots]

@ -551,6 +551,19 @@ class HBPProtocol(DatagramProtocol):
)
return True
def purge_peer_dynamic_tgs(self, peer_id: bytes) -> bool:
"""Monitor/proxy requested TG-4000-equivalent purge for one connected peer."""
if peer_id not in self._peers:
return False
self._apply_tg4000_reset(
peer_id,
1,
"group",
rf_src=peer_id,
stream_id=b"\x00\x00\x00\x00",
)
return True
def _peer_should_receive_dmrd(self, peer_id: bytes, packet: bytes) -> bool:
if peer_id not in self._peers:
return False

@ -52,6 +52,9 @@ class _CaptureStore:
def purge_expired(self, now: float) -> None:
pass
def select_need_reload(self):
return defer.succeed([])
def test_persist_after_register_stores_absolute_expires_from_memory() -> None:
peer_id = _peer_id()
@ -106,6 +109,9 @@ def test_restore_peer_passes_now_as_keyword_to_on_restored() -> None:
def purge_expired(self, now: float) -> None:
pass
def select_need_reload(self):
return defer.succeed([])
uc = DynamicTgUseCases(_Store(), on_restored=on_restored)
out: list[list[int]] = []
d = uc.restore_peer(peer_id, "MASTER-A", sys_cfg)

@ -32,8 +32,8 @@ from adn_server.domain.dmr.bptc import encode_emblc
from adn_server.domain.subscription import TgId
from adn_server.domain.voice_routing import ForwardLeg
from adn_server.infrastructure.acl_router import InMemoryAclRouter
from fakes.subscription_store import InMemorySubscriptionStore
from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder
from fakes.subscription_store import InMemorySubscriptionStore
def _routing(routing_table: dict) -> RoutingUseCases:

@ -0,0 +1,31 @@
"""Domain rules for peer_dynamic_tgs rows."""
from __future__ import annotations
from adn_server.domain.dynamic_tg import DynamicTgEntry, is_persisted_dynamic_row
def test_is_persisted_dynamic_row_rejects_reload_control() -> None:
assert not is_persisted_dynamic_row(
DynamicTgEntry(
int_id=1,
system_name="SYS",
slot=0,
tgid=0,
single_mode=False,
expires_at=None,
updated_at=1.0,
need_reload=True,
)
)
assert is_persisted_dynamic_row(
DynamicTgEntry(
int_id=1,
system_name="SYS",
slot=2,
tgid=7305,
single_mode=True,
expires_at=9_999.0,
updated_at=1.0,
)
)

@ -89,9 +89,22 @@ def test_purge_expired(repo: MysqlDynamicTgRepository) -> None:
assert args == (1_000_000,)
def test_select_need_reload(repo: MysqlDynamicTgRepository) -> None:
repo._pool.runQuery.return_value = succeed([(730039101, "MASTER-A")]) # noqa: SLF001
result: list = []
d = repo.select_need_reload()
d.addCallback(result.append)
from twisted.internet import reactor
reactor.runUntilCurrent()
assert result[0] == [(730039101, "MASTER-A")]
sql = repo._pool.runQuery.call_args[0][0] # noqa: SLF001
assert "need_reload=1" in sql
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)]
[(730039101, "MASTER-A", 2, 7304, 0, None, 1_000_000.0, 0)]
)
result: list = []
d = repo.load_peer(730039101, "MASTER-A")

@ -44,11 +44,13 @@ 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,), None, None, None]
_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("need_reload" in sql for sql in executed)
assert any("idx_need_reload_peer" in sql for sql in executed)
assert any("INSERT IGNORE INTO schema_migrations" in sql for sql in executed)

@ -307,6 +307,24 @@ def test_send_opts_pushes_modified_rows_after_pass() -> None:
assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;"
def test_send_opts_processes_need_reload() -> None:
bridge, _sink, store, _sender = _bridge()
peer = bytes_4(7300444)
purged: list[bytes] = []
class _DynUc:
def process_reload_queue(self, *, try_purge):
try_purge(7300444, "MASTER-A", peer)
from twisted.internet.defer import succeed
return succeed(None)
bridge._dynamic_tg_uc = _DynUc()
bridge._purge_peer_dynamic = lambda pid, _sys: purged.append(pid) or True
_run_deferred(bridge.send_opts())
assert purged == [peer]
def test_session_expired_logs_out() -> None:
bridge, _sink, store, _sender = _bridge()
peer = bytes_4(7300444)

Loading…
Cancel
Save

Powered by TurnKey Linux.