Skip to content
AgentEnv Framework

Registry Plugins

Ship a store in a plugin package, with a content-addressed object store for OpenCiv3's recordings as the example

OpenCiv3 keeps large files as objects: each game's recording, a 184 KB MP4 and 933 KB HTML replay for a 30-turn game, and the 134 MB env image tarball. The plugin ships no store; this page writes one that keeps each payload once.

Stores Are Chosen by Config

A plugin package can replace any of the registry's four stores, but not through an entry point as envs and steps do. There is no store entry-point group: a [stores.<kind>] table names the class with impl, and its config table goes to from_config(**config), which by default calls the constructor.

So installing the package changes no store, and plugin list counts none among what it provides:

agent-env plugin list
PACKAGE             VERSION   PROVIDES                               STATUS
agentenv-framework  0.9.1267  bundle hello                           ok
agentenv-openciv3   0.1.0     3 task steps, 1 CLI command, 1 bundle  ok

Each kind has exactly one store in effect, so it is chosen in config rather than by what is installed. A table that names a missing class, or passes a key the constructor doesn't take, fails the first time the store is built instead of falling back to the default:

Output
agent_env.config.errors.ConfigError: The config table for 'agentenv_openciv3.stores:ContentAddressedObjectStore' has unknown key 'compress'; it takes root, bucket

A Content-Addressed Object Store

Each key maps to a blob named by the sha256 digest of its bytes, through an SQLite index, so two keys with the same bytes share one blob. Trimmed to its write path, it is:

src/agentenv_openciv3/stores.py
from agent_env.store import ObjectAlreadyExistsError, ObjectNotFoundError
from agent_env.store.object_store import DEFAULT_CONTENT_TYPE, ObjectMetadata, ObjectStore

_SCHEMA = """CREATE TABLE IF NOT EXISTS objects (
    key TEXT PRIMARY KEY, digest TEXT NOT NULL, size INTEGER NOT NULL, content_type TEXT, last_modified TEXT NOT NULL)"""


class ContentAddressedObjectStore(ObjectStore):
    """Keys map to sha256-named blobs under root/blobs, through an SQLite index; object URLs are cas://<bucket>/<key>."""

    def put(self, key: str, data: bytes, content_type: str = DEFAULT_CONTENT_TYPE, allow_overwrite: bool = False) -> str:
        digest = hashlib.sha256(data).hexdigest()
        self._store_blob(digest, lambda path: path.write_bytes(data))
        return self._record(key, digest, len(data), content_type, allow_overwrite)

    def put_file(self, key: str, file_path: str, content_type: str = DEFAULT_CONTENT_TYPE) -> str:
        with open(file_path, "rb") as f:
            digest = hashlib.file_digest(f, "sha256").hexdigest()
        self._store_blob(digest, lambda path: shutil.copyfile(file_path, path))
        return self._record(key, digest, os.path.getsize(file_path), content_type, allow_overwrite=False)

    def _blob_path(self, digest: str) -> Path:
        return self._blobs / digest[:2] / digest

    def _store_blob(self, digest: str, write: Callable[[Path], object]) -> None:
        """Write a payload once; a concurrent writer of the same bytes loses the rename harmlessly."""
        path = self._blob_path(digest)
        if path.exists():
            return
        path.parent.mkdir(exist_ok=True)
        fd, staging = tempfile.mkstemp(dir=path.parent)
        os.close(fd)
        try:
            write(Path(staging))
            os.replace(staging, path)
        except BaseException:
            os.unlink(staging)
            raise

    def _record(self, key: str, digest: str, size: int, content_type: str, allow_overwrite: bool) -> str:
        verb = "INSERT OR REPLACE" if allow_overwrite else "INSERT"
        now = datetime.now(timezone.utc).isoformat()
        try:
            with closing(self._connect()) as db, db:
                db.execute(f"{verb} INTO objects VALUES (?, ?, ?, ?, ?)", (key, digest, size, content_type, now))
        except sqlite3.IntegrityError as e:
            raise ObjectAlreadyExistsError(f"Object already exists at {self.object_url(key)}.") from e
        return self.object_url(key)

A second write of the same bytes costs only a hash and an index row. put_file hashes and copies a file in chunks, so an image tarball never sits in memory.

Write-once belongs to the key, not the blob: a plain INSERT on a taken key raises ObjectAlreadyExistsError, and allow_overwrite points the key at another digest. A blob appears only through an atomic rename, and each call opens its own index connection, which keeps it safe when core calls it from several threads.

The rest is the contract every object store meets. It implements the eleven abstract methods, raises ObjectNotFoundError for a missing object and ValueError for a URL it doesn't own, which the inherited owns() relies on.

Select It in Config

The store is named in config.toml, with its root behind an env: reference:

.agentenv/config.toml
[stores.object]
impl = "agentenv_openciv3.stores:ContentAddressedObjectStore"
config = { root = "env:OPENCIV3_RECORDING_ROOT?~/.local/state/agentenv-openciv3/objects", bucket = "openciv3" }

agent-env config show reports it without building it, printing the reference as written:

agent-env config show (excerpt)
object:         ContentAddressedObjectStore  root=env:OPENCIV3_RECORDING_ROOT?~/.local/state/agentenv-openciv3/objects  bucket=openciv3
                from [stores.object]

Objects of @local ids, such as a bundle run's, stay in the per-user local store. The image tarball agent-env openciv3 setup stores lands at:

Output
cas://openciv3/artifacts/docker_image/mcp-server-openciv3/1/mcp-server-openciv3-v1.tar.gz

Prove It With the Conformance Kit

The SDK's object store kit is the 17 CASES in tst/store/object_conformance.py, run by a test of your own. The kit isn't in the wheel, so the test imports it from an agentenv-framework checkout:

tests/test_store_conformance.py
import pytest

from agentenv_openciv3.stores import ContentAddressedObjectStore
from tst.store import object_conformance


@pytest.fixture
def store(tmp_path):
    return ContentAddressedObjectStore(str(tmp_path))


@pytest.mark.parametrize("case", object_conformance.CASES, ids=lambda c: c.__name__)
def test_conformance(case, store):
    case(store, "")


def test_identical_payloads_share_one_blob(store):
    first = store.put("recordings/game-1.mp4", b"same bytes", content_type="video/mp4")
    second = store.put("recordings/game-2.mp4", b"same bytes", content_type="video/mp4")
    assert first != second and store.get(first) == store.get(second) == b"same bytes"
    assert store.blob_count() == 1
    assert store.get_object_metadata("recordings/game-2.mp4").content_type == "video/mp4"

Run it from the checkout root with PYTHONPATH=.:

Terminal
cd <agentenv-framework checkout> && PYTHONPATH=. python -m pytest -p no:cacheprovider -q -rA <path>/tests/test_store_conformance.py

All 17 cases pass, and the dedup test with them:

Output
PASSED ...test_conformance[put_get_roundtrip]
PASSED ...test_conformance[put_is_write_once]
...
PASSED ...test_conformance[missing_object_reads_raise_object_not_found]
PASSED ...test_identical_payloads_share_one_blob
18 passed in 0.55s

What It Doesn't Do

It signs no URLs, so the framework copies objects into a VM sandbox over exec instead of the VM downloading them. What needs a signed URL fails: env snapshots, GitHub image builds, and loading a file artifact into a built-in env that a plugin's provider deployed.

It issues no transfer grants either (supports_transfer_grants stays False), so A2A agents get none from it and the kit's three GRANT_CASES don't apply.

Its blobs live on one host's disk, and nothing removes one, even after an overwrite leaves it unreferenced. Several processes sharing one root have not been tested.

The Other Three Stores

Last updated on

Ask a question · Report an issue

On this page