Merge pull request #98 from pyopower/fix/sub-map-save-on-change

fix(sub_map): save routes within minutes of a change, and atomically
pull/102/head
ce5rpy 4 days ago committed by GitHub
commit ef8d922bc0
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

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

@ -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"]

@ -25,11 +25,19 @@
from __future__ import annotations
import itertools
import os
import pickle
import threading
from collections.abc import Hashable
from pathlib import Path
from ...application.ports import SubMapStore
# A SIGHUP handler runs on the main thread, so it can interrupt a save already in
# progress there: pid and thread alone would give both writers the same temp file.
_tmp_seq = itertools.count()
class PickleSubMapStore(SubMapStore):
"""Persist SUB_MAP as pickle (legacy compatible)."""
@ -49,5 +57,46 @@ 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:
pickle.dump(sub_map, 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(f"{p.name}.tmp.{os.getpid()}.{threading.get_ident()}.{next(_tmp_seq)}")
try:
with open(tmp, "wb") as f:
pickle.dump(sub_map, f)
tmp.replace(p)
finally:
tmp.unlink(missing_ok=True)
# 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

@ -0,0 +1,97 @@
"""SUB_MAP is written when a route changes, not on every tick of the save loop."""
from __future__ import annotations
import pickle
import pytest
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)}
def test_failed_dump_leaves_no_temp_file(tmp_path, monkeypatch):
path = tmp_path / "sub_map.pkl"
PickleSubMapStore().save(str(path), {b"\x00\x00\x01": ("SYS", 1, 1.0)})
def boom(*_a, **_k):
raise pickle.PicklingError("boom")
monkeypatch.setattr(pickle, "dump", boom)
with pytest.raises(pickle.PicklingError):
PickleSubMapStore().save(str(path), {b"\x00\x00\x02": ("SYS", 2, 2.0)})
assert sorted(p.name for p in tmp_path.iterdir()) == ["sub_map.pkl"]
assert PickleSubMapStore().load(str(path)) == {b"\x00\x00\x01": ("SYS", 1, 1.0)}
Loading…
Cancel
Save

Powered by TurnKey Linux.