Skip to Content
DeerFlow

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 evidence

The 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:

MethodReturnsUse
list_changed_runs(cursor=..., limit=...)RunPageDiscover runs created or changed since a cursor
list_run_events(thread_id=..., run_id=..., after_seq=..., limit=...)RunEventPageRead one run’s persisted events, forward from a sequence number
get_run_status(thread_id=..., run_id=...)RunStatusView | NoneRead 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.

FieldMeaning
thread_id, run_idThe run’s identity
statuspending, running, success, error, timeout, or interrupted
created_at, updated_atTimestamps as strings
errorThe error message, when the run failed
stop_reasonWhy the run stopped early, when it did

RunEventView: one persisted event.

FieldMeaning
thread_id, run_idThe run the event belongs to
seqSequence number, increasing within the thread
event_type, categoryThe event’s kind, as the event store recorded it
contentThe event payload, unchanged
metadataEvent metadata, with the legacy auth_token key removed
created_atTimestamp 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=None starts from the beginning.
  • has_more=True means another page is available right now.
  • An empty page means you are caught up. Its next_cursor is 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_cursor only 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 of ValueError. Recover by starting again from None, 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=None starts at the first event. Pass each page’s next_after_seq to continue.
  • seq values are increasing but not contiguous within a run, because the counter is shared by every run in the thread. Always continue from next_after_seq rather 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

SettingRun positions and cursors
database.backend: sqlite or postgresStored in the database. Positions and cursors stay valid across Gateway restarts
database.backend: memoryKept 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.

deerflow_extension_digest/__init__.py
"""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.
  • InvalidRunEvidenceCursor resets the cursor to None instead 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.