Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 4 additions & 32 deletions src/agent_env/bundle/ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@

import hashlib
import json
import os
from collections.abc import Callable, Iterator, Mapping
from contextlib import contextmanager
from dataclasses import dataclass
Expand All @@ -22,26 +21,20 @@
from agent_env.artifact.registry import canonical_type, get_artifact_registry
from agent_env.artifact.store import ARTIFACTS_COLLECTION
from agent_env.config import get_config
from agent_env.config.paths import state_root
from agent_env.env.env import Env
from agent_env.env.registry import get_env_registry
from agent_env.env.store import ENVS_COLLECTION
from agent_env.eval.store import EVALS_COLLECTION
from agent_env.store import Filter, Sort
from agent_env.store.document_store import LocalSqliteDocumentStore
from agent_env.store.local_state import ensure_state_dir
from agent_env.store.local_state import holding_locks
from agent_env.task.store import TASKS_COLLECTION

from .authoring import build_context_files, entry_files
from .parse import Bundle, BundleKind
from .plan import Plan, Write, folder_walk, keeps_base_from_toml
from .resolve import BuiltImage

try:
import fcntl
except ImportError: # Windows: runs of one bundle aren't serialized
fcntl = None

LEDGER_COLLECTION = "bundle_ledger"
# Bumped by hand when what a digest covers changes, or a writer's output changes enough that every
# version it wrote should be written again.
Expand Down Expand Up @@ -173,31 +166,10 @@ def _find(self, collection: str, filter: Filter) -> dict | None:
@contextmanager
def materializing(bundle: Bundle, on_wait: Callable[[], None] | None = None) -> Iterator[None]:
"""Hold a lock on every id the bundle's entries write while its entities are written. Another run that
writes any of those ids, of this bundle or another, waits here, calling ``on_wait`` first. The locks are
taken in one order, so two runs can't each hold one the other waits for. Task runs don't hold them."""
if fcntl is None:
yield
return
ensure_state_dir(state_root() / "locks")
fds, waited = [], False
try:
for path in sorted({_lock_path(entry.id) for entry in bundle.entries}):
fds.append(os.open(path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW, 0o600))
try:
fcntl.flock(fds[-1], fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
if on_wait is not None and not waited:
on_wait()
waited = True
fcntl.flock(fds[-1], fcntl.LOCK_EX)
writes any of those ids, of this bundle or another, waits here, calling ``on_wait`` first. Task runs don't
hold them."""
with holding_locks((entry.id for entry in bundle.entries), on_wait):
yield
finally:
for fd in fds:
os.close(fd)


def _lock_path(id: str) -> Path:
return state_root() / "locks" / f"entity-{hashlib.sha256(id.encode()).hexdigest()[:16]}.lock"


def _tracked(write: Write) -> bool:
Expand Down
361 changes: 361 additions & 0 deletions src/agent_env/bundle/preflight.py

Large diffs are not rendered by default.

29 changes: 20 additions & 9 deletions src/agent_env/bundle/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,9 @@
and the others go on. Each run's sandboxes are torn down as it ends, unless ``keep`` holds them up. Ctrl-C or
SIGTERM cancels the runs: each one that started is marked cancelled and torn down, and a second one stops the
teardown. A bundle's evals run only the bundle's own tasks for now, so one naming a store task is refused
before anything is written. A dry run makes every check the run makes before its first task, and writes and
runs nothing.
before anything is written, and so is a deploy its sandbox provider can't serve (``preflight``). The infra envs a
gateway deploy on the local provider needs are built once the writes are done, before any task starts. A dry run
makes every check the run makes before its first task, and writes, builds and runs nothing.
"""

from __future__ import annotations
Expand All @@ -21,6 +22,7 @@
from enum import Enum
from pathlib import Path

from agent_env.env.bootstrap import InfraBuild, ensure_default_envs
from agent_env.providers import build_sandbox_provider
from agent_env.store.routing import namespace_routing, run_scope
from agent_env.task import Task, record_task_cancelled
Expand All @@ -32,6 +34,7 @@
from .materialize import Materialization, Materialized, materialize
from .parse import BundleEntry, BundleError, BundleKind, parse_bundle
from .plan import Plan, Write, plan_bundle
from .preflight import Preflight, preflight_run
from .resolve import BuiltImage, Reference, resolve_bundle

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -147,6 +150,7 @@ class DryRun:
materialization: Materialization # each write at the version it would leave, and why
runs: tuple[BundleEntry, ...] # the tasks that would run, each once, in the plan's order
skipped: tuple[BundleEntry, ...] # the tasks no eval names, which running every eval leaves out
infra: tuple[InfraBuild, ...] = () # the infra envs the run would build first

def path(self, entry: BundleEntry) -> str:
"""``entry``'s path in the bundle (``tasks/hello.json``)."""
Expand Down Expand Up @@ -191,14 +195,16 @@ def run_bundle(
ends.

Before anything is written, raises RuntimeError when an event loop is already running, ValueError when
``sandbox`` names no provider, and BundleError when the bundle can't be planned. A task that fails its
preflight raises BundleError once the entities it reads are written, before any task is."""
``sandbox`` names no provider, and BundleError when the bundle can't be planned or a deploy can't run where
it would (``preflight``). A task that fails its preflight raises BundleError once the entities it reads are
written, before any task is. Building an infra env the run needs raises what the build raises, before any task
runs."""
_refuse_a_running_loop()
if sandbox:
build_sandbox_provider(sandbox)
say = _progress(on_progress)
with namespace_routing():
plan = _planned(root, tasks, evals, id_root)
plan, preflight = _planned(root, tasks, evals, id_root, sandbox)
materialization = materialize(
plan,
on_wait=lambda: say("waiting for another agent-env run to finish writing this bundle's ids"),
Expand All @@ -207,6 +213,8 @@ def run_bundle(
)
entries = _to_run(plan)
to_run = [(entry, Task.get(entry.id, materialization.version_of("task", entry.id))) for entry in entries]
if preflight.infra:
ensure_default_envs(preflight.infra_kinds, say=say)
with Interrupts() as interrupts:
runs = interrupts.run(_run_all(plan, to_run, model, sandbox, keep, say, interrupts))
_mark_cancelled(runs, interrupts.reason)
Expand Down Expand Up @@ -242,10 +250,10 @@ def dry_run_bundle(
build_sandbox_provider(sandbox)
say = _progress(on_progress)
with namespace_routing():
plan = _planned(root, tasks, evals, id_root)
plan, preflight = _planned(root, tasks, evals, id_root, sandbox)
materialization = materialize(plan, dry_run=True, on_write=lambda done: say(_written(plan, done)))
runs = _to_run(plan)
return DryRun(materialization, runs, _skipped(plan, runs, every=not tasks and not evals))
return DryRun(materialization, runs, _skipped(plan, runs, every=not tasks and not evals), preflight.infra)


def _mark_cancelled(runs: tuple[TaskRun, ...], reason: str) -> None:
Expand All @@ -269,10 +277,13 @@ def _bundle_run(plan: Plan, materialization: Materialization, runs: tuple[TaskRu
return BundleRun(materialization, runs, eval_runs, skipped)


def _planned(root: Path | str, tasks: Sequence[str], evals: Sequence[str], id_root: str | None) -> Plan:
def _planned(root: Path | str, tasks: Sequence[str], evals: Sequence[str], id_root: str | None,
sandbox: str | None) -> tuple[Plan, Preflight]:
plan = plan_bundle(resolve_bundle(parse_bundle(Path(root), id_root=id_root)), tasks=tasks, evals=evals)
refuse_store_tasks(plan)
return plan
to_run = set(_to_run(plan))
return plan, preflight_run(plan, [write.source for write in plan.writes
if write.kind is BundleKind.TASK and write.source.entry in to_run], sandbox)


def _to_run(plan: Plan) -> tuple[BundleEntry, ...]:
Expand Down
39 changes: 5 additions & 34 deletions src/agent_env/cli/env/gateway.py
Original file line number Diff line number Diff line change
@@ -1,17 +1,9 @@
import sys
from pathlib import Path

import click

from agent_env.artifact import DockerImageArtifact
from agent_env.cli.utils import build_platform_option, detect_env_metadata
from agent_env.env import GatewayEnv
from agent_env.utils.docker_build import build_image

_PACKAGE_ROOT = Path(__file__).parent.parent.parent
GATEWAY_DOCKERFILE = _PACKAGE_ROOT / "env" / "gateway" / "Dockerfile"
GATEWAY_CONTEXT = _PACKAGE_ROOT / "env"
GATEWAY_IMAGE_TAG = "env-gateway"
from agent_env.cli.utils import build_platform_option
from agent_env.env.bootstrap import GATEWAY_CONTEXT, GATEWAY_DOCKERFILE, GATEWAY_IMAGE_TAG, put_gateway_env # noqa: F401


@click.group()
Expand All @@ -26,32 +18,11 @@ def gateway():
@build_platform_option
def put(env_id: str, metadata_pairs: tuple[str, ...], build_platform: str):
"""Build and upload a gateway environment."""

click.echo(f"Building gateway Docker image...")
build_image(GATEWAY_DOCKERFILE, GATEWAY_CONTEXT, GATEWAY_IMAGE_TAG, platform=build_platform)

click.echo(f"Creating DockerImageArtifact...")
artifact = DockerImageArtifact.put(
id=f"gateway-{env_id}",
description="Created from agent-env CLI",
image_name=GATEWAY_IMAGE_TAG,
)
click.echo(f"Created artifact: id={artifact.id} version={artifact.version}")

user_metadata = {}
metadata = {}
for pair in metadata_pairs:
if "=" not in pair:
click.echo(f"Invalid metadata format '{pair}', expected key=value", err=True)
sys.exit(1)
key, value = pair.split("=", 1)
user_metadata[key] = value
metadata = detect_env_metadata(GATEWAY_DOCKERFILE, GATEWAY_CONTEXT)
metadata.update(user_metadata)

click.echo(f"Creating GatewayEnv...")
env = GatewayEnv.put(
id=env_id,
docker_image_artifact=artifact,
metadata=metadata if metadata else None,
)
click.echo(f"Created GatewayEnv: id={env.id} version={env.version}")
metadata[key] = value
put_gateway_env(env_id, platform=build_platform, metadata=metadata, say=click.echo)
89 changes: 14 additions & 75 deletions src/agent_env/cli/env/service_db.py
Original file line number Diff line number Diff line change
@@ -1,22 +1,18 @@
"""CLI commands for ServiceDBEnv."""

from pathlib import Path

import click

from agent_env.artifact import DockerImageArtifact
from agent_env.cli.utils import build_platform_option, detect_env_metadata
from agent_env.env.envs.service_db import ServiceDBEnv
from agent_env.cli.utils import build_platform_option
from agent_env.config import get_config
from agent_env.utils.docker_build import build_image

# Path to ServiceDB Dockerfile
SERVICE_DB_DOCKERFILE = Path(__file__).parent.parent.parent / "env" / "envs" / "service_db" / "Dockerfile"
SERVICE_DB_IMAGE_NAME = "agent-env-service-db"
DB_WEB_DOCKERFILE = Path(__file__).parent.parent.parent / "env" / "envs" / "service_db" / "Dockerfile.db-web"
DB_WEB_IMAGE_NAME = "agent-env-db-web"
DB_MCP_DOCKERFILE = Path(__file__).parent.parent.parent / "env" / "envs" / "service_db" / "Dockerfile.db-mcp"
DB_MCP_IMAGE_NAME = "agent-env-db-mcp"
from agent_env.env.bootstrap import ( # noqa: F401
DB_MCP_DOCKERFILE,
DB_MCP_IMAGE_NAME,
DB_WEB_DOCKERFILE,
DB_WEB_IMAGE_NAME,
SERVICE_DB_DOCKERFILE,
SERVICE_DB_IMAGE_NAME,
put_service_db_env,
)


@click.group(name="service-db")
Expand All @@ -31,69 +27,12 @@ def service_db():
@build_platform_option
def put(env_id: str, metadata_pairs: tuple[str, ...], build_platform: str):
"""Build and upload a ServiceDB environment."""
if not env_id:
env_id = get_config().default_service_db_env_id

click.echo(f"Building ServiceDB image from {SERVICE_DB_DOCKERFILE}...")

# Build Docker image
build_image(SERVICE_DB_DOCKERFILE, SERVICE_DB_DOCKERFILE.parent, SERVICE_DB_IMAGE_NAME, platform=build_platform)
click.echo("Docker build successful")

# Create DB DockerImageArtifact
click.echo("Creating DB DockerImageArtifact...")
db_artifact = DockerImageArtifact.put(
id=f"service-db-{env_id}",
description=f"ServiceDB PostgreSQL image for {env_id}",
image_name=SERVICE_DB_IMAGE_NAME,
)
click.echo(f"Created DB DockerImageArtifact: id={db_artifact.id} version={db_artifact.version}")

# Build db-web image (build instead of pull to avoid docker save manifest issues on Apple Silicon)
click.echo(f"Building db-web image from {DB_WEB_DOCKERFILE}...")
build_image(DB_WEB_DOCKERFILE, DB_WEB_DOCKERFILE.parent, DB_WEB_IMAGE_NAME, platform=build_platform)
click.echo("db-web build successful")

# Create db-web DockerImageArtifact
click.echo("Creating db-web DockerImageArtifact...")
db_web_artifact = DockerImageArtifact.put(
id=f"db-web-{env_id}",
description="db-web lightweight web UI for database inspection",
image_name=DB_WEB_IMAGE_NAME,
)
click.echo(f"Created db-web DockerImageArtifact: id={db_web_artifact.id} version={db_web_artifact.version}")

# Build db-mcp image (PostgreSQL MCP server for direct DB access)
click.echo(f"Building db-mcp image from {DB_MCP_DOCKERFILE}...")
build_image(DB_MCP_DOCKERFILE, DB_MCP_DOCKERFILE.parent, DB_MCP_IMAGE_NAME, platform=build_platform)
click.echo("db-mcp build successful")

# Create db-mcp DockerImageArtifact
click.echo("Creating db-mcp DockerImageArtifact...")
db_mcp_artifact = DockerImageArtifact.put(
id=f"db-mcp-{env_id}",
description="db-mcp PostgreSQL MCP server for direct DB access",
image_name=DB_MCP_IMAGE_NAME,
)
click.echo(f"Created db-mcp DockerImageArtifact: id={db_mcp_artifact.id} version={db_mcp_artifact.version}")

user_metadata = {}
metadata = {}
for pair in metadata_pairs:
if "=" not in pair:
click.echo(f"Invalid metadata format '{pair}', expected key=value", err=True)
raise click.Abort()
key, value = pair.split("=", 1)
user_metadata[key] = value
metadata = detect_env_metadata(SERVICE_DB_DOCKERFILE, SERVICE_DB_DOCKERFILE.parent)
metadata.update(user_metadata)

# Create ServiceDBEnv
click.echo("Creating ServiceDBEnv...")
env = ServiceDBEnv.put(
id=env_id,
db_docker_image_artifact=db_artifact,
db_web_docker_image_artifact=db_web_artifact,
db_mcp_docker_image_artifact=db_mcp_artifact,
metadata=metadata if metadata else None,
)
click.echo(f"Created ServiceDBEnv: id={env.id} version={env.version}")
metadata[key] = value
put_service_db_env(env_id or get_config().default_service_db_env_id, platform=build_platform, metadata=metadata,
say=click.echo)
48 changes: 6 additions & 42 deletions src/agent_env/cli/env/website_browser.py
Original file line number Diff line number Diff line change
@@ -1,21 +1,9 @@
import sys
from pathlib import Path

import click

from agent_env.artifact import DockerImageArtifact
from agent_env.cli.utils import build_platform_option, detect_env_metadata
from agent_env.env import MCPServerEnv
from agent_env.env.envs.website_browser import (
PLAYWRIGHT_MCP_VERSION,
WEBSITE_BROWSER_IMAGE_TAG,
WEBSITE_BROWSER_ENVIRONMENT_NAME,
)
from agent_env.utils.docker_build import build_image

_PACKAGE_ROOT = Path(__file__).parent.parent.parent
WEBSITE_BROWSER_DOCKERFILE = _PACKAGE_ROOT / "env" / "envs" / "website_browser" / "Dockerfile"
WEBSITE_BROWSER_CONTEXT = _PACKAGE_ROOT / "env" / "envs" / "website_browser"
from agent_env.cli.utils import build_platform_option
from agent_env.env.bootstrap import WEBSITE_BROWSER_CONTEXT, WEBSITE_BROWSER_DOCKERFILE, put_website_browser_env # noqa: F401


@click.group(name="website-browser")
Expand All @@ -32,36 +20,12 @@ def put(env_id: str, metadata_pairs: tuple[str, ...], build_platform: str):
"""Build and upload a website browser environment."""
from agent_env.config import get_config

if not env_id:
env_id = get_config().default_website_browser_env_id

click.echo("Building website browser Docker image...")
build_image(WEBSITE_BROWSER_DOCKERFILE, WEBSITE_BROWSER_CONTEXT, WEBSITE_BROWSER_IMAGE_TAG, platform=build_platform,
build_args={"PLAYWRIGHT_MCP_VERSION": PLAYWRIGHT_MCP_VERSION})

click.echo("Creating DockerImageArtifact...")
artifact = DockerImageArtifact.put(
id=f"website-browser-{env_id}",
description="Website browser MCP server",
image_name=WEBSITE_BROWSER_IMAGE_TAG,
)
click.echo(f"Created artifact: id={artifact.id} version={artifact.version}")

user_metadata = {}
metadata = {}
for pair in metadata_pairs:
if "=" not in pair:
click.echo(f"Invalid metadata format '{pair}', expected key=value", err=True)
sys.exit(1)
key, value = pair.split("=", 1)
user_metadata[key] = value
metadata = detect_env_metadata(WEBSITE_BROWSER_DOCKERFILE, WEBSITE_BROWSER_CONTEXT)
metadata.update(user_metadata)

click.echo("Creating MCPServerEnv...")
env = MCPServerEnv.put(
id=env_id,
docker_image_artifact=artifact,
environment_name=WEBSITE_BROWSER_ENVIRONMENT_NAME,
metadata=metadata if metadata else None,
)
click.echo(f"Created MCPServerEnv: id={env.id} version={env.version} environment_name={env.environment_name}")
metadata[key] = value
put_website_browser_env(env_id or get_config().default_website_browser_env_id, platform=build_platform,
metadata=metadata, say=click.echo)
Loading
Loading