Skip to content
AgentEnv Framework
RegistryObject Store

Object Store

Where the bytes that registry records point at live, written once under versioned keys

The object store holds the bytes that document store records point at: artifact files, image tarballs and build contexts, collected files, trajectories and env snapshots. It is one of the registry's four stores.

Object URLs

A record holds an object URL, never the bytes: file:// on the local store, s3://<bucket>/ or gs://<bucket>/ on a cloud one. The field keeps its legacy name, such as a file artifact's s3_url, whatever the store.

Keys

An artifact's objects sit at artifacts/<type>/<id>/<version>/<name>, so version 1 of acme-email-file from Creating your artifacts is artifacts/file/acme-email-file/1/emails.json. Other writes use prefixes of their own, among them:

PrefixHolds
artifacts/artifact files, image tarballs and build contexts, collected files
env-snapshots/env snapshots, from agent-env env snapshot
prompt_agent_trajectories/, env_trajectory/agent and env trajectories
agent_snapshots/agent snapshots
github-builds/images built from GitHub

Written once

A write to a key that already holds an object fails with ObjectAlreadyExistsError on every built-in store, so putting a new version never replaces the bytes an earlier one pins.

Signed URLs

When a VM sandbox needs an object, the framework asks the store for a signed URL and the sandbox downloads the object itself. With a store that signs nothing, like the local one, the framework copies the bytes in instead. env snapshot and GitHub image builds upload through a signed URL, so they need a store that signs.

Use a Supported Object Store

Three object stores are built in. A [stores.object] table selects one by its impl and passes its settings under config.

Local Filesystem

The default, which needs no table: each object is a file under ~/.local/state/agent-env/object_store/, and no URL is signed. Name it only to keep the files elsewhere:

.agentenv/config.toml
[stores.object]
impl = "agent_env.store.object_store:LocalFilesystemObjectStore"
config = { root = "/srv/agent-env/objects" }

S3

It takes a bucket and an optional region, signs in through the standard boto3 credential chain, and signs GET and PUT URLs:

.agentenv/config.toml
[stores.object]
impl = "agent_env.store.object_store:S3ObjectStore"
config = { bucket = "<bucket>", region = "<region>" }

share_credentials = true also hands those AWS credentials to the agents and env services the framework deploys, with the identity's whole IAM scope.

Cloud Storage

From the gcp extra. It takes a bucket, an optional project and an optional signing_service_account, and signs in with Application Default Credentials. It signs URLs only with a signer: that service account, or credentials that hold a service account key.

.agentenv/config.toml
[stores.object]
impl = "agent_env.store.object_store.gcs_object_store:GcsObjectStore"
[stores.object.config]
bucket                  = "<bucket>"
signing_service_account = "<name>@<project>.iam.gserviceaccount.com"

agent-env config explain object prints the store in effect and where it came from, without connecting to it:

Output
object
  S3ObjectStore  bucket=<bucket>  region=<region>
  from [stores.object]

Write Your Own Object Store

A store of your own subclasses ObjectStore and implements its eleven abstract methods. A second write to a key must raise ObjectAlreadyExistsError, and a read of a missing object ObjectNotFoundError. Signing is optional: for a store that signs nothing, the framework copies objects into sandboxes itself.

This one keeps objects in an Azure Blob Storage container, which no built-in store covers. Azure itself refuses a write to a name that is taken, and the account key in the connection string signs download URLs:

mycorp/azure_store.py
"""An ObjectStore on Azure Blob Storage: one container, each blob written once."""

from __future__ import annotations

import os
from datetime import datetime, timedelta, timezone

from azure.core.exceptions import ResourceExistsError, ResourceNotFoundError
from azure.storage.blob import (
    BlobSasPermissions,
    BlobServiceClient,
    ContentSettings,
    generate_blob_sas,
)

from agent_env.store import ObjectAlreadyExistsError, ObjectNotFoundError
from agent_env.store.object_store import (
    DEFAULT_CONTENT_TYPE,
    ObjectMetadata,
    ObjectStore,
)

SCHEME = "az://"


class AzureBlobObjectStore(ObjectStore):
    """Object urls are `az://<container>/<key>`. Downloads get SAS URLs; uploads get
    none, because Azure refuses a PUT without an `x-ms-blob-type` header and the
    framework's uploaders send a bare PUT."""

    def __init__(self, connection_string: str, container: str) -> None:
        self._service = BlobServiceClient.from_connection_string(connection_string)
        self._container = container

    def _split(self, object_url: str) -> tuple[str, str]:
        if not object_url.startswith(SCHEME):
            raise ValueError(f"{object_url!r} is not an {SCHEME} url.")
        container, _, key = object_url.removeprefix(SCHEME).partition("/")
        return container, key

    def _blob(self, object_url: str):
        return self._service.get_blob_client(*self._split(object_url))

    def _upload(self, object_url: str, data, content_type: str, overwrite: bool) -> str:
        settings = ContentSettings(content_type=content_type)
        try:
            self._blob(object_url).upload_blob(
                data, overwrite=overwrite, content_settings=settings
            )
        except ResourceExistsError as e:
            raise ObjectAlreadyExistsError(f"An object exists at {object_url}.") from e
        return object_url

    def _download(self, object_url: str):
        try:
            return self._blob(object_url).download_blob()
        except ResourceNotFoundError as e:
            raise ObjectNotFoundError(f"No object at {object_url}.") from e

    def put(
        self,
        key: str,
        data: bytes,
        content_type: str = DEFAULT_CONTENT_TYPE,
        allow_overwrite: bool = False,
    ) -> str:
        return self._upload(self.object_url(key), data, content_type, allow_overwrite)

    def put_file(
        self, key: str, file_path: str, content_type: str = DEFAULT_CONTENT_TYPE
    ) -> str:
        return self.put_file_at(self.object_url(key), file_path, content_type)

    def put_file_at(
        self, object_url: str, file_path: str, content_type: str = DEFAULT_CONTENT_TYPE
    ) -> str:
        with open(file_path, "rb") as f:
            return self._upload(object_url, f, content_type, overwrite=False)

    def get(self, object_url: str) -> bytes:
        return self._download(object_url).readall()

    def download_to_file(self, object_url: str, dest_path: str) -> None:
        stream = self._download(object_url)
        os.makedirs(os.path.dirname(dest_path) or ".", exist_ok=True)
        with open(dest_path, "wb") as f:
            stream.readinto(f)

    def get_object_metadata(self, key: str) -> ObjectMetadata | None:
        return self.get_object_metadata_at(self.object_url(key))

    def get_object_metadata_at(self, object_url: str) -> ObjectMetadata | None:
        try:
            props = self._blob(object_url).get_blob_properties()
        except ResourceNotFoundError:
            return None
        return ObjectMetadata(
            content_type=props.content_settings.content_type,
            size=props.size,
            last_modified=props.last_modified,
            content_encoding=props.content_settings.content_encoding,
        )

    def list(self, prefix: str) -> list[str]:
        urls = self.list_at(self.object_url(prefix))
        return [self.get_object_key(url) for url in urls]

    def list_at(self, url_prefix: str) -> list[str]:
        container, prefix = self._split(url_prefix)
        client = self._service.get_container_client(container)
        names = client.list_blob_names(name_starts_with=prefix)
        return [
            f"{SCHEME}{container}/{name}" for name in names if not name.endswith("/")
        ]

    def object_url(self, key: str) -> str:
        return f"{SCHEME}{self._container}/{key}"

    def get_object_key(self, object_url: str) -> str:
        container, key = self._split(object_url)
        if container != self._container:
            raise ValueError(f"{object_url!r} is not in container {self._container!r}.")
        return key

    def signed_get_url(self, object_url: str, expires_in: int = 3600) -> str | None:
        """A read-only SAS URL, when the connection string holds the account key."""
        credential = self._service.credential
        if not getattr(credential, "account_key", None):
            return None
        container, key = self._split(object_url)
        sas = generate_blob_sas(
            credential.account_name,
            container,
            key,
            account_key=credential.account_key,
            permission=BlobSasPermissions(read=True),
            expiry=datetime.now(timezone.utc) + timedelta(seconds=expires_in),
        )
        return f"{self._blob(object_url).url}?{sas}"

It needs azure-storage-blob, and you select it like a built-in, with its connection string behind a secret: reference:

.agentenv/config.toml
[stores.object]
impl = "mycorp.azure_store:AzureBlobObjectStore"
[stores.object.config]
connection_string = "secret:azure_storage_connection_string"
container         = "agent-env"

It passes all 17 ObjectStore conformance cases against Azurite, Azure's local emulator. 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 prefix:

tests/test_azure_store.py
import uuid

import pytest
from azure.storage.blob import BlobServiceClient

from mycorp.azure_store import AzureBlobObjectStore
from tst.store import object_conformance

AZURITE = "UseDevelopmentStorage=true"


@pytest.fixture(scope="module")
def store():
    container = f"conformance-{uuid.uuid4().hex}"
    BlobServiceClient.from_connection_string(AZURITE).create_container(container)
    return AzureBlobObjectStore(AZURITE, container)


@pytest.mark.parametrize("case", object_conformance.CASES, ids=lambda c: c.__name__)
def test_conformance(case, store):
    case(store, f"{uuid.uuid4().hex}/")

It signs no uploads, so like a store that signs nothing it can't serve env snapshot or GitHub image builds, which upload through a signed URL.

Last updated on

Ask a question · Report an issue

On this page