Skip to content
AgentEnv Framework

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:

Terminal
agent-env env get-instance --id email-mmtflq6t
Output (first lines)
Instance ID: email-mmtflq6t
Env ID: email
Env Version: 1

Env.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:

.agentenv/config.toml
[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:

.agentenv/config.toml
[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:

.agentenv/config.toml
[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:

.agentenv/config.toml
[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:

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

mycorp/postgres_store.py
"""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:

.agentenv/config.toml
[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:

tests/test_postgres_store.py
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

Ask a question · Report an issue

On this page