Run Evidence
Middleware and lifecycle hooks see a run while it happens. A run evidence reader lets an extension look at runs afterwards, from the records the Gateway has already persisted: which runs exist, what state they ended in, and the event stream each one wrote. It is how an extension builds an export, an audit index, or a dashboard without hooking into every call.
The reader is read-only. No method writes to the host.
Getting the reader
The Gateway hands one reader to every service, as ExtensionRuntimeDeps.run_evidence_reader:
class MyService:
async def start(self, deps):
reader = deps.run_evidence_reader
if reader is None:
return # this host does not provide run evidenceThe Gateway always provides it. None means the extension is running on a host that does not implement the reader. Treat that as “unsupported”, never as “no runs”. Service lifetime and ordering are covered in Services and Routes; the reader is usable from the moment start() is called until stop() returns.
The reader has global visibility: it sees every user’s runs and events. Services have no request principal, so the Gateway binds this reader to no user on purpose. Event content is returned as stored. Never return data from it to the caller of a route.
The reader interface
RunEvidenceReader has three async methods, all keyword-only:
| Method | Returns | Use |
|---|---|---|
list_changed_runs(cursor=..., limit=...) | RunPage | Discover runs created or changed since a cursor |
list_run_events(thread_id=..., run_id=..., after_seq=..., limit=...) | RunEventPage | Read one run’s persisted events, forward from a sequence number |
get_run_status(thread_id=..., run_id=...) | RunStatusView | None | Read the authoritative status of one known run |
limit must be an integer from 1 to 2000; anything else raises ValueError. after_seq must be None or a non-negative integer.
Return types
All return types are frozen dataclasses from deerflow_extension_api.
RunStatusView: one run’s lifecycle state.
| Field | Meaning |
|---|---|
thread_id, run_id | The run’s identity |
status | pending, running, success, error, timeout, or interrupted |
created_at, updated_at | Timestamps as strings |
error | The error message, when the run failed |
stop_reason | Why the run stopped early, when it did |
RunEventView: one persisted event.
| Field | Meaning |
|---|---|
thread_id, run_id | The run the event belongs to |
seq | Sequence number, increasing within the thread |
event_type, category | The event’s kind, as the event store recorded it |
content | The event payload, unchanged |
metadata | Event metadata, with the legacy auth_token key removed |
created_at | Timestamp as a string |
RunPage: items (a tuple of RunStatusView), next_cursor, and has_more.
RunEventPage: items (a tuple of RunEventView), next_after_seq, and has_more.
content and metadata are deep copies. The dataclass fields are frozen, but you may modify the nested dicts and lists you receive without affecting host storage.
Discovering changed runs
list_changed_runs returns runs in a stable order, oldest change first. Page through it by passing each page’s next_cursor into the next call:
cursor=Nonestarts from the beginning.has_more=Truemeans another page is available right now.- An empty page means you are caught up. Its
next_cursoris the cursor you passed in, so keep it and poll again later.
Only agent runs appear. Other operations the host records against a thread, such as checkpoint writes, artifact writes, branching, and deletion, are excluded from the feed and from get_run_status.
A run appears again every time it changes. The feed does not contain one entry per run; it contains the current state of each run that changed after your cursor. A run created, started, and finished between two polls shows up once, already finished.
What counts as a change
Each run has a change position that the host advances when the run is created, changes lifecycle state, is cancelled, or has its model name updated. Progress snapshots and lease heartbeats do not advance it, so a long-running run does not flood the feed while it works.
Runs that existed before change tracking was added have position zero. They come first, ordered by run ID.
Cursor rules
The cursor is an opaque string. Store it, pass it back, and do not parse it.
- Replay, never skip. Reusing a cursor is always valid and may return runs you have seen before. A run that changes after you received it comes back with its new state. A run that you have not received yet cannot be skipped.
- Commit after output. Save
next_cursoronly after the work for that page is durable. If you crash in between, the next poll replays the page. Make your processing idempotent, for example by overwriting a per-run record instead of appending to it. - Scope-bound. A cursor belongs to the reader that issued it. Passing a cursor to a reader with a different visibility scope, a malformed cursor, or one from an unsupported version raises
InvalidRunEvidenceCursor, a subclass ofValueError. Recover by starting again fromNone, which replays everything visible.
Deletions
The feed does not report deletions. A deleted run disappears from future pages and from get_run_status, but nothing tells you it is gone. If your extension mirrors runs and must drop deleted ones, periodically call get_run_status for the runs you know about and treat None as deleted.
get_run_status returns None whenever the run is not visible: it does not exist, it was deleted, the thread_id does not match the run, or it is outside the reader’s scope.
Reading a run’s events
list_run_events pages forward through one run’s persisted events:
after_seq=Nonestarts at the first event. Pass each page’snext_after_seqto continue.seqvalues are increasing but not contiguous within a run, because the counter is shared by every run in the thread. Always continue fromnext_after_seqrather than computing the next number.- A run that does not exist or is not visible returns an empty page, never an error, so a caller cannot probe for run IDs outside its scope.
Status always comes from the run store, which is authoritative. Do not infer a run’s final state from its last event.
Storage backends
| Setting | Run positions and cursors |
|---|---|
database.backend: sqlite or postgres | Stored in the database. Positions and cursors stay valid across Gateway restarts |
database.backend: memory | Kept in process memory. All runs and positions are lost on restart |
Events come from the configured run_events store (memory, db, or jsonl) and use its own seq numbering.
With the memory backend, do not persist a cursor across restarts. A new
process numbers changes from zero again, and an old cursor points past all of
them, so the reader returns empty pages until the new process catches up with
the saved position. Start from None after every restart, or keep the cursor
in memory only.
Example: a run digest service
This service polls the reader and records, for every finished run, how many events of each type it persisted. It keeps its cursor in memory, so it rescans from the beginning after a restart. A real exporter would store the cursor next to its output.
"""Count persisted event types for every finished run."""
from __future__ import annotations
import asyncio
import logging
from collections import Counter
from collections.abc import Mapping
from typing import Any
from deerflow_extension_api import (
ExtensionRegistry,
ExtensionRuntimeDeps,
InvalidRunEvidenceCursor,
RunEvidenceReader,
extension,
)
logger = logging.getLogger(__name__)
TERMINAL = {"success", "error", "timeout", "interrupted"}
class RunDigestService:
def __init__(self, interval_seconds: float) -> None:
self.interval_seconds = interval_seconds
self.cursor: str | None = None
self.digests: dict[str, Counter[str]] = {}
self._task: asyncio.Task[None] | None = None
async def start(self, deps: ExtensionRuntimeDeps) -> None:
reader = deps.run_evidence_reader
if reader is None:
logger.warning("run evidence is not available on this host; digest disabled")
return
self._task = asyncio.create_task(self._poll(reader))
async def stop(self) -> None:
if self._task is not None:
self._task.cancel()
await asyncio.gather(self._task, return_exceptions=True)
self._task = None
async def _poll(self, reader: RunEvidenceReader) -> None:
while True:
try:
await self.sync_once(reader)
except Exception:
logger.exception("run digest sync failed; retrying")
await asyncio.sleep(self.interval_seconds)
async def sync_once(self, reader: RunEvidenceReader) -> None:
while True:
try:
page = await reader.list_changed_runs(cursor=self.cursor, limit=100)
except InvalidRunEvidenceCursor:
logger.warning("stored cursor rejected; rescanning from the beginning")
self.cursor = None
continue
for run in page.items:
if run.status in TERMINAL:
self.digests[run.run_id] = await self._count_events(reader, run.thread_id, run.run_id)
# Advance only after this page's output is recorded. Replaying a
# page is harmless because each digest is recomputed, not appended.
self.cursor = page.next_cursor
if not page.has_more:
return
async def _count_events(self, reader: RunEvidenceReader, thread_id: str, run_id: str) -> Counter[str]:
counts: Counter[str] = Counter()
after_seq: int | None = None
while True:
page = await reader.list_run_events(thread_id=thread_id, run_id=run_id, after_seq=after_seq, limit=500)
counts.update(event.event_type for event in page.items)
after_seq = page.next_after_seq
if not page.has_more:
return counts
@extension(api="0.2.0", name="digest")
def install(registry: ExtensionRegistry, config: Mapping[str, Any]) -> None:
registry.service(RunDigestService(float(config.get("interval_seconds", 30))))Points worth copying:
start()only creates the polling task and returns, so Gateway startup is not held up.stop()cancels the task and waits for it, well within the 30-second stop budget.- The digest for a run is recomputed each time the run appears, so replayed pages cause no double counting.
InvalidRunEvidenceCursorresets the cursor toNoneinstead of stopping the service.- An exception in one sync is logged, and the next poll resumes from the last committed cursor.
Common pitfalls
- Treating an empty page as unsupported. Empty means caught up. Unsupported is
run_evidence_reader is None. - Advancing the cursor before the work is saved. A crash then skips a page permanently. Save after.
- Expecting one entry per run. A run reappears each time it changes. Key your output by
run_id. - Waiting for a deletion event. There is none. Reconcile with
get_run_status. - Serving reader data to users. The reader sees every user. Never return its data from a route.