# Document Store (https://www.agentenvframework.com/docs/registry/document-store)

> Where every versioned record, and every record of a deploy or a run, is kept

1. `MCPServerEnv.put` inserts `email` into `envs` as version 1, since no version of that id exists yet.
2. The next `put` reads the highest version and inserts version 2 beside it, leaving version 1 as it was.
3. Two puts race for version 3: one lands, the other hits `DuplicateKeyError` on `(id, version)` and retries as 4.
4. `Env.get` without a version returns the latest, 4. A `deploy_env` step pinned to `env_version: 1` gets 1.
5. The deploy inserts `email-mmtflq6t` into `env_instances`, naming the env and the version it deployed.
6. `Env.from_instance_id` reads that record and loads version 1, not the latest.

Parts of the scene:

- **The framework**: Your code or the CLI: `agent-env env mcp-server put` calls `MCPServerEnv.put`, and a `deploy_env` step calls `Env.get` with its `env_version`.
- **The document store**: SQLite by default, one table per collection in `~/.local/state/agent-env/document_store/documents.db`. MongoDB, DynamoDB and Firestore with MongoDB compatibility (the `gcp` extra) are built in, and `[stores.document]` picks one, or a store you write.
- **envs**: One document per version of each env, which a put inserts and no later put changes. `artifacts`, `a2a_agents`, `tasks`, `task_steps` and `evals` are versioned the same way.
- **The unique index**: Every versioned collection has a unique index on `(id, version)`. A put that loses a race hits `DuplicateKeyError` and tries the next version, five attempts in all (three for an artifact).
- **env_instances**: One record per deploy of an MCP server, multi or website env, under the env id plus eight random lowercase letters and digits, with the `env_id` and `env_version` it deployed. `agent-env env get-instance` reads it back.

The document store keeps every record you register, one document per version, and every record
your deploys and runs write. [The registry](https://www.agentenvframework.com/docs/core-concepts.md#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`](https://www.agentenvframework.com/docs/tasks/important-steps.md#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:

```bash title="Terminal"
agent-env env get-instance --id email-mmtflq6t
```

```text title="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](https://www.agentenvframework.com/docs/tasks/running.md#the-task-run-instance) shows, and each step's context diff in
`task_step_journal`, which lets a [retry](https://www.agentenvframework.com/docs/tasks/task-steps.md#when-a-step-fails) 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:

```toml title=".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:

```toml title=".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:

```toml title=".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:

```toml title=".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:

```text title="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:

```python title="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:

```toml title=".agentenv/config.toml"
[stores.document]
impl = "mycorp.postgres_store:PostgresDocumentStore"
[stores.document.config]
dsn = "secret:postgres_dsn"
```

It passes all 35 `DocumentStore` [conformance cases](https://github.com/scaleapi/agentenv-framework/tree/main/tst/store)
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:

```python title="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.