2026-05-12 14:07:03 +02:00
|
|
|
|
import os
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
import re
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
import sys
|
2026-05-12 14:07:03 +02:00
|
|
|
|
import json
|
|
|
|
|
|
import time
|
|
|
|
|
|
import glob
|
|
|
|
|
|
import logging
|
2026-07-15 15:30:44 +02:00
|
|
|
|
from logging.handlers import RotatingFileHandler
|
2026-07-09 15:38:20 +02:00
|
|
|
|
import urllib.request
|
|
|
|
|
|
import urllib.parse
|
2026-05-12 14:07:03 +02:00
|
|
|
|
import yaml
|
|
|
|
|
|
from datetime import datetime, timezone
|
|
|
|
|
|
from pathlib import Path
|
|
|
|
|
|
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
# Shared liveness logic (thresholds + fresh/stale/dead state machine) lives in
|
|
|
|
|
|
# the control-plane src so the observer and both operator UIs agree on what
|
|
|
|
|
|
# "dead" means. Added to sys.path relative to this file (works both in the
|
|
|
|
|
|
# container, where the repo is mounted at /repo, and in the test tree).
|
|
|
|
|
|
_CP_SRC = Path(__file__).resolve().parent.parent.parent / "services" / "control-plane" / "src"
|
|
|
|
|
|
if str(_CP_SRC) not in sys.path:
|
|
|
|
|
|
sys.path.insert(0, str(_CP_SRC))
|
|
|
|
|
|
from liveness import ( # noqa: E402
|
|
|
|
|
|
compute_liveness,
|
|
|
|
|
|
ttls_for,
|
|
|
|
|
|
liveness_to_status,
|
|
|
|
|
|
FRESH,
|
|
|
|
|
|
STALE,
|
|
|
|
|
|
DEAD,
|
|
|
|
|
|
UNKNOWN,
|
|
|
|
|
|
)
|
|
|
|
|
|
|
2026-06-03 12:26:49 +02:00
|
|
|
|
|
|
|
|
|
|
def _atomic_write_json(path: Path, data) -> None:
|
|
|
|
|
|
"""Write JSON atomically: write to a sibling .tmp, fsync, then os.replace."""
|
|
|
|
|
|
tmp = path.with_suffix(".tmp")
|
|
|
|
|
|
with open(tmp, "w") as f:
|
|
|
|
|
|
json.dump(data, f, indent=2)
|
|
|
|
|
|
f.flush()
|
|
|
|
|
|
os.fsync(f.fileno())
|
|
|
|
|
|
os.replace(tmp, path)
|
|
|
|
|
|
|
2026-06-03 14:29:12 +02:00
|
|
|
|
|
|
|
|
|
|
def _parse_ts(ts) -> float:
|
|
|
|
|
|
"""Return a Unix timestamp float from ts, which may be int/float or an ISO-8601 string.
|
|
|
|
|
|
|
|
|
|
|
|
Events from node-agent use int(time.time()); events from stability-agent / events.py
|
|
|
|
|
|
use ISO format ('2026-06-03T10:30:00Z'). Both appear in incident fields such as
|
|
|
|
|
|
last_occurrence and resolved_at, so any arithmetic on them must go through here.
|
|
|
|
|
|
Returns 0.0 on None or unparseable input so callers can use plain comparisons.
|
|
|
|
|
|
"""
|
|
|
|
|
|
if ts is None:
|
|
|
|
|
|
return 0.0
|
|
|
|
|
|
if isinstance(ts, (int, float)):
|
|
|
|
|
|
return float(ts)
|
|
|
|
|
|
try:
|
|
|
|
|
|
return datetime.fromisoformat(str(ts).replace("Z", "+00:00")).timestamp()
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
return 0.0
|
|
|
|
|
|
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
|
|
|
|
|
|
# Event filenames follow evt-<node>-<unixts>-<type>-<svc>.json (node_agent.py,
|
|
|
|
|
|
# ha_diag/event_emitter.py, and the observer's own _emit_node_transition). The
|
|
|
|
|
|
# <unixts> is the authoritative ordering key for checkpointing — matched the same
|
|
|
|
|
|
# way operator_ui.py::_event_file_ts does (a 9–11 digit run flanked by dashes;
|
|
|
|
|
|
# real epochs are 10 digits and stay so until year 2286).
|
|
|
|
|
|
_EVENT_TS_RE = re.compile(r"-(\d{9,11})-")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _ts_from_event_name(name) -> "int | None":
|
|
|
|
|
|
"""Parse the embedded <unixts> from an event filename or full path.
|
|
|
|
|
|
|
|
|
|
|
|
Returns the int epoch, or None when the name does not carry one (foreign
|
|
|
|
|
|
prefix, events.py naming, junk file) so the caller can fall back to mtime.
|
|
|
|
|
|
"""
|
|
|
|
|
|
m = _EVENT_TS_RE.search(Path(name).stem)
|
|
|
|
|
|
return int(m.group(1)) if m else None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _event_ts_from_path(file_path) -> int:
|
|
|
|
|
|
"""Ordering/checkpoint timestamp for an event file.
|
|
|
|
|
|
|
|
|
|
|
|
Primary: the <unixts> embedded in the filename. Fallback (name doesn't
|
|
|
|
|
|
parse): the file's mtime. NEVER returns 0 for an existing file — a value of
|
|
|
|
|
|
0 would make the file compare as "older than the checkpoint" and be skipped
|
|
|
|
|
|
forever, which is exactly the lexical-path poisoning this fix removes. If
|
|
|
|
|
|
even stat() fails, fall back to now() so the file is treated as new and gets
|
|
|
|
|
|
a chance to be processed (and quarantined if truly unreadable).
|
|
|
|
|
|
"""
|
|
|
|
|
|
ts = _ts_from_event_name(file_path)
|
|
|
|
|
|
if ts is not None:
|
|
|
|
|
|
return ts
|
|
|
|
|
|
try:
|
|
|
|
|
|
return int(os.stat(file_path).st_mtime)
|
|
|
|
|
|
except OSError:
|
|
|
|
|
|
return int(time.time())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _checkpoint_ts_from_value(value) -> int:
|
|
|
|
|
|
"""Coerce a stored checkpoint value into an int epoch (migration helper).
|
|
|
|
|
|
|
|
|
|
|
|
Historical formats of observer_checkpoint.json values:
|
|
|
|
|
|
- int/float → already a timestamp (current format) — kept as-is
|
|
|
|
|
|
- path/name string → the pre-fix lexical checkpoint; parse the embedded
|
|
|
|
|
|
<unixts> out of it so the node resumes near where it
|
|
|
|
|
|
left off instead of reprocessing everything
|
|
|
|
|
|
- anything else / unparseable string → 0 (reprocess all events for that
|
|
|
|
|
|
node). Reprocessing is safe: process_event is idempotent w.r.t.
|
|
|
|
|
|
last_seen/world_state, so re-ingesting duplicates cannot corrupt state —
|
|
|
|
|
|
whereas guessing too high a checkpoint could silently drop events (the
|
|
|
|
|
|
failure mode being fixed). Bias to reprocess, never to skip.
|
|
|
|
|
|
"""
|
|
|
|
|
|
if isinstance(value, bool): # bool is an int subclass — exclude explicitly
|
|
|
|
|
|
return 0
|
|
|
|
|
|
if isinstance(value, (int, float)):
|
|
|
|
|
|
return int(value)
|
|
|
|
|
|
if isinstance(value, str) and value:
|
|
|
|
|
|
ts = _ts_from_event_name(value)
|
|
|
|
|
|
if ts is not None:
|
|
|
|
|
|
return ts
|
|
|
|
|
|
return 0
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
# Constants and Paths
|
|
|
|
|
|
RUNTIME_PATH = os.getenv("RUNTIME_PATH", "/opt/homelab")
|
|
|
|
|
|
EVENTS_DIR = Path(RUNTIME_PATH) / "events"
|
|
|
|
|
|
STATE_DIR = Path(RUNTIME_PATH) / "state"
|
|
|
|
|
|
LOGS_DIR = Path(RUNTIME_PATH) / "logs"
|
|
|
|
|
|
WORLD_DIR = Path(RUNTIME_PATH) / "world"
|
|
|
|
|
|
OBSERVER_STATE_FILE = STATE_DIR / "observer_checkpoint.json"
|
2026-06-12 13:11:15 +02:00
|
|
|
|
FAILED_EVENTS_DIR = STATE_DIR / "observer_failed_events"
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
REPO_ROOT = Path(__file__).parent.parent.parent
|
|
|
|
|
|
INVENTORY_TOPOLOGY = REPO_ROOT / "inventory" / "topology.yaml"
|
|
|
|
|
|
|
2026-07-09 15:38:20 +02:00
|
|
|
|
# --- SHADOW-READ: Prometheus up{} liveness (cutover etap 1) ----------------
|
|
|
|
|
|
# Optional parallel-run source. When PROM_SHADOW_URL is set, the observer ALSO
|
|
|
|
|
|
# queries Prometheus `up{}` each cycle and LOGS any disagreement with its own
|
|
|
|
|
|
# event-driven liveness — but NEVER acts on it. The authoritative liveness stays
|
|
|
|
|
|
# 100% event-driven (compute_liveness in _prune_stale_world). Empty/unset →
|
|
|
|
|
|
# shadow disabled, observer behaves exactly as before (graceful, fail-open).
|
|
|
|
|
|
# Target value (do NOT hardcode — set via env/host override): http://100.95.58.48:9090
|
|
|
|
|
|
PROM_SHADOW_URL = os.environ.get("PROM_SHADOW_URL", "").strip()
|
|
|
|
|
|
PROM_SHADOW_TIMEOUT = int(os.getenv("PROM_SHADOW_TIMEOUT", "5"))
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
# Logging setup
|
|
|
|
|
|
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
|
|
|
|
|
|
logger = logging.getLogger("observer")
|
|
|
|
|
|
|
2026-07-15 15:30:44 +02:00
|
|
|
|
# --- PERSISTENT shadow-mismatch log (cutover etap 2) -----------------------
|
|
|
|
|
|
# SHADOW_LIVENESS_MISMATCH lines are cutover EVIDENCE and must survive a
|
|
|
|
|
|
# `docker rm`/recreate of the observer container. stdout (json-file) logs do
|
|
|
|
|
|
# NOT: an observer recreate on 2026-07-14 destroyed the 07-13/14 mismatch
|
|
|
|
|
|
# material mid-analysis. We therefore ALSO write each mismatch to a
|
|
|
|
|
|
# RotatingFileHandler on a HOST-mounted path so the file outlives the container.
|
|
|
|
|
|
#
|
|
|
|
|
|
# Path is the repo-conventional logs/<service>/ location under RUNTIME_PATH,
|
|
|
|
|
|
# i.e. /opt/homelab/logs/observer/shadow-liveness.log. No dedicated bind-mount
|
|
|
|
|
|
# is required: the base compose already mounts the whole of /opt/homelab into
|
|
|
|
|
|
# the observer, so this file is on the host by construction, and the observer
|
|
|
|
|
|
# (uid 1000) creates the dir itself — avoiding the root-owned bind-source
|
|
|
|
|
|
# ownership footgun a separate /var/log mount would introduce. Overridable via
|
|
|
|
|
|
# SHADOW_LOG_DIR for tests / non-container runs.
|
|
|
|
|
|
SHADOW_LOG_DIR = Path(os.getenv("SHADOW_LOG_DIR", str(LOGS_DIR / "observer")))
|
|
|
|
|
|
SHADOW_LOG_FILENAME = "shadow-liveness.log"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _make_shadow_logger():
|
|
|
|
|
|
"""Build the dedicated PERSISTENT logger for SHADOW_LIVENESS_MISMATCH lines.
|
|
|
|
|
|
|
|
|
|
|
|
Returns a `logging.getLogger("observer.shadow")` wired to a
|
|
|
|
|
|
RotatingFileHandler (5 MiB x 5 backups — mismatches are short lines, that is
|
|
|
|
|
|
weeks of headroom) writing to SHADOW_LOG_DIR/shadow-liveness.log with the
|
|
|
|
|
|
same timestamped format as the main observer log.
|
|
|
|
|
|
|
|
|
|
|
|
propagate=False keeps these records OUT of root/stdout: the call site still
|
|
|
|
|
|
logs the mismatch to stdout via the ordinary `logger` (unchanged), so
|
|
|
|
|
|
without this each mismatch would appear twice in `docker logs`.
|
|
|
|
|
|
|
|
|
|
|
|
FAIL-SAFE: if the directory/file cannot be created or opened (permission,
|
|
|
|
|
|
read-only mount, …) we log ONE warning on the main observer logger and
|
|
|
|
|
|
return a handler-less logger. `.info()` on a handler-less, non-propagating
|
|
|
|
|
|
logger is a silent no-op, so a broken persistent log NEVER raises and NEVER
|
|
|
|
|
|
takes the observer down — the mismatch still reaches stdout at the call site.
|
|
|
|
|
|
Idempotent: any handler from a previous call is dropped first, so
|
|
|
|
|
|
re-instantiation (tests, re-import) never double-writes or pins a stale path.
|
|
|
|
|
|
"""
|
|
|
|
|
|
sl = logging.getLogger("observer.shadow")
|
|
|
|
|
|
sl.setLevel(logging.INFO)
|
|
|
|
|
|
sl.propagate = False
|
|
|
|
|
|
for h in list(sl.handlers):
|
|
|
|
|
|
sl.removeHandler(h)
|
|
|
|
|
|
try:
|
|
|
|
|
|
h.close()
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
try:
|
|
|
|
|
|
os.makedirs(SHADOW_LOG_DIR, exist_ok=True)
|
|
|
|
|
|
handler = RotatingFileHandler(
|
|
|
|
|
|
SHADOW_LOG_DIR / SHADOW_LOG_FILENAME,
|
|
|
|
|
|
maxBytes=5 * 1024 * 1024,
|
|
|
|
|
|
backupCount=5,
|
|
|
|
|
|
)
|
|
|
|
|
|
handler.setFormatter(
|
|
|
|
|
|
logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')
|
|
|
|
|
|
)
|
|
|
|
|
|
sl.addHandler(handler)
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
logger.warning(
|
|
|
|
|
|
"shadow-read: could not open persistent mismatch log at %s (%s) — "
|
|
|
|
|
|
"falling back to stdout only, observer continues",
|
|
|
|
|
|
SHADOW_LOG_DIR / SHADOW_LOG_FILENAME, exc,
|
|
|
|
|
|
)
|
|
|
|
|
|
return sl
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
class Observer:
|
|
|
|
|
|
def __init__(self):
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
# Per-node-directory checkpoint keyed on the last-processed event
|
|
|
|
|
|
# TIMESTAMP (int epoch): {"vps": 1784000000, "piha": 1784000123}.
|
|
|
|
|
|
# A file is "new" iff its event timestamp > the node's checkpoint.
|
|
|
|
|
|
# This replaces the earlier lexical-PATH comparison, which permanently
|
|
|
|
|
|
# poisoned a node the moment a single file with a lexically-larger name
|
|
|
|
|
|
# landed in its dir (e.g. evt-unknown-… > evt-piha-…): every genuinely
|
|
|
|
|
|
# newer event then sorted "before" the checkpoint and was skipped forever.
|
2026-05-27 14:16:58 +02:00
|
|
|
|
self.node_checkpoints: dict = {}
|
2026-05-12 14:07:03 +02:00
|
|
|
|
self.world_state = {
|
|
|
|
|
|
"nodes": {},
|
|
|
|
|
|
"services": {},
|
|
|
|
|
|
"deployments": {},
|
|
|
|
|
|
"incidents": {},
|
|
|
|
|
|
"summary": {
|
2026-05-12 20:59:46 +02:00
|
|
|
|
"last_update": datetime.now(timezone.utc).isoformat(),
|
2026-05-12 14:07:03 +02:00
|
|
|
|
"status": "initializing",
|
|
|
|
|
|
"active_incidents_count": 0
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
self.inventory = self._load_inventory()
|
|
|
|
|
|
self._ensure_dirs()
|
|
|
|
|
|
self._load_checkpoint()
|
2026-07-15 15:30:44 +02:00
|
|
|
|
# Persistent SHADOW_LIVENESS_MISMATCH sink (survives container recreate).
|
|
|
|
|
|
self.shadow_logger = _make_shadow_logger()
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
def _ensure_dirs(self):
|
|
|
|
|
|
WORLD_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
EVENTS_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
LOGS_DIR.mkdir(parents=True, exist_ok=True)
|
2026-06-12 13:11:15 +02:00
|
|
|
|
FAILED_EVENTS_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
|
|
|
|
|
|
def _quarantine_event_file(self, file_path: str, node_dir: str, exc: Exception) -> None:
|
|
|
|
|
|
"""Move an unreadable/unprocessable event out of the hot path."""
|
|
|
|
|
|
src = Path(file_path)
|
|
|
|
|
|
dest_dir = FAILED_EVENTS_DIR / node_dir
|
|
|
|
|
|
dest_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
dest = dest_dir / src.name
|
|
|
|
|
|
if dest.exists():
|
|
|
|
|
|
dest = dest_dir / f"{src.stem}-{int(time.time())}{src.suffix}"
|
|
|
|
|
|
try:
|
|
|
|
|
|
os.replace(src, dest)
|
|
|
|
|
|
logger.error(
|
|
|
|
|
|
"Quarantined bad event for node_dir=%s: %s -> %s (%s: %s)",
|
|
|
|
|
|
node_dir, src, dest, type(exc).__name__, exc,
|
|
|
|
|
|
)
|
|
|
|
|
|
except Exception as move_exc:
|
|
|
|
|
|
logger.error(
|
|
|
|
|
|
"Failed to quarantine bad event for node_dir=%s: %s (%s: %s); move error=%s: %s",
|
|
|
|
|
|
node_dir, src, type(exc).__name__, exc, type(move_exc).__name__, move_exc,
|
|
|
|
|
|
)
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
def _load_inventory(self):
|
|
|
|
|
|
inventory = {"nodes": {}, "services": {}}
|
|
|
|
|
|
try:
|
|
|
|
|
|
if INVENTORY_TOPOLOGY.exists():
|
|
|
|
|
|
with open(INVENTORY_TOPOLOGY, "r") as f:
|
|
|
|
|
|
topo = yaml.safe_load(f)
|
|
|
|
|
|
for node_name, node_info in topo.get("nodes", {}).items():
|
|
|
|
|
|
inventory["nodes"][node_name] = {
|
|
|
|
|
|
"roles": node_info.get("roles", []),
|
|
|
|
|
|
"connectivity": node_info.get("connectivity", {})
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
# Load service assignments from hosts files
|
|
|
|
|
|
hosts_dir = REPO_ROOT / "hosts"
|
|
|
|
|
|
for host_dir in hosts_dir.iterdir():
|
|
|
|
|
|
if host_dir.is_dir():
|
|
|
|
|
|
svc_file = host_dir / "services.yaml"
|
|
|
|
|
|
if svc_file.exists():
|
|
|
|
|
|
with open(svc_file, "r") as f:
|
|
|
|
|
|
svc_data = yaml.safe_load(f)
|
|
|
|
|
|
host_name = svc_data.get("host")
|
|
|
|
|
|
for svc_name, svc_info in svc_data.get("services", {}).items():
|
|
|
|
|
|
if host_name not in inventory["services"]:
|
|
|
|
|
|
inventory["services"][host_name] = {}
|
|
|
|
|
|
inventory["services"][host_name][svc_name] = {
|
|
|
|
|
|
"role": svc_info.get("role"),
|
|
|
|
|
|
"exposure": svc_info.get("exposure")
|
|
|
|
|
|
}
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to load inventory: {e}")
|
|
|
|
|
|
return inventory
|
|
|
|
|
|
|
|
|
|
|
|
def _load_checkpoint(self):
|
|
|
|
|
|
if OBSERVER_STATE_FILE.exists():
|
|
|
|
|
|
try:
|
|
|
|
|
|
with open(OBSERVER_STATE_FILE, "r") as f:
|
|
|
|
|
|
checkpoint = json.load(f)
|
2026-05-27 14:16:58 +02:00
|
|
|
|
|
|
|
|
|
|
if "node_checkpoints" in checkpoint:
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
# Per-directory checkpoints. Values may be int epochs (current
|
|
|
|
|
|
# format) OR pre-fix path strings — coerce every value to an
|
|
|
|
|
|
# int timestamp so a checkpoint file written by the old
|
|
|
|
|
|
# lexical-path observer migrates transparently on first start.
|
|
|
|
|
|
raw = checkpoint["node_checkpoints"] or {}
|
|
|
|
|
|
self.node_checkpoints = {
|
|
|
|
|
|
node: _checkpoint_ts_from_value(val)
|
|
|
|
|
|
for node, val in raw.items()
|
|
|
|
|
|
}
|
|
|
|
|
|
if any(not isinstance(v, (int, float)) for v in raw.values()):
|
|
|
|
|
|
logger.info(
|
|
|
|
|
|
"Migrated path-based node_checkpoints → timestamps: %s",
|
|
|
|
|
|
self.node_checkpoints,
|
|
|
|
|
|
)
|
2026-05-27 14:16:58 +02:00
|
|
|
|
elif "last_processed_file" in checkpoint:
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
# Migrate the very old single-file checkpoint: extract node dir
|
|
|
|
|
|
# from the path and the timestamp from the filename.
|
2026-05-27 14:16:58 +02:00
|
|
|
|
old = checkpoint["last_processed_file"]
|
|
|
|
|
|
if old:
|
|
|
|
|
|
try:
|
|
|
|
|
|
node_dir = Path(old).relative_to(EVENTS_DIR).parts[0]
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
self.node_checkpoints = {node_dir: _checkpoint_ts_from_value(old)}
|
2026-05-27 14:16:58 +02:00
|
|
|
|
logger.info(f"Migrated old checkpoint → node_checkpoints: {self.node_checkpoints}")
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass # Bad path — start fresh
|
|
|
|
|
|
|
|
|
|
|
|
self._load_world_from_disk()
|
2026-05-12 14:07:03 +02:00
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to load checkpoint: {e}")
|
|
|
|
|
|
|
|
|
|
|
|
def _load_world_from_disk(self):
|
|
|
|
|
|
# Optional: Load existing state to resume faster
|
|
|
|
|
|
files = {
|
|
|
|
|
|
"nodes": WORLD_DIR / "nodes.json",
|
|
|
|
|
|
"services": WORLD_DIR / "services.json",
|
|
|
|
|
|
"deployments": WORLD_DIR / "deployments.json",
|
|
|
|
|
|
"incidents": WORLD_DIR / "incidents.json",
|
|
|
|
|
|
"summary": WORLD_DIR / "runtime-summary.json"
|
|
|
|
|
|
}
|
|
|
|
|
|
for key, path in files.items():
|
|
|
|
|
|
if path.exists():
|
|
|
|
|
|
try:
|
|
|
|
|
|
with open(path, "r") as f:
|
|
|
|
|
|
self.world_state[key] = json.load(f)
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to load {key} state: {e}")
|
|
|
|
|
|
|
|
|
|
|
|
def _save_checkpoint(self):
|
|
|
|
|
|
try:
|
2026-06-03 12:26:49 +02:00
|
|
|
|
_atomic_write_json(OBSERVER_STATE_FILE, {"node_checkpoints": self.node_checkpoints})
|
2026-05-12 14:07:03 +02:00
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to save checkpoint: {e}")
|
|
|
|
|
|
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
def _emit_node_transition(self, node_name, prev, new, node_info, now):
|
|
|
|
|
|
"""Write an event when a node crosses a liveness boundary.
|
|
|
|
|
|
|
|
|
|
|
|
Makes liveness transitions visible (panel event feed) and actionable
|
|
|
|
|
|
(the supervisor routes node_offline/node_stale to an alert). Without
|
|
|
|
|
|
this the outage would be silent — exactly the failure mode being fixed.
|
|
|
|
|
|
|
|
|
|
|
|
Recovery (stale/dead → fresh) emits node_online so a node coming back is
|
|
|
|
|
|
as visible as it going down.
|
|
|
|
|
|
|
|
|
|
|
|
The affected node is the top-level ``node`` (so the supervisor can route
|
|
|
|
|
|
on it) AND ``payload.affected_node`` — note events.py::emit_event would
|
|
|
|
|
|
tag node=gethostname() (the observer host, vps), so the observer writes
|
|
|
|
|
|
the event directly instead of using that helper.
|
|
|
|
|
|
|
|
|
|
|
|
``source: "observer"`` marks the event so run_once skips re-ingesting it
|
|
|
|
|
|
(re-ingestion would reset the node's last_seen and resurrect it).
|
|
|
|
|
|
"""
|
|
|
|
|
|
if new == DEAD:
|
|
|
|
|
|
etype, severity = "node_offline", "high"
|
|
|
|
|
|
elif new == STALE:
|
|
|
|
|
|
etype, severity = "node_stale", "warning"
|
|
|
|
|
|
elif new == FRESH and prev in (STALE, DEAD):
|
|
|
|
|
|
etype, severity = "node_online", "info"
|
|
|
|
|
|
else:
|
|
|
|
|
|
return # not a transition we alert on
|
|
|
|
|
|
|
|
|
|
|
|
ts = int(now)
|
|
|
|
|
|
event_id = f"evt-{node_name}-{ts}-{etype}-node"
|
|
|
|
|
|
age = int(now - _parse_ts(node_info.get("last_seen")))
|
|
|
|
|
|
event = {
|
|
|
|
|
|
"id": event_id,
|
|
|
|
|
|
"timestamp": ts,
|
|
|
|
|
|
"date": datetime.now(timezone.utc).isoformat(),
|
|
|
|
|
|
"type": etype,
|
|
|
|
|
|
"severity": severity,
|
|
|
|
|
|
"node": node_name,
|
|
|
|
|
|
"service": None,
|
|
|
|
|
|
"source": "observer",
|
|
|
|
|
|
"message": f"Node {node_name} liveness {prev} -> {new} (last_seen {age}s ago)",
|
|
|
|
|
|
"payload": {
|
|
|
|
|
|
"affected_node": node_name,
|
|
|
|
|
|
"from": prev,
|
|
|
|
|
|
"to": new,
|
|
|
|
|
|
"last_seen": node_info.get("last_seen"),
|
|
|
|
|
|
"age_secs": age,
|
|
|
|
|
|
},
|
|
|
|
|
|
}
|
|
|
|
|
|
try:
|
|
|
|
|
|
node_dir = EVENTS_DIR / node_name
|
|
|
|
|
|
node_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
_atomic_write_json(node_dir / f"{event_id}.json", event)
|
|
|
|
|
|
logger.warning("Node %s transition %s -> %s (emitted %s)",
|
|
|
|
|
|
node_name, prev, new, etype)
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
logger.error("Failed to emit node transition for %s: %s", node_name, exc)
|
|
|
|
|
|
|
2026-07-09 15:38:20 +02:00
|
|
|
|
def _query_prometheus_liveness(self):
|
|
|
|
|
|
"""SHADOW-READ (cutover etap 1): Prometheus up{} → {node_name: up_bool}.
|
|
|
|
|
|
|
|
|
|
|
|
Parallel-run only: the result is compared against — and logged next to —
|
|
|
|
|
|
the authoritative event-driven liveness, but NEVER changes it. See
|
|
|
|
|
|
PROM_SHADOW_URL. Returns {} when shadow is disabled OR on ANY error, so
|
|
|
|
|
|
the observer can never crash or change behaviour because of this call:
|
|
|
|
|
|
every failure is fail-open (info/warning, never error, never raised).
|
|
|
|
|
|
|
|
|
|
|
|
Each `up` series is keyed by its `node` label (Prometheus fleet-node
|
|
|
|
|
|
targets carry node:vps/piha/solaria/lustro). Series without a node label
|
|
|
|
|
|
(e.g. the prometheus self-scrape up{job="prometheus"}) are ignored.
|
|
|
|
|
|
"""
|
|
|
|
|
|
url = PROM_SHADOW_URL
|
|
|
|
|
|
if not url:
|
|
|
|
|
|
return {} # shadow disabled — graceful no-op, observer unchanged
|
|
|
|
|
|
query_url = url.rstrip("/") + "/api/v1/query?" + urllib.parse.urlencode({"query": "up"})
|
|
|
|
|
|
try:
|
|
|
|
|
|
with urllib.request.urlopen(query_url, timeout=PROM_SHADOW_TIMEOUT) as resp:
|
|
|
|
|
|
payload = json.loads(resp.read().decode("utf-8"))
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
# Fail-open: Prometheus down / timeout / bad JSON → no shadow data.
|
|
|
|
|
|
# Deliberately NOT logger.error — shadow-read must never look like a
|
|
|
|
|
|
# critical observer failure.
|
|
|
|
|
|
logger.warning(
|
|
|
|
|
|
"shadow-read: Prometheus query failed (%s: %s) — fail-open, "
|
|
|
|
|
|
"event liveness unaffected", type(exc).__name__, exc,
|
|
|
|
|
|
)
|
|
|
|
|
|
return {}
|
|
|
|
|
|
result: dict = {}
|
|
|
|
|
|
try:
|
|
|
|
|
|
for series in payload.get("data", {}).get("result", []):
|
|
|
|
|
|
node = series.get("metric", {}).get("node")
|
|
|
|
|
|
if not node:
|
|
|
|
|
|
continue # e.g. up{job="prometheus"} — no node label to map
|
|
|
|
|
|
value = series.get("value", [None, None])[1]
|
|
|
|
|
|
result[node] = (value == "1")
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
logger.warning(
|
|
|
|
|
|
"shadow-read: could not parse Prometheus response (%s: %s) — "
|
|
|
|
|
|
"fail-open", type(exc).__name__, exc,
|
|
|
|
|
|
)
|
|
|
|
|
|
return {}
|
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
|
|
def _shadow_compare_liveness(self, node_name, event_liveness, node_info, prom_map, now):
|
|
|
|
|
|
"""SHADOW-READ comparison (cutover etap 1): log event-vs-Prometheus
|
|
|
|
|
|
liveness disagreement WITHOUT changing anything.
|
|
|
|
|
|
|
|
|
|
|
|
Maps both sources onto a shared up/down axis:
|
|
|
|
|
|
- event DEAD ≈ prom down
|
|
|
|
|
|
- event FRESH or STALE ≈ prom up (STALE = "seen recently, just
|
|
|
|
|
|
ageing" — still counts as up for this comparison; documented
|
|
|
|
|
|
assumption of etap 1)
|
|
|
|
|
|
|
|
|
|
|
|
A node Prometheus doesn't know (no series, e.g. chelsty-infra — not
|
|
|
|
|
|
scraped) is NOT a mismatch: missing data is not disagreement (debug only).
|
|
|
|
|
|
This method is read-only w.r.t. liveness/status and must never raise.
|
|
|
|
|
|
"""
|
|
|
|
|
|
if node_name not in prom_map:
|
|
|
|
|
|
logger.debug("shadow-read: no prom data for node=%s", node_name)
|
|
|
|
|
|
return
|
|
|
|
|
|
prom_up = prom_map[node_name]
|
|
|
|
|
|
event_up = event_liveness in (FRESH, STALE)
|
|
|
|
|
|
age = now - _parse_ts(node_info.get("last_seen"))
|
|
|
|
|
|
if event_up != prom_up:
|
|
|
|
|
|
logger.info(
|
|
|
|
|
|
"SHADOW_LIVENESS_MISMATCH node=%s event=%s prom=%s last_seen_age=%.0fs",
|
|
|
|
|
|
node_name, event_liveness, "up" if prom_up else "down", age,
|
|
|
|
|
|
)
|
2026-07-15 15:30:44 +02:00
|
|
|
|
# ALSO write to the persistent, host-mounted log so this evidence
|
|
|
|
|
|
# survives a container recreate (stdout json-file logs do not).
|
|
|
|
|
|
# Handler-less (fail-safe) → silent no-op; never affects the above.
|
|
|
|
|
|
self.shadow_logger.info(
|
|
|
|
|
|
"SHADOW_LIVENESS_MISMATCH node=%s event=%s prom=%s last_seen_age=%.0fs",
|
|
|
|
|
|
node_name, event_liveness, "up" if prom_up else "down", age,
|
|
|
|
|
|
)
|
2026-07-09 15:38:20 +02:00
|
|
|
|
else:
|
|
|
|
|
|
logger.debug(
|
|
|
|
|
|
"shadow-read agree node=%s event=%s prom=%s last_seen_age=%.0fs",
|
|
|
|
|
|
node_name, event_liveness, "up" if prom_up else "down", age,
|
|
|
|
|
|
)
|
|
|
|
|
|
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
def _prune_stale_world(self):
|
|
|
|
|
|
"""Remove world-state entries for nodes absent from the topology inventory.
|
|
|
|
|
|
|
|
|
|
|
|
Root cause this guards against: when NODE_NAME env var is unset, node_agent.py
|
|
|
|
|
|
falls back to socket.gethostname(), which inside a Docker container returns the
|
|
|
|
|
|
12-char hex container ID (e.g. 'be17cb6eb0f6') instead of the canonical host name
|
|
|
|
|
|
('vps'). The observer ingests those events and creates ghost entries that never
|
|
|
|
|
|
expire on their own.
|
|
|
|
|
|
|
|
|
|
|
|
Also ages out resolved incidents older than 7 days to keep world state lean.
|
|
|
|
|
|
"""
|
|
|
|
|
|
known_nodes = set(self.inventory["nodes"].keys())
|
|
|
|
|
|
if not known_nodes:
|
|
|
|
|
|
# Inventory failed to load — don't prune to avoid wiping valid state.
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
stale_nodes = [n for n in list(self.world_state["nodes"].keys())
|
|
|
|
|
|
if n not in known_nodes]
|
|
|
|
|
|
for n in stale_nodes:
|
|
|
|
|
|
logger.info(f"Pruning stale node from world state: {n}")
|
|
|
|
|
|
del self.world_state["nodes"][n]
|
|
|
|
|
|
|
|
|
|
|
|
stale_svcs = [k for k in list(self.world_state["services"].keys())
|
|
|
|
|
|
if k.split("/")[0] in stale_nodes]
|
|
|
|
|
|
for k in stale_svcs:
|
|
|
|
|
|
logger.info(f"Pruning stale service from world state: {k}")
|
|
|
|
|
|
del self.world_state["services"][k]
|
|
|
|
|
|
|
2026-05-27 15:41:13 +02:00
|
|
|
|
# Prune ghost service keys whose service-name portion is a hash-prefixed
|
|
|
|
|
|
# Docker stale-state artifact (e.g. "9e36297651e7_control-plane-observer").
|
|
|
|
|
|
# These are created when node-agent incorrectly uses c.name instead of the
|
|
|
|
|
|
# compose label, and accumulate on every container rebuild.
|
|
|
|
|
|
# Pattern: <node>/<12hexchars>_<real-name>
|
|
|
|
|
|
ghost_svcs = [
|
|
|
|
|
|
k for k in list(self.world_state["services"].keys())
|
|
|
|
|
|
if len(k.split("/", 1)) == 2
|
|
|
|
|
|
and len(k.split("/", 1)[1]) > 13
|
|
|
|
|
|
and k.split("/", 1)[1][12] == "_"
|
|
|
|
|
|
and all(ch in "0123456789abcdef" for ch in k.split("/", 1)[1][:12])
|
|
|
|
|
|
]
|
|
|
|
|
|
for k in ghost_svcs:
|
|
|
|
|
|
logger.info(f"Pruning ghost (hash-prefixed) service key from world state: {k}")
|
|
|
|
|
|
del self.world_state["services"][k]
|
|
|
|
|
|
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
now = time.time()
|
2026-06-03 12:26:49 +02:00
|
|
|
|
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
# --- Authoritative node liveness (fresh / stale / dead) -------------
|
|
|
|
|
|
# Status is derived from how long ago the node last reported, NOT from
|
|
|
|
|
|
# the last status event. This runs every cycle (including cycles with
|
|
|
|
|
|
# no new events — see run_once), so a node that silently stops sending
|
|
|
|
|
|
# heartbeats transitions on its own without any node_offline event ever
|
|
|
|
|
|
# being emitted. This is the fix for the "dead node shown NOMINAL"
|
|
|
|
|
|
# class of silent outage.
|
2026-07-09 15:38:20 +02:00
|
|
|
|
#
|
|
|
|
|
|
# SHADOW-READ (cutover etap 1): fetch Prometheus up{} ONCE per cycle
|
|
|
|
|
|
# (before the node loop, not per-node) for comparison-only logging. {}
|
|
|
|
|
|
# when disabled/unreachable — fail-open, event liveness is authoritative.
|
|
|
|
|
|
prom_liveness_map = self._query_prometheus_liveness()
|
2026-06-17 19:26:01 +02:00
|
|
|
|
for node_name, node_info in self.world_state["nodes"].items():
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
roles = (node_info.get("roles")
|
|
|
|
|
|
or self.inventory["nodes"].get(node_name, {}).get("roles", []))
|
|
|
|
|
|
liveness = compute_liveness(
|
|
|
|
|
|
node_info.get("last_seen"), now=now, ttls=ttls_for(node_name, roles)
|
|
|
|
|
|
)
|
|
|
|
|
|
if liveness == UNKNOWN:
|
|
|
|
|
|
# No last_seen yet — don't guess a node dead. Leave status as-is.
|
2026-06-17 19:26:01 +02:00
|
|
|
|
continue
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
prev = node_info.get("liveness")
|
|
|
|
|
|
node_info["liveness"] = liveness
|
|
|
|
|
|
new_status = liveness_to_status(liveness)
|
|
|
|
|
|
if new_status:
|
|
|
|
|
|
node_info["status"] = new_status
|
|
|
|
|
|
# Emit on a real transition only. prev is None on the very first
|
|
|
|
|
|
# classification after (re)start — that is a baseline, not an event.
|
|
|
|
|
|
if prev is not None and prev != liveness:
|
|
|
|
|
|
self._emit_node_transition(node_name, prev, liveness, node_info, now)
|
2026-07-09 15:38:20 +02:00
|
|
|
|
# SHADOW-READ (cutover etap 1): compare (log-only) the event-derived
|
|
|
|
|
|
# `liveness` computed above against Prometheus up{}. Purely additive —
|
|
|
|
|
|
# does NOT read back or mutate node_info["liveness"]/["status"], and
|
|
|
|
|
|
# the authoritative decision above is already committed.
|
|
|
|
|
|
self._shadow_compare_liveness(node_name, liveness, node_info, prom_liveness_map, now)
|
2026-06-17 19:26:01 +02:00
|
|
|
|
|
2026-06-03 14:29:12 +02:00
|
|
|
|
try:
|
|
|
|
|
|
# Collect incident_ids currently referenced by any service entry.
|
|
|
|
|
|
linked_ids: set = {
|
|
|
|
|
|
svc.get("incident_id")
|
|
|
|
|
|
for svc in self.world_state["services"].values()
|
|
|
|
|
|
if svc.get("incident_id")
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
# Case 1 — service is healthy but still points at an active incident.
|
|
|
|
|
|
# process_event already calls _resolve_incident on service_healthy events,
|
|
|
|
|
|
# but if the observer restarted with on-disk state where the link was
|
|
|
|
|
|
# intact (inconsistency from a pre-atomic-write crash), it may not get
|
|
|
|
|
|
# resolved until the next service_healthy event is processed. Resolve
|
|
|
|
|
|
# immediately — a healthy service cannot have an ongoing incident.
|
|
|
|
|
|
for svc_key, svc in self.world_state["services"].items():
|
|
|
|
|
|
if svc.get("status") != "healthy":
|
|
|
|
|
|
continue
|
2026-06-03 12:26:49 +02:00
|
|
|
|
inc_id = svc.get("incident_id")
|
2026-06-03 14:29:12 +02:00
|
|
|
|
if not inc_id:
|
|
|
|
|
|
continue
|
|
|
|
|
|
inc = self.world_state["incidents"].get(inc_id, {})
|
|
|
|
|
|
if inc.get("status") == "active":
|
|
|
|
|
|
logger.info(
|
|
|
|
|
|
f"Auto-resolving incident {inc_id} for {svc_key}: "
|
|
|
|
|
|
f"service is healthy"
|
|
|
|
|
|
)
|
|
|
|
|
|
inc["status"] = "resolved"
|
|
|
|
|
|
inc["resolved_at"] = now
|
|
|
|
|
|
svc["incident_id"] = None
|
|
|
|
|
|
linked_ids.discard(inc_id)
|
|
|
|
|
|
|
|
|
|
|
|
# Case 2 — orphaned active incident: no service entry links to it and
|
|
|
|
|
|
# last_occurrence is older than 5 minutes (guard against creation races).
|
|
|
|
|
|
# These are the stale records left behind when on-disk state was
|
|
|
|
|
|
# inconsistent: the service entry had incident_id cleared but incidents.json
|
|
|
|
|
|
# still had the record as "active".
|
|
|
|
|
|
for inc_id, inc in self.world_state["incidents"].items():
|
|
|
|
|
|
if inc.get("status") != "active":
|
|
|
|
|
|
continue
|
|
|
|
|
|
if inc_id in linked_ids:
|
|
|
|
|
|
continue
|
|
|
|
|
|
age = now - _parse_ts(inc.get("last_occurrence"))
|
|
|
|
|
|
if age > 300: # 5-minute guard
|
|
|
|
|
|
logger.info(
|
|
|
|
|
|
f"Auto-resolving orphaned incident {inc_id} "
|
|
|
|
|
|
f"(service={inc.get('service')}, node={inc.get('node')}): "
|
|
|
|
|
|
f"no service references it, age={int(age)}s"
|
|
|
|
|
|
)
|
|
|
|
|
|
inc["status"] = "resolved"
|
|
|
|
|
|
inc["resolved_at"] = now
|
|
|
|
|
|
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
logger.error(f"Error during incident auto-resolve in _prune_stale_world: {exc}")
|
2026-06-03 12:26:49 +02:00
|
|
|
|
|
|
|
|
|
|
# Remove resolved incidents older than 7 days.
|
2026-06-03 14:29:12 +02:00
|
|
|
|
# Use _parse_ts so ISO-string resolved_at values are handled correctly.
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
stale_incidents = [
|
|
|
|
|
|
k for k, v in self.world_state["incidents"].items()
|
|
|
|
|
|
if v.get("status") == "resolved"
|
2026-06-03 14:29:12 +02:00
|
|
|
|
and now - _parse_ts(v.get("resolved_at")) > 7 * 86400
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
]
|
|
|
|
|
|
for k in stale_incidents:
|
|
|
|
|
|
del self.world_state["incidents"][k]
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
def _save_world(self):
|
|
|
|
|
|
self.world_state["summary"]["last_update"] = datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
active_incidents = [
|
|
|
|
|
|
k for k, v in self.world_state["incidents"].items() if v.get("status") == "active"
|
|
|
|
|
|
]
|
|
|
|
|
|
self.world_state["summary"]["active_incidents_count"] = len(active_incidents)
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
self.world_state["summary"]["node_count"] = len(self.world_state["nodes"])
|
|
|
|
|
|
self.world_state["summary"]["service_count"] = len(self.world_state["services"])
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
if active_incidents:
|
|
|
|
|
|
self.world_state["summary"]["status"] = "degraded"
|
|
|
|
|
|
else:
|
|
|
|
|
|
self.world_state["summary"]["status"] = "nominal"
|
|
|
|
|
|
|
|
|
|
|
|
files = {
|
|
|
|
|
|
"nodes.json": self.world_state["nodes"],
|
|
|
|
|
|
"services.json": self.world_state["services"],
|
|
|
|
|
|
"deployments.json": self.world_state["deployments"],
|
|
|
|
|
|
"incidents.json": self.world_state["incidents"],
|
2026-06-03 12:26:49 +02:00
|
|
|
|
"recommendations.json": [],
|
2026-05-12 14:07:03 +02:00
|
|
|
|
"runtime-summary.json": self.world_state["summary"]
|
|
|
|
|
|
}
|
|
|
|
|
|
for filename, data in files.items():
|
|
|
|
|
|
try:
|
2026-06-03 12:26:49 +02:00
|
|
|
|
_atomic_write_json(WORLD_DIR / filename, data)
|
2026-05-12 14:07:03 +02:00
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to save {filename}: {e}")
|
|
|
|
|
|
|
|
|
|
|
|
def process_event(self, event):
|
|
|
|
|
|
etype = event.get("type")
|
|
|
|
|
|
node = event.get("node")
|
|
|
|
|
|
service = event.get("service")
|
|
|
|
|
|
severity = event.get("severity")
|
|
|
|
|
|
timestamp = event.get("timestamp")
|
|
|
|
|
|
cid = event.get("correlation_id")
|
|
|
|
|
|
payload = event.get("payload", {})
|
|
|
|
|
|
|
|
|
|
|
|
# 1. Update Node State
|
|
|
|
|
|
if node not in self.world_state["nodes"]:
|
|
|
|
|
|
self.world_state["nodes"][node] = {
|
|
|
|
|
|
"status": "unknown",
|
|
|
|
|
|
"last_seen": None,
|
|
|
|
|
|
"roles": self.inventory["nodes"].get(node, {}).get("roles", [])
|
|
|
|
|
|
}
|
|
|
|
|
|
self.world_state["nodes"][node]["last_seen"] = timestamp
|
feat(node-agent): implement health monitor and safe cleanup policy
scripts/monitor/health-monitor.sh (new):
- Standalone bash health monitor: disk/RAM/CPU checks + docker container health
- Per-node-type cleanup policy enforced:
lte_node (chelsty-infra, chelsty-ha): NO cleanup, no docker ops
sd_card (piha, saturn): dangling images + containers, rate-limited once/24h
ai_node (solaria): dangling + containers + build cache, NEVER -a
standard (vps): dangling + containers + build cache + CP filesystem rotation
- VPS filesystem rotation: completed/failed actions >7d, deploy logs >30d,
events >3d AND past observer checkpoint
- Emits structured JSON events (node_health, disk_pressure, high_memory, high_cpu,
containers_not_running, healthcheck_failed)
services/node-agent/ (new):
- Python daemon (node_agent.py): same policy as bash script, Docker SDK
for container checks and cleanup, /proc for system metrics
- Optional event shipping to VPS via rsync+SSH (VPS_EVENTS_HOST env var)
- Dockerfile: python:3.11-slim + openssh-client + rsync + docker>=6.0
- docker-compose.yml: mounts docker socket, /opt/homelab, repo read-only
observer.py:
- Handle node_health: update node status + disk/mem/cpu metrics, clear disk_pressure
- Handle disk_pressure: record severity on node, clear when healthy
- Handle high_memory / high_cpu: record pressure level for correlation
supervisor.py:
- Add NO_DISK_CLEANUP_NODES = {chelsty-infra, chelsty-ha}
- reconcile() step 3: generate disk_cleanup actions for nodes with high disk pressure
- _generate_disk_cleanup_recommendation(): stable ID disk-cleanup-{node},
checks all active states, risk=guarded (operator approval required)
executor.py:
- Handle disk_cleanup action type via _execute_disk_cleanup()
- Commands come from action payload; safety gate rejects any command touching
/opt/homelab/data/, /opt/homelab/config/, /opt/homelab/state/, or rm -rf /
hosts/*/services.yaml:
- Rename stability-agent -> node-agent on piha, vps, solaria, chelsty-infra
- Add node-agent to chelsty-ha (previously missing)
- Add cleanup policy notes to LTE node comments
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:15:06 +02:00
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
if etype == "node_online":
|
|
|
|
|
|
self.world_state["nodes"][node]["status"] = "online"
|
|
|
|
|
|
elif etype == "node_offline":
|
|
|
|
|
|
self.world_state["nodes"][node]["status"] = "offline"
|
|
|
|
|
|
|
feat(node-agent): implement health monitor and safe cleanup policy
scripts/monitor/health-monitor.sh (new):
- Standalone bash health monitor: disk/RAM/CPU checks + docker container health
- Per-node-type cleanup policy enforced:
lte_node (chelsty-infra, chelsty-ha): NO cleanup, no docker ops
sd_card (piha, saturn): dangling images + containers, rate-limited once/24h
ai_node (solaria): dangling + containers + build cache, NEVER -a
standard (vps): dangling + containers + build cache + CP filesystem rotation
- VPS filesystem rotation: completed/failed actions >7d, deploy logs >30d,
events >3d AND past observer checkpoint
- Emits structured JSON events (node_health, disk_pressure, high_memory, high_cpu,
containers_not_running, healthcheck_failed)
services/node-agent/ (new):
- Python daemon (node_agent.py): same policy as bash script, Docker SDK
for container checks and cleanup, /proc for system metrics
- Optional event shipping to VPS via rsync+SSH (VPS_EVENTS_HOST env var)
- Dockerfile: python:3.11-slim + openssh-client + rsync + docker>=6.0
- docker-compose.yml: mounts docker socket, /opt/homelab, repo read-only
observer.py:
- Handle node_health: update node status + disk/mem/cpu metrics, clear disk_pressure
- Handle disk_pressure: record severity on node, clear when healthy
- Handle high_memory / high_cpu: record pressure level for correlation
supervisor.py:
- Add NO_DISK_CLEANUP_NODES = {chelsty-infra, chelsty-ha}
- reconcile() step 3: generate disk_cleanup actions for nodes with high disk pressure
- _generate_disk_cleanup_recommendation(): stable ID disk-cleanup-{node},
checks all active states, risk=guarded (operator approval required)
executor.py:
- Handle disk_cleanup action type via _execute_disk_cleanup()
- Commands come from action payload; safety gate rejects any command touching
/opt/homelab/data/, /opt/homelab/config/, /opt/homelab/state/, or rm -rf /
hosts/*/services.yaml:
- Rename stability-agent -> node-agent on piha, vps, solaria, chelsty-infra
- Add node-agent to chelsty-ha (previously missing)
- Add cleanup policy notes to LTE node comments
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:15:06 +02:00
|
|
|
|
elif etype == "node_health":
|
|
|
|
|
|
# Regular heartbeat from node-agent; updates resource metrics.
|
|
|
|
|
|
# Clears disk_pressure if disk is now healthy (< warn threshold).
|
|
|
|
|
|
self.world_state["nodes"][node]["status"] = "online"
|
|
|
|
|
|
self.world_state["nodes"][node].update({
|
|
|
|
|
|
"disk_usage_pct": payload.get("disk_pct"),
|
|
|
|
|
|
"mem_usage_pct": payload.get("mem_pct"),
|
|
|
|
|
|
"cpu_usage_pct": payload.get("cpu_pct"),
|
|
|
|
|
|
})
|
|
|
|
|
|
if (payload.get("disk_pct") or 0) < 75:
|
|
|
|
|
|
self.world_state["nodes"][node].pop("disk_pressure", None)
|
|
|
|
|
|
|
|
|
|
|
|
elif etype == "disk_pressure":
|
|
|
|
|
|
# Emitted when disk usage crosses 75 % (medium) or 85 % (high).
|
|
|
|
|
|
# The supervisor reads disk_pressure to generate disk_cleanup actions.
|
|
|
|
|
|
self.world_state["nodes"][node]["disk_pressure"] = severity
|
|
|
|
|
|
self.world_state["nodes"][node]["disk_usage_pct"] = payload.get("usage_pct")
|
|
|
|
|
|
|
|
|
|
|
|
elif etype == "high_memory":
|
|
|
|
|
|
# Memory pressure observation; recorded on the node for correlation.
|
|
|
|
|
|
# No automated action — operator decides if a container restart helps.
|
|
|
|
|
|
self.world_state["nodes"][node]["memory_pressure"] = severity
|
|
|
|
|
|
self.world_state["nodes"][node]["mem_usage_pct"] = payload.get("usage_pct")
|
|
|
|
|
|
|
|
|
|
|
|
elif etype == "high_cpu":
|
|
|
|
|
|
# CPU pressure observation; recorded for visibility.
|
|
|
|
|
|
self.world_state["nodes"][node]["cpu_pressure"] = severity
|
|
|
|
|
|
self.world_state["nodes"][node]["cpu_usage_pct"] = payload.get("usage_pct")
|
|
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
# 2. Update Service State
|
|
|
|
|
|
if service and service != "all":
|
|
|
|
|
|
svc_key = f"{node}/{service}"
|
|
|
|
|
|
if svc_key not in self.world_state["services"]:
|
|
|
|
|
|
self.world_state["services"][svc_key] = {
|
|
|
|
|
|
"node": node,
|
|
|
|
|
|
"service": service,
|
|
|
|
|
|
"status": "unknown",
|
|
|
|
|
|
"last_check": None,
|
|
|
|
|
|
"incident_id": None
|
|
|
|
|
|
}
|
|
|
|
|
|
self.world_state["services"][svc_key]["last_check"] = timestamp
|
|
|
|
|
|
|
|
|
|
|
|
if etype == "service_recovered":
|
|
|
|
|
|
self.world_state["services"][svc_key]["status"] = "healthy"
|
|
|
|
|
|
self._resolve_incident(svc_key, timestamp)
|
2026-05-27 14:49:56 +02:00
|
|
|
|
elif etype == "service_healthy":
|
|
|
|
|
|
# Positive confirmation from node-agent that a managed container
|
|
|
|
|
|
# is running. This keeps services.json populated so the supervisor
|
|
|
|
|
|
# can correctly detect drift (absent entry = never reported = unknown,
|
2026-05-27 15:20:19 +02:00
|
|
|
|
# not the same as confirmed missing).
|
|
|
|
|
|
# Also resolve any active incident — if a service that had been
|
|
|
|
|
|
# unhealthy/crashing is now confirmed healthy, the incident is over.
|
2026-05-27 14:49:56 +02:00
|
|
|
|
self.world_state["services"][svc_key]["status"] = "healthy"
|
2026-05-27 15:20:19 +02:00
|
|
|
|
self._resolve_incident(svc_key, timestamp)
|
2026-07-14 20:21:03 +02:00
|
|
|
|
elif etype in ["service_unhealthy", "healthcheck_failed", "containers_not_running"]:
|
|
|
|
|
|
# containers_not_running: node-agent (node_agent.py) — and the
|
|
|
|
|
|
# stability-agent — report a *managed* container that has exited /
|
|
|
|
|
|
# dead or is crash-looping (restarting past RestartCount threshold).
|
|
|
|
|
|
# Treat it exactly like the other hard-failure signals: mark the
|
|
|
|
|
|
# service unhealthy and open an incident. The incident's
|
|
|
|
|
|
# trigger_type is the event type ("containers_not_running"), which
|
|
|
|
|
|
# the supervisor already recognises in CONTAINER_RESTART_TRIGGERS and
|
|
|
|
|
|
# remediates with a low-risk container_restart (vs. a full redeploy).
|
|
|
|
|
|
#
|
|
|
|
|
|
# Before this branch existed, containers_not_running fell through
|
|
|
|
|
|
# this if/elif chain entirely: node-agent saw the dead container and
|
|
|
|
|
|
# emitted the event, the observer bumped last_check but left status at
|
|
|
|
|
|
# its last "healthy" value and created NO incident. The supervisor
|
|
|
|
|
|
# then saw no drift and generated no action — so a dead / crash-looping
|
|
|
|
|
|
# container produced ZERO operator alerts. The monitoring loop was
|
|
|
|
|
|
# silently broken (matches the "action queue empty despite failures"
|
|
|
|
|
|
# symptom). (node_agent.py's own comment claims this event "rides the
|
|
|
|
|
|
# existing, supervisor-wired remediation path" — that path only exists
|
|
|
|
|
|
# once the observer opens the incident here.)
|
2026-05-12 14:07:03 +02:00
|
|
|
|
self.world_state["services"][svc_key]["status"] = "unhealthy"
|
|
|
|
|
|
self._handle_incident(svc_key, event)
|
2026-07-14 20:21:03 +02:00
|
|
|
|
elif etype in ["container_restarting", "container_state_unexpected"]:
|
|
|
|
|
|
# Intentionally observational — parity with node_agent.check_containers,
|
|
|
|
|
|
# which emits these as low/medium severity and explicitly does NOT wire
|
|
|
|
|
|
# them to any supervisor trigger. Causes: a transient post-deploy
|
|
|
|
|
|
# restart still below the crash-loop threshold, an operator `docker
|
|
|
|
|
|
# pause`, or an unhandled-but-not-fault docker state.
|
|
|
|
|
|
#
|
|
|
|
|
|
# We deliberately do NOT set status=unhealthy (that would make the
|
|
|
|
|
|
# supervisor generate remediation for a transient blip — a redeploy,
|
|
|
|
|
|
# since there is no CONTAINER_RESTART_TRIGGERS incident) and do NOT
|
|
|
|
|
|
# open an incident (noise). But the event must not vanish, so we
|
|
|
|
|
|
# record a lightweight observational trace on the service entry; the
|
|
|
|
|
|
# raw event also stays visible in the /events feed. If a restart is a
|
|
|
|
|
|
# real crash-loop it escalates on its own: node-agent re-emits it as
|
|
|
|
|
|
# containers_not_running once RestartCount crosses the threshold, and
|
|
|
|
|
|
# the branch above then alarms.
|
|
|
|
|
|
self.world_state["services"][svc_key]["last_observation"] = {
|
|
|
|
|
|
"type": etype,
|
|
|
|
|
|
"severity": severity,
|
|
|
|
|
|
"timestamp": timestamp,
|
|
|
|
|
|
"message": event.get("message"),
|
|
|
|
|
|
}
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
# 3. Update Deployment State
|
|
|
|
|
|
if etype.startswith("deployment_") and cid:
|
|
|
|
|
|
if cid not in self.world_state["deployments"]:
|
|
|
|
|
|
self.world_state["deployments"][cid] = {
|
|
|
|
|
|
"node": node,
|
|
|
|
|
|
"service": service,
|
|
|
|
|
|
"status": "unknown",
|
|
|
|
|
|
"started_at": None,
|
|
|
|
|
|
"finished_at": None,
|
|
|
|
|
|
"events": []
|
|
|
|
|
|
}
|
|
|
|
|
|
self.world_state["deployments"][cid]["events"].append({
|
|
|
|
|
|
"type": etype,
|
|
|
|
|
|
"timestamp": timestamp,
|
|
|
|
|
|
"payload": payload
|
|
|
|
|
|
})
|
|
|
|
|
|
if etype == "deployment_started":
|
|
|
|
|
|
self.world_state["deployments"][cid]["status"] = "in_progress"
|
|
|
|
|
|
self.world_state["deployments"][cid]["started_at"] = timestamp
|
|
|
|
|
|
elif etype == "deployment_completed":
|
|
|
|
|
|
self.world_state["deployments"][cid]["status"] = "completed"
|
|
|
|
|
|
self.world_state["deployments"][cid]["finished_at"] = timestamp
|
|
|
|
|
|
elif etype == "deployment_failed":
|
|
|
|
|
|
self.world_state["deployments"][cid]["status"] = "failed"
|
|
|
|
|
|
self.world_state["deployments"][cid]["finished_at"] = timestamp
|
|
|
|
|
|
# Deployment failure often creates an incident
|
|
|
|
|
|
self._handle_deployment_failure(event)
|
|
|
|
|
|
|
|
|
|
|
|
def _handle_incident(self, svc_key, event):
|
|
|
|
|
|
# Correlation: collapse repeated failures for the same service on the same node
|
|
|
|
|
|
active_incident = self.world_state["services"][svc_key].get("incident_id")
|
|
|
|
|
|
|
|
|
|
|
|
if active_incident and active_incident in self.world_state["incidents"]:
|
|
|
|
|
|
incident = self.world_state["incidents"][active_incident]
|
|
|
|
|
|
if incident["status"] == "active":
|
|
|
|
|
|
incident["last_occurrence"] = event["timestamp"]
|
|
|
|
|
|
incident["occurrence_count"] = incident.get("occurrence_count", 1) + 1
|
|
|
|
|
|
incident["events"].append(event["timestamp"])
|
|
|
|
|
|
return
|
|
|
|
|
|
|
|
|
|
|
|
# Create new incident
|
|
|
|
|
|
incident_id = f"inc-{int(time.time())}-{event.get('node')}-{event.get('service')}"
|
|
|
|
|
|
self.world_state["incidents"][incident_id] = {
|
|
|
|
|
|
"id": incident_id,
|
|
|
|
|
|
"node": event.get("node"),
|
|
|
|
|
|
"service": event.get("service"),
|
|
|
|
|
|
"status": "active",
|
|
|
|
|
|
"severity": event.get("severity"),
|
2026-05-27 12:42:03 +02:00
|
|
|
|
# trigger_type records the event type that opened this incident so that
|
|
|
|
|
|
# the supervisor can choose the appropriate remediation action
|
|
|
|
|
|
# (e.g. container_restart for containers_not_running / mqtt_unreachable
|
|
|
|
|
|
# vs. a full redeploy for other causes).
|
|
|
|
|
|
"trigger_type": event.get("type"),
|
2026-05-12 14:07:03 +02:00
|
|
|
|
"started_at": event.get("timestamp"),
|
|
|
|
|
|
"last_occurrence": event.get("timestamp"),
|
|
|
|
|
|
"occurrence_count": 1,
|
|
|
|
|
|
"events": [event["timestamp"]],
|
|
|
|
|
|
"correlation_id": event.get("correlation_id")
|
|
|
|
|
|
}
|
|
|
|
|
|
self.world_state["services"][svc_key]["incident_id"] = incident_id
|
|
|
|
|
|
|
|
|
|
|
|
def _resolve_incident(self, svc_key, timestamp):
|
|
|
|
|
|
incident_id = self.world_state["services"][svc_key].get("incident_id")
|
|
|
|
|
|
if incident_id and incident_id in self.world_state["incidents"]:
|
|
|
|
|
|
if self.world_state["incidents"][incident_id]["status"] == "active":
|
|
|
|
|
|
self.world_state["incidents"][incident_id]["status"] = "resolved"
|
|
|
|
|
|
self.world_state["incidents"][incident_id]["resolved_at"] = timestamp
|
|
|
|
|
|
self.world_state["services"][svc_key]["incident_id"] = None
|
|
|
|
|
|
|
|
|
|
|
|
def _handle_deployment_failure(self, event):
|
|
|
|
|
|
# Specific logic for deployment failures
|
|
|
|
|
|
svc_key = f"{event.get('node')}/{event.get('service')}"
|
|
|
|
|
|
self._handle_incident(svc_key, event)
|
|
|
|
|
|
|
|
|
|
|
|
# Link diagnostics if available in payload
|
|
|
|
|
|
incident_id = self.world_state["services"][svc_key].get("incident_id")
|
|
|
|
|
|
if incident_id and incident_id in self.world_state["incidents"]:
|
|
|
|
|
|
payload = event.get("payload", {})
|
|
|
|
|
|
if "diagnostics_file" in payload:
|
|
|
|
|
|
self.world_state["incidents"][incident_id]["diagnostics_ref"] = payload["diagnostics_file"]
|
|
|
|
|
|
elif "error" in payload:
|
|
|
|
|
|
self.world_state["incidents"][incident_id]["last_error"] = payload["error"]
|
|
|
|
|
|
|
|
|
|
|
|
def run_once(self):
|
2026-05-12 20:59:46 +02:00
|
|
|
|
# Update heartbeat
|
|
|
|
|
|
heartbeat_file = STATE_DIR / "observer.heartbeat"
|
|
|
|
|
|
try:
|
|
|
|
|
|
heartbeat_file.touch()
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.error(f"Failed to touch heartbeat file: {e}")
|
|
|
|
|
|
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
# Collect all event files grouped by node directory. A file is "new"
|
|
|
|
|
|
# when its event TIMESTAMP (from the filename, mtime fallback) is greater
|
|
|
|
|
|
# than the node's checkpoint timestamp — never a lexical path compare, so
|
|
|
|
|
|
# a lexically-smaller-but-newer filename (evt-piha-… after a stray
|
|
|
|
|
|
# evt-unknown-…) can no longer poison a node into skipping every event.
|
|
|
|
|
|
all_files = glob.glob(str(EVENTS_DIR / "**" / "*.json"), recursive=True)
|
2026-05-27 14:16:58 +02:00
|
|
|
|
|
2026-05-12 14:07:03 +02:00
|
|
|
|
new_files = []
|
2026-05-27 14:16:58 +02:00
|
|
|
|
for file_path in all_files:
|
2026-05-12 14:07:03 +02:00
|
|
|
|
try:
|
2026-05-27 14:16:58 +02:00
|
|
|
|
node_dir = str(Path(file_path).relative_to(EVENTS_DIR).parts[0])
|
|
|
|
|
|
except (IndexError, ValueError):
|
|
|
|
|
|
node_dir = "__unknown__"
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
ev_ts = _event_ts_from_path(file_path)
|
|
|
|
|
|
last_for_node = self.node_checkpoints.get(node_dir, 0)
|
|
|
|
|
|
if ev_ts > last_for_node:
|
|
|
|
|
|
new_files.append((ev_ts, node_dir, file_path))
|
|
|
|
|
|
|
|
|
|
|
|
# Process oldest-first (tie-break on path for determinism) so the
|
|
|
|
|
|
# checkpoint advances monotonically in time and a mid-batch crash resumes
|
|
|
|
|
|
# from the right place.
|
|
|
|
|
|
new_files.sort(key=lambda t: (t[0], t[2]))
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
if not new_files:
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
# Even if no new events, prune stale entries and refresh summary freshness.
|
|
|
|
|
|
self._prune_stale_world()
|
2026-05-12 20:59:46 +02:00
|
|
|
|
self._save_world()
|
2026-05-12 14:07:03 +02:00
|
|
|
|
return
|
|
|
|
|
|
|
2026-05-27 14:16:58 +02:00
|
|
|
|
logger.info(f"Processing {len(new_files)} new events across "
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
f"{len({n for _, n, _ in new_files})} node(s)")
|
|
|
|
|
|
for ev_ts, node_dir, file_path in new_files:
|
2026-05-12 14:07:03 +02:00
|
|
|
|
try:
|
|
|
|
|
|
with open(file_path, "r") as f:
|
|
|
|
|
|
event = json.load(f)
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
# Skip events the observer emitted itself (node liveness
|
|
|
|
|
|
# transitions). Re-ingesting them would call process_event,
|
|
|
|
|
|
# which sets last_seen = event timestamp and would resurrect a
|
|
|
|
|
|
# node we just declared dead. They are still consumed by the
|
|
|
|
|
|
# supervisor (alerting) and the panel event feed.
|
|
|
|
|
|
if event.get("source") == "observer":
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
if ev_ts > self.node_checkpoints.get(node_dir, 0):
|
|
|
|
|
|
self.node_checkpoints[node_dir] = ev_ts
|
feat(observer): 3-state node liveness (fresh/stale/dead) + transitions + read-time net
Fixes the "dead node shown NOMINAL" silent outage: node status was set only by
events and never expired, so a node that crashed/lost connectivity stayed
"online" forever (chelsty-infra was online for 16d, piha ~6d). The only thing
that flipped status to offline was a node_offline event, which an unreachable
node can never emit.
Now node status is derived from freshness (now - last_seen), recomputed every
observer cycle (incl. cycles with no new events):
- always-on: fresh <=180s, stale 180-600s, dead >600s (3x the 60s heartbeat)
- remote/LTE (chelsty-*): fresh <=900s, stale 900-3600s, dead >3600s
Thresholds + tier logic live in ONE shared helper, services/control-plane/src/
liveness.py, imported by the observer and both operator UIs (bind-mounted into
the agent-system webui image). No 3x copy.
Transitions are not silent: the observer emits node_stale / node_offline /
node_online (recovery) events tagged source=observer (skipped on re-ingest so
they never reset last_seen), routed by the supervisor to alert_only actions.
Read-time safety net: both UIs recompute liveness from last_seen at request
time, so a stalled observer still surfaces dead nodes. Services inherit their
node's liveness (cascade, variant B) without mutating services.json.
Replaces the earlier binary NODE_OFFLINE_TTL_SECS flip.
Tests: liveness unit tests, observer 3-state + transitions/recovery/baseline +
self-event skip, operator_ui read-time net + cascade, supervisor node-event
routing. 89 passed. docker compose config valid for both stacks.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 20:07:25 +02:00
|
|
|
|
continue
|
|
|
|
|
|
self.process_event(event)
|
fix(observer): checkpoint by timestamp not lexical path — lexically-smaller-but-newer events were silently skipped forever (poisoned node)
Per-node checkpoint now stores the last-processed event TIMESTAMP (int epoch)
instead of a file path compared lexically. A file is "new" iff its timestamp
(parsed from evt-<node>-<unixts>-<type>-<svc>.json, mtime fallback) exceeds the
node's checkpoint; processing is ordered by timestamp, not path.
Root cause (PIHA dead ~34d, 2026-07-12): a stray evt-unknown-<ts>-… file landed
in events/piha/, lexically greater than every evt-piha-… name. The lexical
checkpoint pinned there, so every genuinely newer piha event sorted "before" it
and was skipped forever. Event backlog grew to 7344 files, last_seen frozen,
shadow-read logged false SHADOW_LIVENESS_MISMATCH event=dead prom=up.
- _event_ts_from_path: filename epoch, mtime fallback; NEVER returns 0 for an
existing file (0 == "older than checkpoint" == the poison).
- _checkpoint_ts_from_value: graceful migration of pre-fix path-string
checkpoints (and the older last_processed_file format) to int epochs;
unparseable → 0 (reprocess all — safe, process_event is idempotent on
last_seen/world_state; bias to reprocess, never to skip).
- Preserved: quarantine of bad events, observer-source re-ingest guard.
- Regression tests (test_incident_lifecycle.py section 9): lexically-smaller-
but-newer processed, unparseable name falls back to mtime (not wedged),
ts-not-path ordering, both checkpoint-format migrations, helper units.
Separate bug filed in backlog (not fixed here): ha-diag-agent emits node=
"unknown" events (config.py node_name default) into another node's dir when
NODE_NAME reaches the compose volume path but not the app env — the source of
the poison file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 15:55:38 +02:00
|
|
|
|
# Advance per-node checkpoint by timestamp (only forward).
|
|
|
|
|
|
if ev_ts > self.node_checkpoints.get(node_dir, 0):
|
|
|
|
|
|
self.node_checkpoints[node_dir] = ev_ts
|
2026-05-12 14:07:03 +02:00
|
|
|
|
except Exception as e:
|
2026-06-12 13:11:15 +02:00
|
|
|
|
logger.error(
|
|
|
|
|
|
"Error processing node_dir=%s file=%s (%s: %s)",
|
|
|
|
|
|
node_dir, file_path, type(e).__name__, e,
|
|
|
|
|
|
)
|
|
|
|
|
|
self._quarantine_event_file(file_path, node_dir, e)
|
2026-05-12 14:07:03 +02:00
|
|
|
|
|
|
|
|
|
|
self._save_checkpoint()
|
fix(observer+operator-ui): fix stale world state, dict→list API, event time filter
Root cause of stale data:
- node_agent.py falls back to socket.gethostname() when NODE_NAME is unset.
Inside a Docker container this returns the 12-char container ID (e.g.
'be17cb6eb0f6'), not the host name. Observer ingested those events and
created ghost entries in world/nodes.json that never expired.
observer.py:
- _prune_stale_world(): removes node/service/incident entries for nodes absent
from topology inventory; called on every run_once() cycle (both new-events
and idle paths). Resolved incidents older than 7 days are also aged out.
- _save_world(): now writes node_count and service_count to runtime-summary.json
so the Dashboard's System Overview cards show real numbers instead of undefined.
operator_ui.py:
- current_nodes/services/deployments/incidents(): the observer stores world state
as keyed dicts; the frontend calls .map() which requires an array. All four
functions now convert the dict to a properly-shaped list. Each item has the
fields the Nodes, Services, Topology, Deployments, and Correlation views expect
(hostname, health, capabilities, desired_state, dependencies, etc.).
- current_incidents(): synthesises a human-readable 'message' field from node +
service + trigger_type (observer does not store one; dashboard showed undefined).
- current_events(): adds a 24 h time filter (EVENTS_MAX_AGE_HOURS env var,
default 24). Without this, every event file ever written was returned,
including events from ghost-node deploys.
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-27 13:51:03 +02:00
|
|
|
|
self._prune_stale_world()
|
2026-05-12 14:07:03 +02:00
|
|
|
|
self._save_world()
|
|
|
|
|
|
|
|
|
|
|
|
def loop(self, interval=5):
|
|
|
|
|
|
logger.info("Starting observer loop")
|
|
|
|
|
|
while True:
|
|
|
|
|
|
self.run_once()
|
|
|
|
|
|
time.sleep(interval)
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == "__main__":
|
|
|
|
|
|
import sys
|
|
|
|
|
|
observer = Observer()
|
|
|
|
|
|
if "--run-once" in sys.argv:
|
|
|
|
|
|
observer.run_once()
|
|
|
|
|
|
else:
|
|
|
|
|
|
observer.loop()
|