fix: retry alias JSON downloads on a short poll instead of waiting a full STALE_DAYS cycle

pull/73/head
Rodrigo Pérez 1 week ago
parent 699af5d4a1
commit f6cc0fe281

@ -64,6 +64,11 @@ ALIASES:
TGID_URL: https://servers.adn.systems/talkgroup_ids.json
LOCAL_SUBSCRIBER_FILE: local_subcriber_ids.json
STALE_DAYS: 1
# How often (seconds) the server checks whether alias files need (re)downloading.
# Independent of STALE_DAYS: a file already fresh is skipped (no network call), but a
# failed/missing download is retried on this short cadence instead of waiting a full
# STALE_DAYS cycle. Optional, defaults to 900 (15 minutes).
POLL_INTERVAL_SEC: 900
SUB_MAP_FILE: ""
SERVER_ID_URL: https://servers.adn.systems/server_ids.tsv
SERVER_ID_FILE: server_ids.tsv

@ -171,6 +171,75 @@ def _wire_monitor_downlink_ctx(
report_factory.set_downlink_ctx_for_system(_ctx_for)
_DEFAULT_ALIAS_POLL_INTERVAL_SEC = 900.0
def _resolve_alias_poll_interval(aliases_cfg: dict[str, Any], logger: logging.Logger) -> float:
"""ALIASES.POLL_INTERVAL_SEC, defaulting to 900s. Never raises — a bad value (missing,
non-numeric, zero/negative) falls back to the default with a warning instead of
aborting server startup."""
raw = aliases_cfg.get("POLL_INTERVAL_SEC")
if raw is None or raw == "":
return _DEFAULT_ALIAS_POLL_INTERVAL_SEC
try:
value = float(raw)
except (TypeError, ValueError):
logger.warning(
"(ALIAS) invalid ALIASES.POLL_INTERVAL_SEC=%r, using default %gs",
raw,
_DEFAULT_ALIAS_POLL_INTERVAL_SEC,
)
return _DEFAULT_ALIAS_POLL_INTERVAL_SEC
if value <= 0:
logger.warning(
"(ALIAS) ALIASES.POLL_INTERVAL_SEC=%s must be positive, using default %gs",
raw,
_DEFAULT_ALIAS_POLL_INTERVAL_SEC,
)
return _DEFAULT_ALIAS_POLL_INTERVAL_SEC
return value
def _log_alias_health(config: Any, systems_cfg: dict[str, Any], logger: logging.Logger) -> None:
"""Warn loudly when peer/subscriber/server ID tables are empty and would fail-closed.
validate_id() and validate_obp_source_server_id() (udp_hbp.py) reject any ID that is
not found in _PEER_IDS/_SUB_IDS/_LOCAL_SUBSCRIBER_IDS whenever ALLOW_UNREG_ID is False,
and OBP source-server checks reject on an empty _SERVER_IDS when VALIDATE_SERVER_IDS is
True. An empty table under those flags means every registration/relay is silently
rejected — this makes that failure mode visible in the logs instead of only showing up
as "users can't connect" reports.
"""
peer_ids = config.get("_PEER_IDS") or {}
sub_ids = config.get("_SUB_IDS") or {}
local_sub_ids = config.get("_LOCAL_SUBSCRIBER_IDS") or {}
server_ids = config.get("_SERVER_IDS") or {}
logger.info(
"(ALIAS) dictionary sizes: peer_ids=%d subscriber_ids=%d local_subscriber_ids=%d "
"talkgroup_ids=%d server_ids=%d",
len(peer_ids),
len(sub_ids),
len(local_sub_ids),
len(config.get("_TG_IDS") or {}),
len(server_ids),
)
enforces_reg = any(
sys_cfg.get("ENABLED", True) and not sys_cfg.get("ALLOW_UNREG_ID", True)
for sys_cfg in systems_cfg.values()
)
if enforces_reg and not peer_ids and not sub_ids and not local_sub_ids:
logger.error(
"(ALIAS) peer_ids/subscriber_ids/local_subscriber_ids are ALL EMPTY while at "
"least one system has ALLOW_UNREG_ID: false — every peer/hotspot registration "
"will be rejected until a download succeeds or a .bak file is available"
)
if config.get("GLOBAL", {}).get("VALIDATE_SERVER_IDS") and not server_ids:
logger.error(
"(ALIAS) server_ids is EMPTY while GLOBAL.VALIDATE_SERVER_IDS is true — all "
"OpenBridge traffic with a 4-5 digit source server will be rejected"
)
def _looping_errback(logger: logging.Logger, failure):
"""Errback for LoopingCalls (legacy loopingErrHandle). Stops reactor to avoid memory leaks."""
try:
@ -210,6 +279,7 @@ def run_peer_server(
config["_LOCAL_SUBSCRIBER_IDS"] = local_subscriber_ids
config["_SERVER_IDS"] = server_ids
config["CHECKSUMS"] = checksums
_log_alias_health(config, config.get("SYSTEMS", {}), logger)
# SUB_MAP (shared mutable; used by SubMapTrimmer and shutdown)
aliases_cfg = config.get("ALIASES", {})
@ -472,26 +542,55 @@ def run_peer_server(
)
).start(66).addErrback(_looping_errback, logger)
# Alias reload (STALE_DAYS -> seconds)
alias_interval = (aliases_cfg.get("STALE_DAYS") or 1) * 86400
# Alias reload: poll every ALIAS_POLL_INTERVAL_SEC (default 900s / 15 min), independent
# from STALE_DAYS. STALE_DAYS only controls how old a cached file must be before
# try_download re-fetches it (see alias_loader.try_download); polling on a short, fixed
# cadence means a failed download (e.g. first boot with no cached files yet) is retried
# within minutes instead of waiting up to a full STALE_DAYS cycle. Once files are fresh,
# each tick is a cheap mtime check with no network call.
alias_poll_interval = _resolve_alias_poll_interval(aliases_cfg, logger)
alias_reload_state: dict[str, Any] = {"generation": 0, "pending": False}
def alias_reload_loop():
logger.debug("(ALIAS) starting alias thread")
alias_reload_state["generation"] += 1
my_generation = alias_reload_state["generation"]
if alias_reload_state["pending"]:
logger.warning(
"(ALIAS) previous alias download still running after %gs, starting a new "
"attempt (gen %s); its result will be discarded if it arrives late",
alias_poll_interval,
my_generation,
)
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)
d.addCallback(
lambda loaded: DefaultAliasLoader.merge_reload_into_config(
config, alias_loader, *loaded
)
)
d.addErrback(
lambda failure: logger.warning(
"(ALIAS) alias reload failed: %s", failure.getErrorMessage()
def _apply(loaded):
if my_generation != alias_reload_state["generation"]:
logger.info(
"(ALIAS) discarding alias download result from a superseded attempt "
"(gen %s, current gen %s)",
my_generation,
alias_reload_state["generation"],
)
return
alias_reload_state["pending"] = False
DefaultAliasLoader.merge_reload_into_config(config, alias_loader, *loaded)
_log_alias_health(config, config.get("SYSTEMS", {}), logger)
def _fail(failure):
if my_generation == alias_reload_state["generation"]:
alias_reload_state["pending"] = False
logger.warning(
"(ALIAS) alias reload failed (gen %s): %s", my_generation, failure.getErrorMessage()
)
)
task.LoopingCall(alias_reload_loop).start(alias_interval).addErrback(_looping_errback, logger)
d.addCallback(_apply)
d.addErrback(_fail)
task.LoopingCall(alias_reload_loop).start(alias_poll_interval).addErrback(_looping_errback, logger)
# SubMapTrimmer (3600s) + save
def sub_map_trimmer_loop():

@ -41,8 +41,25 @@ from ...application.ports import AliasLoader
logger = logging.getLogger(__name__)
def try_download(path: Path, file_name: str, url: str, stale_sec: float) -> str:
"""Legacy try_download: download file from url if missing or older than stale_sec. Returns result message."""
def try_download(
path: Path,
file_name: str,
url: str,
stale_sec: float,
*,
max_attempts: int = 3,
first_timeout: float = 30,
retry_timeout: float = 10,
retry_delay_sec: float = 3,
) -> str:
"""Legacy try_download: download file from url if missing or older than stale_sec.
Retries up to `max_attempts` times on network failure (timeout, connection refused,
DNS failure, etc). The first attempt uses `first_timeout`; retries use the shorter
`retry_timeout` so a fully unreachable host can't stall a caller for `max_attempts *
first_timeout` (this runs synchronously during startup, and off-thread on the
periodic reload — see alias_reload_loop in peer_server.py). Returns a result message.
"""
if not url:
return f"ID ALIAS MAPPER: '{file_name}' URL empty, not downloaded"
full = path / file_name
@ -54,14 +71,30 @@ def try_download(path: Path, file_name: str, url: str, stale_sec: float) -> str:
file_old = True
if not file_old and file_exists:
return f"ID ALIAS MAPPER: '{file_name}' is current, not downloaded"
try:
ctx = ssl.create_default_context()
ctx.check_hostname = False
ctx.verify_mode = ssl.CERT_NONE
with urlopen(url, context=ctx, timeout=30) as response:
data = response.read()
except OSError as e:
return f"ID ALIAS MAPPER: '{file_name}' could not be downloaded due to an IOError: {e}"
data: bytes | None = None
last_error: OSError | None = None
for attempt in range(1, max_attempts + 1):
timeout = first_timeout if attempt == 1 else retry_timeout
try:
ctx = ssl.create_default_context()
ctx.check_hostname = False
ctx.verify_mode = ssl.CERT_NONE
with urlopen(url, context=ctx, timeout=timeout) as response:
data = response.read()
last_error = None
break
except OSError as e:
last_error = e
if attempt < max_attempts:
logger.warning(
"(ALIAS) ID ALIAS MAPPER: '%s' download attempt %d/%d failed (%s), retrying in %gs",
file_name, attempt, max_attempts, e, retry_delay_sec,
)
if retry_delay_sec:
time.sleep(retry_delay_sec)
if last_error is not None:
return f"ID ALIAS MAPPER: '{file_name}' could not be downloaded due to an IOError: {last_error}"
if not data or data == b"{}":
return f"ID ALIAS MAPPER: '{file_name}' file not written because downloaded data is empty"
try:
@ -72,6 +105,21 @@ def try_download(path: Path, file_name: str, url: str, stale_sec: float) -> str:
return f"ID ALIAS MAPPER: '{file_name}' successfully downloaded"
_DOWNLOAD_FAILURE_MARKERS = (
"could not be downloaded",
"could not be written",
"file not written because",
)
def _log_download_result(result: str) -> None:
"""Log a try_download result at a level a monitoring/alerting pipeline can filter on."""
if any(marker in result for marker in _DOWNLOAD_FAILURE_MARKERS):
logger.error("(ALIAS) %s", result)
else:
logger.info("(ALIAS) %s", result)
def _blake2bsum(file_path: Path) -> str:
"""Blake2b hex digest of file (legacy blake2bsum)."""
h = hashlib.blake2b()
@ -102,7 +150,7 @@ class DefaultAliasLoader(AliasLoader):
if aliases.get("TRY_DOWNLOAD"):
if aliases.get("CHECKSUM_FILE") and aliases.get("CHECKSUM_URL"):
result = try_download(path, aliases["CHECKSUM_FILE"], aliases.get("CHECKSUM_URL", ""), stale_sec)
logger.info("(ALIAS) %s", result)
_log_download_result(result)
for key, url_key in [
("PEER_FILE", "PEER_URL"),
("SUBSCRIBER_FILE", "SUBSCRIBER_URL"),
@ -112,7 +160,7 @@ class DefaultAliasLoader(AliasLoader):
url = aliases.get(url_key)
if url and aliases.get(key):
result = try_download(path, aliases[key], url, stale_sec)
logger.info("(ALIAS) %s", result)
_log_download_result(result)
checksums = self._load_checksums(path, aliases.get("CHECKSUM_FILE"))
peer_file = aliases.get("PEER_FILE", "peer_ids.json")
sub_file = aliases.get("SUBSCRIBER_FILE", "subscriber_ids.json")
@ -225,7 +273,10 @@ class DefaultAliasLoader(AliasLoader):
def _load_verified(target: Path) -> dict[int, str]:
if not target.is_file():
return {}
# Raise (not return {}) so the caller falls back to .bak below instead of
# silently ending up with an empty dictionary when the primary file is
# simply missing (e.g. first boot with the download still failing).
raise FileNotFoundError(f"'{target.name}' file does not exist")
if expected_checksum:
if _blake2bsum(target) != expected_checksum:
raise ValueError("bad checksum")
@ -239,7 +290,7 @@ class DefaultAliasLoader(AliasLoader):
loaded_from_primary = True
except Exception as e:
logger.error(
"(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of %s file. not updating.: %s",
"(ALIAS) ID ALIAS MAPPER: problem loading %s file (%s), falling back to .bak",
name,
e,
)
@ -252,8 +303,15 @@ class DefaultAliasLoader(AliasLoader):
name,
f,
)
else:
logger.warning(
"(ALIAS) ID ALIAS MAPPER: no .bak available for %s, dictionary will be empty",
name,
)
if result:
logger.info("(ALIAS) ID ALIAS MAPPER: %s dictionary is available", name)
else:
logger.warning("(ALIAS) ID ALIAS MAPPER: %s dictionary is empty", name)
if loaded_from_primary and full.is_file():
try:
shutil.copy(full, bak)
@ -302,16 +360,20 @@ class DefaultAliasLoader(AliasLoader):
loaded_from_primary = False
try:
if expected_checksum and full.is_file():
if _blake2bsum(full) != expected_checksum:
raise ValueError("bad checksum")
# Raise (not just skip) on a missing primary file too, so the except block
# below falls back to .bak instead of silently ending up with an empty dict
# (e.g. first boot with the download still failing).
if not full.is_file():
raise FileNotFoundError(f"'{file_name}' file does not exist")
if expected_checksum and _blake2bsum(full) != expected_checksum:
raise ValueError("bad checksum")
result = self._load_server_tsv(path, file_name)
if full.is_file() and not result:
if not result:
raise ValueError("empty server_ids")
loaded_from_primary = bool(result)
loaded_from_primary = True
except Exception as e:
logger.error(
"(ALIAS) ID ALIAS MAPPER: problem with blake2bsum of server_ids file: %s",
"(ALIAS) ID ALIAS MAPPER: problem loading server_ids file (%s), falling back to .bak",
e,
)
if bak.is_file():
@ -324,6 +386,8 @@ class DefaultAliasLoader(AliasLoader):
)
if result:
logger.info("(ALIAS) ID ALIAS MAPPER: server_ids dictionary is available")
else:
logger.warning("(ALIAS) ID ALIAS MAPPER: server_ids dictionary is empty")
if loaded_from_primary and full.is_file():
try:
shutil.copy(full, bak)

@ -0,0 +1,54 @@
# ADN DMR Peer Server - alias dictionary health check surfaces fail-closed rejection risk
from __future__ import annotations
import logging
from adn_server.infrastructure.bootstrap.peer_server import _log_alias_health
logger = logging.getLogger(__name__)
def test_warns_when_id_tables_empty_and_registration_enforced(caplog) -> None:
config = {"_PEER_IDS": {}, "_SUB_IDS": {}, "_LOCAL_SUBSCRIBER_IDS": {}, "GLOBAL": {}}
systems_cfg = {"SYSTEM": {"ENABLED": True, "ALLOW_UNREG_ID": False}}
with caplog.at_level(logging.ERROR):
_log_alias_health(config, systems_cfg, logger)
assert any("ALL EMPTY" in r.message for r in caplog.records)
def test_no_warning_when_unregistered_ids_allowed(caplog) -> None:
config = {"_PEER_IDS": {}, "_SUB_IDS": {}, "_LOCAL_SUBSCRIBER_IDS": {}, "GLOBAL": {}}
systems_cfg = {"SYSTEM": {"ENABLED": True, "ALLOW_UNREG_ID": True}}
with caplog.at_level(logging.ERROR):
_log_alias_health(config, systems_cfg, logger)
assert not any("ALL EMPTY" in r.message for r in caplog.records)
def test_no_warning_when_id_tables_populated(caplog) -> None:
config = {"_PEER_IDS": {1: "X"}, "_SUB_IDS": {}, "_LOCAL_SUBSCRIBER_IDS": {}, "GLOBAL": {}}
systems_cfg = {"SYSTEM": {"ENABLED": True, "ALLOW_UNREG_ID": False}}
with caplog.at_level(logging.ERROR):
_log_alias_health(config, systems_cfg, logger)
assert not any("ALL EMPTY" in r.message for r in caplog.records)
def test_warns_when_server_ids_empty_and_validation_enabled(caplog) -> None:
config = {
"_PEER_IDS": {1: "X"},
"_SUB_IDS": {2: "Y"},
"_LOCAL_SUBSCRIBER_IDS": {},
"_SERVER_IDS": {},
"GLOBAL": {"VALIDATE_SERVER_IDS": True},
}
with caplog.at_level(logging.ERROR):
_log_alias_health(config, {}, logger)
assert any("server_ids is EMPTY" in r.message for r in caplog.records)
def test_disabled_system_does_not_trigger_warning(caplog) -> None:
config = {"_PEER_IDS": {}, "_SUB_IDS": {}, "_LOCAL_SUBSCRIBER_IDS": {}, "GLOBAL": {}}
systems_cfg = {"SYSTEM": {"ENABLED": False, "ALLOW_UNREG_ID": False}}
with caplog.at_level(logging.ERROR):
_log_alias_health(config, systems_cfg, logger)
assert not any("ALL EMPTY" in r.message for r in caplog.records)

@ -0,0 +1,44 @@
# ADN DMR Peer Server - ALIASES.POLL_INTERVAL_SEC resolution never aborts startup
from __future__ import annotations
import logging
from adn_server.infrastructure.bootstrap.peer_server import (
_DEFAULT_ALIAS_POLL_INTERVAL_SEC,
_resolve_alias_poll_interval,
)
logger = logging.getLogger(__name__)
def test_missing_key_defaults_to_900() -> None:
assert _resolve_alias_poll_interval({}, logger) == _DEFAULT_ALIAS_POLL_INTERVAL_SEC
def test_empty_string_defaults_to_900() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": ""}, logger) == 900.0
def test_valid_value_is_used() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": 60}, logger) == 60.0
def test_valid_string_value_is_parsed() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": "120"}, logger) == 120.0
def test_non_numeric_value_falls_back_without_raising() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": "900s"}, logger) == 900.0
def test_zero_falls_back_without_raising() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": 0}, logger) == 900.0
def test_negative_falls_back_without_raising() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": -5}, logger) == 900.0
def test_none_type_falls_back_without_raising() -> None:
assert _resolve_alias_poll_interval({"POLL_INTERVAL_SEC": None}, logger) == 900.0

@ -4,6 +4,7 @@ from __future__ import annotations
import json
from pathlib import Path
from unittest.mock import patch
from adn_server.infrastructure.persistence.alias_loader import DefaultAliasLoader, try_download
@ -17,12 +18,80 @@ def test_try_download_failure_does_not_erase_existing_file(tmp_path: Path) -> No
file_name = "subscriber_ids.json"
full = tmp_path / file_name
full.write_bytes(b'{"subscribers":[{"id":7300391,"callsign":"CE5RPY"}]}')
# stale_sec=0 forces download attempt; bad URL simulates selfcare down
result = try_download(tmp_path, file_name, "http://127.0.0.1:1/nope.json", stale_sec=0)
# stale_sec=0 forces download attempt; bad URL simulates selfcare down.
# max_attempts=1 keeps the test fast; retry behavior is covered separately below.
result = try_download(
tmp_path, file_name, "http://127.0.0.1:1/nope.json", stale_sec=0, max_attempts=1
)
assert "could not be downloaded" in result or "IOError" in result
assert full.read_bytes().startswith(b"{")
def test_try_download_retries_and_recovers_from_transient_failure(tmp_path: Path) -> None:
file_name = "peer_ids.json"
calls: list[float] = []
def _flaky_urlopen(url, context=None, timeout=None):
calls.append(timeout)
if len(calls) < 3:
raise OSError("connection refused")
return _FakeResponse(b'{"peers":[{"id":1,"callsign":"X"}]}')
with patch(
"adn_server.infrastructure.persistence.alias_loader.urlopen", side_effect=_flaky_urlopen
):
result = try_download(
tmp_path,
file_name,
"https://example.invalid/peer_ids.json",
stale_sec=0,
max_attempts=3,
retry_delay_sec=0,
)
assert "successfully downloaded" in result
assert len(calls) == 3
# first attempt uses the longer timeout, retries use the shorter one
assert calls[0] == 30
assert calls[1] == 10
def test_try_download_gives_up_after_max_attempts(tmp_path: Path) -> None:
file_name = "peer_ids.json"
calls: list[float] = []
def _always_fails(url, context=None, timeout=None):
calls.append(timeout)
raise OSError("connection refused")
with patch(
"adn_server.infrastructure.persistence.alias_loader.urlopen", side_effect=_always_fails
):
result = try_download(
tmp_path,
file_name,
"https://example.invalid/peer_ids.json",
stale_sec=0,
max_attempts=3,
retry_delay_sec=0,
)
assert "could not be downloaded" in result
assert len(calls) == 3
class _FakeResponse:
def __init__(self, data: bytes) -> None:
self._data = data
def read(self) -> bytes:
return self._data
def __enter__(self) -> "_FakeResponse":
return self
def __exit__(self, *exc: object) -> None:
return None
def test_merge_reload_keeps_previous_sub_ids_on_empty_reload() -> None:
loader = DefaultAliasLoader()
config = {
@ -56,3 +125,28 @@ def test_load_id_dict_with_backup_uses_bak_on_checksum_mismatch(tmp_path: Path)
"subscriber_ids",
)
assert loaded.get(7300391) == "GOOD"
def test_load_id_dict_with_backup_uses_bak_when_primary_missing(tmp_path: Path) -> None:
loader = DefaultAliasLoader()
file_name = "subscriber_ids.json"
# No primary file at all (e.g. first boot, download never succeeded), only a .bak
# from a previous successful run.
_write_subscriber_file(tmp_path, f"{file_name}.bak", 7300391, "GOOD")
loaded = loader._load_id_dict_with_backup(
tmp_path,
file_name,
None,
"subscriber_ids",
)
assert loaded.get(7300391) == "GOOD"
def test_load_server_tsv_with_backup_uses_bak_when_primary_missing(tmp_path: Path) -> None:
loader = DefaultAliasLoader()
file_name = "server_ids.tsv"
(tmp_path / f"{file_name}.bak").write_text(
"OPB Net ID\tCountry\n1234\tChile\n", encoding="utf-8"
)
loaded = loader._load_server_tsv_with_backup(tmp_path, file_name, None)
assert loaded.get("1234") == "Chile"

Loading…
Cancel
Save

Powered by TurnKey Linux.