From 91db6829f4c41d88d0fc1cc8a71bd265670aa0ed Mon Sep 17 00:00:00 2001 From: oskar Date: Wed, 26 Aug 2026 21:08:04 +0200 Subject: [PATCH] fix(control-plane): unique container_restart action_id, no more history overwrite MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _generate_recommendation() built container_restart ids as the bare container-restart--. Two DIFFERENT incidents for the same node+service (e.g. a generic containers_not_running restart, later followed — after recovery and recurrence — by an unrelated restart for the same service) produced the identical id. Once the first action reached cancelled/completed/failed, the second action's own transition into that same directory silently overwrote the first one's history file. This is exactly what happened 2026-08-26 to a shadow-mode HA-websocket restart colliding with an unrelated 08-06 entry (docs/sessions/2026-08-26.md) — worked around by hand-renaming the file that session. Fix: suffix the id with the triggering incident's started_at — container-restart--- — NOT time.time() at generation call time. reconcile() calls _generate_recommendation() on every loop iteration while the drift persists, and the pending/ approved/running existence check immediately below is what makes that idempotent; it only works if repeated calls for the SAME ongoing incident produce the SAME id. started_at is fixed for an incident's whole life (observer._handle_incident only bumps last_occurrence/occurrence_count on repeat occurrences — see COMMIT-1-adjacent code) and changes only when a genuinely new incident opens for that service, which is exactly "same id while ongoing, different id on recurrence". When the incident record is missing/unlinked, fall back to the bare pre-fix id (container-restart--, no suffix) — NOT time.time(). This is not just a malformed-data corner case: observer._prune_stale_world Case 3 (commit 71a7af5) clears a service's incident_id after 24h of event silence even while the drift is still ongoing, so a live restarting service can naturally hit this path. time.time() would mint a new action_id — and a new pending file — on every single reconcile() tick, which is the exact non-idempotency this commit exists to fix, just via a different trigger. The bare id can't distinguish same-incident from different-incident recurrences the way the suffixed id can, but it is stable across calls, which is what the dedup check actually needs. Scope: only the generic CONTAINER_RESTART_TRIGGERS path (_generate_recommendation). Left unchanged, deliberately: - redeploy-- ids — no observed collision, out of scope for this fix (flagged as a latent follow-up below). - The HA-specific container-restart--homeassistant id used by _generate_ha_container_restart / _generate_ha_shadow_alert / _cancel_ha_container_restart: these three functions rely on an exact-match lookup of that fixed id (cooldown check via _ha_action_recently_completed, and the cancel path finding the specific pending file to move) — adding a suffix there would break both without a broader refactor to prefix-glob lookups. - alert-ha-*/alert-node-* ids: _ha_action_recently_completed also exact-matches these for cooldown dedup; a suffix would defeat cooldown entirely (every occurrence would look "new"). node-agent idempotency gate confirmed unaffected: _already_processed() in node_agent.py does a full-string action_id match against processed-actions/.done, guarding against RE-processing the exact same dispatched action file (e.g. a duplicate rsync delivery) — not against a new action_id for a new occurrence of the same service. A suffixed id is legitimately a new action to node-agent, which is the correct behavior (a genuine new incident should actually restart the container again). Tests: test_supervisor_action_id_uniqueness.py covers (1) repeated _generate_recommendation() calls for the same ongoing incident produce the same id and do not duplicate the pending file, (2) a new incident after the old one completed gets a different id and does not overwrite the old completed record, (3) fallback to the bare pre-fix id when the incident record is missing, (4) that bare fallback id is stable across repeated calls for the same missing-record drift — no duplicate pending file, same as case (1) but for the no-incident path, (5) redeploy ids stay bare. Updated test_observer_container_events.py's end-to-end assertion to match by prefix instead of exact filename. Full control-plane suite: 184 passed; node-agent suite: 70 passed (unchanged, confirming the idempotency gate needed no code change). Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_012pjmfPfrF5UYHqki2YwvdG --- services/control-plane/src/supervisor.py | 65 ++++++- .../tests/test_observer_container_events.py | 11 +- .../test_supervisor_action_id_uniqueness.py | 181 ++++++++++++++++++ 3 files changed, 249 insertions(+), 8 deletions(-) create mode 100644 services/control-plane/tests/test_supervisor_action_id_uniqueness.py diff --git a/services/control-plane/src/supervisor.py b/services/control-plane/src/supervisor.py index 6de3dc9..b97aa77 100644 --- a/services/control-plane/src/supervisor.py +++ b/services/control-plane/src/supervisor.py @@ -16,6 +16,23 @@ def _atomic_write_json(path: Path, data) -> None: os.fsync(f.fileno()) os.replace(tmp, path) + +def _parse_ts(ts) -> float: + """Return a Unix timestamp float from ts (int/float, or an ISO-8601 string + as stability-agent / events.py write it). Mirrors observer.observer._parse_ts + — duplicated rather than imported to keep supervisor.py's dependency on the + observer module (an unrelated component with its own import-time sys.path + manipulation) at zero. Returns 0.0 on None/unparseable input.""" + if ts is None: + return 0.0 + if isinstance(ts, (int, float)): + return float(ts) + try: + from datetime import datetime + return datetime.fromisoformat(str(ts).replace("Z", "+00:00")).timestamp() + except Exception: + return 0.0 + # Constants and Paths RUNTIME_PATH = os.getenv("RUNTIME_PATH", "/opt/homelab") WORLD_DIR = Path(RUNTIME_PATH) / "world" @@ -433,13 +450,51 @@ class Supervisor: service = drift["service"] trigger_type = drift.get("trigger_type") - # Choose action type first so we can build the stable, deterministic ID. - # Stable IDs mean reconcile is truly idempotent: the same drift always - # produces the same filename, so we never create duplicates even across - # restarts of the supervisor. + # Choose action type first so we can build the ID. + # + # container_restart IDs carry a suffix so two DIFFERENT incidents for + # the same node+service never collide in cancelled/completed/failed + # (2026-08-26: a generic containers_not_running restart and a later + # ha-diag-agent shadow-mode restart both used the bare + # container-restart-piha-homeassistant id and overwrote each other's + # history — worked around manually that session, see + # docs/sessions/2026-08-26.md). + # + # The suffix is the triggering incident's started_at, NOT time.time() + # at generation time: reconcile() calls _generate_recommendation on + # every loop iteration while the drift persists, and the pending/ + # approved/running existence check below is what makes that + # idempotent — it only works if repeated calls for the SAME ongoing + # incident produce the SAME action_id. started_at stays fixed for the + # life of one incident (observer._handle_incident only bumps + # last_occurrence/occurrence_count on repeat occurrences) and changes + # only when a new incident is opened for that service — exactly the + # cases we want "same id" and "different id" for, respectively. + # + # No incident record → bare id (the pre-fix format, no suffix at + # all), NOT time.time(). A missing/unlinked incident_id is not just + # a malformed-data edge case: observer._prune_stale_world Case 3 + # (commit 71a7af5) clears service.incident_id after 24h of event + # silence even while the underlying drift is still ongoing, so this + # path is hit by a live, still-restarting service. time.time() would + # mint a new action_id — and a new pending file — on every single + # reconcile() tick, defeating the dedup check below entirely. The + # bare id has no incident to distinguish "same" from "different" + # occurrences by, but it is at least stable across calls, which is + # what idempotency here actually requires. if trigger_type in CONTAINER_RESTART_TRIGGERS: - action_id = f"container-restart-{node}-{service}" + incident_id = self.actual_state["services"].get(drift["svc_key"], {}).get("incident_id") + incident = self.actual_state["incidents"].get(incident_id, {}) if incident_id else {} + started_ts = int(_parse_ts(incident.get("started_at"))) + if started_ts: + action_id = f"container-restart-{node}-{service}-{started_ts}" + else: + action_id = f"container-restart-{node}-{service}" else: + # redeploy IDs stay bare (node-service) — out of scope for this + # fix (see commit message: no observed collision here yet), and + # _cancel_resolved_pending_actions/_ha_action_recently_completed + # do not key off redeploy ids so nothing here depends on it. action_id = f"redeploy-{node}-{service}" # Skip if an action for this ID is already live in any active state diff --git a/services/control-plane/tests/test_observer_container_events.py b/services/control-plane/tests/test_observer_container_events.py index c2e80a0..9814997 100644 --- a/services/control-plane/tests/test_observer_container_events.py +++ b/services/control-plane/tests/test_observer_container_events.py @@ -178,10 +178,15 @@ def test_supervisor_generates_container_restart_from_observer_output( sup = Supervisor() sup.reconcile() - action_id = "container-restart-piha-paperless" - action_path = actions / "pending" / f"{action_id}.json" - assert action_path.exists(), "supervisor did not generate a container_restart" + # action_id carries a per-incident timestamp suffix (2026-08-26 fix: two + # different incidents for the same node/service must never collide in + # cancelled/completed/failed — see supervisor._generate_recommendation), + # so match by prefix rather than an exact id. + matches = list((actions / "pending").glob("container-restart-piha-paperless-*.json")) + assert len(matches) == 1, "supervisor did not generate exactly one container_restart" + action_path = matches[0] action = json.loads(action_path.read_text()) + assert action["action_id"] == action_path.stem assert action["type"] == "container_restart" assert action["node"] == "piha" assert action["service"] == "paperless" diff --git a/services/control-plane/tests/test_supervisor_action_id_uniqueness.py b/services/control-plane/tests/test_supervisor_action_id_uniqueness.py new file mode 100644 index 0000000..703b62e --- /dev/null +++ b/services/control-plane/tests/test_supervisor_action_id_uniqueness.py @@ -0,0 +1,181 @@ +"""action_id uniqueness for container_restart (2026-08-26 fix). + +Before this fix, _generate_recommendation() built container_restart action +ids as the bare `container-restart--` — no timestamp, no +incident reference. Two DIFFERENT incidents for the same node+service (e.g. a +generic containers_not_running restart followed, after recovery and a later +recurrence, by another restart) produced the identical id. Once the first +action reached cancelled/completed/failed, the second action's move into that +same directory overwrote the first one's history file outright — this is +exactly what happened on 2026-08-26 to a shadow-mode HA-websocket restart +colliding with an unrelated 08-06 entry (docs/sessions/2026-08-26.md), worked +around manually that session by hand-renaming the file. + +The fix anchors the id's suffix to the triggering incident's `started_at` +instead of wall-clock time-of-generation, so: + - repeated reconcile() calls for the SAME ongoing incident keep producing + the SAME action_id (required for the pending/approved/running dedup + check in _generate_recommendation to still work — see its comment), and + - a NEW incident (after the old one resolved) gets a NEW id, so its + eventual cancelled/completed/failed file cannot collide with the old + incident's. +""" +from __future__ import annotations + +import json +import sys +from pathlib import Path + +import pytest + +sys.path.insert(0, str(Path(__file__).parent.parent / "src")) +import supervisor as supervisor_module # noqa: E402 +from supervisor import Supervisor, CONTAINER_RESTART_TRIGGERS # noqa: E402 + + +@pytest.fixture +def sup(tmp_path, monkeypatch): + actions = tmp_path / "actions" + events = tmp_path / "events" + world = tmp_path / "world" + repo = tmp_path / "repo" + for d in (actions, events, world, repo / "inventory", repo / "hosts"): + d.mkdir(parents=True, exist_ok=True) + monkeypatch.setattr(supervisor_module, "ACTIONS_DIR", actions) + monkeypatch.setattr(supervisor_module, "EVENTS_DIR", events) + monkeypatch.setattr(supervisor_module, "WORLD_DIR", world) + monkeypatch.setattr(supervisor_module, "REPO_ROOT", repo) + s = Supervisor() + # _generate_recommendation only reads actual_state — bypass disk-backed + # _load_desired_state()/_load_actual_state() and set it directly. + s.actual_state = {"services": {}, "nodes": {}, "incidents": {}} + return s + + +def _drift(node, service, trigger_type="containers_not_running"): + return { + "type": "unhealthy_service", + "svc_key": f"{node}/{service}", + "node": node, + "service": service, + "status": "unhealthy", + "trigger_type": trigger_type, + } + + +def _seed_incident(sup, node, service, incident_id, started_at, status="active"): + svc_key = f"{node}/{service}" + sup.actual_state["services"][svc_key] = { + "node": node, "service": service, "status": "unhealthy", + "incident_id": incident_id, + } + sup.actual_state["incidents"][incident_id] = { + "id": incident_id, "status": status, + "trigger_type": "containers_not_running", + "started_at": started_at, + "last_occurrence": started_at, + } + + +def _pending(tmp_path): + return list((tmp_path / "actions" / "pending").glob("*.json")) + + +def test_repeated_calls_for_same_incident_produce_same_action_id_no_duplicate(sup, tmp_path): + """Simulates reconcile() calling _generate_recommendation() on every loop + iteration while the drift persists: must not spam a new pending action + each time.""" + _seed_incident(sup, "piha", "paperless", "inc-1000-piha-paperless", started_at=1000) + drift = _drift("piha", "paperless") + + sup._generate_recommendation(drift) + sup._generate_recommendation(drift) + sup._generate_recommendation(drift) + + pending = _pending(tmp_path) + assert len(pending) == 1, f"expected exactly one pending action, got {[p.name for p in pending]}" + assert pending[0].name == "container-restart-piha-paperless-1000.json" + + +def test_new_incident_after_old_completed_gets_different_action_id(sup, tmp_path): + """A second, later incident for the same node/service must not collide + with the first incident's already-completed action file.""" + _seed_incident(sup, "piha", "paperless", "inc-1000-piha-paperless", started_at=1000) + drift = _drift("piha", "paperless") + sup._generate_recommendation(drift) + + first_action_path = tmp_path / "actions" / "pending" / "container-restart-piha-paperless-1000.json" + assert first_action_path.exists() + + # First incident's action ran to completion (moved out of pending/ by the + # executor in real life; simulate that here). + completed_dir = tmp_path / "actions" / "completed" + completed_dir.mkdir(parents=True, exist_ok=True) + first_action = json.loads(first_action_path.read_text()) + first_action["status"] = "completed" + (completed_dir / first_action_path.name).write_text(json.dumps(first_action)) + first_action_path.unlink() + + # A NEW incident recurs later for the same node+service (old one resolved + # in between — new incident_id, new started_at). + _seed_incident(sup, "piha", "paperless", "inc-2000-piha-paperless", started_at=2000) + sup._generate_recommendation(drift) + + second_action_path = tmp_path / "actions" / "pending" / "container-restart-piha-paperless-2000.json" + assert second_action_path.exists() + # The first incident's completed record must be untouched — not overwritten. + assert json.loads((completed_dir / "container-restart-piha-paperless-1000.json").read_text())["status"] == "completed" + + +def test_fallback_to_bare_id_when_incident_record_missing(sup, tmp_path): + """Malformed/missing incident data (incident_id set on the service but no + matching record in incidents.json) must still produce a usable + (non-crashing) action_id, not block remediation. The fallback must be the + bare pre-fix id — NOT a time.time() suffix — since this path is hit + naturally (not just on malformed data): observer._prune_stale_world + Case 3 (commit 71a7af5) clears service.incident_id after 24h of event + silence while the drift is still ongoing, so a time.time() suffix would + mint a new id on every reconcile() tick forever.""" + svc_key = "piha/paperless" + sup.actual_state["services"][svc_key] = { + "node": "piha", "service": "paperless", "status": "unhealthy", + "incident_id": "inc-missing", + } + # Deliberately no matching entry in sup.actual_state["incidents"]. + drift = _drift("piha", "paperless") + + sup._generate_recommendation(drift) + + pending = _pending(tmp_path) + assert len(pending) == 1 + assert pending[0].name == "container-restart-piha-paperless.json" + + +def test_fallback_bare_id_stable_across_repeated_calls(sup, tmp_path): + """Same missing-incident-record scenario, but simulating reconcile() + calling _generate_recommendation() on every loop iteration while the + drift persists: must not spam a new pending action each time, exactly + like the has-an-incident-record case above.""" + svc_key = "piha/paperless" + sup.actual_state["services"][svc_key] = { + "node": "piha", "service": "paperless", "status": "unhealthy", + "incident_id": "inc-missing", + } + drift = _drift("piha", "paperless") + + sup._generate_recommendation(drift) + sup._generate_recommendation(drift) + sup._generate_recommendation(drift) + + pending = _pending(tmp_path) + assert len(pending) == 1, f"expected exactly one pending action, got {[p.name for p in pending]}" + assert pending[0].name == "container-restart-piha-paperless.json" + + +def test_redeploy_action_id_stays_bare(sup, tmp_path): + """Non-container_restart drift (redeploy path) is out of scope for this + fix and keeps its existing bare node-service id.""" + drift = _drift("piha", "outline", trigger_type="service_unhealthy") + sup._generate_recommendation(drift) + + assert (tmp_path / "actions" / "pending" / "redeploy-piha-outline.json").exists()