Review round 3: honest replay reminder, one-shot exactly-once vs re-enable, once-only pruning

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-08-18 04:58:21 +03:00
parent 74c11db4bb
commit 627e517ee7
6 changed files with 186 additions and 29 deletions

View file

@ -77,6 +77,28 @@ async def api_schedules_upsert(request: Request) -> JSONResponse:
enabled = _enabled_value(body)
if isinstance(enabled, str):
return json_error(enabled, 400)
completed_at = ""
if trigger.get("type") == "once":
# Exactly-once vs re-enable: a one-shot that already fired (non-empty
# completed_at) cannot be re-armed by flipping enabled back on — that
# would silently re-run the consumed task. Re-arming requires a NEW
# trigger.run_at, which clears the consumed receipt; a disable/edit that
# keeps the same run_at carries the receipt forward so GC still sees it.
from supervisor.queue import list_scheduled_tasks
wanted = str(body.get("id") or "").strip()
existing = next(
(item for item in list_scheduled_tasks(request_drive_root(request)).get("tasks") or []
if isinstance(item, dict) and str(item.get("id") or "") == wanted), None)
prev = (existing or {}).get("trigger")
prev = prev if isinstance(prev, dict) else {}
if (existing is not None and str(existing.get("completed_at") or "")
and str(prev.get("run_at") or "") == trigger["run_at"]):
if enabled:
return json_error(
"this one-shot schedule already fired; re-arming it requires a new "
"trigger.run_at (a fresh run_at clears completed_at)", 400)
completed_at = str(existing.get("completed_at"))
record = {
"id": str(body.get("id") or "").strip(),
"name": str(body.get("name") or body.get("id") or "scheduled-task").strip(),
@ -86,6 +108,8 @@ async def api_schedules_upsert(request: Request) -> JSONResponse:
"trigger": trigger,
"task": task,
}
if completed_at:
record["completed_at"] = completed_at
from supervisor.queue import upsert_scheduled_task
return JSONResponse({"ok": True, "schedule": upsert_scheduled_task(record, drive_root=request_drive_root(request))})

View file

@ -470,7 +470,9 @@ def force_plan_decision(
not_required = {"required": False, "allow": True, "status": "not_required"}
if bool(getattr(ctx, "is_ephemeral_turn", False)):
return not_required
from ouroboros.task_results import load_plan_review_state, plan_review_gate_projection
from ouroboros.task_results import (
current_plan_review_wave, load_plan_review_state, plan_review_gate_projection,
)
try:
root = _canonical_root(ctx)
@ -496,6 +498,12 @@ def force_plan_decision(
"self_opened": not bool(metadata.get("force_plan")),
**plan_review_gate_projection(state, effective, hard_rail=hard_rail),
}
if decision.get("reviewer_slots_degraded"):
# The reminder's replay promise is conditional on the recorded wave's
# structural health epoch (empty epoch = a re-dispatch is PAID), so the
# epoch fact rides the decision for plan_review_reminder.
decision["degraded_health_epoch"] = (
(current_plan_review_wave(state) or {}).get("health_epoch") or "")
if hurry_armed and str(enforcement or "").lower() == "blocking":
# Attribution only (task detail); the durable state and the configured
# global enforcement are byte-identical before/after.
@ -516,13 +524,19 @@ def plan_review_reminder(decision: Dict[str, Any]) -> str:
"plan_task with your goal, plan and spec to start a fresh review before finalizing."
)
if decision.get("reviewer_slots_degraded"):
# B2: facts, never a retry coach (P5). The recorded wave carries each slot's typed
# state; an identical envelope replays it free; a changed spec buys the next cycle.
# B2: facts, never a retry coach (P5). The replay promise is CONDITIONAL —
# wording SSOT: plan_render._degraded_replay_note (a free replay exists only
# while a non-empty recorded epoch and the reviewer roster stand; an
# empty-epoch wave re-dispatches a paid panel). Lazy import; plan_render
# never imports owner_hurry, so no cycle.
from ouroboros.tools.plan_render import _degraded_replay_note
note = _degraded_replay_note({"health_epoch": decision.get("degraded_health_epoch")})
return (
f"{tag} Blocking plan review is OPEN with no parseable reviewer quorum (DEGRADED). "
"The recorded wave lists each reviewer slot's typed state and reset time; an "
"identical envelope replays that recorded result at no cost, and a changed spec "
"starts the next paid cycle. Implementation stays held while the review is open."
f"The recorded wave lists each reviewer slot's typed state and reset time; {note}; "
"a changed spec starts the next paid cycle. Implementation stays held while the "
"review is open."
)
if outcome == "REVIEW_REQUIRED":
return (

View file

@ -146,9 +146,8 @@ from supervisor.task_admission import ( # noqa: E402,F401 - public queue API
reserve_task_admission,
)
# Variant A off-loop worker reaper lives in supervisor/task_reaper.py (extracted for
# module size). Re-export the thin names the enforce path and tests use; monkeypatching
# these queue-module names still works because the enforce path references them here.
# Variant A off-loop worker reaper lives in supervisor/task_reaper.py (module size); re-export
# the thin names the enforce path and tests use — monkeypatching these queue names still works.
from supervisor.task_reaper import ( # noqa: E402,F401 — re-exported for enforce path + tests
ensure_reaper_started as _ensure_reaper_started,
reap_queue as _reap_queue,
@ -212,9 +211,8 @@ def enqueue_task(
task_id = str(t.get("id") or "").strip()
reserved_token = str(ADMISSION_RESERVATIONS.get(task_id) or "")
if reserved_token and admission_token != reserved_token:
# A reservation owns this id until its request either enqueues or
# releases it. Tokenless internal callers and competing ingress
# requests must not be able to consume/collide with that id.
# A reservation owns this id until its request either enqueues or releases it.
# Tokenless internal callers and competing ingress must not consume/collide with it.
t["_admission_blocked"] = "admission_reservation_owned"
return t
if require_worker_pool:
@ -581,8 +579,11 @@ def check_scheduled_tasks() -> None:
now = now_utc.astimezone(tz)
expr = ""
if trigger_type == "once":
# One-shot entry (B2b W=A): fires once at/after run_at, then is
# marked done below — the same admission path as cron schedules.
# One-shot (B2b W=A): fires once at/after run_at via the same admission path
# as cron, then is marked done below. A consumed receipt (non-empty completed_at)
# NEVER re-fires even re-enabled from UI; re-arm = gateway upsert, fresh run_at.
if record.get("completed_at"):
continue
due, once_error = _once_due(trigger, tz, now)
if once_error:
changed = _record_last_error(record, once_error) or changed
@ -936,9 +937,8 @@ def restore_pending_from_snapshot(max_age_sec: int = 900) -> int:
# Never resurrect a terminal/cancelled task as a ghost pending entry.
# AR2-10 (§8-A1): the intent projection is consulted UNDER the queue lock at
# restore — the "no active intent" read and the enqueue form one serialized step
# against assignment/drop, the same invariant the pre-assignment consult keeps.
# Boot-time and contention-free; _queue_lock is an RLock, so enqueue_task's own
# acquisition stays re-entrant.
# against assignment/drop (same invariant as the pre-assignment consult). Boot-time
# and contention-free; _queue_lock is an RLock, so enqueue_task stays re-entrant.
with _queue_lock:
skip_revival = False
try:

View file

@ -68,21 +68,24 @@ def record_last_error(record: Dict[str, Any], message: str) -> bool:
def prune_consumed_once_records(tasks: list, cutoff_epoch: float) -> tuple[list, int]:
"""``(kept, pruned_count)`` — drop CONSUMED one-shot records (``enabled=False``
+ ``completed_at``) whose completion is older than the unified GC retention
cutoff (epoch seconds; ``retention.age_cutoff``). The consumed record is a
durable receipt, not a standing schedule, so it ages out like every other
disposable runtime artifact. ENABLED records are never pruned here (an
owner-disabled cron row has no ``completed_at`` and is kept too); an
unparseable ``completed_at`` is kept, conservatively."""
"""``(kept, pruned_count)`` — drop CONSUMED one-shot records (``trigger.type ==
"once"`` + ``enabled=False`` + ``completed_at``) whose completion is older than
the unified GC retention cutoff (epoch seconds; ``retention.age_cutoff``). The
consumed one-shot is a durable receipt, not a standing schedule, so it ages out
like every other disposable runtime artifact. ONLY one-shots are pruned: a
disabled CRON row is a standing schedule the owner may re-enable, and is kept
even if it carries a stray ``completed_at``. ENABLED records are never pruned;
an unparseable ``completed_at`` is kept, conservatively."""
kept, pruned = [], 0
for record in tasks:
if (isinstance(record, dict) and not record.get("enabled", True)
and record.get("completed_at")):
done = parse_schedule_time(record.get("completed_at"), datetime.timezone.utc)
if done is not None and done.timestamp() < float(cutoff_epoch):
pruned += 1
continue
trigger = record.get("trigger") if isinstance(record.get("trigger"), dict) else {}
if str(trigger.get("type") or "") == "once":
done = parse_schedule_time(record.get("completed_at"), datetime.timezone.utc)
if done is not None and done.timestamp() < float(cutoff_epoch):
pruned += 1
continue
kept.append(record)
return kept, pruned

View file

@ -193,6 +193,40 @@ def test_effort_only_roster_change_lapses_replay_authority(harness, monkeypatch)
assert _state(harness)["cycles_paid"] == 2
def test_degraded_reminder_promises_free_replay_only_with_structural_epoch(harness, monkeypatch):
"""Round-3: the user-turn DEGRADED reminder mirrors plan_render's
_degraded_replay_note (the wording SSOT) — a wave WITH a recorded structural
epoch is promised the free replay with its conditions (unchanged epoch +
roster); an EMPTY-epoch wave re-dispatches a paid panel, so the old
unconditional "replays ... at no cost" promise must not appear."""
from ouroboros.owner_hurry import force_plan_decision, plan_review_reminder
# Empty epoch: every slot dies at dispatch time, invisible to the snapshot.
_patch_health(monkeypatch, lambda slots: {})
harness.install({"s1": "", "s2": "", "s3": ""})
ctx = harness.make_ctx()
out = _call(ctx)
assert _control(out) == {"outcome": "DEGRADED", "closed": False}
decision = force_plan_decision(ctx, {}, enforcement="blocking")
assert decision["reviewer_slots_degraded"] and decision["degraded_health_epoch"] == ""
reminder = plan_review_reminder(decision)
assert "re-dispatches a fresh panel" in reminder
assert "no cost" not in reminder and "no further cost" not in reminder
# Non-empty epoch: structural snapshot evidence recorded — the free replay is
# promised together with its conditions (epoch + roster stand).
_patch_health(monkeypatch, lambda slots: dict(_DEAD_PANEL))
harness.install({"s1": CLEAN})
ctx2 = harness.make_ctx(task_id="task-epoch")
out2 = _call(ctx2)
assert _control(out2) == {"outcome": "DEGRADED", "closed": False}
decision2 = force_plan_decision(ctx2, {}, enforcement="blocking")
assert decision2["reviewer_slots_degraded"] and decision2["degraded_health_epoch"]
reminder2 = plan_review_reminder(decision2)
assert "at no further cost" in reminder2
assert "reviewer roster stand" in reminder2
assert "re-dispatches a fresh panel" not in reminder2
def test_replay_decision_config_failure_keeps_replay_but_logs_loudly(caplog):
"""Review fix 4 (accepted-partial): a configuration-resolution failure keeps the
recorded free replay (fail-open) but is logged as a WARNING with the exception

View file

@ -106,6 +106,26 @@ def test_once_schedule_survives_a_refused_admission_and_retries(tmp_path, monkey
assert record.get("last_error") == ""
def test_re_enabled_completed_once_never_refires(tmp_path):
"""Round-3 exactly-once: a consumed one-shot (non-empty completed_at) must not
fire again even when the owner flips enabled back on from the UI — re-arming
goes through the gateway upsert with a fresh run_at, never a bare toggle."""
queue, pending = _queue(tmp_path)
fired = datetime.datetime(2020, 1, 1, tzinfo=UTC).isoformat()
queue.upsert_scheduled_task({
"id": "fu-rearmed", "name": "Follow-up", "enabled": True, # UI re-enable
"completed_at": fired, "last_task_id": "t-already-ran",
"trigger": {"type": "once", "run_at": "2000-01-01T00:00:00+00:00"}, # long due
"task": {"type": "task", "text": "must not run twice"},
})
queue.check_scheduled_tasks()
queue.check_scheduled_tasks()
assert pending == []
record = queue.list_scheduled_tasks(tmp_path)["tasks"][0]
assert record["completed_at"] == fired # receipt untouched
assert record["last_task_id"] == "t-already-ran"
def test_once_schedule_with_invalid_run_at_records_a_typed_error(tmp_path):
queue, pending = _queue(tmp_path)
queue.upsert_scheduled_task({
@ -270,6 +290,61 @@ def test_schedules_gateway_accepts_and_validates_once_triggers(tmp_path):
assert unknown.status_code == 400
def test_gateway_rearm_of_completed_once_requires_a_fresh_run_at(tmp_path):
"""Round-3 exactly-once vs re-enable: re-enabling a CONSUMED one-shot through
the Schedules upsert without a NEW run_at is a 400; supplying a fresh run_at
re-arms it (completed_at cleared) and it fires exactly once; a disable that
keeps the same run_at carries the receipt forward for GC."""
from starlette.applications import Starlette
from starlette.routing import Route
from starlette.testclient import TestClient
from ouroboros.gateway.schedules import api_schedules_upsert
from supervisor import queue
queue.init(tmp_path, 600, 1800)
pending: list = []
queue.init_queue_refs(pending, {}, {"value": 0})
fired = datetime.datetime(2020, 1, 1, tzinfo=UTC).isoformat()
queue.upsert_scheduled_task({
"id": "fu-done", "name": "Follow-up", "enabled": False, "completed_at": fired,
"trigger": {"type": "once", "run_at": "2000-01-01T00:00:00+00:00"},
"task": {"type": "task", "text": "resume"},
})
app = Starlette(routes=[Route("/api/schedules", endpoint=api_schedules_upsert, methods=["POST"])])
app.state.drive_root = tmp_path
client = TestClient(app)
# Bare re-enable with the SAME run_at: refused with a clear re-arm message.
refused = client.post("/api/schedules", json={
"id": "fu-done", "enabled": True,
"trigger": {"type": "once", "run_at": "2000-01-01T00:00:00+00:00"},
"task": {"type": "task", "text": "resume"},
})
assert refused.status_code == 400 and "run_at" in refused.json()["error"]
# Disable/edit keeping the same run_at: allowed, receipt carried forward.
kept = client.post("/api/schedules", json={
"id": "fu-done", "enabled": False,
"trigger": {"type": "once", "run_at": "2000-01-01T00:00:00+00:00"},
"task": {"type": "task", "text": "resume"},
})
assert kept.status_code == 200
assert kept.json()["schedule"]["completed_at"] == fired
# A fresh run_at re-arms: completed_at cleared, and the record fires ONCE.
rearmed = client.post("/api/schedules", json={
"id": "fu-done", "enabled": True,
"trigger": {"type": "once", "run_at": "2000-02-01T00:00:00+00:00"},
"task": {"type": "task", "text": "resume"},
})
assert rearmed.status_code == 200
assert "completed_at" not in rearmed.json()["schedule"]
queue.check_scheduled_tasks()
queue.check_scheduled_tasks()
assert len(pending) == 1
record = queue.list_scheduled_tasks(tmp_path)["tasks"][0]
assert record["enabled"] is False and record["completed_at"]
def test_scheduled_tasks_digest_projects_run_at_for_once_records(tmp_path):
"""Review fix 9: the context digest shows a one-shot's fire instant (run_at)
instead of an empty-string cron; cron records keep their cron field."""
@ -325,9 +400,16 @@ def test_consumed_once_records_are_pruned_past_gc_retention(tmp_path):
"trigger": {"type": "cron", "expr": "0 3 * * *"},
"task": {"type": "task", "text": "paused"},
})
queue.upsert_scheduled_task({ # round-3: disabled CRON with a stray old completed_at
"id": "disabled-cron-stamped", "enabled": False, "completed_at": old,
"trigger": {"type": "cron", "expr": "0 4 * * *"},
"task": {"type": "task", "text": "paused, once ran"},
})
queue.check_scheduled_tasks()
ids = {r["id"] for r in queue.list_scheduled_tasks(tmp_path)["tasks"]}
assert ids == {"consumed-fresh", "enabled-future", "disabled-cron"}
# Only the aged-out CONSUMED ONE-SHOT is pruned; a disabled cron row is a
# standing schedule the owner may re-enable, even when it carries completed_at.
assert ids == {"consumed-fresh", "enabled-future", "disabled-cron", "disabled-cron-stamped"}
assert pending == []