Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
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: 15 additions & 21 deletions src/agent_env/artifact/artifacts/docker_image.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,8 +86,12 @@ def put(

store = get_artifact_store()
version = store.next_version(id)
# Objects are written once, so each attempt writes under a prefix of its own: one that stopped
# before its document was written doesn't block the next.
prefix = store.attempt_prefix("docker_image", id)

config = get_config()
objects = config.get_object_store_to_write(prefix, id)
image_store = config.get_image_store_for(id)
repository = image_repository(id)
image_ref = image_store.image_ref(repository, f"v{version}")
Expand Down Expand Up @@ -115,19 +119,14 @@ def put(
stderr = save_proc.stderr.read().decode() if save_proc.stderr else ""
raise RuntimeError(f"docker save {image_ref} failed: {stderr}")

tar_gz_s3_url = store.put_object_file(
artifact_type="docker_image",
id=id,
version=version,
object_name=f"{fs_safe(id)}-v{version}.tar.gz",
file_path=str(tmp_path),
content_type="application/gzip",
tar_gz_object_url = objects.put_file_at(
f"{prefix}{fs_safe(id)}-v{version}.tar.gz", str(tmp_path), "application/gzip"
)
finally:
if tmp_path.exists():
tmp_path.unlink()

build_context_s3_url = None
build_context_object_url = None
if build_context_path:
context_dir = Path(build_context_path)
paths_to_include = _get_dockerfile_copy_sources(context_dir, dockerfile_path)
Expand All @@ -139,13 +138,8 @@ def put(
full_path = context_dir / rel_path
if full_path.exists():
tar.add(str(full_path), arcname=rel_path)
build_context_s3_url = store.put_object_file(
artifact_type="docker_image",
id=id,
version=version,
object_name="build-context.tar.gz",
file_path=str(ctx_tmp_path),
content_type="application/gzip",
build_context_object_url = objects.put_file_at(
f"{prefix}build-context.tar.gz", str(ctx_tmp_path), "application/gzip"
)
finally:
if ctx_tmp_path.exists():
Expand All @@ -155,8 +149,8 @@ def put(
id,
description=description,
image_name=image_ref,
tar_gz_s3_url=tar_gz_s3_url,
build_context_s3_url=build_context_s3_url,
tar_gz_s3_url=tar_gz_object_url,
build_context_s3_url=build_context_object_url,
)

@classmethod
Expand Down Expand Up @@ -297,7 +291,7 @@ def _signed_put(key: str) -> tuple[str, str]:
)
return url, put

tar_gz_s3_url, image_put_url = await asyncio.to_thread(_signed_put, f"github-builds/{id}/{image_tag}.tar.gz")
tar_gz_object_url, image_put_url = await asyncio.to_thread(_signed_put, f"github-builds/{id}/{image_tag}.tar.gz")
await sandbox.exec_script(f'curl -fsSL -X PUT --upload-file /tmp/image.tar.gz "{image_put_url}"')

log("upload_context", "Uploading build context...", 80)
Expand All @@ -306,7 +300,7 @@ def _signed_put(key: str) -> tuple[str, str]:
copy_sources = _parse_copy_sources(dockerfile_content, df_rel)
tar_paths = " ".join(shlex.quote(p) for p in copy_sources)
await sandbox.exec_script(f"tar czf /tmp/build-context.tar.gz -C {context_abs} {tar_paths}")
build_context_s3_url, context_put_url = await asyncio.to_thread(
build_context_object_url, context_put_url = await asyncio.to_thread(
_signed_put, f"github-builds/{id}/{image_tag}-context.tar.gz"
)
await sandbox.exec_script(f'curl -fsSL -X PUT --upload-file /tmp/build-context.tar.gz "{context_put_url}"')
Expand All @@ -319,8 +313,8 @@ def _signed_put(key: str) -> tuple[str, str]:
id=id,
description=f"Built from GitHub: {dockerfile_github_url}",
image_name=image_ref,
tar_gz_s3_url=tar_gz_s3_url,
build_context_s3_url=build_context_s3_url,
tar_gz_s3_url=tar_gz_object_url,
build_context_s3_url=build_context_object_url,
)
logger.info(f"put_from_github: created artifact id={artifact.id} version={artifact.version}")

Expand Down
10 changes: 10 additions & 0 deletions src/agent_env/bundle/authoring.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,16 @@ def entry_files(bundle: Bundle, entry: BundleEntry) -> dict[str, Path]:
return _Walk(bundle, entry).files()


def build_context_files(bundle: Bundle, entry: BundleEntry) -> dict[str, Path]:
"""What an image built from an entry's folder is built from: its ``entry_files`` and its own toml, which
the Dockerfile can copy too."""
files = entry_files(bundle, entry)
toml = entry.path / CONFIG_FILES[entry.kind]
if toml.is_file():
files[toml.name] = toml
return dict(sorted(files.items()))


def entry_file(bundle: Bundle, entry: BundleEntry) -> tuple[str, Path]:
"""The one file in an entry's folder, as ``(name, path)``. Raises BundleError otherwise."""
files = entry_files(bundle, entry)
Expand Down
21 changes: 14 additions & 7 deletions src/agent_env/bundle/ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@
from agent_env.store.local_state import ensure_state_dir
from agent_env.task.store import TASKS_COLLECTION

from .authoring import entry_files
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
Expand Down Expand Up @@ -118,10 +118,12 @@ def digest(self, write: Write, needs: Mapping[tuple[str, str], int]) -> Digest |
tell it all, so it is written every run."""
if not _tracked(write):
return None
inputs = {"type": _type(write), "config": _sha256(_canonical(write.source.config)),
config = {"dockerfile": write.source.dockerfile} if isinstance(write.source, BuiltImage) else write.source.config
inputs = {"type": _type(write), "config": _sha256(_canonical(config)),
"files": {}, "needs": {}, "store_refs": {}}
if write.kind is BundleKind.ARTIFACT:
for key, path in entry_files(self._plan.bundle.bundle, write.source.entry).items():
files = build_context_files if isinstance(write.source, BuiltImage) else entry_files
for key, path in files(self._plan.bundle.bundle, write.source.entry).items():
inputs["files"][key] = _file_sha256(path)
elif write.kind in (BundleKind.ENV, BundleKind.AGENT):
# An env's or agent's document records the versions of what it references, so one written anew, or
Expand Down Expand Up @@ -199,10 +201,13 @@ def _lock_path(id: str) -> Path:


def _tracked(write: Write) -> bool:
"""Whether the ledger can list everything ``write`` is made from. Not yet for a built image or a skill,
nor for a type with a ``from_toml`` of its own, which may read its folder in ways it can't see. An agent
is made from its agent.toml alone: its image is a reference."""
if isinstance(write.source, BuiltImage) or write.kind is BundleKind.SKILL:
"""Whether the ledger can list everything ``write`` is made from. Not yet for a skill, nor for a type with
a ``from_toml`` of its own, which may read its folder in ways it can't see. An agent is made from its
agent.toml alone: its image is a reference. A built image is made from its build context,
``build_context_files``; what the build fetches (its base image, packages) isn't an input."""
if isinstance(write.source, BuiltImage):
return True
if write.kind is BundleKind.SKILL:
return False
if write.kind is BundleKind.ARTIFACT:
return folder_walk(get_artifact_registry().get(_type(write))) is not None
Expand All @@ -213,6 +218,8 @@ def _tracked(write: Write) -> bool:


def _type(write: Write) -> str:
if isinstance(write.source, BuiltImage):
return "docker_image"
type_ = write.source.entry.type
return canonical_type(type_) if write.kind is BundleKind.ARTIFACT else type_

Expand Down
75 changes: 68 additions & 7 deletions src/agent_env/bundle/materialize.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@

Materializing first refuses every write this release has no writer for, so nothing is written for a bundle
that can't be written whole. It needs the CLI's namespace routing, which sends ``@local`` writes to the
``@local`` namespace's store. Holding the bundle's lock, it writes the entities in the plan's order, then
``@local`` namespace's store. Holding the bundle's lock, it writes the entities in the plan's order, an image
built from an entry's Dockerfile just before the entry, with ``docker build`` on this machine; then it
builds and preflights every task before writing any of them, and writes the evals last, since they name the
tasks. Each write goes through the ledger, so one whose inputs haven't changed reuses the version the bundle
last wrote. A reused task keeps whatever its steps took from config when it was first written, such as a
Expand All @@ -13,22 +14,28 @@
from __future__ import annotations

import copy
import shutil
import tempfile
from collections.abc import Callable, Iterator
from contextlib import contextmanager, nullcontext
from dataclasses import dataclass
from pathlib import Path
from typing import Any

from agent_env.a2a_agent import A2AAgent
from agent_env.artifact.artifacts.docker_image import DockerImageArtifact
from agent_env.artifact.registry import canonical_type, get_artifact_registry
from agent_env.entity_refs import EntityRef, RefRole, ref_sites
from agent_env.eval import Eval, EvalTask
from agent_env.store.ids import image_repository
from agent_env.store.routing import namespace_routing_enabled
from agent_env.task import Task
from agent_env.task_step.registry import get_task_step_registry
from agent_env.task_step.task_step import TaskStep
from agent_env.utils.docker_build import DockerBuildError, build_image

from ._fs import relative, with_article
from .authoring import AuthoringContext
from .authoring import AuthoringContext, build_context_files
from .ledger import Ledger, materializing
from .parse import BundleError, BundleKind
from .plan import Plan, Write, folder_walk
Expand Down Expand Up @@ -66,11 +73,13 @@ def materialize(
*,
dry_run: bool = False,
on_wait: Callable[[], None] | None = None,
on_build: Callable[[Write], None] | None = None,
on_write: Callable[[Materialized], None] | None = None,
) -> Materialization:
"""Write ``plan``'s entities, tasks and evals. The first write that fails stops it, and every earlier
one stays: they're in the ledger, so the next run reuses them. ``on_wait`` is called when another run
holds a lock this one needs, and ``on_write`` after each write, reused or not.
holds a lock this one needs, ``on_build`` before an image is built, and ``on_write`` after each write,
reused or not.

``dry_run`` checks and preflights what a run would, and writes nothing: no entity, ledger row or lock.
Each write gets the version it would reuse or the store's next, which another run can take first. A step
Expand Down Expand Up @@ -99,8 +108,9 @@ def through_ledger(write: Write, write_fn: Callable[[], int]) -> None:
on_write(done[_key(write)])

with nullcontext() if dry_run else materializing(plan.bundle.bundle, on_wait):
_refuse_builds_without_docker(plan, ledger) # once another run writing these ids is done
for write in entities:
through_ledger(write, lambda: _WRITERS[write.kind](plan, write))
through_ledger(write, lambda: _write_entity(plan, write, on_build))
unwritten = {_key(item.write) for item in done.values() if not item.reused} if dry_run else set()
built, problems, unchecked = {}, [], []
for write in tasks:
Expand All @@ -118,6 +128,48 @@ def through_ledger(write: Write, write_fn: Callable[[], int]) -> None:
return Materialization(plan, tuple(done[_key(write)] for write in plan.writes), tuple(unchecked))


def _write_entity(plan: Plan, write: Write, on_build: Callable[[Write], None] | None) -> int:
if not isinstance(write.source, BuiltImage):
return _WRITERS[write.kind](plan, write)
if on_build is not None:
on_build(write)
return _write_built_image(plan, write)


def _write_built_image(plan: Plan, write: Write) -> int:
"""Build the image an entry's Dockerfile describes and write it as a docker_image artifact: pushed to the
image store, saved as a tarball, and its build context kept for installing it into a running container.
The build context is a copy of the files the ledger hashes, ``build_context_files``, so an image the
ledger reuses was built from what it hashed."""
image = write.source
tag = f"{image_repository(write.id)}:bundle"
with tempfile.TemporaryDirectory(prefix="agent-env-build-") as staged:
context = Path(staged)
for key, path in build_context_files(plan.bundle.bundle, image.entry).items():
(context / key).parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(path, context / key)
Comment thread
earakely-scale marked this conversation as resolved.
try:
build_image(context / image.dockerfile, context, tag, platform=None)
except DockerBuildError as e:
raise BundleError([f"{_path(plan, write)}: {_tail(str(e))}"]) from None
try:
return DockerImageArtifact.put(
id=write.id, description=f"built from {_path(plan, write)}/{image.dockerfile}", image_name=tag,
build_context_path=str(context),
).version
except RuntimeError as e: # the image store, docker push or docker save
raise BundleError([f"{_path(plan, write)}: {e}"]) from None


_BUILD_OUTPUT_TAIL_LINES = 40


def _tail(message: str) -> str:
"""A failed build's message, its first line and the end of docker's output."""
first, _, output = message.partition("\n")
return "\n".join([first, *output.splitlines()[-_BUILD_OUTPUT_TAIL_LINES:]])


def _write_artifact(plan: Plan, write: Write) -> int:
entry = write.source.entry
cls = get_artifact_registry()[canonical_type(entry.type)]
Expand Down Expand Up @@ -166,11 +218,20 @@ def _refuse_unwritable(plan: Plan) -> None:
raise BundleError(problems)


def _refuse_builds_without_docker(plan: Plan, ledger: Ledger) -> None:
"""An image the ledger will reuse needs no docker, so only the ones it would build are refused."""
if shutil.which("docker") is not None:
return
problems = [f"{_path(plan, write)}: building its image from {write.source.dockerfile} needs docker, and it isn't "
"on PATH" for write in plan.writes
if isinstance(write.source, BuiltImage) and not ledger.check(write, {}).unchanged]
Comment thread
earakely-scale marked this conversation as resolved.
if problems:
raise BundleError(problems)


def _unwritable(write: Write) -> str | None:
"""What ``write`` would write, when this release has no writer for it."""
if isinstance(write.source, BuiltImage):
return f"an image built from {write.source.dockerfile}"
if write.kind is BundleKind.TASK:
if isinstance(write.source, BuiltImage) or write.kind is BundleKind.TASK:
return None
if write.kind not in _WRITERS:
return with_article(write.kind.value.removesuffix("s"))
Expand Down
4 changes: 2 additions & 2 deletions src/agent_env/bundle/resolve.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
from agent_env.env.registry import get_env_registry
from agent_env.eval.eval import Eval
from agent_env.plugins import _registration
from agent_env.store.ids import LOCAL_PREFIX, validate_local_id
from agent_env.store.ids import LOCAL_PREFIX, derive_id, validate_local_id
from agent_env.task_step.registry import get_task_step_registry
from agent_env.task_step.task_step import TaskStep, attach_retry_config, dependencies

Expand Down Expand Up @@ -230,7 +230,7 @@ def _image(self, entry: BundleEntry, config: dict, ref: EntityRef, references: l
self._problem(entry, f"{ref.path}: there is no {dockerfile!r} in this folder to build")
return True
role = f"{entry.kind.store}_image" if ref.path == "image" else ref.path
image = BuiltImage(f"{entry.id}__{role}", entry, dockerfile)
image = BuiltImage(derive_id(entry.id, role), entry, dockerfile)
config[ref.path] = image.id
self.built.append(image)
references.append(Reference(EntityKind.ARTIFACT, image.id, None, image, ref.path, ref.artifact_type))
Expand Down
13 changes: 10 additions & 3 deletions src/agent_env/bundle/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,8 @@
from ._fs import relative
from .materialize import Materialization, Materialized, materialize
from .parse import BundleEntry, BundleError, BundleKind, parse_bundle
from .plan import Plan, plan_bundle
from .resolve import Reference, resolve_bundle
from .plan import Plan, Write, plan_bundle
from .resolve import BuiltImage, Reference, resolve_bundle

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -202,6 +202,7 @@ def run_bundle(
materialization = materialize(
plan,
on_wait=lambda: say("waiting for another agent-env run to finish writing this bundle's ids"),
on_build=lambda write: say(f"{_label(plan, write)}: building with docker, which can take minutes"),
on_write=lambda done: say(_written(plan, done)),
)
entries = _to_run(plan)
Expand Down Expand Up @@ -409,9 +410,15 @@ def _task_run_command(ref: Reference) -> str:


def _written(plan: Plan, done: Materialized) -> str:
what = f"{_path(plan, done.write.source.entry)}: v{done.version}"
what = f"{_label(plan, done.write)}: v{done.version}"
return f"{what}, unchanged" if done.reused else f"{what} ({'; '.join(done.reasons)})"


def _label(plan: Plan, write: Write) -> str:
"""How a write is named: its entry's folder, and for an image built from it, which image."""
path = _path(plan, write.source.entry)
return f"{path} ({write.source.dockerfile} image)" if isinstance(write.source, BuiltImage) else path


def _path(plan: Plan, entry: BundleEntry) -> str:
return relative(plan.bundle.bundle.root, entry.path)
Loading
Loading