feat: best-effort mirror (dual-send) of heartbeats + events
Optional second destination so each heartbeat and event is posted to a mirror server (the office NAS) alongside the primary — dual-write to de-risk the migration. Default off (blank mirror URLs); existing installs unchanged. - event_forwarder: mirror_reachable() fast-fail probe + mirror_forward_pass() (reliable event mirror with its OWN state file, total exception isolation, never raises into the primary path). - series4_ingest: mirror_api_url/mirror_sfm_url/mirror_sfm_state_file config; best-effort heartbeat mirror; isolated event-mirror pass after the primary. - settings dialog: new 'Mirror' tab. - 9 new tests incl. the isolation invariant (down mirror = fast no-op that never touches primary state). 42 passing. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+5
-1
@@ -24,5 +24,9 @@
|
||||
"sfm_http_timeout": 60,
|
||||
"sfm_state_file": "",
|
||||
"sfm_max_forwards_per_pass": 500,
|
||||
"sfm_max_event_age_days": 365
|
||||
"sfm_max_event_age_days": 365,
|
||||
|
||||
"mirror_api_url": "",
|
||||
"mirror_sfm_url": "",
|
||||
"mirror_sfm_state_file": ""
|
||||
}
|
||||
|
||||
@@ -692,6 +692,74 @@ def forward_pending(
|
||||
return counts
|
||||
|
||||
|
||||
# ── Mirror (dual-send) ────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def mirror_reachable(base_url: str, timeout: float = 3.0) -> bool:
|
||||
"""Quick liveness probe of a mirror SFM base URL.
|
||||
|
||||
Run before a mirror forward pass so an unreachable mirror can't stall
|
||||
the watcher loop on a per-event HTTP timeout for every pending file.
|
||||
Probes ``<base>/health`` (the SFM server exposes it). ANY HTTP
|
||||
response — even 404/500 — means the server is up, so proceed; only a
|
||||
connection-level failure (refused / timed out / DNS) counts as down.
|
||||
Best-effort: never raises.
|
||||
"""
|
||||
if not base_url:
|
||||
return False
|
||||
url = base_url.rstrip("/") + "/health"
|
||||
try:
|
||||
with urllib.request.urlopen(
|
||||
urllib.request.Request(url, method="GET"), timeout=timeout
|
||||
):
|
||||
return True
|
||||
except urllib.error.HTTPError:
|
||||
return True # got a status back → server is alive
|
||||
except Exception:
|
||||
return False # refused / timeout / DNS / socket → down
|
||||
|
||||
|
||||
def mirror_forward_pass(
|
||||
watch_dir: str,
|
||||
mirror_url: str,
|
||||
mirror_state: ForwardState,
|
||||
*,
|
||||
reachable_fn=mirror_reachable,
|
||||
reachable_timeout: float = 3.0,
|
||||
**forward_kwargs: Any,
|
||||
) -> Optional[Dict[str, int]]:
|
||||
"""One best-effort forwarding pass against a *mirror* SFM server.
|
||||
|
||||
The dual-send entry point. Wraps :func:`forward_pending` with two
|
||||
properties that keep the mirror from ever harming the primary path:
|
||||
|
||||
1. **Fast-fail reachability guard** — probe ``mirror_url`` first; if
|
||||
it's down, return ``None`` immediately rather than let
|
||||
:func:`forward_pending` block on a per-event timeout for every
|
||||
pending file.
|
||||
2. **Total exception isolation** — any error (probe, forward, state
|
||||
I/O) is swallowed and reported as ``None``. This function NEVER
|
||||
raises, so the caller's already-completed primary forward is
|
||||
untouched.
|
||||
|
||||
The mirror keeps its OWN ``ForwardState`` file, so its idempotency
|
||||
and retry are independent of the primary's. Nothing is lost while
|
||||
the mirror is down — skipped events stay pending and are delivered
|
||||
on a later pass once it returns.
|
||||
|
||||
Returns the :func:`forward_pending` counts dict on a completed pass,
|
||||
or ``None`` if the mirror was unreachable or the pass raised.
|
||||
"""
|
||||
try:
|
||||
if not mirror_url:
|
||||
return None
|
||||
if not reachable_fn(mirror_url, reachable_timeout):
|
||||
return None
|
||||
return forward_pending(watch_dir, mirror_url, mirror_state, **forward_kwargs)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
# ── Seed-state mode (skip historical backfill on first deploy) ────────────────
|
||||
|
||||
|
||||
|
||||
@@ -69,6 +69,15 @@ def load_config(config_path: str) -> Dict[str, Any]:
|
||||
"sfm_state_file": "", # blank → <log_dir>/thor_forwarded.json
|
||||
"sfm_max_forwards_per_pass": 500,
|
||||
"sfm_max_event_age_days": 365,
|
||||
|
||||
# Mirror (dual-send) — best-effort second destination, default OFF.
|
||||
# See docs/mirror-dual-send-design.md. Heartbeat-mirror fires iff
|
||||
# api_url AND mirror_api_url are set; event-mirror runs iff SFM
|
||||
# forwarding is on AND mirror_sfm_url is set. Own state file →
|
||||
# idempotency independent of the primary.
|
||||
"mirror_api_url": "",
|
||||
"mirror_sfm_url": "",
|
||||
"mirror_sfm_state_file": "", # blank → <log_dir>/thor_forwarded_mirror.json
|
||||
}
|
||||
|
||||
with open(config_path, "r", encoding="utf-8") as f:
|
||||
@@ -329,6 +338,12 @@ def run_watcher(state: Dict[str, Any], stop_event: threading.Event) -> None:
|
||||
sfm_state_path = str(cfg.get("sfm_state_file", "")).strip() or \
|
||||
os.path.join(state["log_dir"], "thor_forwarded.json")
|
||||
|
||||
# Mirror (dual-send) config
|
||||
MIRROR_API_URL = str(cfg.get("mirror_api_url", "")).strip()
|
||||
MIRROR_SFM_URL = str(cfg.get("mirror_sfm_url", "")).strip()
|
||||
mirror_state_path = str(cfg.get("mirror_sfm_state_file", "")).strip() or \
|
||||
os.path.join(state["log_dir"], "thor_forwarded_mirror.json")
|
||||
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[cfg] THORDATA_PATH={} SCAN_INTERVAL={}s API_INTERVAL={}s API={} SFM={}".format(
|
||||
THORDATA_PATH, SCAN_INTERVAL, API_INTERVAL, bool(API_URL),
|
||||
@@ -362,6 +377,24 @@ def run_watcher(state: Dict[str, Any], stop_event: threading.Event) -> None:
|
||||
else:
|
||||
state["sfm_status"] = "disabled"
|
||||
|
||||
# Mirror event-forward state — independent of the primary's. Created
|
||||
# only when the primary forwarder is on AND a mirror URL is set.
|
||||
mirror_state_obj: Optional[event_forwarder.ForwardState] = None
|
||||
if sfm_state_obj is not None and MIRROR_SFM_URL:
|
||||
try:
|
||||
mirror_state_obj = event_forwarder.ForwardState(mirror_state_path)
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] event mirror ready url={} state={} known={}".format(
|
||||
MIRROR_SFM_URL, mirror_state_path, mirror_state_obj.count()))
|
||||
print("[MIRROR] event mirror ready url={}".format(MIRROR_SFM_URL))
|
||||
except Exception as exc:
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] event mirror init failed: {}".format(exc))
|
||||
mirror_state_obj = None
|
||||
if MIRROR_API_URL:
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] heartbeat mirror enabled url={}".format(MIRROR_API_URL))
|
||||
|
||||
state["last_forward"] = None
|
||||
state["last_forward_counts"] = None
|
||||
|
||||
@@ -419,6 +452,13 @@ def run_watcher(state: Dict[str, Any], stop_event: threading.Event) -> None:
|
||||
payload["log_tail"] = _read_log_tail(log_file, 25)
|
||||
response = send_api_payload(payload, API_URL, API_TIMEOUT)
|
||||
last_api_ts = now_ts
|
||||
# Best-effort heartbeat mirror — never affects primary.
|
||||
if MIRROR_API_URL:
|
||||
try:
|
||||
send_api_payload(payload, MIRROR_API_URL, API_TIMEOUT)
|
||||
except Exception as exc:
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] heartbeat post failed: {}".format(exc))
|
||||
if response is not None:
|
||||
state["api_status"] = "ok"
|
||||
if response.get("update_available"):
|
||||
@@ -459,6 +499,25 @@ def run_watcher(state: Dict[str, Any], stop_event: threading.Event) -> None:
|
||||
print(msg)
|
||||
log_message(log_file, ENABLE_LOGGING, msg)
|
||||
|
||||
# Best-effort event mirror — own state, fast-fail guard,
|
||||
# fully isolated. Runs after the primary each tick.
|
||||
if mirror_state_obj is not None:
|
||||
m_counts = event_forwarder.mirror_forward_pass(
|
||||
THORDATA_PATH, MIRROR_SFM_URL, mirror_state_obj,
|
||||
max_age_days=SFM_MAX_AGE_DAYS,
|
||||
quiescence_seconds=SFM_QUIESCENCE,
|
||||
missing_report_grace_seconds=SFM_GRACE,
|
||||
timeout=SFM_HTTP_TIMEOUT,
|
||||
max_per_pass=SFM_MAX_PER_PASS,
|
||||
logger=lambda m: log_message(log_file, ENABLE_LOGGING, "[mirror] " + m),
|
||||
)
|
||||
if m_counts is None:
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] forward skipped (unreachable/error): {}".format(MIRROR_SFM_URL))
|
||||
else:
|
||||
log_message(log_file, ENABLE_LOGGING,
|
||||
"[mirror] pass forwarded={forwarded} errors={errors}".format(**m_counts))
|
||||
|
||||
except Exception as e:
|
||||
err = "[loop-error] {}".format(e)
|
||||
print(err)
|
||||
|
||||
@@ -712,5 +712,134 @@ class TestForwardPending(unittest.TestCase):
|
||||
self.assertEqual(len(_FakeImportHandler.received), 2)
|
||||
|
||||
|
||||
# ── Mirror (dual-send) ───────────────────────────────────────────────────────
|
||||
|
||||
|
||||
class TestMirrorReachable(unittest.TestCase):
|
||||
|
||||
def test_empty_url_is_unreachable(self):
|
||||
self.assertFalse(ef.mirror_reachable("", timeout=1.0))
|
||||
|
||||
def test_dead_port_is_unreachable_and_fast(self):
|
||||
# Nothing listens on 127.0.0.1:1 → connection refused → False,
|
||||
# and it must NOT hang (this is the whole point of the guard).
|
||||
t0 = time.time()
|
||||
self.assertFalse(ef.mirror_reachable("http://127.0.0.1:1", timeout=2.0))
|
||||
self.assertLess(time.time() - t0, 2.5)
|
||||
|
||||
def test_any_http_response_counts_as_reachable(self):
|
||||
# The fake server 501s on GET /health, but it's UP — so reachable.
|
||||
server, base = _start_fake_server()
|
||||
try:
|
||||
self.assertTrue(ef.mirror_reachable(base, timeout=2.0))
|
||||
finally:
|
||||
server.shutdown()
|
||||
server.server_close()
|
||||
|
||||
|
||||
class TestMirrorForwardPass(unittest.TestCase):
|
||||
"""The dual-send entry point: best-effort, isolated, own state."""
|
||||
|
||||
def setUp(self):
|
||||
_FakeImportHandler.received = []
|
||||
self.server, self.base_url = _start_fake_server()
|
||||
|
||||
def tearDown(self):
|
||||
self.server.shutdown()
|
||||
self.server.server_close()
|
||||
|
||||
def _one_event(self, root: Path) -> None:
|
||||
unit = _make_thordata(root, "Project A", "UM11719")
|
||||
_make_event(unit, "UM11719_20231219163444.IDFW", age_seconds=200, content=b"binary")
|
||||
_make_txt(unit, "UM11719_20231219163444.IDFW", age_seconds=100, content=b"report")
|
||||
|
||||
def test_empty_mirror_url_is_noop(self):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
self._one_event(root)
|
||||
mstate = ef.ForwardState(str(root / "mirror.json"))
|
||||
self.assertIsNone(ef.mirror_forward_pass(str(root), "", mstate, max_age_days=30))
|
||||
self.assertEqual(len(_FakeImportHandler.received), 0)
|
||||
|
||||
def test_skips_when_unreachable_without_posting(self):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
self._one_event(root)
|
||||
mstate = ef.ForwardState(str(root / "mirror.json"))
|
||||
result = ef.mirror_forward_pass(
|
||||
str(root), self.base_url, mstate,
|
||||
reachable_fn=lambda url, timeout=3.0: False, # force "down"
|
||||
max_age_days=30, quiescence_seconds=5,
|
||||
missing_report_grace_seconds=60, timeout=5.0,
|
||||
)
|
||||
self.assertIsNone(result)
|
||||
self.assertEqual(mstate.count(), 0) # nothing forwarded
|
||||
self.assertEqual(len(_FakeImportHandler.received), 0) # never POSTed
|
||||
|
||||
def test_forwards_when_reachable_using_its_own_state(self):
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
self._one_event(root)
|
||||
mstate = ef.ForwardState(str(root / "mirror.json"))
|
||||
counts = ef.mirror_forward_pass(
|
||||
str(root), self.base_url, mstate,
|
||||
reachable_fn=lambda url, timeout=3.0: True,
|
||||
max_age_days=30, quiescence_seconds=5,
|
||||
missing_report_grace_seconds=60, timeout=5.0,
|
||||
)
|
||||
self.assertIsNotNone(counts)
|
||||
self.assertEqual(counts["forwarded"], 1)
|
||||
self.assertEqual(mstate.count(), 1)
|
||||
self.assertEqual(len(_FakeImportHandler.received), 1)
|
||||
|
||||
def test_never_raises_when_forward_blows_up(self):
|
||||
# Even if forward_pending raises, the mirror swallows it → None.
|
||||
def _boom(*a, **k):
|
||||
raise RuntimeError("boom")
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
self._one_event(root)
|
||||
mstate = ef.ForwardState(str(root / "mirror.json"))
|
||||
orig = ef.forward_pending
|
||||
ef.forward_pending = _boom
|
||||
try:
|
||||
result = ef.mirror_forward_pass(
|
||||
str(root), self.base_url, mstate,
|
||||
reachable_fn=lambda url, timeout=3.0: True,
|
||||
max_age_days=30,
|
||||
)
|
||||
finally:
|
||||
ef.forward_pending = orig
|
||||
self.assertIsNone(result)
|
||||
|
||||
def test_down_mirror_leaves_primary_state_untouched(self):
|
||||
# The invariant: a real down mirror (via the real reachability
|
||||
# probe) is a fast no-op and the primary forward is unaffected.
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
self._one_event(root)
|
||||
|
||||
primary_state = ef.ForwardState(str(root / "primary.json"))
|
||||
pcounts = ef.forward_pending(
|
||||
str(root), self.base_url, primary_state,
|
||||
max_age_days=30, quiescence_seconds=5,
|
||||
missing_report_grace_seconds=60, timeout=5.0,
|
||||
)
|
||||
self.assertEqual(pcounts["forwarded"], 1)
|
||||
primary_snapshot = primary_state.count()
|
||||
|
||||
mstate = ef.ForwardState(str(root / "mirror.json"))
|
||||
t0 = time.time()
|
||||
result = ef.mirror_forward_pass(
|
||||
str(root), "http://127.0.0.1:1", mstate, # nothing listens
|
||||
max_age_days=30, quiescence_seconds=5,
|
||||
missing_report_grace_seconds=60, timeout=5.0,
|
||||
)
|
||||
self.assertIsNone(result)
|
||||
self.assertLess(time.time() - t0, 3.5) # fast-failed, didn't hang
|
||||
self.assertEqual(mstate.count(), 0)
|
||||
self.assertEqual(primary_state.count(), primary_snapshot) # untouched
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -47,6 +47,11 @@ DEFAULTS = {
|
||||
"sfm_state_file": "",
|
||||
"sfm_max_forwards_per_pass": 500,
|
||||
"sfm_max_event_age_days": 365,
|
||||
|
||||
# Mirror (dual-send) — best-effort second destination, default OFF.
|
||||
"mirror_api_url": "",
|
||||
"mirror_sfm_url": "",
|
||||
"mirror_sfm_state_file": "",
|
||||
}
|
||||
|
||||
|
||||
@@ -220,6 +225,14 @@ class SettingsDialog:
|
||||
self.var_sfm_max_age_days = tk.StringVar(value=str(v.get("sfm_max_event_age_days", 365)))
|
||||
self.var_sfm_state_file = tk.StringVar(value=str(v.get("sfm_state_file", "")))
|
||||
|
||||
# Mirror (dual-send) — best-effort second destination
|
||||
raw_mirror = str(v.get("mirror_api_url", ""))
|
||||
if raw_mirror.endswith(_suffix):
|
||||
raw_mirror = raw_mirror[:-len(_suffix)]
|
||||
self.var_mirror_api_url = tk.StringVar(value=raw_mirror)
|
||||
self.var_mirror_sfm_url = tk.StringVar(value=str(v.get("mirror_sfm_url", "")))
|
||||
self.var_mirror_sfm_state_file = tk.StringVar(value=str(v.get("mirror_sfm_state_file", "")))
|
||||
|
||||
# ── UI construction ───────────────────────────────────────────────────────
|
||||
|
||||
def _build_ui(self):
|
||||
@@ -245,6 +258,7 @@ class SettingsDialog:
|
||||
self._build_tab_scanning(nb)
|
||||
self._build_tab_logging(nb)
|
||||
self._build_tab_forwarding(nb)
|
||||
self._build_tab_mirror(nb)
|
||||
self._build_tab_updates(nb)
|
||||
|
||||
btn_frame = tk.Frame(outer)
|
||||
@@ -484,6 +498,49 @@ class SettingsDialog:
|
||||
finally:
|
||||
self._sfm_test_btn.config(state="normal")
|
||||
|
||||
def _build_tab_mirror(self, nb):
|
||||
f = self._tab_frame(nb, "Mirror")
|
||||
|
||||
intro = (
|
||||
"Optional dual-send: post each heartbeat and event to a SECOND\n"
|
||||
"(\"mirror\") server in addition to the primary. Best-effort — the\n"
|
||||
"mirror can never delay or fail the primary. Leave blank to disable."
|
||||
)
|
||||
tk.Label(f, text=intro, justify="left", fg="#1a5276", wraplength=420).grid(
|
||||
row=0, column=0, columnspan=2, sticky="w", padx=(8, 8), pady=(6, 8)
|
||||
)
|
||||
|
||||
_add_label_entry(f, 1, "Mirror Terra-View URL", self.var_mirror_api_url)
|
||||
_add_label_entry(f, 2, "Mirror SFM URL", self.var_mirror_sfm_url)
|
||||
|
||||
def browse_mirror_state():
|
||||
p = filedialog.asksaveasfilename(
|
||||
title="Select Mirror State File",
|
||||
defaultextension=".json",
|
||||
filetypes=[("JSON files", "*.json"), ("All files", "*.*")],
|
||||
initialfile=os.path.basename(
|
||||
self.var_mirror_sfm_state_file.get() or "thor_forwarded_mirror.json"),
|
||||
initialdir=os.path.dirname(self.var_mirror_sfm_state_file.get() or "C:\\"),
|
||||
)
|
||||
if p:
|
||||
self.var_mirror_sfm_state_file.set(p.replace("/", "\\"))
|
||||
|
||||
_add_label_browse_entry(f, 3, "Mirror State File", self.var_mirror_sfm_state_file,
|
||||
browse_mirror_state)
|
||||
|
||||
help_text = (
|
||||
"Mirror Terra-View URL → base like http://10.0.0.x:8001 (heartbeats).\n"
|
||||
"Mirror SFM URL → base like http://10.0.0.x:8200 (events).\n"
|
||||
"The event mirror keeps its OWN state file, so nothing is lost while\n"
|
||||
"the mirror is down — pending events are delivered once it returns.\n"
|
||||
"A quick reachability check skips the mirror cleanly when unreachable,\n"
|
||||
"so the primary is never slowed.\n"
|
||||
"State file blank → defaults to <log_dir>\\thor_forwarded_mirror.json."
|
||||
)
|
||||
tk.Label(f, text=help_text, justify="left", fg="#555555", wraplength=420).grid(
|
||||
row=4, column=0, columnspan=2, sticky="w", padx=(8, 8), pady=(8, 4)
|
||||
)
|
||||
|
||||
def _build_tab_updates(self, nb):
|
||||
f = self._tab_frame(nb, "Updates")
|
||||
|
||||
@@ -599,6 +656,12 @@ class SettingsDialog:
|
||||
sfm_url = ""
|
||||
sfm_url = sfm_url.rstrip("/") # event_forwarder adds the endpoint path
|
||||
|
||||
# Mirror (dual-send) — best-effort second destination.
|
||||
mirror_api_url = self.var_mirror_api_url.get().strip()
|
||||
if mirror_api_url:
|
||||
mirror_api_url = mirror_api_url.rstrip("/") + "/api/series4/heartbeat"
|
||||
mirror_sfm_url = self.var_mirror_sfm_url.get().strip().rstrip("/")
|
||||
|
||||
values = {
|
||||
"thordata_path": self.var_thordata_path.get().strip(),
|
||||
"scan_interval": int_values["Scan Interval"],
|
||||
@@ -623,6 +686,10 @@ class SettingsDialog:
|
||||
"sfm_max_forwards_per_pass": int_values["Max Forwards Per Pass"],
|
||||
"sfm_max_event_age_days": int_values["Max Event Age (days)"],
|
||||
"sfm_state_file": self.var_sfm_state_file.get().strip(),
|
||||
|
||||
"mirror_api_url": mirror_api_url,
|
||||
"mirror_sfm_url": mirror_sfm_url,
|
||||
"mirror_sfm_state_file": self.var_mirror_sfm_state_file.get().strip(),
|
||||
}
|
||||
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user