perf(aliases): stop each alias reload from raising the server's memory

- Parse the subscriber file once per reload, turning each record into a
  profile as the decoder finishes it, instead of parsing it twice and
  building the whole object tree first.
- Profiles are plain tuples with interned names: a dict each doubled the
  memory, and a NamedTuple stays tracked by the garbage collector.
- One copy of subscriber_ids, and the default local subscriber file (the
  subscriber file itself) is not parsed again.
perf/alias-memory
Rodrigo Pérez 11 hours ago
parent bec05936f9
commit 1aa602a934

@ -83,14 +83,10 @@ def format_talker_alias_text(config: dict[str, Any], rf_src: bytes) -> str:
settings = talker_alias_settings(config) settings = talker_alias_settings(config)
template = settings["format"] template = settings["format"]
rid = int_id(rf_src) rid = int_id(rf_src)
profiles = config.get("_SUB_PROFILES", {}) callsign, fname, surname, talker_alias = config.get("_SUB_PROFILES", {}).get(rid) or (None, "", "", "")
profile = profiles.get(rid, {}) callsign = callsign or config.get("_SUB_IDS", {}).get(rid) or ""
sub_ids = config.get("_SUB_IDS", {}) if talker_alias:
callsign = profile.get("callsign") or sub_ids.get(rid) or "" text = talker_alias
fname = profile.get("fname") or ""
surname = profile.get("surname") or ""
if profile.get("talker_alias"):
text = str(profile["talker_alias"])
else: else:
try: try:
text = template.format( text = template.format(

@ -54,6 +54,12 @@ class TgId:
return self.value 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] Slot = Literal[1, 2]
"""Timeslot 1 or 2.""" """Timeslot 1 or 2."""

@ -32,6 +32,7 @@ import logging
import os import os
import shutil import shutil
import ssl import ssl
import sys
import threading import threading
import time import time
from collections.abc import Callable from collections.abc import Callable
@ -40,6 +41,7 @@ from typing import Any
from urllib.request import urlopen from urllib.request import urlopen
from ...application.ports import AliasLoader from ...application.ports import AliasLoader
from ...domain.value_objects import SubscriberProfile
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -213,14 +215,16 @@ class DefaultAliasLoader(AliasLoader):
peer_ids = self._load_with_backup( peer_ids = self._load_with_backup(
path, peer_file, checksums.get("peer_ids"), "peer_ids", self._load_id_json, path, peer_file, checksums.get("peer_ids"), "peer_ids", self._load_id_json,
) )
subscriber_ids = self._load_with_backup( subscriber_ids = self._subscriber_ids(self._load_with_backup(
path, sub_file, checksums.get("subscriber_ids"), "subscriber_ids", self._load_id_json, path, sub_file, checksums.get("subscriber_ids"), "subscriber_ids", self._load_subscriber_json,
) ))
talkgroup_ids = self._load_with_backup( talkgroup_ids = self._load_with_backup(
path, tgid_file, checksums.get("talkgroup_ids"), "talkgroup_ids", self._load_id_json, path, tgid_file, checksums.get("talkgroup_ids"), "talkgroup_ids", self._load_id_json,
) )
local_subscriber_ids = self._load_id_json( local_file = aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json")
path / 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( server_ids = self._load_with_backup(
path, server_file, checksums.get("server_ids"), "server_ids", path, server_file, checksums.get("server_ids"), "server_ids",
@ -257,10 +261,10 @@ class DefaultAliasLoader(AliasLoader):
_keep("_PEER_IDS", peer_ids, "peer_ids") _keep("_PEER_IDS", peer_ids, "peer_ids")
if subscriber_ids: if subscriber_ids:
sub = dict(subscriber_ids) # In place: a copy would keep a second 300k-entry table alive.
sub[900999] = "D-APRS" subscriber_ids[900999] = "D-APRS"
sub[4294967295] = "SC" subscriber_ids[4294967295] = "SC"
config["_SUB_IDS"] = sub config["_SUB_IDS"] = subscriber_ids
if profiles is not None: if profiles is not None:
config["_SUB_PROFILES"] = profiles config["_SUB_PROFILES"] = profiles
elif isinstance(alias_loader, DefaultAliasLoader): elif isinstance(alias_loader, DefaultAliasLoader):
@ -407,55 +411,59 @@ class DefaultAliasLoader(AliasLoader):
pass pass
return out return out
def load_subscriber_profiles(self, config: dict[str, Any]) -> dict[int, dict[str, str]]: def load_subscriber_profiles(self, config: dict[str, Any]) -> dict[int, SubscriberProfile]:
"""Load {id: {callsign, fname, surname, talker_alias?}} from subscriber JSON files.""" """{id: profile} from the subscriber file, overlaid with the local subscriber file."""
aliases = config.get("ALIASES", {}) aliases = config.get("ALIASES", {})
path = Path(aliases.get("PATH", "./data/")).resolve() path = Path(aliases.get("PATH", "./data/")).resolve()
sub_file = aliases.get("SUBSCRIBER_FILE", "subscriber_ids.json") sub_file = aliases.get("SUBSCRIBER_FILE", "subscriber_ids.json")
local_file = aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json") local_file = aliases.get("LOCAL_SUBSCRIBER_FILE", "subscriber_ids.json")
files = [path / sub_file, path / local_file] # load_aliases has just parsed this file: this is a cache hit, not a second parse.
stamps = [_file_stamp(f) for f in files] 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") remembered = self._parsed.get("_profiles")
if remembered is not None and remembered[0] == stamps: if remembered is not None and remembered[0] is main and remembered[1] == stamp:
return remembered[1] return remembered[2]
profiles: dict[int, dict[str, str]] = {} local = self._load_subscriber_json(path / local_file)
for file_path in files: profiles = {**main, **local} if local else main
self._merge_subscriber_profiles(file_path, profiles) self._parsed["_profiles"] = (main, stamp, profiles)
self._parsed["_profiles"] = (stamps, profiles)
return 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(): 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: try:
with open(file_path, "r", encoding="utf-8") as f: with open(file_path, "r", encoding="utf-8") as f:
data = json.load(f) json.load(f, object_hook=record)
except (json.JSONDecodeError, OSError): except (json.JSONDecodeError, OSError):
return return {}
if not isinstance(data, dict): return out
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

@ -53,7 +53,7 @@ def talker_alias_config() -> dict[str, Any]:
config["GLOBAL"]["TALKER_ALIAS_SEND_DMRA"] = True config["GLOBAL"]["TALKER_ALIAS_SEND_DMRA"] = True
rid = 3120001 rid = 3120001
config["_SUB_PROFILES"] = { config["_SUB_PROFILES"] = {
rid: {"callsign": "CE5RPY", "fname": "Rodrigo", "surname": "Perez"}, rid: ("CE5RPY", "Rodrigo", "Perez", ""),
} }
config["_SUB_IDS"] = {rid: "CE5RPY"} config["_SUB_IDS"] = {rid: "CE5RPY"}
return config return config

@ -0,0 +1,102 @@
# ADN DMR Peer Server - tests infrastructure alias reload memory
#
# 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
###############################################################################
"""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

@ -133,13 +133,12 @@ def test_subscriber_profiles_are_not_rebuilt_for_unchanged_files(tmp_path: Path)
first = loader.load_subscriber_profiles(cfg) first = loader.load_subscriber_profiles(cfg)
assert 7300391 in first assert 7300391 in first
merges = [] parses = []
original = loader._merge_subscriber_profiles original = loader._load_subscriber_json
loader._merge_subscriber_profiles = lambda p, out: merges.append(p) # type: ignore[assignment] loader._load_subscriber_json = lambda p: (parses.append(p), original(p))[1] # type: ignore[assignment]
assert loader.load_subscriber_profiles(cfg) is first 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") _write(tmp_path, "subscriber_ids.json", 7300392, "CE5ABC")
os.utime(tmp_path / "subscriber_ids.json", (2_000_000_000, 2_000_000_000)) os.utime(tmp_path / "subscriber_ids.json", (2_000_000_000, 2_000_000_000))
assert 7300392 in loader.load_subscriber_profiles(cfg) assert 7300392 in loader.load_subscriber_profiles(cfg)

Loading…
Cancel
Save

Powered by TurnKey Linux.