diff --git a/src/adn_server/application/talker_alias_use_cases.py b/src/adn_server/application/talker_alias_use_cases.py index 220dc26..9605000 100644 --- a/src/adn_server/application/talker_alias_use_cases.py +++ b/src/adn_server/application/talker_alias_use_cases.py @@ -83,14 +83,10 @@ def format_talker_alias_text(config: dict[str, Any], rf_src: bytes) -> str: settings = talker_alias_settings(config) template = settings["format"] rid = int_id(rf_src) - profiles = config.get("_SUB_PROFILES", {}) - profile = profiles.get(rid, {}) - sub_ids = config.get("_SUB_IDS", {}) - callsign = profile.get("callsign") or sub_ids.get(rid) or "" - fname = profile.get("fname") or "" - surname = profile.get("surname") or "" - if profile.get("talker_alias"): - text = str(profile["talker_alias"]) + callsign, fname, surname, talker_alias = config.get("_SUB_PROFILES", {}).get(rid) or (None, "", "", "") + callsign = callsign or config.get("_SUB_IDS", {}).get(rid) or "" + if talker_alias: + text = talker_alias else: try: text = template.format( diff --git a/src/adn_server/domain/value_objects.py b/src/adn_server/domain/value_objects.py index 960672a..98bbdf9 100644 --- a/src/adn_server/domain/value_objects.py +++ b/src/adn_server/domain/value_objects.py @@ -54,6 +54,12 @@ class TgId: return self.value +# (callsign, fname, surname, talker_alias). The subscriber file holds over 300k of these: +# a dict each costs twice the memory, and a NamedTuple stays tracked by the garbage +# collector, adding ~40 ms to every full collection on the reactor thread. +SubscriberProfile = tuple[str | None, str, str, str] + + Slot = Literal[1, 2] """Timeslot 1 or 2.""" diff --git a/src/adn_server/infrastructure/persistence/alias_loader.py b/src/adn_server/infrastructure/persistence/alias_loader.py index 227d0ac..7470630 100644 --- a/src/adn_server/infrastructure/persistence/alias_loader.py +++ b/src/adn_server/infrastructure/persistence/alias_loader.py @@ -32,6 +32,7 @@ import logging import os import shutil import ssl +import sys import threading import time from collections.abc import Callable @@ -40,6 +41,7 @@ from typing import Any from urllib.request import urlopen from ...application.ports import AliasLoader +from ...domain.value_objects import SubscriberProfile logger = logging.getLogger(__name__) @@ -213,14 +215,16 @@ class DefaultAliasLoader(AliasLoader): peer_ids = self._load_with_backup( path, peer_file, checksums.get("peer_ids"), "peer_ids", self._load_id_json, ) - subscriber_ids = self._load_with_backup( - path, sub_file, checksums.get("subscriber_ids"), "subscriber_ids", self._load_id_json, - ) + subscriber_ids = self._subscriber_ids(self._load_with_backup( + path, sub_file, checksums.get("subscriber_ids"), "subscriber_ids", self._load_subscriber_json, + )) talkgroup_ids = self._load_with_backup( path, tgid_file, checksums.get("talkgroup_ids"), "talkgroup_ids", self._load_id_json, ) - local_subscriber_ids = self._load_id_json( - path / aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json") + local_file = aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json") + # The default points at the subscriber file itself: don't parse 300k records twice. + local_subscriber_ids = ( + subscriber_ids if local_file == sub_file else self._load_id_json(path / local_file) ) server_ids = self._load_with_backup( path, server_file, checksums.get("server_ids"), "server_ids", @@ -257,10 +261,10 @@ class DefaultAliasLoader(AliasLoader): _keep("_PEER_IDS", peer_ids, "peer_ids") if subscriber_ids: - sub = dict(subscriber_ids) - sub[900999] = "D-APRS" - sub[4294967295] = "SC" - config["_SUB_IDS"] = sub + # In place: a copy would keep a second 300k-entry table alive. + subscriber_ids[900999] = "D-APRS" + subscriber_ids[4294967295] = "SC" + config["_SUB_IDS"] = subscriber_ids if profiles is not None: config["_SUB_PROFILES"] = profiles elif isinstance(alias_loader, DefaultAliasLoader): @@ -407,55 +411,59 @@ class DefaultAliasLoader(AliasLoader): pass return out - def load_subscriber_profiles(self, config: dict[str, Any]) -> dict[int, dict[str, str]]: - """Load {id: {callsign, fname, surname, talker_alias?}} from subscriber JSON files.""" + def load_subscriber_profiles(self, config: dict[str, Any]) -> dict[int, SubscriberProfile]: + """{id: profile} from the subscriber file, overlaid with the local subscriber file.""" aliases = config.get("ALIASES", {}) path = Path(aliases.get("PATH", "./data/")).resolve() sub_file = aliases.get("SUBSCRIBER_FILE", "subscriber_ids.json") local_file = aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json") - files = [path / sub_file, path / local_file] - stamps = [_file_stamp(f) for f in files] + # load_aliases has just parsed this file: this is a cache hit, not a second parse. + main = self._load_with_backup(path, sub_file, None, "subscriber_ids", self._load_subscriber_json) + if local_file == sub_file: + return main + stamp = _file_stamp(path / local_file) remembered = self._parsed.get("_profiles") - if remembered is not None and remembered[0] == stamps: - return remembered[1] - profiles: dict[int, dict[str, str]] = {} - for file_path in files: - self._merge_subscriber_profiles(file_path, profiles) - self._parsed["_profiles"] = (stamps, profiles) + if remembered is not None and remembered[0] is main and remembered[1] == stamp: + return remembered[2] + local = self._load_subscriber_json(path / local_file) + profiles = {**main, **local} if local else main + self._parsed["_profiles"] = (main, stamp, profiles) return profiles - def _merge_subscriber_profiles(self, file_path: Path, out: dict[int, dict[str, str]]) -> None: + def _subscriber_ids(self, profiles: dict[int, SubscriberProfile]) -> dict[int, str]: + remembered = self._parsed.get("_sub_ids") + if remembered is not None and remembered[0] is profiles: + return remembered[1] + ids = {rid: p[0] for rid, p in profiles.items() if p[0] is not None} + self._parsed["_sub_ids"] = (profiles, ids) + return ids + + def _load_subscriber_json(self, file_path: Path) -> dict[int, SubscriberProfile]: if not file_path.is_file(): - return + return {} + out: dict[int, SubscriberProfile] = {} + + # Each record becomes a profile as soon as the decoder finishes it, so the + # file's whole tree (eight strings per subscriber) never exists at once. + def record(obj: dict[str, Any]) -> Any: + if "id" not in obj: + return obj + try: + rid = int(obj["id"]) + except (TypeError, ValueError): + return None + callsign = obj.get("callsign") + out[rid] = ( + None if callsign is None else str(callsign), + sys.intern(str(obj.get("fname") or "")), + sys.intern(str(obj.get("surname") or "")), + str(obj.get("talker_alias") or ""), + ) + return None + try: with open(file_path, "r", encoding="utf-8") as f: - data = json.load(f) + json.load(f, object_hook=record) except (json.JSONDecodeError, OSError): - return - if not isinstance(data, dict): - return - if "count" in data: - data = {k: v for k, v in data.items() if k != "count"} - for _key, val in data.items(): - if not isinstance(val, list): - continue - for record in val: - if not isinstance(record, dict) or "id" not in record: - continue - try: - rid = int(record["id"]) - except (ValueError, TypeError): - continue - entry: dict[str, str] = {} - if record.get("callsign"): - entry["callsign"] = str(record["callsign"]) - if record.get("fname"): - entry["fname"] = str(record["fname"]) - if record.get("surname"): - entry["surname"] = str(record["surname"]) - if record.get("talker_alias"): - entry["talker_alias"] = str(record["talker_alias"]) - if entry: - prev = out.get(rid, {}) - prev.update(entry) - out[rid] = prev + return {} + return out diff --git a/tests/harness/scenarios.py b/tests/harness/scenarios.py index a5fd774..1dcbf47 100644 --- a/tests/harness/scenarios.py +++ b/tests/harness/scenarios.py @@ -53,7 +53,7 @@ def talker_alias_config() -> dict[str, Any]: config["GLOBAL"]["TALKER_ALIAS_SEND_DMRA"] = True rid = 3120001 config["_SUB_PROFILES"] = { - rid: {"callsign": "CE5RPY", "fname": "Rodrigo", "surname": "Perez"}, + rid: ("CE5RPY", "Rodrigo", "Perez", ""), } config["_SUB_IDS"] = {rid: "CE5RPY"} return config diff --git a/tests/infrastructure/test_alias_memory.py b/tests/infrastructure/test_alias_memory.py new file mode 100644 index 0000000..56a3f9f --- /dev/null +++ b/tests/infrastructure/test_alias_memory.py @@ -0,0 +1,102 @@ +# ADN DMR Peer Server - tests infrastructure alias reload memory +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# 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 +############################################################################### + +"""Every real change to the 300k-subscriber file raised the server's RSS by ~130 MB.""" + +from __future__ import annotations + +import gc +import json +from pathlib import Path + +from adn_server.application.talker_alias_use_cases import format_talker_alias_text +from adn_server.infrastructure.persistence.alias_loader import DefaultAliasLoader + + +def _config(path: Path, local_file: str = "local_subscriber_ids.json") -> dict: + (path / "subscriber_ids.json").write_text(json.dumps({"count": 2, "results": [ + {"id": "7300391", "callsign": "CE5RPY", "fname": "Rodrigo", "surname": "Perez", "city": "Talca"}, + {"id": "7300392", "callsign": "CE5ABC", "fname": "Rodrigo", "surname": "Soto", "city": "Talca"}, + ]}), encoding="utf-8") + (path / local_file).write_text(json.dumps({"results": [ + {"id": 7300392, "callsign": "CE5ABC", "fname": "Juan", "surname": "Soto"}, + ]}), encoding="utf-8") + return {"GLOBAL": {"TALKER_ALIAS_FORMAT": "{callsign} {fname}"}, "ALIASES": { + "PATH": str(path), "SUBSCRIBER_FILE": "subscriber_ids.json", "LOCAL_SUBSCRIBER_FILE": local_file, + }} + + +def _reload(loader: DefaultAliasLoader, config: dict) -> tuple: + loaded = loader.load_aliases(config) + DefaultAliasLoader.merge_reload_into_config( + config, loader, *loaded, profiles=loader.load_subscriber_profiles(config) + ) + return loaded + + +def test_the_subscriber_file_is_parsed_once_per_reload(tmp_path: Path) -> None: + config = _config(tmp_path) + loader = DefaultAliasLoader() + parsed = [] + original = loader._load_subscriber_json + loader._load_subscriber_json = lambda p: (parsed.append(p.name), original(p))[1] # type: ignore[assignment] + + _reload(loader, config) + + assert parsed.count("subscriber_ids.json") == 1 + + +def test_ids_and_profiles_come_out_of_that_one_parse(tmp_path: Path) -> None: + config = _config(tmp_path) + _reload(DefaultAliasLoader(), config) + + assert config["_SUB_IDS"][7300391] == "CE5RPY" + assert config["_SUB_IDS"][900999] == "D-APRS" + assert format_talker_alias_text(config, (7300391).to_bytes(3, "big")) == "CE5RPY Rodrigo" + # The local subscriber file still overrides the downloaded one. + assert format_talker_alias_text(config, (7300392).to_bytes(3, "big")) == "CE5ABC Juan" + + +def test_the_live_subscriber_ids_are_not_a_second_copy(tmp_path: Path) -> None: + config = _config(tmp_path) + loader = DefaultAliasLoader() + subscriber_ids = _reload(loader, config)[1] + + assert config["_SUB_IDS"] is subscriber_ids + + +def test_profiles_stay_out_of_the_garbage_collector(tmp_path: Path) -> None: + """300k tracked objects add ~40 ms to every full collection, on the reactor thread.""" + config = _config(tmp_path) + _reload(DefaultAliasLoader(), config) + gc.collect() + + assert not any(gc.is_tracked(p) for p in config["_SUB_PROFILES"].values()) + + +def test_the_default_local_file_is_not_parsed_again(tmp_path: Path) -> None: + """LOCAL_SUBSCRIBER_FILE defaults to the subscriber file itself.""" + config = _config(tmp_path, local_file="subscriber_ids.json") + loader = DefaultAliasLoader() + loader._load_id_json = lambda p: {} if p.name != "subscriber_ids.json" else 1 / 0 # type: ignore[assignment] + + peer_ids, subscriber_ids, _tg, local_ids, _srv, _chk = loader.load_aliases(config) + + assert local_ids is subscriber_ids diff --git a/tests/infrastructure/test_alias_reload_skips_unchanged.py b/tests/infrastructure/test_alias_reload_skips_unchanged.py index 948b5fe..db505b9 100644 --- a/tests/infrastructure/test_alias_reload_skips_unchanged.py +++ b/tests/infrastructure/test_alias_reload_skips_unchanged.py @@ -133,13 +133,12 @@ def test_subscriber_profiles_are_not_rebuilt_for_unchanged_files(tmp_path: Path) first = loader.load_subscriber_profiles(cfg) assert 7300391 in first - merges = [] - original = loader._merge_subscriber_profiles - loader._merge_subscriber_profiles = lambda p, out: merges.append(p) # type: ignore[assignment] + parses = [] + original = loader._load_subscriber_json + loader._load_subscriber_json = lambda p: (parses.append(p), original(p))[1] # type: ignore[assignment] assert loader.load_subscriber_profiles(cfg) is first - assert merges == [] + assert parses == [] - loader._merge_subscriber_profiles = original # type: ignore[assignment] _write(tmp_path, "subscriber_ids.json", 7300392, "CE5ABC") os.utime(tmp_path / "subscriber_ids.json", (2_000_000_000, 2_000_000_000)) assert 7300392 in loader.load_subscriber_profiles(cfg)