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:
| Prefix | Holds |
|---|---|
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:
[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:
[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.
[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:
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:
"""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:
[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:
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