P7-1: fence surviving RUNNING snapshot rows with a durable cancel intent

Restore read only snap['pending'], so a window closed on live work left both
roots stored as their pre-restart status forever: the startup re-persist then
overwrote the snapshot and the one owner line named the task that was NOT lost.

Restore is the last holder of the pre-restart running list, but it must not be
a second terminal writer racing a worker that outlived SIGTERM. It now walks
snap['running'] before the stale early return and mints one durable cancel
intent per surviving row (reason='server_shutdown', source='snapshot_restore'),
skipping rows that are terminal, cancel-requested, already owned by an active
intent, or about to be revived as pending; an unreadable cancel authority mints
nothing and is disclosed instead. The existing custody path then claims, kills,
reconciles and writes the terminal with its own text, and the task-done seam
expires the open quiz and closes the paired owner wait (owner Q11=A).

The fenced ids are counted as terminalized_running in the existing restore
ledger row and returned through a keyword-only out-parameter, so the int return
and its asserts are unchanged. The boot notice fires on restored OR fenced work
and states an INTENT: custody writes each terminal result at least a watchdog
window (10s) later. The live startup order (restore, then kill_workers) is
unchanged and the minting is idempotent across boots.
This commit is contained in:
Ouroboros 2026-09-12 10:25:56 +03:00
parent 927ee5fb07
commit bc7037bd71
7 changed files with 288 additions and 24 deletions

View file

@ -1014,7 +1014,7 @@ Each Chat instance handles `open` by resynchronizing archive-aware durable histo
A headless task is ADDRESSED when it is admitted, not when it is displayed (`log_addressing.ingress_chat_id`). A registered project's run has exactly ONE destination: an explicit `chat_id` may only agree with that thread, and any other value — the hidden partition included — is refused with a typed 400 rather than honoured or silently overridden, because a run addressed away from its room is the one shape that puts a card in Main whose project holds none of its work. Without a Project, ordinary API tasks default to `HIDDEN_CHAT_ID` (0). The confirmed browser Publish flow explicitly carries `source="web"` and `WEB_UI_CHAT_ID` to request Main; source is caller-declared addressing on the existing owner API, not a new authentication proof. Other non-Project conversation addresses remain refused. A run scoped to a REGISTERED, active project is admitted into that project's thread (dialogue, children, attachments and answer in the room the owner already has; Main still receives the one host-stamped completion row), and Main is told it finished only when its work is actually in that room — addressed there at admission or BOUND to the project; Registration alone does not qualify. Every other run stays in the hidden partition, silent in every chat, read back through the terminal, `--result-json-out`, the chat-blind Logs panel and `GET /api/tasks/<id>`; a reserved but inactive project keeps its chat acceptable so the queue's lifecycle fence refuses with its own typed reason, and a project deleted mid-run keeps its reserved chat. The run is also NAMED at admission and chat promotion, without a new model call: a caller-supplied `title` (`ouroboros run --title` or the top-level contract field; `metadata.title` is refused with a 400 like `metadata.project_id`) is authorship and fills both `title` and `suggested_name`; otherwise the request's first line, stripped of markdown and capped at the project-name length, fills `suggested_name` ALONE, so a truncated prompt never outranks a real name coined later, and a `task_named` frame is broadcast on admission so the live card is never born showing its status phrase as a title — the client buffers a `task_named` that arrives before the card's record exists (`web/modules/chat.js`), so frame order does not matter.
`queue_snapshot.json` is an atomic recovery and diagnostic projection, not a second scheduler. It carries pending and running rows, acceptance and root-budget fences, resident/active/parked worker counts, assignable capacity, and any pool-disabled reason. Startup restores a recent snapshot into an otherwise-empty pending queue and never resurrects ordinary RUNNING work. A selected native owner-wait handoff remains eligible beyond snapshot age only through its current waiting source and acknowledged planned-restart transaction. Terminal tasks stay terminal, a task with an active durable cancel intent (or a legacy cancel-requested latch file) is left for cancellation custody, descendants below an accepted or sealed root finalize as cancelled, and malformed durable fence evidence fails closed. Snapshot capture copies the live containers under the queue lock because concurrent HTTP mutation can otherwise crash the supervisor mid-iteration.
`queue_snapshot.json` is an atomic recovery and diagnostic projection, not a second scheduler. It carries pending and running rows, acceptance and root-budget fences, resident/active/parked worker counts, assignable capacity, and any pool-disabled reason. Startup restores a recent snapshot into an otherwise-empty pending queue and never resurrects ordinary RUNNING work: instead it FENCES every surviving RUNNING row with a durable cancel intent (`reason='server_shutdown'`) and lets the ordinary cancellation custody path terminalize it, expire its open quiz and close the paired owner wait, so a window closed on live work ends honestly instead of leaving a ghost; the restore ledger row names them as `terminalized_running` and the boot notice states the intent, because custody writes each terminal result a watchdog window later. A selected native owner-wait handoff remains eligible beyond snapshot age only through its current waiting source and acknowledged planned-restart transaction. Terminal tasks stay terminal, a task with an active durable cancel intent (or a legacy cancel-requested latch file) is left for cancellation custody, descendants below an accepted or sealed root finalize as cancelled, and malformed durable fence evidence fails closed. Snapshot capture copies the live containers under the queue lock because concurrent HTTP mutation can otherwise crash the supervisor mid-iteration.
Pooled completion separates a finished file-save attempt from publishable terminal truth. `worker_process.worker_main` calls `headless.prepare_terminal_task_files` before its own non-ephemeral buffered `task_done`, after blocking post-task work; earlier answer/metrics frames retain their order. The worker stays occupied until that first attempt ends. Its private integer `_files_prepared_attempt` identifies the attempt, not save success, and never becomes a durable result or public event field. `events_task_done` re-reads CURRENT and uses `headless.terminal_task_files_ready`: a split drive requires the existing child-bound copyback projection and body, not merely an early terminal post-task checkpoint; pending refs may remain, but workspace artifact finalization cannot still be pending. Queue removal, slot release, accounting and project/evolution hooks stay with the normal event owner.

View file

@ -648,7 +648,8 @@ def _run_supervisor(settings: dict) -> None:
_migrate_startup_cancel_latches(DATA_DIR)
prior_worker_pids = _startup_worker_pids(DATA_DIR)
restored_pending = restore_pending_from_snapshot()
interrupted_running: list = []
restored_pending = restore_pending_from_snapshot(terminalized=interrupted_running)
kill_workers(preserve_pending=True)
spawn_workers(max_workers)
persist_queue_snapshot(reason="startup")
@ -670,11 +671,23 @@ def _run_supervisor(settings: dict) -> None:
_prune_delegated_snapshots()
if restored_pending > 0:
if restored_pending > 0 or interrupted_running:
st_boot = load_state()
if st_boot.get("owner_chat_id"):
send_with_budget(int(st_boot["owner_chat_id"]),
f"♻️ Restored pending queue from snapshot: {restored_pending} tasks.")
# The second clause states an INTENT, not an outcome: restore only
# fences an interrupted task with a durable cancel intent, and
# cancellation custody writes its terminal result a watchdog
# window later (task_lifecycle._INTENT_WATCHDOG_MIN_AGE_SEC).
notice = ["♻️"]
if restored_pending > 0:
notice.append(f"Restored pending queue from snapshot: {restored_pending} tasks.")
if interrupted_running:
count = len(interrupted_running)
notice.append(
f"Cancelling {count} task{'' if count == 1 else 's'} that "
f"{'was' if count == 1 else 'were'} still running when the server stopped."
)
send_with_budget(int(st_boot["owner_chat_id"]), " ".join(notice))
_startup_retired_settings_notice(settings)
auto_resume_after_restart()

View file

@ -12,7 +12,7 @@ import json
import logging
import pathlib
import time
from typing import Optional
from typing import Any, Optional
from ouroboros.contracts.schema_versions import SCHEMA_VERSION_KEY
from ouroboros.utils import utc_now_iso
@ -210,8 +210,90 @@ def parse_iso_to_ts(iso_ts: str) -> Optional[float]:
return None
def restore_pending_from_snapshot(max_age_sec: int = 900) -> int:
"""Restore recent pending tasks from queue snapshot."""
def _fence_snapshot_running_rows(rows: Any, *, restored_ids: "set[str]") -> "list[str]":
"""Fence every RUNNING row that survived the shutdown with a durable cancel intent.
Restore is the last holder of the pre-restart running list, but it is NOT a
terminal writer: minting the intent hands each row to the one settle owner,
which claims it, kills a worker that outlived SIGTERM, reconciles, and only
then writes the terminal with its own text — expiring an open quiz and
closing the paired owner wait through the task-done seam. A second writer
here would race that surviving worker; an intent cannot. An UNREADABLE
cancel authority mints nothing: the unknown is disclosed, never fenced.
Returns the fenced task ids.
"""
from ouroboros.cancel_intents import has_active_intent, request_cancel
from ouroboros.task_results import (
_TRULY_TERMINAL_STATUSES, STATUS_CANCEL_REQUESTED, load_task_result,
)
fenced: list[str] = []
unreadable: list[str] = []
for row in rows if isinstance(rows, list) else []:
task_id = str(row.get("id") or "") if isinstance(row, dict) else ""
if not task_id or task_id in restored_ids:
continue
try:
stored = load_task_result(_queue().DRIVE_ROOT, task_id, strict=True) or {}
status = str(stored.get("status") or "")
if status in _TRULY_TERMINAL_STATUSES or status == STATUS_CANCEL_REQUESTED:
continue
if has_active_intent(_queue().DRIVE_ROOT, task_id, strict=True):
continue # cancellation custody already owns this row
intent = request_cancel(
_queue().DRIVE_ROOT, task_id,
reason="server_shutdown", source="snapshot_restore",
)
except Exception:
unreadable.append(task_id)
log.warning("Snapshot restore left running row %s unfenced: its cancel "
"authority is unreadable", task_id, exc_info=True)
continue
if not intent.get("already_settled"):
fenced.append(task_id)
if unreadable:
_queue().append_jsonl(
_queue().DRIVE_ROOT / "logs" / "supervisor.jsonl",
{"ts": utc_now_iso(), "type": "queue_restore_running_fence_unreadable",
"task_ids": unreadable},
)
return fenced
def _record_queue_restore(
*, restored: int = 0, skipped_terminal: int = 0,
cancel_authority_holds: Optional[list] = None, blocked_admission: Optional[list] = None,
invalid_task_depth: Optional[list] = None, terminalized_running: Optional[list] = None,
) -> None:
"""The one durable row a restore leaves: what it revived, what it left to
cancellation custody, and which surviving RUNNING rows it fenced. A stale
snapshot with nothing to revive still records the fences it minted."""
if not (restored or skipped_terminal or blocked_admission or terminalized_running):
return
_queue().append_jsonl(
_queue().DRIVE_ROOT / "logs" / "supervisor.jsonl",
{
"ts": utc_now_iso(),
"type": "queue_restored_from_snapshot",
"restored_pending": restored,
"skipped_terminal": skipped_terminal,
"cancel_authority_holds": list(cancel_authority_holds or []),
"blocked_admission": list(blocked_admission or []),
"invalid_task_depth": list(invalid_task_depth or []),
"terminalized_running": list(terminalized_running or []),
},
)
def restore_pending_from_snapshot(
max_age_sec: int = 900, *, terminalized: Optional[list] = None,
) -> int:
"""Restore recent pending tasks from queue snapshot.
Returns the number of PENDING rows revived. ``terminalized`` collects the ids
of surviving RUNNING rows fenced with a cancel intent, so the caller can name
them without changing what the returned count means.
"""
if _queue().PENDING:
return 0
try:
@ -243,7 +325,17 @@ def restore_pending_from_snapshot(max_age_sec: int = 900) -> int:
if (not task.get("_owner_wait_resume") and not stale)
or (task.get("_owner_wait_resume") and restore_owner_wait_allowed(_queue().DRIVE_ROOT, task))
]
# The pre-restart RUNNING rows are read HERE, before the stale gate: this
# is the last moment the list exists, and a stale snapshot is exactly the
# case where nothing else will ever settle them.
fenced_running = _fence_snapshot_running_rows(
snap.get("running"),
restored_ids={str(task.get("id") or "") for task in snapshot_pending},
)
if terminalized is not None:
terminalized.extend(fenced_running)
if stale and not snapshot_pending:
_record_queue_restore(terminalized_running=fenced_running)
return 0
snapshot_pending, pending_by_id, restored = restore_terminalization_retry_rows(
snapshot_pending, pending=_queue().PENDING, running=_queue().RUNNING,
@ -429,18 +521,12 @@ def restore_pending_from_snapshot(max_age_sec: int = 900) -> int:
"root_task_ids": sorted(fenced_roots),
},
)
if restored > 0 or skipped_terminal > 0 or blocked_restore:
_queue().append_jsonl(
_queue().DRIVE_ROOT / "logs" / "supervisor.jsonl",
{
"ts": utc_now_iso(),
"type": "queue_restored_from_snapshot",
"restored_pending": restored,
"skipped_terminal": skipped_terminal,
"cancel_authority_holds": cancel_authority_holds,
"blocked_admission": blocked_restore, "invalid_task_depth": invalid_depth_restore,
},
)
_record_queue_restore(
restored=restored, skipped_terminal=skipped_terminal,
cancel_authority_holds=cancel_authority_holds,
blocked_admission=blocked_restore, invalid_task_depth=invalid_depth_restore,
terminalized_running=fenced_running,
)
from supervisor.queue_transitions import sweep_orphaned_budget_fences
sweep_orphaned_budget_fences(

View file

@ -829,3 +829,135 @@ def test_steer_refusal_removes_the_just_staged_attachments(tmp_path, monkeypatch
from ouroboros.owner_mailbox import drain_owner_messages
assert drain_owner_messages(tmp_path, "steer-stage") == []
def _restart_snapshot(qenv, monkeypatch, *, running: list, pending: list = (), ts: str = ""):
"""Write a queue snapshot the way a shutdown left it and point restore at it."""
from ouroboros.utils import utc_now_iso
state_dir = qenv.drive / "state"
state_dir.mkdir(parents=True, exist_ok=True)
path = state_dir / "queue_snapshot.json"
path.write_text(json.dumps({
"ts": ts or utc_now_iso(),
"pending": [{"task": task} for task in pending],
"running": running,
"acceptance_fences": [],
"budget_root_fences": [],
}), encoding="utf-8")
monkeypatch.setattr(qenv.q, "QUEUE_SNAPSHOT_PATH", path, raising=False)
return path
def test_snapshot_restore_fences_a_surviving_running_row_for_custody(qenv, monkeypatch):
"""Q11=A: work the window killed ends as Cancelled — through the ONE settle
owner. Restore only mints the durable intent; the terminal write, the kill and
the reconcile stay with cancellation custody a watchdog window later."""
import time
from supervisor import task_lifecycle
write_task_result(qenv.drive, "interrupted-root", STATUS_RUNNING, chat_id=1)
_restart_snapshot(qenv, monkeypatch, running=[
{"id": "interrupted-root", "task": {"id": "interrupted-root", "chat_id": 1}},
])
fenced: list = []
assert qenv.q.restore_pending_from_snapshot(terminalized=fenced) == 0
assert fenced == ["interrupted-root"]
# Restore is NOT a terminal writer: the row is untouched, the intent is durable.
assert load_task_result(qenv.drive, "interrupted-root")["status"] == STATUS_RUNNING
intent = ci.active_intent(qenv.drive, "interrupted-root")
assert intent and intent["reason"] == "server_shutdown"
assert intent["source"] == "snapshot_restore"
rows = [json.loads(line) for line in
(qenv.drive / "logs" / "supervisor.jsonl").read_text(encoding="utf-8").splitlines()]
restore_rows = [row for row in rows if row["type"] == "queue_restored_from_snapshot"]
assert restore_rows[-1]["terminalized_running"] == ["interrupted-root"]
assert restore_rows[-1]["restored_pending"] == 0
# A second boot before custody ran re-reads the same row and mints nothing new.
second: list = []
assert qenv.q.restore_pending_from_snapshot(terminalized=second) == 0
assert second == []
assert ci.active_intent(qenv.drive, "interrupted-root")["request_id"] == intent["request_id"]
# The existing watchdog half terminalizes it; the text is custody's own.
outcomes = task_lifecycle.sweep_cancel_intents(now=time.time() + 60)
assert outcomes["interrupted-root"] == "cancelled"
settled = load_task_result(qenv.drive, "interrupted-root")
assert settled["status"] == STATUS_CANCELLED
assert settled["result"].startswith("Task cancelled")
assert ci.active_intent(qenv.drive, "interrupted-root") is None
def test_snapshot_restore_leaves_owned_and_terminal_running_rows_alone(qenv, monkeypatch):
"""The fence is for rows nothing else owns: an active intent belongs to
cancellation custody, and a task that finished stays finished."""
write_task_result(qenv.drive, "already-owned", STATUS_RUNNING, chat_id=1)
owned = ci.request_cancel(qenv.drive, "already-owned", reason="owner pressed Stop")
write_task_result(qenv.drive, "finished", STATUS_COMPLETED, chat_id=1, result="done")
write_task_result(qenv.drive, "revived", "scheduled", chat_id=1)
_restart_snapshot(
qenv, monkeypatch,
running=[
{"id": "already-owned", "task": {"id": "already-owned", "chat_id": 1}},
{"id": "finished", "task": {"id": "finished", "chat_id": 1}},
{"id": "", "task": {}},
],
pending=[{"id": "revived", "chat_id": 1, "type": "chat"}],
)
fenced: list = []
assert qenv.q.restore_pending_from_snapshot(terminalized=fenced) == 1
assert fenced == []
assert [task["id"] for task in qenv.q.PENDING] == ["revived"]
assert ci.active_intent(qenv.drive, "already-owned")["request_id"] == owned["request_id"]
assert ci.active_intent(qenv.drive, "already-owned")["reason"] == "owner pressed Stop"
assert ci.active_intent(qenv.drive, "finished") is None
assert load_task_result(qenv.drive, "finished")["status"] == STATUS_COMPLETED
def test_snapshot_restore_fence_expires_the_open_quiz_and_closes_the_owner_wait(qenv, monkeypatch):
"""A restart that kills a task with an open question reuses the EXISTING
expiry state: no new 'lost to a restart' quiz state is invented."""
import time
from ouroboros import owner_quiz
from supervisor import queue_transitions, task_lifecycle, workers
events: list = []
monkeypatch.setattr(workers, "get_event_q",
lambda: types.SimpleNamespace(put=events.append), raising=False)
task_id = "asked-and-interrupted"
write_task_result(qenv.drive, task_id, STATUS_RUNNING, chat_id=1)
owner_quiz.record_asked(
qenv.drive, task_id, quiz_id="q1", question="Which folder?",
options=["A", "B"], wait_for_answer=True,
)
write_task_result(
qenv.drive, task_id, STATUS_RUNNING,
owner_wait={"state": "waiting", "quiz_id": "q1"},
)
_restart_snapshot(qenv, monkeypatch, running=[
{"id": task_id, "owner_wait": {"state": "waiting"}, "task": {"id": task_id, "chat_id": 1}},
])
fenced: list = []
assert qenv.q.restore_pending_from_snapshot(terminalized=fenced) == 0
assert fenced == [task_id]
assert owner_quiz.quiz_states(qenv.drive, task_id)["q1"]["state"] == "open"
assert task_lifecycle.sweep_cancel_intents(now=time.time() + 60)[task_id] == "cancelled"
assert load_task_result(qenv.drive, task_id)["status"] == STATUS_CANCELLED
# Custody publishes the terminal event; the supervisor's task-done seam
# (events_task_done -> reconcile_terminal_task_projections) is what closes
# the per-task owner-control projections, exactly as for any other terminal.
done = [event for event in events if event["type"] == "task_done"]
assert len(done) == 1 and done[0]["status"] == STATUS_CANCELLED
queue_transitions.reconcile_terminal_task_projections(qenv.drive, task_id)
assert owner_quiz.quiz_states(qenv.drive, task_id)["q1"]["state"] == "expired_terminal"
assert load_task_result(qenv.drive, task_id)["owner_wait"]["state"] == "expired_terminal"

View file

@ -213,4 +213,4 @@ def test_the_supervisor_boot_calls_the_notice_after_the_queue_restore():
source = (pathlib.Path(__file__).resolve().parents[1] / "server.py").read_text(encoding="utf-8")
body = source.split("def _run_supervisor(settings: dict) -> None:", 1)[1]
assert "_startup_retired_settings_notice(settings)" in body.split("\ndef ", 1)[0]
assert body.index("restore_pending_from_snapshot()") < body.index("_startup_retired_settings_notice(settings)")
assert body.index("restore_pending_from_snapshot(") < body.index("_startup_retired_settings_notice(settings)")

View file

@ -296,7 +296,7 @@ def test_supervisor_startup_restores_queue_before_worker_reset():
import server
source = inspect.getsource(server._run_supervisor)
restore = source.index("restored_pending = restore_pending_from_snapshot()")
restore = source.index("restored_pending = restore_pending_from_snapshot(")
reset = source.index("kill_workers(preserve_pending=True)")
spawn = source.index("spawn_workers(max_workers)")
assert restore < reset < spawn
@ -580,6 +580,16 @@ class _Recorder:
self.stop = _FakeStopEvent()
self.restart = None
self.ready = None
# Snapshot restore reports revived PENDING rows in its return value and
# the RUNNING rows it fenced for cancellation custody through the
# keyword-only out-parameter; a test can script either.
self.restored_pending = 0
self.fenced_running: list = []
def restore(self, *_args, terminalized=None, **_kwargs):
if terminalized is not None:
terminalized.extend(self.fenced_running)
return self.restored_pending
def _supervisor_harness(monkeypatch, tmp_path, steps):
@ -669,7 +679,7 @@ def _supervisor_harness(monkeypatch, tmp_path, steps):
"persist_queue_snapshot", "cancel_task_by_id", "queue_deep_self_review_task",
"sort_pending", "check_scheduled_tasks"):
monkeypatch.setattr(queue_pkg, name, noop)
monkeypatch.setattr(queue_pkg, "restore_pending_from_snapshot", lambda: 0)
monkeypatch.setattr(queue_pkg, "restore_pending_from_snapshot", rec.restore)
for name in ("init", "spawn_workers", "kill_workers", "assign_tasks", "ensure_workers_healthy",
"auto_resume_after_restart"):
monkeypatch.setattr(workers_mod, name, noop)
@ -687,6 +697,28 @@ def _run(rec, server):
return rec
@pytest.mark.parametrize("restored, fenced, expected", [
(0, ["a"], "♻️ Cancelling 1 task that was still running when the server stopped."),
(2, ["a", "b"],
"♻️ Restored pending queue from snapshot: 2 tasks. "
"Cancelling 2 tasks that were still running when the server stopped."),
])
def test_boot_notice_states_the_restart_cancellations_as_an_intent(
monkeypatch, tmp_path, restored, fenced, expected):
"""A window closed on live work: the boot line names how many tasks the
restart interrupted and says they are BEING cancelled. It cannot claim they
ENDED: restore only minted the durable intent, and cancellation custody
writes each terminal result a watchdog window later."""
import server
rec = _supervisor_harness(monkeypatch, tmp_path, ["stop"])
rec.restored_pending = restored
rec.fenced_running = list(fenced)
_run(rec, server)
assert [text for _chat, text in rec.alerts if text.startswith("♻️")] == [expected]
def test_exception_while_stopping_exits_quietly_without_counting_a_crash(monkeypatch, tmp_path, caplog):
"""The graceful-shutdown race: the teardown sets the stop flag, then the tick
meets the torn-down bus (BrokenPipe). That is not a crash: no owner alarm,

View file

@ -314,7 +314,8 @@ def test_real_supervisor_orders_custody_recovery_before_prune(roots, monkeypatch
order = []
monkeypatch.setattr(server, "_migrate_startup_cancel_latches", lambda root: order.append("migrate"))
monkeypatch.setattr(server, "_startup_worker_pids", lambda root: order.append("capture-pids") or {777})
monkeypatch.setattr(queue, "restore_pending_from_snapshot", lambda: order.append("restore") or 0)
monkeypatch.setattr(queue, "restore_pending_from_snapshot",
lambda **_kw: order.append("restore") or 0)
monkeypatch.setattr(workers, "kill_workers", lambda **k: order.append("kill"))
monkeypatch.setattr(workers, "spawn_workers", lambda n: order.append("spawn"))
monkeypatch.setattr(server, "_startup_custody_sweep", lambda: order.append("custody"))