Document Store
Where every versioned record, and every record of a deploy or a run, is kept
The document store keeps every record you register, one document per version, and every record your deploys and runs write. The registry shows it beside the other three stores.
What it holds
Versioned records live in envs, artifacts, a2a_agents, tasks, task_steps and evals.
Deploys and runs write beside them, in collections such as env_instances, a2a_agent_instances,
env_state_instances, task_instances, task_step_journal and runs.
Versions
A put reads the highest version of its id and inserts the next one, so it never changes an earlier
version. A unique index on (id, version) settles a race: the losing put gets
DuplicateKeyError and tries the next version, five attempts in all (three for an artifact).
A read without a version gets the latest, and a read with one gets exactly that version.
A deploy_env step pins one with env_version.
Instances and runs
Deploying an MCP server, multi or website env inserts a record into env_instances under an
instance id such as email-mmtflq6t, naming the env and the version it deployed:
agent-env env get-instance --id email-mmtflq6tInstance ID: email-mmtflq6t
Env ID: email
Env Version: 1Env.from_instance_id reconnects through this record, so it loads the deployed version, not the
latest. A task run records itself in task_instances, as
The task run instance shows, and each step's context diff in
task_step_journal, which lets a retry roll a failed
span back.
Use a Supported Document Store
Four document stores are built in. A [stores.document] table selects one by its impl and passes
its settings under config.
SQLite
The default, which needs no table: one table per collection in
~/.local/state/agent-env/document_store/documents.db. Name it only to keep the file elsewhere:
[stores.document]
impl = "agent_env.store.document_store:LocalSqliteDocumentStore"
config = { path = "/srv/agent-env/documents.db" }MongoDB
It takes a connection uri and a database, and pings the server when the store is built:
[stores.document]
impl = "agent_env.store.document_store:MongoDocumentStore"
[stores.document.config]
uri = "secret:mongodb_uri"
database = "agent_env"DynamoDB
Each collection gets its own on-demand table, {table_prefix}{collection}, created on first use.
The store signs in through the standard boto3 credential chain:
[stores.document]
impl = "agent_env.store.document_store:DynamoDbDocumentStore"
config = { table_prefix = "agentenv_", region = "<region>" }Firestore
Firestore with MongoDB compatibility, from the gcp extra. It takes the database's host and
name, and signs in with Application Default Credentials:
[stores.document]
impl = "agent_env.store.document_store.firestore_mongo_document_store:FirestoreMongoDocumentStore"
config = { host = "<uid>.<location>.firestore.goog", database = "<database>" }agent-env config explain document prints the store in effect and where it came from, without
connecting to it:
document
MongoDocumentStore uri=secret:mongodb_uri database=agent_env
from [stores.document]Write Your Own Document Store
A store of your own subclasses DocumentStore and implements its nine abstract methods.
agent_env.store.document_store.evaluation runs filters, sorts and updates in Python, so a backend
only has to keep whole documents, enforce unique indexes and make each write atomic, as the
built-in SQLite and DynamoDB stores do.
This one keeps every collection in one PostgreSQL table, a JSONB row per document. Each write locks its collection before it reads the document it changes, so of two writers racing on one document only one wins, and a unique index is a partial index on that collection:
"""A DocumentStore on PostgreSQL: one JSONB row per document, filtered in Python."""
import copy
import hashlib
import json
import threading
from contextlib import contextmanager
from datetime import datetime
from typing import Optional
from psycopg2 import errors, sql
from psycopg2.pool import ThreadedConnectionPool
from agent_env.store.document_store import (
DocumentStore,
DuplicateKeyError,
Filter,
Sort,
UpdateSpec,
)
from agent_env.store.document_store import evaluation
SELECT_ROWS = "SELECT pk, doc FROM documents WHERE collection = %s ORDER BY pk"
def _dumps(doc: dict) -> str:
def default(value):
if isinstance(value, datetime):
return value.isoformat()
raise TypeError(f"{type(value).__name__} is not JSON-serializable")
return json.dumps(doc, default=default)
class PostgresDocumentStore(DocumentStore):
"""Every collection shares one `documents` table. A unique index is a partial
expression index on the collection, and each write holds a lock on its
collection until it commits."""
def __init__(self, dsn: str, max_connections: int = 10) -> None:
self._pool = ThreadedConnectionPool(1, max_connections, dsn)
self._slots = threading.BoundedSemaphore(max_connections)
with self._transaction() as cur:
cur.execute("SELECT pg_advisory_xact_lock(hashtext('documents'))")
cur.execute(
"CREATE TABLE IF NOT EXISTS documents (pk bigserial PRIMARY KEY,"
" collection text NOT NULL, doc jsonb NOT NULL)"
)
cur.execute(
"CREATE INDEX IF NOT EXISTS documents_collection"
" ON documents (collection, pk)"
)
@contextmanager
def _transaction(self):
with self._slots:
conn = self._pool.getconn()
try:
with conn, conn.cursor() as cur: # commits, or rolls back on an error
yield cur
finally:
self._pool.putconn(conn, close=bool(conn.closed))
def _matching(self, collection: str, filter: Filter) -> list[dict]:
with self._transaction() as cur:
cur.execute(SELECT_ROWS, (collection,))
return [doc for _, doc in cur.fetchall() if evaluation.matches(doc, filter)]
def _first_locked(self, cur, collection: str, filter: Filter):
"""The first match, read after taking the collection's write lock, so no
other write can change it before this transaction commits."""
cur.execute("SELECT pg_advisory_xact_lock(hashtext(%s))", (collection,))
cur.execute(SELECT_ROWS, (collection,))
rows = cur.fetchall()
return next((row for row in rows if evaluation.matches(row[1], filter)), None)
def _write(self, cur, query: str, params: tuple) -> None:
try:
cur.execute(query, params)
except errors.UniqueViolation as e:
raise DuplicateKeyError(str(e)) from e
def _insert(self, cur, collection: str, doc: dict) -> None:
query = "INSERT INTO documents (collection, doc) VALUES (%s, %s)"
self._write(cur, query, (collection, _dumps(doc)))
def _overwrite(self, cur, pk: int, doc: dict) -> None:
query = "UPDATE documents SET doc = %s WHERE pk = %s"
self._write(cur, query, (_dumps(doc), pk))
def find_one(
self, collection: str, filter: Filter, sort: Optional[Sort] = None
) -> Optional[dict]:
docs = evaluation.sort_docs(self._matching(collection, filter), sort)
return docs[0] if docs else None
def query(
self,
collection: str,
filter: Filter,
sort: Optional[Sort] = None,
limit: Optional[int] = None,
offset: Optional[int] = None,
) -> list[dict]:
docs = evaluation.sort_docs(self._matching(collection, filter), sort)
docs = docs[offset:] if offset else docs
return docs[:limit] if limit else docs
def count(self, collection: str, filter: Filter) -> int:
return len(self._matching(collection, filter))
def insert(self, collection: str, doc: dict) -> None:
with self._transaction() as cur:
self._insert(cur, collection, doc)
def update(
self, collection: str, filter: Filter, update: UpdateSpec, upsert: bool = False
) -> int:
doc = self.update_one_and_get(collection, filter, update, upsert=upsert)
return 0 if doc is None else 1
def update_one_and_get(
self,
collection: str,
filter: Filter,
update: UpdateSpec,
return_after: bool = True,
upsert: bool = False,
) -> Optional[dict]:
with self._transaction() as cur:
match = self._first_locked(cur, collection, filter)
if match is None:
if not upsert:
return None
doc = evaluation.synthesize(filter, update)
self._insert(cur, collection, doc)
return doc if return_after else None
pk, doc = match
before = copy.deepcopy(doc)
evaluation.apply_update(doc, update)
self._overwrite(cur, pk, doc)
return doc if return_after else before
def replace(
self, collection: str, filter: Filter, doc: dict, upsert: bool = False
) -> int:
with self._transaction() as cur:
match = self._first_locked(cur, collection, filter)
if match is not None:
self._overwrite(cur, match[0], doc)
elif upsert:
self._insert(cur, collection, doc)
else:
return 0
return 1
def delete(self, collection: str, filter: Filter) -> int:
with self._transaction() as cur:
match = self._first_locked(cur, collection, filter)
if match is None:
return 0
cur.execute("DELETE FROM documents WHERE pk = %s", (match[0],))
return 1
def ensure_index(
self,
collection: str,
fields: list[str],
unique: bool = False,
ttl_seconds: Optional[int] = None,
) -> None:
"""Builds unique indexes only: reads filter in Python, so no other index
would ever be used."""
if not unique:
return
digest = hashlib.sha256(f"{collection}:{fields}".encode()).hexdigest()
name = f"documents_{digest[:16]}"
keys = sql.SQL(", ").join(
sql.SQL("(doc #>> {})").format(sql.Literal(field.split(".")))
for field in fields
)
create = sql.SQL(
"CREATE UNIQUE INDEX IF NOT EXISTS {} ON documents ({})"
" WHERE collection = {}"
).format(sql.Identifier(name), keys, sql.Literal(collection))
with self._transaction() as cur:
cur.execute("SELECT pg_advisory_xact_lock(hashtext(%s))", (name,))
cur.execute(create)It needs psycopg2, and you select it like a built-in, with its DSN behind a secret: reference:
[stores.document]
impl = "mycorp.postgres_store:PostgresDocumentStore"
[stores.document.config]
dsn = "secret:postgres_dsn"It passes all 35 DocumentStore conformance cases
against PostgreSQL 16, including the ones where eight threads race for one document. The suite
ships in the agentenv-framework repository's tst/, not in the wheel, so run it from a clone.
Each case gets a fresh collection:
import uuid
import pytest
from mycorp.postgres_store import PostgresDocumentStore
from tst.store import conformance
@pytest.fixture(scope="module")
def store():
return PostgresDocumentStore("postgresql://postgres@localhost:5432/postgres")
@pytest.mark.parametrize("case", conformance.CASES, ids=lambda c: c.__name__)
def test_conformance(case, store):
case(store, f"coll_{uuid.uuid4().hex}")Every read loads the whole collection, which suits a registry's records. A store for millions of documents would push its filters down into SQL instead.
Last updated on