Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 49 additions & 3 deletions s16code/events/outbox.py
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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

Expand Down Expand Up @@ -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."""
Expand All @@ -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:
Expand Down
31 changes: 31 additions & 0 deletions tests/test_autonomous_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"}