diff --git a/s16code/events/outbox.py b/s16code/events/outbox.py index 9d27482..b5bcd16 100644 --- a/s16code/events/outbox.py +++ b/s16code/events/outbox.py @@ -1,4 +1,12 @@ -"""A tiny idempotency outbox for capability side effects.""" +"""A tiny idempotency outbox for capability side effects. + +A record left ``started`` means the process died between dispatching a side +effect and learning whether it landed. :meth:`ActionOutbox.execute` refuses to +guess and parks the run on ``outbox.reconcile``. Nothing emits that event: it is +a question for a person, answered with :meth:`resolve` (or by completing the +handle through ``POST /v1/agent/completions``) once they have checked whether +the effect actually happened. :meth:`uncertain` is how they find them. +""" from __future__ import annotations @@ -7,6 +15,7 @@ import os import tempfile import threading +from datetime import UTC, datetime from pathlib import Path from typing import Any, Awaitable, Callable @@ -53,6 +62,39 @@ def _decode(value: dict[str, Any]) -> dict[str, Any] | Deferred: return Deferred(value["handle"], value["event_type"], value.get("metadata", {})) return value["value"] + def uncertain(self) -> list[dict[str, Any]]: + """Every effect whose outcome a crash left unknown, oldest first. + + A run parked on ``outbox.reconcile`` waits for an event this package + never sends. That is deliberate -- repeating an effect that may already + have landed is worse than waiting -- but it only works if the operator + can see what is waiting and how long it has been. + """ + with self._lock: + stuck: list[dict[str, Any]] = [] + for path in sorted(self.root.glob("*.json")): + try: + record = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + continue + if record.get("status") != "started": + continue + stuck.append({"key": path.stem, "started_at": record.get("started_at", ""), + "attempt": record.get("attempt", 1), "event_type": "outbox.reconcile"}) + return sorted(stuck, key=lambda item: item["started_at"]) + + async def resolve(self, key: str, receipt: dict[str, Any]) -> None: + """Record what a person found out about an uncertain effect. + + Closing the loop belongs here rather than in the caller: the receipt has + to be stored in the same shape ``execute`` returns, or the next call + re-parks on a question that has already been answered. + """ + with self._lock: + self._write(key, {"status": "completed", "receipt": self._encode(receipt), + "resolved_by": "operator", + "resolved_at": datetime.now(UTC).isoformat()}) + async def execute(self, key: str, operation: Callable[[], Awaitable[dict[str, Any] | Deferred]]) -> dict[str, Any] | Deferred: """Reuse a receipt; never blindly repeat an operation left in-flight by a crash.""" @@ -63,10 +105,14 @@ async def execute(self, key: str, if record["status"] == "completed": return self._decode(record["receipt"]) if record["status"] == "started": - return Deferred(key, "outbox.reconcile", {"uncertain": True}) + return Deferred(key, "outbox.reconcile", + {"uncertain": True, "started_at": record.get("started_at", "")}) if record["status"] == "failed": raise RuntimeError(record["error"]) - self._write(key, {"status": "started"}) + # Timestamped so an operator can tell an effect that is in flight + # from one abandoned by a crash. Without it every parked record looks + # equally fresh and equally hopeless. + self._write(key, {"status": "started", "started_at": datetime.now(UTC).isoformat()}) try: result = await operation() except Exception as error: diff --git a/tests/test_autonomous_events.py b/tests/test_autonomous_events.py index 34a2fdd..9f23173 100644 --- a/tests/test_autonomous_events.py +++ b/tests/test_autonomous_events.py @@ -114,3 +114,34 @@ async def effect(): uncertain = await outbox.execute("crash-key", effect) assert uncertain.handle == "crash-key" and uncertain.event_type == "outbox.reconcile" assert calls == 1 + + +@pytest.mark.asyncio +async def test_outbox_can_list_what_a_crash_left_uncertain(tmp_path): + """A run parked on an uncertain effect has to be findable. + + `execute()` parks on `Deferred(key, "outbox.reconcile")` when it meets a + record left `started` by a crash. Nothing in the package ever emits that + event, so the only way out is an operator completing the handle by hand -- + and they cannot, because nothing lists the stuck keys or says when they got + stuck. The run waits forever and no surface admits it. + """ + outbox = ActionOutbox(tmp_path) + + async def effect(): + return {"external_id": "receipt-1"} + + assert outbox.uncertain() == [] + + outbox._write("crash-key", {"status": "started"}) + parked = await outbox.execute("crash-key", effect) + assert parked.event_type == "outbox.reconcile" + + stuck = outbox.uncertain() + assert [item["key"] for item in stuck] == ["crash-key"] + assert stuck[0]["event_type"] == "outbox.reconcile" + + # Resolving it the documented way clears it from the list. + await outbox.resolve("crash-key", {"external_id": "confirmed-by-operator"}) + assert outbox.uncertain() == [] + assert await outbox.execute("crash-key", effect) == {"external_id": "confirmed-by-operator"}