From 915aacfa82dde1953687ce477ef3b75a9f42a9e9 Mon Sep 17 00:00:00 2001 From: yo Date: Wed, 23 Sep 2026 23:09:02 +0200 Subject: [PATCH] fix(sub_map): save routes within minutes of a change, and atomically SUB_MAP was written once an hour. The shutdown save only runs if SIGTERM reaches the server, which it does not under a Docker entrypoint that stays PID 1, so a restart dropped up to an hour of routes. Private calls to those units then go nowhere (FORWARD: []) until each radio transmits again. The loop now ticks every 300s and writes only when a route (system, slot, peer) changed or an entry's timestamp moved to a new hour, so the 24h trim still sees fresh times after a restart. An idle master does not write at all. The pickle is written to a .tmp and renamed. A crash mid-dump used to leave a truncated file, which load() reads as an empty map. Co-Authored-By: Claude Opus 5.5 --- .../infrastructure/bootstrap/peer_server.py | 19 +++-- .../infrastructure/persistence/__init__.py | 4 +- .../persistence/sub_map_store.py | 42 +++++++++- .../test_sub_map_save_on_change.py | 80 +++++++++++++++++++ 4 files changed, 134 insertions(+), 11 deletions(-) create mode 100644 tests/infrastructure/test_sub_map_save_on_change.py diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 00650cc..cb0610d 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -79,7 +79,7 @@ from adn_server.infrastructure.config_normalizer import ( ) from adn_server.infrastructure.config_reload import BindSpec, reload_server_config from adn_server.infrastructure.logging_config import reopen_file_handlers -from adn_server.infrastructure.persistence import PickleSubMapStore +from adn_server.infrastructure.persistence import PickleSubMapStore, SubMapSaver from adn_server.infrastructure.persistence.alias_loader import DefaultAliasLoader from adn_server.infrastructure.persistence.database_config import database_settings from adn_server.infrastructure.persistence.dynamic_tg_repository import MysqlDynamicTgRepository @@ -295,6 +295,7 @@ def run_peer_server( sub_map_store = PickleSubMapStore() sub_map = sub_map_store.load(sub_map_path) config["_SUB_MAP"] = sub_map + sub_map_saver = SubMapSaver(sub_map_store, sub_map_path, sub_map) # Generator: expand MASTER systems with GENERATOR > 1 into SYSTEM-0, SYSTEM-1, ... (legacy) _expand_generator(config, logger) @@ -624,7 +625,9 @@ def run_peer_server( task.LoopingCall(alias_reload_loop).start(alias_poll_interval).addErrback(_looping_errback, logger) - # SubMapTrimmer (3600s) + save + # SubMapTrimmer + save. Every 300s, but it writes only when a route changed: a + # restart used to lose up to an hour of routes (Docker rarely delivers the + # signal that triggers the shutdown save), leaving private calls with no target. def sub_map_trimmer_loop(): logger.debug("(SUBSCRIBER) Subscriber Map trimmer loop started") now = time.time() @@ -633,12 +636,12 @@ def run_peer_server( sub_map.pop(k, None) if aliases_cfg.get("SUB_MAP_FILE"): try: - sub_map_store.save(sub_map_path, sub_map) - logger.info("(SUBSCRIBER) Writing SUB_MAP to disk") + if sub_map_saver.save_if_changed(): + logger.info("(SUBSCRIBER) Writing SUB_MAP to disk") except Exception as e: logger.warning("(SUBSCRIBER) Cannot write SUB_MAP to file: %s", e) - task.LoopingCall(sub_map_trimmer_loop).start(3600).addErrback(_looping_errback, logger) + task.LoopingCall(sub_map_trimmer_loop).start(300).addErrback(_looping_errback, logger) # Kill switch + shutdown (legacy kill_server every 5s; SIGTERM/SIGINT trigger _KILL_SERVER) config.setdefault("GLOBAL", {})["_KILL_SERVER"] = False @@ -653,7 +656,7 @@ def run_peer_server( if reactor.running: reactor.stop() if aliases_cfg.get("SUB_MAP_FILE"): - sub_map_store.save(sub_map_path, sub_map) + sub_map_saver.save() try: keys_store.save(keys_path, keys) logger.info("(KEYS) saved system keys to keystore") @@ -668,7 +671,7 @@ def run_peer_server( """On reactor shutdown: save SUB_MAP and keys.""" if aliases_cfg.get("SUB_MAP_FILE"): try: - sub_map_store.save(sub_map_path, config["_SUB_MAP"]) + sub_map_saver.save() logger.info("(SUBSCRIBER) Writing SUB_MAP to disk (shutdown)") except Exception as e: logger.warning("(SUBSCRIBER) Cannot write SUB_MAP on shutdown: %s", e) @@ -826,7 +829,7 @@ def run_peer_server( logger.info("(CONFIG-RELOAD) SIGHUP received, scheduling reload") if aliases_cfg.get("SUB_MAP_FILE"): try: - sub_map_store.save(sub_map_path, sub_map) + sub_map_saver.save() logger.info("(SUBSCRIBER) Writing SUB_MAP to disk (SIGHUP)") except Exception as e: logger.warning("(SUBSCRIBER) Cannot write SUB_MAP to file: %s", e) diff --git a/src/adn_server/infrastructure/persistence/__init__.py b/src/adn_server/infrastructure/persistence/__init__.py index 50f004e..6fc71db 100644 --- a/src/adn_server/infrastructure/persistence/__init__.py +++ b/src/adn_server/infrastructure/persistence/__init__.py @@ -22,6 +22,6 @@ ############################################################################### from .keys_store import JsonKeysStore -from .sub_map_store import PickleSubMapStore +from .sub_map_store import PickleSubMapStore, SubMapSaver -__all__ = ["PickleSubMapStore", "JsonKeysStore"] +__all__ = ["PickleSubMapStore", "SubMapSaver", "JsonKeysStore"] diff --git a/src/adn_server/infrastructure/persistence/sub_map_store.py b/src/adn_server/infrastructure/persistence/sub_map_store.py index fe71f59..6ff27ab 100644 --- a/src/adn_server/infrastructure/persistence/sub_map_store.py +++ b/src/adn_server/infrastructure/persistence/sub_map_store.py @@ -25,7 +25,9 @@ from __future__ import annotations +import os import pickle +from collections.abc import Hashable from pathlib import Path from ...application.ports import SubMapStore @@ -49,5 +51,43 @@ class PickleSubMapStore(SubMapStore): """Save SUB_MAP to pickle file.""" p = Path(path) p.parent.mkdir(parents=True, exist_ok=True) - with open(p, "wb") as f: + # Write aside and rename: a crash mid-dump must not leave a truncated + # file, which load() would turn into an empty map. + tmp = p.with_name(p.name + ".tmp") + with open(tmp, "wb") as f: pickle.dump(sub_map, f) + os.replace(tmp, p) + + +# SUB_MAP is rewritten on every frame, but only the route of an entry (system, +# slot, peer) matters after a restart. The timestamp only feeds the 24h trim, so +# it counts as changed once per hour, as often as the old hourly save wrote it. +_TIME_BUCKET_S = 3600 + + +def _route_snapshot(sub_map: dict) -> dict[bytes, Hashable]: + return {k: (v[0], v[1], v[3] if len(v) > 3 else None, int(v[2] // _TIME_BUCKET_S)) for k, v in sub_map.items()} + + +class SubMapSaver: + """Write SUB_MAP to disk when its routes changed since the last write.""" + + def __init__(self, store: SubMapStore, path: str, sub_map: dict): + self._store = store + self._path = path + self._sub_map = sub_map + self._saved = _route_snapshot(sub_map) + + def save(self) -> None: + """Write unconditionally (shutdown, SIGHUP).""" + self._store.save(self._path, self._sub_map) + self._saved = _route_snapshot(self._sub_map) + + def save_if_changed(self) -> bool: + """Write only if a route changed; True when it wrote.""" + snapshot = _route_snapshot(self._sub_map) + if snapshot == self._saved: + return False + self._store.save(self._path, self._sub_map) + self._saved = snapshot + return True diff --git a/tests/infrastructure/test_sub_map_save_on_change.py b/tests/infrastructure/test_sub_map_save_on_change.py new file mode 100644 index 0000000..f14c5f8 --- /dev/null +++ b/tests/infrastructure/test_sub_map_save_on_change.py @@ -0,0 +1,80 @@ +"""SUB_MAP is written when a route changes, not on every tick of the save loop.""" + +from __future__ import annotations + +import pickle + +from adn_server.infrastructure.persistence import PickleSubMapStore, SubMapSaver + + +class _CountingStore(PickleSubMapStore): + def __init__(self): + self.writes = 0 + + def save(self, path, sub_map): + self.writes += 1 + super().save(path, sub_map) + + +def _saver(tmp_path, sub_map): + store = _CountingStore() + return store, SubMapSaver(store, str(tmp_path / "sub_map.pkl"), sub_map) + + +def test_loaded_map_is_not_rewritten(tmp_path): + sub_map = {b"\x00\x00\x01": ("MASTER-1", 2, 1000.0, b"\x00\x00\x00\x01")} + store, saver = _saver(tmp_path, sub_map) + assert saver.save_if_changed() is False + assert store.writes == 0 + + +def test_same_route_newer_frame_is_not_a_change(tmp_path): + sub_map = {b"\x00\x00\x01": ("MASTER-1", 2, 1000.0, b"\x00\x00\x00\x01")} + store, saver = _saver(tmp_path, sub_map) + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 2, 1030.0, b"\x00\x00\x00\x01") + assert saver.save_if_changed() is False + + +def test_new_unit_moved_unit_and_trim_are_changes(tmp_path): + sub_map = {} + store, saver = _saver(tmp_path, sub_map) + + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 2, 1000.0, b"\x00\x00\x00\x01") + assert saver.save_if_changed() is True + + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 2, 1010.0, b"\x00\x00\x00\x02") # other hotspot + assert saver.save_if_changed() is True + + del sub_map[b"\x00\x00\x01"] + assert saver.save_if_changed() is True + assert saver.save_if_changed() is False + assert store.writes == 3 + + +def test_timestamp_is_refreshed_once_per_hour(tmp_path): + # Otherwise a unit that never moves keeps its first saved time and a + # restart would trim it while it is still active. + sub_map = {b"\x00\x00\x01": ("MASTER-1", 2, 3600.0, None)} + store, saver = _saver(tmp_path, sub_map) + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 2, 7199.0, None) + assert saver.save_if_changed() is False + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 2, 7200.0, None) + assert saver.save_if_changed() is True + + +def test_forced_save_resets_the_baseline(tmp_path): + sub_map = {} + store, saver = _saver(tmp_path, sub_map) + sub_map[b"\x00\x00\x01"] = ("MASTER-1", 1, 1000.0) # legacy 3-tuple + saver.save() + assert saver.save_if_changed() is False + assert store.writes == 1 + assert store.load(str(tmp_path / "sub_map.pkl")) == sub_map + + +def test_save_leaves_no_temp_file_and_replaces_atomically(tmp_path): + path = tmp_path / "sub_map.pkl" + path.write_bytes(pickle.dumps({b"old": ("X", 1, 0.0)})) + PickleSubMapStore().save(str(path), {b"new": ("Y", 2, 1.0)}) + assert [p.name for p in tmp_path.iterdir()] == ["sub_map.pkl"] + assert PickleSubMapStore().load(str(path)) == {b"new": ("Y", 2, 1.0)}