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.
pull/94/head
Rodrigo Pérez 7 days ago
parent 59bb95afdb
commit 3bbc174973

@ -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):

@ -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:

@ -0,0 +1,127 @@
# ADN DMR Peer Server - tests infrastructure alias reload skips unchanged files
#
# 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
###############################################################################
"""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)
Loading…
Cancel
Save

Powered by TurnKey Linux.