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:
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 okEach 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:
agent_env.config.errors.ConfigError: The config table for 'agentenv_openciv3.stores:ContentAddressedObjectStore' has unknown key 'compress'; it takes root, bucketA 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:
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:
[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:
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:
cas://openciv3/artifacts/docker_image/mcp-server-openciv3/1/mcp-server-openciv3-v1.tar.gzProve 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:
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=.:
cd <agentenv-framework checkout> && PYTHONPATH=. python -m pytest -p no:cacheprovider -q -rA <path>/tests/test_store_conformance.pyAll 17 cases pass, and the dedup test with them:
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.55sWhat 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
- A document store: Write Your Own Document Store.
- An image store: Write Your Own Image Store.
- A secret store: Write Your Own Secret Store.
Last updated on