From 3bbc17497362b9806e4ce473c27f0c863171d8b3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Tue, 22 Sep 2026 00:52:21 -0300 Subject: [PATCH] perf(alias): reload alias files only when they change The reload loop ticks every 900s so a failed download is retried soon, but STALE_DAYS replaces the files about once a day. Every tick in between re-read them: blake2b over 50MB, a JSON parse, a 50MB copy to .bak, and a 300k-entry profile rebuild, all to produce the same dicts. Parsed results are now kept against each file's (mtime_ns, size, inode) and the backup is written only when the primary has been verified, which is the point where it is worth keeping as a fallback. A primary that fails drops its remembered parse so the next tick looks at the file again. The profile build also ran in merge_reload_into_config, on the reactor thread, at over a second per cycle. It is built in the thread pool now and handed in. A tick with nothing new: 2568ms -> 0.2ms. --- .../infrastructure/bootstrap/peer_server.py | 21 ++- .../persistence/alias_loader.py | 46 ++++++- .../test_alias_reload_skips_unchanged.py | 127 ++++++++++++++++++ 3 files changed, 183 insertions(+), 11 deletions(-) create mode 100644 tests/infrastructure/test_alias_reload_skips_unchanged.py diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 831461a..00650cc 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -586,11 +586,18 @@ def run_peer_server( ) alias_reload_state["pending"] = True logger.debug("(ALIAS) starting alias thread (gen %s)", my_generation) - # Downloads run in the thread pool (blocking HTTP would stall hotspot pings); - # the config swap is applied back on the reactor thread. - d = threads.deferToThread(alias_loader.load_aliases, config) - - def _apply(loaded): + # Downloads and parsing run in the thread pool; only the config swap comes + # back to the reactor thread. + def _load_off_reactor(): + loaded = alias_loader.load_aliases(config) + if isinstance(alias_loader, DefaultAliasLoader): + return loaded, alias_loader.load_subscriber_profiles(config) + return loaded, None + + d = threads.deferToThread(_load_off_reactor) + + def _apply(result): + loaded, profiles = result if my_generation != alias_reload_state["generation"]: logger.info( "(ALIAS) discarding alias download result from a superseded attempt " @@ -600,7 +607,9 @@ def run_peer_server( ) return alias_reload_state["pending"] = False - DefaultAliasLoader.merge_reload_into_config(config, alias_loader, *loaded) + DefaultAliasLoader.merge_reload_into_config( + config, alias_loader, *loaded, profiles=profiles + ) _log_alias_health(config, config.get("SYSTEMS", {}), logger) def _fail(failure): diff --git a/src/adn_server/infrastructure/persistence/alias_loader.py b/src/adn_server/infrastructure/persistence/alias_loader.py index 1e0250e..d996c38 100644 --- a/src/adn_server/infrastructure/persistence/alias_loader.py +++ b/src/adn_server/infrastructure/persistence/alias_loader.py @@ -151,8 +151,24 @@ def _atomic_copy(src: Path, dst: Path) -> None: tmp.unlink(missing_ok=True) +def _file_stamp(file_path: Path) -> tuple[int, int, int] | None: + """Identifies a file's contents: try_download replaces, so the inode moves too.""" + try: + stat = file_path.stat() + except OSError: + return None + return (stat.st_mtime_ns, stat.st_size, stat.st_ino) + + class DefaultAliasLoader(AliasLoader): - """Load aliases from JSON files and optional downloads. Legacy mk_aliases.""" + """Load aliases from JSON files and optional downloads. Legacy mk_aliases. + + The reload loop ticks far more often than STALE_DAYS replaces the files, so + parsed results are kept against each file's stamp. + """ + + def __init__(self) -> None: + self._parsed: dict[str, tuple[Any, Any]] = {} def load_aliases( self, @@ -220,8 +236,13 @@ class DefaultAliasLoader(AliasLoader): local_subscriber_ids: dict[int, str], server_ids: dict[str, str], checksums: dict[str, str], + profiles: dict[int, dict[str, str]] | None = None, ) -> None: - """Apply alias reload without wiping in-memory tables on partial download failure.""" + """Apply alias reload without wiping in-memory tables on partial download failure. + + This runs on the reactor thread, so ``profiles`` lets the caller build the + 300k of them off it instead. + """ def _keep(key: str, new_val: dict, label: str) -> None: if new_val: config[key] = new_val @@ -238,7 +259,9 @@ class DefaultAliasLoader(AliasLoader): sub[900999] = "D-APRS" sub[4294967295] = "SC" config["_SUB_IDS"] = sub - if isinstance(alias_loader, DefaultAliasLoader): + if profiles is not None: + config["_SUB_PROFILES"] = profiles + elif isinstance(alias_loader, DefaultAliasLoader): config["_SUB_PROFILES"] = alias_loader.load_subscriber_profiles(config) elif config.get("_SUB_IDS"): logger.warning( @@ -295,6 +318,10 @@ class DefaultAliasLoader(AliasLoader): """Legacy mk_aliases peer/subscriber/tgid load with .bak fallback.""" full = path / file_name bak = path / f"{file_name}.bak" + stamp = _file_stamp(full) + remembered = self._parsed.get(file_name) + if stamp is not None and remembered is not None and remembered[0] == stamp: + return remembered[1] result: dict[int, str] = {} loaded_from_primary = False @@ -316,6 +343,7 @@ class DefaultAliasLoader(AliasLoader): result = _load_verified(full) loaded_from_primary = True except Exception as e: + self._parsed.pop(file_name, None) logger.error( "(ALIAS) ID ALIAS MAPPER: problem loading %s file (%s), falling back to .bak", name, @@ -348,6 +376,8 @@ class DefaultAliasLoader(AliasLoader): name, g, ) + if stamp is not None: + self._parsed[file_name] = (stamp, result) return result def _load_id_json(self, file_path: Path) -> dict[int, str]: @@ -431,9 +461,15 @@ class DefaultAliasLoader(AliasLoader): 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] + 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_name in (sub_file, local_file): - self._merge_subscriber_profiles(path / file_name, profiles) + for file_path in files: + self._merge_subscriber_profiles(file_path, profiles) + self._parsed["_profiles"] = (stamps, profiles) return profiles def _merge_subscriber_profiles(self, file_path: Path, out: dict[int, dict[str, str]]) -> None: diff --git a/tests/infrastructure/test_alias_reload_skips_unchanged.py b/tests/infrastructure/test_alias_reload_skips_unchanged.py new file mode 100644 index 0000000..6a6f6a2 --- /dev/null +++ b/tests/infrastructure/test_alias_reload_skips_unchanged.py @@ -0,0 +1,127 @@ +# ADN DMR Peer Server - tests infrastructure alias reload skips unchanged files +# +# 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 +############################################################################### + +"""The reload loop ticks far more often than the files behind it change. + +STALE_DAYS keeps the download to about once a day while the loop runs every few +minutes so a failed download is retried soon. Every tick in between reads the +same bytes, so it must not pay for them again. +""" + +from __future__ import annotations + +import json +import os +from pathlib import Path + +from adn_server.infrastructure.persistence.alias_loader import DefaultAliasLoader + + +def _write(path: Path, name: str, rid: int, callsign: str) -> None: + (path / name).write_text( + json.dumps({"subscribers": [{"id": rid, "callsign": callsign}]}), encoding="utf-8" + ) + + +def _load(loader: DefaultAliasLoader, path: Path, name: str) -> dict: + return loader._load_id_dict_with_backup(path, name, None, "subscriber_ids") + + +def test_an_unchanged_file_is_not_parsed_again(tmp_path: Path) -> None: + _write(tmp_path, "subscriber_ids.json", 7300391, "CE5RPY") + loader = DefaultAliasLoader() + first = _load(loader, tmp_path, "subscriber_ids.json") + + parses = [] + original = loader._load_id_json + loader._load_id_json = lambda p: (parses.append(p), original(p))[1] # type: ignore[assignment] + again = _load(loader, tmp_path, "subscriber_ids.json") + + assert parses == [] + assert again == first + + +def test_an_unchanged_file_is_not_copied_to_bak_again(tmp_path: Path) -> None: + """The backup is 50MB in production; rewriting it every tick is pure disk wear.""" + _write(tmp_path, "subscriber_ids.json", 7300391, "CE5RPY") + loader = DefaultAliasLoader() + _load(loader, tmp_path, "subscriber_ids.json") + bak = tmp_path / "subscriber_ids.json.bak" + assert bak.is_file() + + os.utime(bak, (1_000_000_000, 1_000_000_000)) + _load(loader, tmp_path, "subscriber_ids.json") + + assert bak.stat().st_mtime == 1_000_000_000 + + +def test_a_file_that_changed_is_parsed_again(tmp_path: Path) -> None: + _write(tmp_path, "subscriber_ids.json", 7300391, "CE5RPY") + loader = DefaultAliasLoader() + assert _load(loader, tmp_path, "subscriber_ids.json") == {7300391: "CE5RPY"} + + _write(tmp_path, "subscriber_ids.json", 7300392, "CE5ABC") + os.utime(tmp_path / "subscriber_ids.json", (2_000_000_000, 2_000_000_000)) + + assert _load(loader, tmp_path, "subscriber_ids.json") == {7300392: "CE5ABC"} + + +def test_a_bad_primary_is_retried_rather_than_remembered(tmp_path: Path) -> None: + """Falling back to .bak must not cache that result: the next tick has to look + at the primary again, which is how a repaired download gets picked up. + """ + name = "subscriber_ids.json" + _write(tmp_path, name, 7300391, "CE5RPY") + loader = DefaultAliasLoader() + _load(loader, tmp_path, name) + + (tmp_path / name).write_text("{ not json", encoding="utf-8") + assert _load(loader, tmp_path, name) == {7300391: "CE5RPY"} # from .bak + + _write(tmp_path, name, 7300392, "CE5ABC") + assert _load(loader, tmp_path, name) == {7300392: "CE5ABC"} + + +def test_subscriber_profiles_are_not_rebuilt_for_unchanged_files(tmp_path: Path) -> None: + """This one runs to over a second on 300k subscribers, so it is the tick cost + that matters most. + """ + _write(tmp_path, "subscriber_ids.json", 7300391, "CE5RPY") + cfg = { + "ALIASES": { + "PATH": str(tmp_path), + "SUBSCRIBER_FILE": "subscriber_ids.json", + "LOCAL_SUBSCRIBER_FILE": "subscriber_ids.json", + } + } + loader = DefaultAliasLoader() + 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] + assert loader.load_subscriber_profiles(cfg) is first + assert merges == [] + + 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)