Step Persistence
StepPersistence saves a snapshot of an agent run at every settled step and when the run fails, so you can resume the run, or fork it from an earlier step, in the same process or another one, from memory, files, SQLite, MongoDB, or your own store. Alongside the snapshots it keeps an append-only trail of step events and a ledger of tool side effects, so after a crash you can tell which tool calls completed and which may or may not have run. It is also the persistence layer for orchestrators that delegate to sub-agents, for example continuing a delegate’s investigation with a follow-up question.
A snapshot holds the run’s message history, not everything around it: capability state outside the messages, workspace files, and resuming from the middle of a step are tracked separately (see pydantic-ai-harness issues #149 and #196). For recovery inside a step, run the agent on durable execution, which StepPersistence works alongside.
While Pydantic AI Harness is on 0.x releases, the API may change between minor releases; when it does, deprecation warnings and release-note migration guidance tell you (or your agent) exactly how to upgrade. See the version policy.
- Append-only step events. Every interesting boundary (run start/end, model request, tool call, failure) appends a
StepEvent. A run killed mid-tool-call still leaves a usable event trail. - Continuable snapshots. A
ContinuableSnapshotis saved at settled node boundaries, and a failing run saves its live at-failure history. Each snapshot carries astate:completewhen everyToolCallParthas a matching result,interruptedwhen the capture holds unsettled tool work (e.g. a crash mid-tool-cycle).latest_snapshotandcontinue_runreturn onlycompletesnapshots unless the caller passesinclude_interrupted=True. Pass the snapshot’smessagesback toAgent.run(message_history=...)to continue or fork. - Tool-effect ledger. Every tool call’s lifecycle (
started,completed,failed) is recorded against(run_id, tool_call_id). After a crash, a tool with astartedrecord and no terminal update should be treated asunknown_after_crash: the side effect may or may not have happened. - Lineage metadata.
conversation_id(sequence) andparent_run_id(hierarchy) are independent axes. See Three-level identity.
import asyncio
from pydantic_ai import Agent
from pydantic_ai_harness import StepPersistence
from pydantic_ai_harness.step_persistence import InMemoryStepStore
store = InMemoryStepStore()
librarian = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='code_librarian')],
)
async def main():
await librarian.run('Find ThinkingPartDelta and confirm the callable allowance')
asyncio.run(main())
That is the whole setup. run_id is always per-Agent.run call, matching pydantic_ai’s RunContext.run_id. For multi-turn logical grouping use conversation_id= — that is the pydantic_ai-native primitive for it (see Three-level identity).
run_id resolution per call:
- Explicit
run_id='libr-1'becomes the id for this one call. This suits single-shot use cases (a deterministic id for testing, replay, debugging, or a one-off scripted run). Reusing one capability instance with the same explicitrun_idacross multiple.run()calls raisesValueErrorinbefore_run. The tool-effect ledger is keyed by(run_id, tool_call_id)and providers reuse deterministic tool-call ids, so a silent collision would erase theunknown_after_crashsignal. Useconversation_id=for multi-turn grouping instead. agent_nameset,run_idunset derives a path-safe base64url encoding of the complete(agent_name, ctx.run_id)pair. The encoding is injective withinFileStepStore’s 200-character limit, so replay addresses the same stored run without collisions between distinct accepted context ids. A longer derived id raisesValueErrorbefore backend selection, including with the memory, SQLite, and Mongo stores.- Neither set uses
ctx.run_idunchanged. A missing context run id raisesRuntimeErrorbecause inventing one would disconnect replayed writes.
The run record, events, and snapshots carry agent_name. When it is unset, they record the running agent’s name instead (which Pydantic AI infers from the variable name when Agent(name=...) is not passed). That fallback does not feed the run_id derivation above, so the store key stays ctx.run_id and runs can still be looked up by the id passed to or returned from Agent.run.
StepPersistence has the stable capability id step_persistence, so it can be attached alongside a Pydantic AI durability capability without passing id=. Pass an explicit id only when the same agent has more than one StepPersistence instance.
Six store boundaries are durable operations: registration identity, run registration, event append, snapshot save, tool-effect start, and tool-effect completion or failure. Every persisted timestamp is read inside one of those operations, so replay uses the journaled value instead of reading the workflow wall clock again. The journaled registration identity makes a retried registration idempotent while a distinct reuse of the same run id still fails.
Events and snapshots written by the capability carry deterministic per-run idempotency keys. Every built-in store suppresses a key it already applied, while records created directly with idempotency_key=None retain append behavior. Snapshot keys use a per-run save sequence together with step_index and state. Replay produces the same sequence, while distinct snapshots at the same step and state retain their write order and newer history.
The orchestrator pattern — one logical agent serving many turns — uses conversation_id, not a shared run_id:
import asyncio
from pydantic_ai import Agent
from pydantic_ai_harness import StepPersistence
from pydantic_ai_harness.step_persistence import InMemoryStepStore
store = InMemoryStepStore()
orchestrator = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='orchestrator')],
)
async def main():
for turn in turns:
await orchestrator.run(turn, conversation_id='orch-conv')
# All turns of this orchestrator, chronological:
records = await store.list_runs(conversation_id='orch-conv')
asyncio.run(main())
The capability mirrors pydantic_ai’s identity stack:
| Concept | Definition | Granularity |
|---|---|---|
conversation_id | The dialogue. Resolved by pydantic_ai from the conversation_id= argument to Agent.run, or the most recent conversation_id on message_history, or a fresh UUID7. | sequence of runs |
run_id | One Agent.run invocation. | one step in the sequence |
step_index | Graph-node count within a run (ctx.run_step). | one node within one run |
StepEvent.conversation_id and RunRecord.conversation_id are populated from ctx.conversation_id. So three .run() calls sharing one conversation_id produce three distinct run_ids, all queryable as a group:
import asyncio
async def main():
runs = await store.list_runs(conversation_id='conv-abc') # 3 records, chronological
asyncio.run(main())
pydantic_ai already has message_history= for “carry on with this prior context”. StepPersistence does not introduce a parallel mechanism. It exposes one helper that loads the most recent settled snapshot:
import asyncio
from pydantic_ai import Agent
from pydantic_ai_harness import StepPersistence
from pydantic_ai_harness.step_persistence import InMemoryStepStore, continue_run
store = InMemoryStepStore()
librarian = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='code_librarian')],
)
async def main():
# Earlier: tag the first turn with a conversation id so the follow-up can find it.
await librarian.run(
'Find ThinkingPartDelta and confirm the callable allowance',
conversation_id='libr-conv',
)
# Later (possibly a different process):
prior_run = (await store.list_runs(conversation_id='libr-conv'))[-1].run_id
history = await continue_run(store, run_id=prior_run)
await librarian.run(
'Read _apply_provider_details_delta and check the path',
message_history=history,
conversation_id='libr-conv', # keep the conversation grouping
)
asyncio.run(main())
fork_run(store, run_id=...) returns the same shape but is intended when the caller wants a branched logical run from that snapshot point (the new run gets a fresh run_id and probably a fresh conversation_id).
By default continue_run returns the messages of the latest complete snapshot for that run_id — a point whose tool work was fully settled when captured. Snapshots are written at these boundaries:
- after every
CallToolsNodewhose tool calls all returned — the pending tool-return request is folded in, so the point is durable the moment the tool completes, before the next model request is even sent, - at
after_run, when the run ended past that boundary (a run that reached no boundary at all, or anAgent.run_streamwhose closing response lands after the last one), and - when a run fails: the live history at failure time is saved, whatever its shape — a model request that raises after a clean tool cycle produces a
completesnapshot; a crash mid-tool-cycle produces aninterruptedone carrying every completed cycle.
An interrupted snapshot is sendable, but not necessarily safe: a pending tool call may be re-executed or closed out with a synthesized interrupted return, and neither says whether the original side effect happened. That is the tool-effect ledger’s job. Which one happens depends on how you continue. Resuming without a new prompt executes the pending calls. With a new prompt, calls are closed out if some of their batch already returned; if none did, the run raises UserError rather than abandon calls that could still be answered, so call repair_messages on the history first to close them out yourself. So the default read path skips interrupted snapshots; pass include_interrupted=True to continue_run / fork_run / latest_snapshot after checking list_unresolved_tool_effects. If no matching snapshot exists, continue_run raises LookupError.
parent_run_id is a lineage label, not a functional dependency. It does two things:
- Every
StepEventandRunRecordcarries it, so you can filter and group. store.list_runs(parent_run_id='orch-1')returns every delegate run pointing at that orchestrator.
It is auto-inferred for in-process delegation: when an orchestrator’s tool synchronously calls a delegate’s Agent.run(...), the delegate’s StepPersistence picks up the orchestrator’s run_id via a ContextVar that the orchestrator’s wrap_run set. No threading required:
import asyncio
from pydantic_ai import Agent
from pydantic_ai_harness import StepPersistence
from pydantic_ai_harness.step_persistence import InMemoryStepStore
store = InMemoryStepStore()
orchestrator = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='orchestrator')],
)
librarian = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='code_librarian')],
)
@orchestrator.tool_plain
async def ask_librarian(question: str) -> str:
result = await librarian.run(question) # parent_run_id auto-filled
return result.output
async def main():
# Tag the orchestrator turn so the lookup below can find its run_id.
await orchestrator.run(
'Where is ThinkingPartDelta defined?',
conversation_id='orch-conv',
)
# All librarian runs now point at the orchestrator's run_id:
orch_run_id = (await store.list_runs(conversation_id='orch-conv'))[-1].run_id
delegates = await store.list_runs(parent_run_id=orch_run_id)
asyncio.run(main())
Set parent_run_id= explicitly to override (for example, cross-process delegation where ContextVars do not propagate).
parent_run_id is distinct from conversation_id. The orchestrator and delegate usually live in different conversations (the orchestrator talks to a user; the delegate talks to itself). But they share a parent-child link.
list_runs returns matches sorted by started_at ascending across all backends — pick the most recent with [-1].
import asyncio
async def main():
# Every delegate of one orchestrator run (chronological)
delegates = await store.list_runs(parent_run_id='orch-3f2a')
# Every run in one dialogue (multi-turn conversation across many .run() calls)
turns = await store.list_runs(conversation_id='conv-abc')
latest_turn = turns[-1]
# Filters combine (AND):
focused = await store.list_runs(
parent_run_id='orch-3f2a',
conversation_id='libr-conv',
)
# Detail per run:
events = await store.list_events(run_id=delegates[0].run_id)
snapshot = await store.latest_snapshot(run_id=delegates[0].run_id)
unresolved = await store.list_unresolved_tool_effects(run_id=delegates[0].run_id)
asyncio.run(main())
import asyncio
async def main():
# An earlier delegate run died mid-investigation.
events = await store.list_events(run_id='libr-3f2a')
unresolved = await store.list_unresolved_tool_effects(run_id='libr-3f2a')
for record in unresolved:
# status == 'started' with no terminal update -- unknown_after_crash.
print(f'tool {record.tool_name} ({record.tool_call_id}) may or may not have run')
print(f' idempotency_key={record.idempotency_key} '
f'effect_summary={record.effect_summary}')
# Decide whether to resume or branch:
history = await continue_run(store, run_id='libr-3f2a')
# If the unresolved tools were read-only and safe to redo:
await librarian.run('continue investigating', message_history=history,
conversation_id='libr-conv')
# If side effects might have happened and the orchestrator wants a fresh attempt:
history = await fork_run(store, run_id='libr-3f2a')
# ... pass to a new delegate run with a different agent_name / conversation_id.
# To resume from the interrupted frontier itself (the crashed cycle included),
# after checking the unresolved effects above:
history = await continue_run(store, run_id='libr-3f2a', include_interrupted=True)
asyncio.run(main())
Side-effect deduplication is the orchestrator’s responsibility. Tools that write external state should annotate their in-flight ToolEffectRecord via annotate_tool_effect:
from pydantic_ai import RunContext
from pydantic_ai_harness.step_persistence import annotate_tool_effect
@orchestrator.tool
async def set_label(ctx: RunContext[Deps], issue: int, label: str) -> str:
await annotate_tool_effect(
store,
ctx,
idempotency_key=f'issue-{issue}::label::{label}',
effect_summary=f'set label {label!r} on issue #{issue}',
)
await github.set_label(issue, label) # the actual side effect
return 'ok'
The helper reads the active run_id from the StepPersistence ContextVar and tool_call_id / tool_name from ctx, then merges the metadata into the prior record. It is a no-op when called outside a step-persistence-wrapped tool call. after_tool_execute preserves both fields when it writes the terminal completed / failed entry.
StepPersistence.compaction_transcript_handle() exposes the current run_id to compaction receipts. It is an identifier for this store’s persisted run history, not a promise that the pre-compaction transcript remains available: snapshots can already contain compacted history and configured retention can delete older snapshots.
InMemoryStepStore— process-local; great for tests.FileStepStore(directory)— directory layout under<directory>/<run_id>/:run.json—RunRecord(lineage)events.jsonl— append-onlyStepEventstool_effects.jsonl— append-onlyToolEffectRecords, scoped to this runsnapshot-keys.jsonl— replay-suppression keys retained independently of snapshot pruningsnapshots/{seq}.json—ContinuableSnapshots, named by a per-run monotonic counter (notstep_index, which would collide when the samerun_idis reused acrossAgent.runcalls, sincectx.run_stepresets to 0 each call).
SqliteStepStore(database='runs.db')— single SQLite file with tablesruns,events,snapshots,snapshot_idempotency_keys,tool_effects, and a siblingmediatable for externalized blobs (see Persisting media below). WAL mode is enabled;tool_effectsupserts per(run_id, tool_call_id)so the latest state wins; snapshots useAUTOINCREMENT seqto mirrorFileStepStore._next_snapshot_seq. Databases created before the snapshotstatecolumn existed gain it automatically on open (existing rows read ascomplete). Passconnection=instead ofdatabase=to share asqlite3.Connectionwith the rest of your application; the connection must be opened withcheck_same_thread=Falsebecause hook calls are dispatched onto a worker thread.MongoStepStore(client= or db_url=, database=...)— MongoDB collectionsruns,events,snapshots,snapshot_idempotency_keys,tool_effects, andcounters(atomic$incallocates the monotonicseq). Run registration uses an atomic insert byruns._id = run_id; duplicate ids raiseValueError. Needs themongodbextra (which installspymongo>=4.17.0); pass a sharedAsyncMongoClientasclient=, or a connection string asdb_url=(the store then owns the client — callawait store.aclose()to release it). Individual parts at or abovemedia_threshold_bytesexternalize by default to aMongoMediaStoreon the same client. That is a per-value offload, not an aggregate cap: a snapshot of many below-threshold parts can still exceed MongoDB’s 16 MiB document limit and fail on insert, so lower the threshold if that is a risk for your workload.
All implement the same async StepStore protocol, so capability hooks never block the event loop on the file/sqlite backends (I/O is dispatched via anyio.to_thread); the Mongo backend is natively async.
FileStepStore validates run_id against [A-Za-z0-9_.-]{1,200} (and rejects ..) to prevent path traversal. Callers passing user-controlled IDs should still sanitise first.
The store issues createIndex on its first write, for ten indexes: conversation_id and parent_run_id (both sparse) plus started_at on runs; (run_id, seq) and unique keyed (run_id, idempotency_key) on events; (run_id, seq), unique keyed (run_id, idempotency_key), and (run_id, state, seq) on snapshots; and a unique (run_id, tool_call_id) plus (run_id, status) on tool_effects. The idempotency indexes include only documents whose key is a string, so None retains append behavior. Its default MongoMediaStore adds one more, described on the media page. Three consequences worth knowing before pointing the store at an existing deployment:
- The connecting user needs the privilege to create indexes. A restricted Atlas role without it fails on the first write, not at construction.
- The unique index build fails if an existing
tool_effectscollection already holds duplicate(run_id, tool_call_id)pairs. - Index builds against already-populated collections cost time and I/O on that first call.
RunRecord.metadata and StepEvent.metadata are stored as nested documents, so their keys become BSON field names: keys containing . or starting with �IC4� and �IC3� are stored as nested documents, so their keys become BSON field names: keys containing �IC2� or starting with need [MongoDB 5.0 or later](https://www.mongodb.com/docs/manual/core/dot-dollar-considerations/), and a key containing a NULL byte is rejected by the BSON encoder before it reaches the server. CI exercises both Mongo backends against mongo:8`.
Install MongoDB support:
pip install "pydantic-ai-harness[mongodb]"
uv add "pydantic-ai-harness[mongodb]"
Each step writes a new full-history snapshot keyed by an incrementing seq, and nothing is pruned by default. Within one long Agent.run the snapshot count equals the number of settled tool-call steps, so a long single run pays a growing storage cost.
All four stores — InMemoryStepStore, FileStepStore, SqliteStepStore, and MongoStepStore — accept an opt-in max_snapshots_per_run: int | None (default None, unbounded — byte-for-byte the prior behavior). When set to N >= 1, each save_snapshot prunes the run down to a retain set:
- the newest
Nsnapshots byseq, - the newest snapshot overall (serves
latest_snapshot(include_interrupted=True)), - the newest
completesnapshot (serves the default read path).
The last two keep both read modes correct even when the newest N snapshots are all interrupted and the newest resumable complete sits below that window, so the retain set can exceed N. from_spec(..., max_snapshots_per_run=N) forwards the bound to the store it constructs (backend='memory', 'file', or 'sqlite'; a Mongo store is built directly, not from a spec).
from pydantic_ai_harness.step_persistence import FileStepStore
store = FileStepStore('runs', max_snapshots_per_run=8)
Pruning a snapshot never deletes its externalized media: blobs are content-addressed and may be shared across snapshots and runs, so orphaned-blob GC is out of scope (see the non-goals below). Age-based (TTL) expiry is out of scope too — it belongs at whole-run granularity, not per snapshot.
Bounded retention discards older per-step snapshots, including pre-compaction ones. Any downstream that reconstructs history by unioning a run’s retained snapshots — snapshot search or a compaction receipt keyed on run_id — can only see what is retained. With a tight bound (for example max_snapshots_per_run=1) the older, pre-compaction states are gone, so treat the bound as a hard limit on how far back such recovery can reach. Leave the bound at None, or set it high enough to cover the history you need to recover, when historical reconstruction matters.
BinaryContent payloads (images, audio, documents, video) inlined as base64 inside a snapshot would balloon every file or row containing the message; a large text part (e.g. a big tool-return string) does the same and can push a MongoStepStore snapshot past MongoDB’s 16 MiB document cap (#440). The file, sqlite, and mongo backends externalize any BinaryContent.data, and any part whose string content is at or above 64 KiB, through a configured MediaStore, leaving a URI reference in the snapshot. The same media_threshold_bytes governs binary and text alike; there is no separate text knob. Round-trip is transparent: latest_snapshot(...).messages[*] returns the original BinaryContent bytes and text.
Text externalization is not Mongo-only and has no opt-out short of media_store=None: the walker is shared, so an existing FileStepStore or SqliteStepStore deployment starts writing blobs for large text parts as well as binary ones from this release on. Snapshots written before it still restore — the reader recognises the older binary marker shape. This compatibility is upgrade-only: a release that predates text externalization treats every marker as binary, so it cannot validate a snapshot containing an externalized text marker. Keep a current reader for persisted snapshots that contain those markers.
Reserved-key escaping is a second marker-format generation with the same rule for these stores: a payload using the marker format’s namespaced keys is moved into a versioned reserved mapping (the __harness_external_escaped_keys__ stash, stamped with the format version under __harness_external_marker_format__), and the current reader moves those values back to their own keys. Compatibility the other way is upgrade-only. A reader that predates the escaping format re-inlines the externalized field correctly, but it leaves both reserved keys sitting in the restored payload rather than removing them. A marker carrying both, stamped with a version this reader does not know, is rejected rather than restored with the reserved values stripped: restore_media raises ValueError, and latest_snapshot surfaces it to the caller for the file, sqlite, and mongo stores. list_snapshots is different: each store treats the failed snapshot as unparsable, skips it, and logs the error, so an unknown version shows up as a missing snapshot rather than an exception. That rejection is the version gate and is intended, but store users have to anticipate it. Keep a current reader for persisted snapshots that contain escaped markers.
| StepStore | Default media_store | Where blobs live |
|---|---|---|
InMemoryStepStore | not applicable | bytes stay in the in-memory snapshot |
FileStepStore | DiskMediaStore(<root>/media/) | <root>/media/<sha256>.bin |
SqliteStepStore | SqliteMediaStore(database=<same db>) | sibling media table in the same DB |
MongoStepStore | MongoMediaStore(client=<same client>) | sibling media + media_chunks collections |
Override the destination by passing your own MediaStore:
from pydantic_ai_harness.media import S3MediaStore
from pydantic_ai_harness.step_persistence import FileStepStore
store = FileStepStore(
'runs',
media_store=S3MediaStore(
bucket='my-bucket',
endpoint='https://<account>.r2.cloudflarestorage.com',
region='auto',
access_key_id=...,
secret_access_key=...,
),
media_threshold_bytes=64 * 1024, # raise or lower if you want
)
Opt out entirely (keep bytes inline in the snapshot JSON/row):
from pydantic_ai_harness.step_persistence import FileStepStore, SqliteStepStore
FileStepStore('runs', media_store=None)
SqliteStepStore(database='runs.db', media_store=None)
URIs are media+sha256://<hex>, content-addressed. The same blob written through any MediaStore resolves the same way, so dedup is automatic and moving the underlying storage is a one-line swap. The shipped implementations are:
DiskMediaStore(directory)— one file per blob at<directory>/<sha256>.bin.SqliteMediaStore(database=...)orSqliteMediaStore(connection=...)— one row per blob (INSERT OR IGNOREfor content-addressed dedup).S3MediaStore(bucket=, endpoint=, region=, access_key_id=, secret_access_key=)— path-style URLs plus handrolled SigV4. Compatible with AWS S3, Cloudflare R2 (region='auto'), MinIO, and other S3-compatible providers. PUT/GET/HEAD only — no multipart, lifecycle, or listing in v1.MongoMediaStore(client= or db_url=, database=...)— MongoDB, needs themongodbextra. Each blob is sha256-addressed chunks across amediamanifest document and a siblingmedia_chunkscollection (manual chunking rather than GridFS, so dedup is preserved — see the media page), so a blob larger than one BSON document still stores and reads back. Chunking bounds the document, not memory: there is no streaming API, so each blob is held whole in process memory on bothputandget. The manifest holdsMediaContext.metadatainline and is not chunked, so keep per-blob metadata small.collection=renames both collections andchunk_size_bytes=(default 8 MiB) sets the split size.
Each store accepts a public_url= callable that turns the canonical media+sha256://<hex> URI into a URL the model can fetch directly. The forthcoming MediaExternalizer capability will use this to swap BinaryContent parts for ImageUrl / AudioUrl / other URL parts before the model sees the message, letting providers fetch big media over the wire without re-encoding bytes into the request body.
Static base URL (public R2 bucket, CDN):
from pydantic_ai_harness.media import S3MediaStore, make_static_public_url
store = S3MediaStore(
bucket='my-bucket',
endpoint='https://<acc>.r2.cloudflarestorage.com',
region='auto',
access_key_id=..., secret_access_key=...,
key_prefix='media/',
public_url=make_static_public_url('https://pub-abc.r2.dev', key_prefix='media/'),
)
Presigned or rotating-signature URL — pass any async callable that takes (uri, MediaContext):
from pydantic_ai_harness.media import MediaContext, S3MediaStore
async def presign(uri: str, ctx: MediaContext) -> str:
key = 'media/' + uri.removeprefix('media+sha256://') + '.bin'
return await my_signer.generate(key, ttl=3600, content_type=ctx.media_type)
store = S3MediaStore(..., public_url=presign)
Every MediaStore method (put, get, exists, public_url, get_metadata) and both user-supplied callables (PublicUrlResolver, KeyStrategy) accept a MediaContext:
from collections.abc import Mapping
from dataclasses import dataclass, field
@dataclass(frozen=True, kw_only=True)
class MediaContext:
media_type: str | None = None # e.g. 'image/png'
filename: str | None = None # original filename, when known
metadata: Mapping[str, str] = field(default_factory=dict) # user-supplied tags
All fields default; new fields are added non-breakingly as use cases emerge. Pass what you have, ignore the rest.
Persistence by store. get_metadata(uri) round-trips the user-supplied metadata mapping on all four stores. media_type is also persisted but is not part of what get_metadata returns (it is stored for the byte payload itself, for example as the Content-Type).
SqliteMediaStorewritesmetadatato a JSON column andmedia_typeto a dedicated column.S3MediaStoresendsmetadataas signedx-amz-meta-*headers (ASCII alphanumeric plus dash key names) andmedia_typeasContent-Type;get_metadatareads thex-amz-meta-*values back from the HEAD response.DiskMediaStorewrites a sidecar JSON file (<resolved>.meta.json) alongside each blob, atomic via tmp plus rename. Sidecars are absent only when the put carried no metadata.MongoMediaStorewritesmetadataas a JSON string andmedia_typeas a dedicated field on the blob’s manifest document (themediacollection by default);get_metadatadecodes the JSON string back. Because the mapping is one JSON string rather than nested fields, metadata keys are not subject to BSON field-name rules here.
Default is <sha256>.bin. DiskMediaStore and S3MediaStore accept overrides to fit existing layouts; SqliteMediaStore and MongoMediaStore do not (the digest is their primary key, so a user-chosen key would either break dedup or be a no-op — use table= / collection= to move the rows or documents):
from pydantic_ai_harness.media import DiskMediaStore, MediaContext
def by_media_type(uri: str, ctx: MediaContext) -> str:
digest = uri.removeprefix('media+sha256://')
ext = {'image/png': '.png', 'image/jpeg': '.jpg'}.get(ctx.media_type or '', '.bin')
return f'images/{digest}{ext}'
store = DiskMediaStore('runs', key_strategy=by_media_type)
Caveat: if your strategy depends on context.media_type (for example, to pick an extension), get(uri) and exists(uri) will not find the blob unless the same context is supplied at read time. For pure path-organisation strategies (no context dependency) the constraint does not apply.
DiskMediaStore rejects strategies that produce absolute paths or paths containing .. segments, to prevent escaping the store directory.
Separately, all four stores accept a public_url= resolver, useful when a CDN, local HTTP server, or signed-URL service fronts the bytes. Without it public_url(...) returns None (the model never sees a URL unless a resolver is configured and it returns a string).
pydantic_ai providers transparently download bytes from a URL when the target model does not natively accept that URL type, so emitting a URL is always safe: you only ever lose wire savings, never correctness.
DynamoDB, Postgres, Redis, GCS, and other backends are out of scope for this release. Write your own StepStore (about ten methods on a Protocol) or your own MediaStore (five methods: put, get, exists, public_url, get_metadata) and pass it via store= / media_store=. Please open an issue if you ship one — we want to feed the eventual shared adapter layer with N >= 3 real implementations before abstracting.
pydantic_ai_harness.step_persistence.conversations provides
SqliteConversationStore, ConversationSummary, and SavedConversation for
multi-turn applications. A conversation head is separate from per-run checkpoints:
it includes accepted prompts and between-run edits such as compaction. Do not
reconstruct it by concatenating overlapping run snapshots.
save(summary=..., messages=...) compares the supplied content revision and
returns the committed summary. A stale writer or a deleted session raises
ConversationConflict. get(conversation_id=...) restores messages through the
same media format used by step snapshots. listing(query=..., limit=..., offset=...)
returns summaries without loading messages; search matches saved user/assistant
text and metadata using Unicode case folding, including text entries within
multimodal prompts. Search does not include tool output, reasoning, or discarded
pre-compaction history. Unknown metadata schema versions are rejected.
Metadata naming uses a separate version. name(source=..., title=..., ...) cannot
overwrite a newer content revision, newer name, or a manual title. Naming does not
change the activity timestamp. delete(source=...) removes the conversation and
associated run records from the same SQLite database atomically, retaining shared
media. It is not secure erasure. A local live PID marks an unfinished conversation
as busy; this is not a distributed lease and the database must not be shared
between hosts. PID reuse is conservatively treated as busy.
The database is created owner-only where supported. Contents are not encrypted. There is no automatic conversation TTL or media garbage collection.
pydantic_ai_harness.step_persistence.naming is deprecated and emits a
HarnessDeprecationWarning on import. Its naming prompt, queue bounds, and
failure policy are CLAI’s resume-browser policy rather than a Harness primitive,
so CLAI now owns its own copy. Copy the helpers you use into your application;
the module will be removed in a future release. The conversation store above is
unaffected.
Until then, the module provides a tool-free naming agent
and SessionNamer, a worker owned by the application’s task group. submit(id)
coalesces jobs in a queue bounded to ten sessions. run() processes one job at a
time until its owner cancels it. backfill(entries) considers up to ten newest
entries. Naming failures are logged at debug level and leave existing metadata
usable; cancellation propagates. Applications must join the worker before closing
its model clients or storage dependencies.
Names consist of a short title, subtitle, and up to four tags. The model receives
the prior title/detail plus a bounded 2,400-character current conversation tail.
This is deliberately not a message-index cursor: compaction and recovery can
replace the list. Generated names become eligible again after 16 content
revisions. Manual names are not changed. Naming requests have a 60-second worker
timeout and the provided generate_name helper allows at most two model requests
and 250 output tokens. The helper’s agent is named session_namer, has no tools,
and does not inherit the foreground agent’s capabilities.
Core’s agent spans attribute auxiliary model calls to session_namer; no second
span hierarchy is emitted. Successful naming response token counts are stored
separately from foreground history, including results rejected as stale while the
session still exists. Failed or timed-out requests may incur provider usage not
available to the application. Monetary pricing of auxiliary calls is not included
in retained-history cost. Applications choose the naming model and disclose the
additional provider requests to their users.
Set capture_frontier=True to save accepted request histories before model
requests and the model response frontier before tool execution. The default is
False to preserve existing checkpoint frequency. CLAI enables it. A first model
request failure can then retain its prompt, and a process killed mid-tool-cycle
can retain the proposed calls and arguments even before a cycle settles.
inspect_recovery(store=..., run_id=...) in
pydantic_ai_harness.step_persistence.recovery returns the newest and settled
snapshots, unresolved effects, and names of recorded completed/failed tools.
It does not infer that an effect is safe to replay.
These are still message checkpoints, not graph-state checkpoints. Snapshots at
unsettled frontiers are interrupted and remain off the default read path.
after_run compares final content, not only message count, to catch same-length
or shortened history rewrites. Put the recorder before capabilities whose
after_run transforms history: core runs after-hooks in reverse order. Snapshot
message values are copied before storage so later mutations cannot alter a saved
in-memory checkpoint through shared references.
SnapshotSaved is a typed capability event emitted after a checkpoint write
completes. It carries persistence_run_id, conversation_id, step_index, and
state. Subscribe using core’s hooks.on.event(SnapshotSaved), which a CLAI plugin
returns from get_capabilities. Store writes are the source of truth; notifications may
repeat during durable replay and observer failures cannot undo committed writes.
Automatic execution recovery is not implemented. Two core contracts should be addressed before promising it:
on_run_errorshould expose authoritative post-cleanup history. Today Harness stashes a live list reference from node/request hooks because the outer error context can reference the start-of-run list. That depends on core continuing to mutate the working history in place. Core’s cancellation result APIs are useful to callers, but do not establish the same contract for every error hook.- An awaited checkpoint boundary should expose normalized results as individual
tools settle, including accompanying user content, retries, and parallel
siblings.
after_tool_executesees raw results before all normalization;after_node_runsees a settled batch.FunctionToolResultEventexposes a normalized result, but observing a stream is not an atomic commit of that result with the tool-effect ledger and the execution frontier.
A hard kill during a parallel batch can therefore leave a completed effect with
no persisted result. A started effect is unknown after a crash, and even a
failed tool may have made partial external changes. Returning to an older
complete checkpoint does not undo those changes. Tools with external effects
need their own idempotency/reconciliation strategy. No Harness event can make an
external side effect atomic with a local SQLite write.
Tests cover a real subprocess kill, early-request failure, final-history rewrite, revision conflicts, and bounded/cancelled naming. The kill test confirms that frontier capture survives without error hooks; it is not an exactly-once execution guarantee.
- It does not restore capability per-run state, graph-node state, retry counters, or in-flight streaming responses.
- It does not deduplicate replayed side effects automatically. Tools that write artifacts, labels, PRs, or external state should call
annotate_tool_effect(store, ctx, ...)(see Failure recovery) so the orchestrator can decide whether replay is safe. - It does not prune events, and by default does not prune snapshots. Retention is the caller’s responsibility; snapshot growth can be bounded opt-in with
max_snapshots_per_run(see Bounding snapshot growth). - It does not garbage-collect externalized media. Pruning a snapshot leaves its content-addressed blobs in place, since they may be shared across snapshots and runs.
- It does not emit OpenTelemetry spans. pydantic_ai’s
Instrumentationcapability already spansagent run/chat/running tooland populatesgen_ai.agent.name,gen_ai.agent.call.id,gen_ai.conversation.idvia baggage. A future change may add step-persistence attributes to the active span; that is tracked as a follow-up issue.
Bases: AbstractCapability[AgentDepsT]
Append-only step log + continuable snapshots + tool-effect ledger.
The capability emits a StepEvent at every interesting boundary
(run/model-request/tool-call start, completion, failure), records a
ToolEffectRecord per tool call so the orchestrator can decide whether
replay is safe, and saves a ContinuableSnapshot at every settled
CallToolsNode boundary — folding in the pending tool-return request, so
the point is durable the moment the tool completes — plus a fallback save
at after_run when the run ends past that boundary. A run that fails
saves the live at-failure history (see on_run_error), classified by its
tool-work state: complete when every tool call is resolved,
interrupted otherwise.
A run that crashes between before_tool_execute and after_tool_execute
leaves a visible event trail, a started tool-effect record (the
unknown_after_crash signal), and an interrupted snapshot carrying
every completed cycle. The default latest_snapshot / continue_run
read path only returns complete snapshots; pass
include_interrupted=True to resume from the interrupted frontier after
consulting list_unresolved_tool_effects.
from pydantic_ai import Agent
from pydantic_ai_harness.step_persistence import StepPersistence, InMemoryStepStore
store = InMemoryStepStore()
librarian = Agent(
'openai:gpt-5',
capabilities=[StepPersistence(store=store, agent_name='code_librarian')],
)
await librarian.run('Find ThinkingPartDelta and confirm the callable allowance')
Use continue_run(store, run_id=...) / fork_run(store, run_id=...)
to load a prior snapshot, then pass the result to
Agent.run(..., message_history=...).
Backend that records events, snapshots, and tool effects.
Type: StepStore Default: field(default_factory=InMemoryStepStore)
Logical agent name (e.g. code_librarian, reproducer).
Recorded on the run, its events, and its snapshots. When set, it is also
encoded into the context-derived run_id so store inspection identifies both
the agent and the durable run.
When unset, the running agent’s name is recorded instead, but the run_id
stays ctx.run_id: deriving the store key from Agent.name would change the
key of runs that callers look up by their context run id.
Type: str | None Default: None
Identifier for this one Agent.run call.
run_id is per-call, matching pydantic_ai.RunContext.run_id. For
multi-turn logical grouping use conversation_id on Agent.run(...) —
that is the pyai-native primitive for it.
Resolution order (materialised in for_run):
- Explicit value → used as-is. Single-shot use cases:
deterministic id for testing, replay, debugging. Reusing the
capability across multiple
.run()calls with the same explicitrun_idraisesValueErrorinbefore_run— the tool-effect ledger keys on(run_id, tool_call_id)and providers reuse deterministic tool-call ids, so a silent collision would erase theunknown_after_crashsignal. Useconversation_id=onAgent.runfor multi-turn grouping. agent_nameset,run_idunset -> path-safe encoding of both values. Reusing the capability instance yields distinct ids because the agent graph assigns each run its own id.- Neither set ->
ctx.run_idper.run().
Type: str | None Default: None
Run that spawned this one.
Auto-inferred from the enclosing StepPersistence wrap_run scope —
when an orchestrator’s tool synchronously calls a delegate’s
Agent.run(...), the delegate picks up the orchestrator’s run_id
here without manual threading. Set explicitly to override (e.g. for
cross-process delegation where ContextVars do not propagate).
Type: str | None Default: None
Free-form metadata stored on the RunRecord and on each event.
Type: dict[str, str] Default: field(default_factory=_empty_metadata)
Also checkpoint accepted inputs and model tool-call frontiers before execution.
These additional writes make first-request failures and process kills before a settled tool cycle inspectable. They do not make side effects safe to replay.
Type: bool Default: False
@classmethod
def from_spec(cls, *args: Any, **kwargs: Any) -> StepPersistence[Any]
Construct from a serialised spec.
Supports backend='memory' (default), backend='file' (with
directory), or backend='sqlite' (with database). Raises
ValueError for any other backend value — silently falling
back to in-memory storage would turn a typo into accidental
non-durability.
max_snapshots_per_run (default None, unbounded) is forwarded to
the constructed store to bound per-run snapshot growth.
StepPersistence[Any]
def compaction_transcript_handle() -> str | None
Retrieval handle to this run’s transcript, for compaction receipts.
Satisfies the compaction TranscriptHandleProvider protocol structurally (no import
coupling). A compaction strategy discovers this capability via RunContext.capabilities
and records its run id. Returns None before for_run has materialised the id.
@async
def for_run(ctx: RunContext[AgentDepsT]) -> AbstractCapability[AgentDepsT]
Materialise run_id and parent_run_id for this Agent.run call.
Reads the contextvar set by any enclosing StepPersistence.wrap_run
before the local run overwrites it, so a delegate’s parent_run_id
ends up pointing at its orchestrator’s run_id.
A separate ContextVar is needed because pydantic_ai’s own
cross-run signals (RUN_ID_BAGGAGE_KEY via OTel baggage,
RunContext.run_id, and _CURRENT_RUN_CONTEXT) are single-slot:
the inner Instrumentation.wrap_run overwrites them before any
nested capability sees the parent. The harness-local contextvar
lets us snapshot the parent here, before the local wrap_run
rebinds it.
AbstractCapability[AgentDepsT]
@async
def wrap_run(
ctx: RunContext[AgentDepsT],
*,
handler: WrapRunHandler,
) -> AgentRunResult[Any]
Push this run’s id onto the contextvar so nested delegates can read it.
@async
def before_run(ctx: RunContext[AgentDepsT]) -> None
Register run lineage and emit run_started.
Reject reuse by a distinct registration — the
tool-effect ledger keys on (run_id, tool_call_id) and providers
reuse deterministic tool-call ids, so a second Agent.run with
the same run_id would silently collide. A journaled registration id
distinguishes a retry of this durable operation from a new run.
@async
def after_run(
ctx: RunContext[AgentDepsT],
*,
result: AgentRunResult[Any],
) -> AgentRunResult[Any]
Emit run_completed, saving a final snapshot only as a fallback.
When a terminal CallToolsNode already saved the final history via
after_node_run it carries the correct step_index, whereas by
after_run ctx.run_step is reset to 0 — so re-saving would both
duplicate the tail and stamp a misleading step_index. We save only
when the final content differs from the newest boundary snapshot,
including same-length rewrites. The comparison uses per-run copied
state, not an external store read which could change under durable replay.
That covers a run which reached no provider-valid boundary at all, and
Agent.run_stream, which ends through SetFinalResult rather than a
terminal CallToolsNode and appends its closing response after the last
boundary — leaving after_run the only hook that sees the full run.
@async
def on_run_error(
ctx: RunContext[AgentDepsT],
*,
error: BaseException,
) -> AgentRunResult[Any]
Persist the live at-failure history as the run’s last resume point, then emit run_failed.
The single error-path save site: reads the list reference stashed by
after_node_run (see _stash_live_history), whose content at this
point is the full history the run had built when it failed — including
a failing model request’s payload and any partial tool returns captured
by the graph during unwind. Nothing is compared against the store:
the live history is by definition the newest state, so an earlier
boundary snapshot is simply superseded, and a history a sticky
processor trimmed is persisted as trimmed — exactly what the next
request would have sent.
The history is saved whenever it contains a model response (a bare
prompt equals restarting the run), classified complete when every
tool call is resolved and interrupted otherwise. Interrupted
snapshots stay off the default latest_snapshot read path.
@async
def on_model_request_error(
ctx: RunContext[AgentDepsT],
*,
request_context: ModelRequestContext,
error: Exception,
) -> ModelResponse
Emit model_request_failed and re-raise.
No snapshot is saved here: the failing request’s payload already sits
in the live history (the graph appends the request before sending), so
on_run_error’s save covers it. A failure the model layer recovers
from (retry, fallback) needs no rescue at all.
@async
def after_node_run(
ctx: RunContext[AgentDepsT],
*,
node: AgentNode[AgentDepsT],
result: NodeResult[AgentDepsT],
) -> NodeResult[AgentDepsT]
Save a continuable snapshot after a settled CallToolsNode, and refresh the live-history stash.
At that boundary every tool call from the preceding ModelRequestNode
has a matching tool return, so the history is provider-valid. The
returned ModelRequestNode carries those returns and is not yet in
ctx.messages, so its request is folded in before validation —
without it a worker killed right after a completed tool call would
leave no resume point at all (#373). is_provider_valid doubles as a
defense in case a custom node reshapes history. after_run compares the
saved content with the final history, including same-length rewrites
by other capabilities.
This save is the durable one: it lands in the store while the run is
still healthy, so it survives a hard kill that fires no hook. The
error path (on_run_error) only rescues histories that a raise unwinds
through.
Every node boundary also re-stashes the live message list so that
on_run_error can persist the at-failure history when a later node
raises before its own after_node_run fires. The stash holds that list
by reference, so the snapshot candidate rebinds to a new list rather
than appending to it — an append would leak result.request into the
history the error path later reads, duplicating it once the graph
appends the request itself.
NodeResult[AgentDepsT]