diff --git a/src/agent_env/artifact/artifacts/docker_image.py b/src/agent_env/artifact/artifacts/docker_image.py index fc86794a..f2bb5e0d 100644 --- a/src/agent_env/artifact/artifacts/docker_image.py +++ b/src/agent_env/artifact/artifacts/docker_image.py @@ -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}") @@ -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) @@ -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(): @@ -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 @@ -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) @@ -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}"') @@ -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}") diff --git a/src/agent_env/bundle/authoring.py b/src/agent_env/bundle/authoring.py index 8cd43299..521953ee 100644 --- a/src/agent_env/bundle/authoring.py +++ b/src/agent_env/bundle/authoring.py @@ -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) diff --git a/src/agent_env/bundle/ledger.py b/src/agent_env/bundle/ledger.py index aeb1de3f..570dc4fa 100644 --- a/src/agent_env/bundle/ledger.py +++ b/src/agent_env/bundle/ledger.py @@ -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 @@ -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 @@ -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 @@ -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_ diff --git a/src/agent_env/bundle/materialize.py b/src/agent_env/bundle/materialize.py index d13e1481..c4fefd60 100644 --- a/src/agent_env/bundle/materialize.py +++ b/src/agent_env/bundle/materialize.py @@ -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 @@ -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 @@ -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 @@ -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: @@ -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) + 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)] @@ -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] + 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")) diff --git a/src/agent_env/bundle/resolve.py b/src/agent_env/bundle/resolve.py index db9ffe1d..21f96bac 100644 --- a/src/agent_env/bundle/resolve.py +++ b/src/agent_env/bundle/resolve.py @@ -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 @@ -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)) diff --git a/src/agent_env/bundle/run.py b/src/agent_env/bundle/run.py index 1f533748..e87e4b55 100644 --- a/src/agent_env/bundle/run.py +++ b/src/agent_env/bundle/run.py @@ -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__) @@ -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) @@ -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) diff --git a/tst/integration/cli/run_bundle_built_agents_test.py b/tst/integration/cli/run_bundle_built_agents_test.py new file mode 100644 index 00000000..78e26674 --- /dev/null +++ b/tst/integration/cli/run_bundle_built_agents_test.py @@ -0,0 +1,204 @@ +"""``agent-env run`` on a bundle whose agents are built from their folders' Dockerfiles, with real docker builds and +the agents deployed on the local sandbox: the shapes a Dockerfile can take in an agent folder, a rerun that builds +nothing, an edit that rebuilds only its own agent, and a build that fails. Needs a Docker daemon and the local registry.""" + +import json +import logging +import shutil + +import pytest +from click.testing import CliRunner + +from agent_env.artifact.store import reset_artifact_store +from agent_env.cli import cli +from agent_env.config import configure, reset_config +from agent_env.store.object_store.local.grant_server import grant_server +from tst.util.a2a_test_agent import AGENT_DIR, PROTOCOL_DIR + +pytestmark = [pytest.mark.integration, pytest.mark.int_test_slow] + +BASE = (AGENT_DIR / "Dockerfile").read_text().splitlines()[0] # the echo agent's pinned base image +NO_MODEL = {"LITELLM_API_KEY": "unused", "LITELLM_BASE_URL": "http://unused.invalid"} + +# The echo agent, replying with what its build put in it instead of an echo: a greeting (the GREETING env var, else +# greeting.txt beside it) and which of the files a build might leave out it can see. +AGENT_PY = (AGENT_DIR / "agent.py").read_text().replace( + 'reply = f"Echo: {prompt}"', + 'greeting = os.environ.get("GREETING") or (HERE / "greeting.txt").read_text().strip()\n' + ' seen = [name for name in ("notes.md", "secret.txt", "__pycache__") if (HERE / name).exists()]\n' + ' reply = f"{greeting} sees {seen}"', +).replace("from typing import Any\n", "import os\nfrom pathlib import Path\nfrom typing import Any\n\nHERE = Path(__file__).parent\n") +assert "HERE = Path" in AGENT_PY and "sees {seen}" in AGENT_PY + +# The steps every agent's image starts with, so the build cache shares them. +PREFIX = f'''{BASE} +WORKDIR /app +COPY agentenv-protocol/ ./agentenv-protocol/ +RUN pip install --no-cache-dir './agentenv-protocol[agent]' +''' + + +class Link(str): + """A symlink to this path, relative to the link.""" + + +def _toml(**fields): + lines = [f"{key} = {value}" for key, value in fields.items() if key != "env"] + env = {**NO_MODEL, **fields.get("env", {})} + return "\n".join([*lines, "[default_env_vars]", *(f'{key} = "{value}"' for key, value in env.items())]) + "\n" + + +AGENTS = { + # A Dockerfile at the folder's root and no agent.toml: the deploy step passes the model variables. + "plain": { + "Dockerfile": PREFIX + 'COPY agent.py greeting.txt ./\nCMD ["python", "agent.py"]\n', + "greeting.txt": "plain-v1\n", + }, + # `COPY .`: the whole folder, less what .dockerignore names. A link is copied as its target, and what the OS or + # Python leaves behind stays out of the build. + "whole": { + "Dockerfile": PREFIX + 'COPY . ./\nCMD ["python", "agent.py"]\n', + ".dockerignore": "secret.txt\n", + "data/greeting.txt": "whole-v1\n", + "greeting.txt": Link("data/greeting.txt"), + "notes.md": "kept\n", + "secret.txt": "left out by .dockerignore\n", + "__pycache__/agent.cpython-312.pyc": "left out of the build\n", + ".DS_Store": "left out of the build\n", + "agent.toml": _toml(), + }, + # A Dockerfile in a subfolder, named by agent.toml; the build context is still the agent's folder. agent.toml's + # env vars reach the container. + "subdir": { + "docker/Dockerfile": PREFIX + 'COPY agent.py ./\nCMD ["python", "agent.py"]\n', + "agent.toml": _toml(image='{ dockerfile = "docker/Dockerfile" }', env={"GREETING": "subdir-from-toml"}), + }, + # A multi-stage build under another name. + "staged": { + "Dockerfile.agent": f'{BASE} AS greeting\nRUN echo staged-v1 > /greeting.txt\n\n' + PREFIX + + 'COPY agent.py ./\nCOPY --from=greeting /greeting.txt ./greeting.txt\n' + + 'CMD ["python", "agent.py"]\n', + "agent.toml": _toml(image='{ dockerfile = "Dockerfile.agent" }'), + }, +} + +# How the run names each agent's image: its folder, and the Dockerfile it's built from. +IMAGES = {"plain": "agents/plain (Dockerfile image)", "whole": "agents/whole (Dockerfile image)", + "subdir": "agents/subdir (docker/Dockerfile image)", "staged": "agents/staged (Dockerfile.agent image)"} + +REPLIES = { + "plain": "plain-v1 sees []", + "whole": "whole-v1 sees ['notes.md']", + "subdir": "subdir-from-toml sees []", + "staged": "staged-v1 sees []", +} + + +def _task(name, reply): + deploy = {"id": "deploy", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": name, "sandbox_type": "local"} + if name == "plain": + deploy["env_vars"] = NO_MODEL + return json.dumps([ + deploy, + {"id": "ask", "type": "prompt_agent", "prompt": "what did your build put in you?", "prompt_id": "ask", + "depends_on": ["deploy"]}, + {"id": "check", "type": "agent_prompt_response_verifier", "prompt_id": "ask", "verifier_id": "reply", + "criteria": [{"type": "response_contains", "needles": [reply]}], "depends_on": ["ask"]}, + ]) + + +def _agent_folder(root, name, files): + folder = root / "agents" / name + for path, text in {"agent.py": AGENT_PY, **files}.items(): + (folder / path).parent.mkdir(parents=True, exist_ok=True) + if isinstance(text, Link): + (folder / path).symlink_to(text) + else: + (folder / path).write_text(text) + shutil.copytree(PROTOCOL_DIR, folder / "agentenv-protocol", + ignore=shutil.ignore_patterns("__pycache__", "*.pyc", "tests", "examples", "*.egg-info")) + return folder + + +@pytest.fixture +def state(monkeypatch, tmp_path): + """Local stores under this test's folder. Whatever ran, no sandbox work folder may be left behind.""" + sandboxes = tmp_path / "sandboxes" + # HOME stays: docker's credential helper (the macOS keychain, for one) can hang a build's base-image lookup + # when HOME moves. + monkeypatch.chdir(tmp_path) + monkeypatch.setenv("XDG_STATE_HOME", str(tmp_path / "state")) + monkeypatch.setenv("AGENT_ENV_DOCUMENT_STORE", "local") + monkeypatch.setenv("AGENT_ENV_OBJECT_STORE", "local") + monkeypatch.setenv("AGENT_ENV_LOCAL_SANDBOX_DIR", str(sandboxes)) + configure() + reset_artifact_store() + logging.disable(logging.CRITICAL) # pytest's live logging would take CliRunner's stdout + try: + yield tmp_path + assert not sandboxes.exists() or not any(sandboxes.iterdir()), "a run left a sandbox work folder" + finally: + logging.disable(logging.NOTSET) + # The agents moved their trajectories through this process's grant server, which serves a certificate from + # this test's state root; a later test's agents trust another root's CA, so it must start afresh. + grant_server(None, None).close() + reset_artifact_store() + reset_config() + + +def _run(root): + return CliRunner().invoke(cli, ["run", str(root)]) + + +def test_each_shape_of_agent_dockerfile_builds_runs_and_is_rebuilt_only_when_its_folder_changes(state): + root = state / "agents-bundle" + for name, files in AGENTS.items(): + _agent_folder(root, name, files) + (root / "tasks").mkdir(parents=True, exist_ok=True) + (root / f"tasks/{name}.json").write_text(_task(name, REPLIES[name])) + + first = _run(root) + + assert first.exit_code == 0, first.output + for name in AGENTS: + assert f"{IMAGES[name]}: building with docker" in first.output, first.output + assert f"{IMAGES[name]}: v1 (new)" in first.output, first.output + assert f"tasks/{name}.json v1: passed" in first.output, first.output + + again = _run(root) + + assert again.exit_code == 0, again.output + assert "building with docker" not in again.output, again.output + for name in AGENTS: + assert f"{IMAGES[name]}: v1, unchanged" in again.output, again.output + assert f"tasks/{name}.json v1: passed" in again.output, again.output + + (root / "agents/plain/greeting.txt").write_text("plain-v2\n") + (root / "tasks/plain.json").write_text(_task("plain", "plain-v2 sees []")) + (root / "agents/whole/__pycache__/agent.cpython-312.pyc").write_text("recompiled\n") + (root / "agents/whole/.DS_Store").write_text("moved\n") + + edited = _run(root) + + assert edited.exit_code == 0, edited.output + assert f"{IMAGES['plain']}: v2 (files changed: greeting.txt)" in edited.output, edited.output + assert "agents/plain: v2 (" in edited.output, edited.output + for name in ("whole", "subdir", "staged"): + assert f"{IMAGES[name]}: v1, unchanged" in edited.output, edited.output + assert "tasks/plain.json v2: passed" in edited.output, edited.output + assert edited.output.count("building with docker") == 1, edited.output + + +def test_a_dockerfile_that_fails_to_build_is_one_problem_with_dockers_output_and_writes_no_agent(state): + root = state / "broken-bundle" + _agent_folder(root, "broken", {"Dockerfile": f"{BASE}\nRUN echo about to fail && exit 3\n"}) + (root / "tasks").mkdir() + (root / "tasks/broken.json").write_text(_task("broken", "never")) + + result = _run(root) + + assert result.exit_code == 1, result.output + assert "Error: agents/broken: docker build of " in result.output, result.output + assert "about to fail" in result.output, result.output + assert "Traceback" not in result.output, result.output + assert "agents/broken: v1" not in result.output, result.output diff --git a/tst/unit/bundle/ledger_test.py b/tst/unit/bundle/ledger_test.py index 7ca7b332..fa647abe 100644 --- a/tst/unit/bundle/ledger_test.py +++ b/tst/unit/bundle/ledger_test.py @@ -409,14 +409,9 @@ def test_the_ledger_never_touches_the_configured_store(bundle_dir, cli_routing): ({"artifacts/greeting/artifact.toml": 'type = "own_file_ledger_test"\n'}, None, GREETING), ({"envs/own/env.toml": 'type = "own_env_ledger_test"\n'}, {"id": "own", "type": "deploy_env", "env_id": "own"}, f"{ROOT}/own"), - ({"envs/imaged/env.toml": 'type = "imaged_ledger_test"\n', "envs/imaged/Dockerfile": "FROM scratch\n"}, - {"id": "imaged", "type": "deploy_env", "env_id": "imaged"}, f"{ROOT}/imaged__env_image"), ({"skills/pdf/SKILL.md": "---\nname: pdf\n---\n"}, {"id": "pdf", "type": "load_artifact", "env_id": "tickets", "artifact_id": "pdf"}, f"{ROOT}/pdf"), - ({"agents/solver/Dockerfile": "FROM scratch\n"}, - {"id": "agent", "type": "deploy_agent", "env_ids": ["tickets"], "a2a_agent_id": "solver"}, - f"{ROOT}/solver__agent_image"), -], ids=["artifact-with-own-from_toml", "env-with-own-from_toml", "built-image", "skill", "agent-image"]) +], ids=["artifact-with-own-from_toml", "env-with-own-from_toml", "skill"]) def test_a_write_whose_inputs_arent_tracked_is_written_every_run(bundle_dir, files, step, id): for rel, text in files.items(): (bundle_dir / rel).parent.mkdir(parents=True, exist_ok=True) @@ -428,6 +423,26 @@ def test_a_write_whose_inputs_arent_tracked_is_written_every_run(bundle_dir, fil assert _run(bundle_dir)[id].reasons == UNTRACKED +def test_a_built_image_is_made_from_every_file_of_its_folder_the_toml_too(bundle_dir): + (bundle_dir / "agents/solver").mkdir(parents=True) + (bundle_dir / "agents/solver/Dockerfile").write_text("FROM scratch\nCOPY run.sh /\n") + (bundle_dir / "agents/solver/run.sh").write_text("echo hi\n") + (bundle_dir / "tasks/t.json").write_text( + _steps({"id": "agent", "type": "deploy_agent", "env_ids": ["tickets"], "a2a_agent_id": "solver"})) + image, agent = f"{ROOT}/solver__agent_image", f"{ROOT}/solver" + + assert _run(bundle_dir)[image].reasons == ("new",) + assert _run(bundle_dir)[image].unchanged + (bundle_dir / "agents/solver/agent.toml").write_text('[metadata]\ndefault_model = "m"\n') + retoml = _run(bundle_dir) + (bundle_dir / "agents/solver/run.sh").write_text("echo bye\n") + rebuilt = _run(bundle_dir) + + assert retoml[image].reasons == ("files added: agent.toml",) and not retoml[agent].unchanged + assert rebuilt[image].reasons == ("files changed: run.sh",) + assert not rebuilt[agent].unchanged + + @pytest.mark.parametrize("other", [ "the-same-folder", "another-folder-with-its-id-root", "another-bundle-declaring-its-id", ]) diff --git a/tst/unit/bundle/materialize_test.py b/tst/unit/bundle/materialize_test.py index 9480c042..6a082343 100644 --- a/tst/unit/bundle/materialize_test.py +++ b/tst/unit/bundle/materialize_test.py @@ -1,10 +1,12 @@ """Materializing a planned bundle: what gets written, what is reused, and what is refused before any write.""" +import contextlib import json import subprocess import sys import textwrap import time +from pathlib import Path from typing import Literal import pytest @@ -26,6 +28,7 @@ from agent_env.env.env import Env from agent_env.eval import Eval, EvalTask from agent_env.store import Filter +from agent_env.store.ids import image_repository from agent_env.store.routing import disable_namespace_routing from agent_env.task import Task from agent_env.task_step.task_step import TaskStep @@ -199,11 +202,9 @@ def test_what_has_no_writer_yet_is_refused_before_anything_is_written(bundle_dir ]) assert sorted(_problems(lambda: _run(bundle_dir, dry_run))) == [ - "agents/solver: writing an image built from Dockerfile isn't supported yet", "artifacts/base-mcp: writing a docker_image artifact isn't supported yet", "artifacts/snap: writing an environment artifact isn't supported yet", "envs/imaged: writing an env isn't supported yet", - "envs/imaged: writing an image built from Dockerfile isn't supported yet", "envs/tickets: writing an env isn't supported yet", "skills/pdf: writing a skill isn't supported yet", ] @@ -310,6 +311,198 @@ def test_another_bundles_entities_are_read_like_store_ids_and_an_agent_pins_the_ assert A2AAgent.get(f"{ROOT}/solver").docker_image_artifact.version == 2 +@pytest.fixture +def docker_on_path(monkeypatch): + monkeypatch.setattr(materialize_module.shutil, "which", lambda name: f"/usr/bin/{name}") + + +@pytest.fixture +def builds(monkeypatch, docker_on_path): + """Stands in for docker: each build is recorded with the files of its context, a link marked with a + trailing ``@``, and the image written as a docker_image document.""" + calls = [] + + def build(dockerfile, context, tag, *, platform): + calls.append({"build": (dockerfile.relative_to(context).as_posix(), _listing(context), tag, platform)}) + + def put(id, *, description, image_name, build_context_path=None, dockerfile_path=None): + calls[-1]["put"] = (id, image_name, _listing(build_context_path), dockerfile_path) + return get_artifact_store().put_document(DockerImageArtifact( + id=id, description=description, image_name=image_name, tar_gz_s3_url=f"file:///{id}.tar.gz")) + + monkeypatch.setattr(materialize_module, "build_image", build) + monkeypatch.setattr(materialize_module.DockerImageArtifact, "put", put) + return calls + + +def _listing(folder): + return sorted(path.relative_to(folder).as_posix() + ("@" if path.is_symlink() else "") + for path in Path(folder).rglob("*") if not path.is_dir() or path.is_symlink()) + + +def test_an_agent_folder_with_a_dockerfile_is_built_and_the_agent_written_over_its_image(bundle_dir, builds): + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\nCOPY run.sh /\n", "agents/solver/run.sh": "echo hi\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + image = f"{ROOT}/solver__agent_image" + announced = [] + + first = materialize(plan_of(bundle_dir), on_build=lambda write: announced.append(write.id)) + + tag = f"{image_repository(image)}:bundle" + assert builds == [{"build": ("Dockerfile", ["Dockerfile", "run.sh"], tag, None), + "put": (image, tag, ["Dockerfile", "run.sh"], None)}] + assert announced == [image] + assert {id: summary[:2] for id, summary in _summary(first).items() if "solver" in id} == { + image: (1, False), f"{ROOT}/solver": (1, False)} + agent = A2AAgent.get(f"{ROOT}/solver") + assert (agent.docker_image_artifact.id, agent.docker_image_artifact.version) == (image, 1) + + again = materialize(plan_of(bundle_dir), on_build=lambda write: announced.append(write.id)) + + assert len(builds) == 1 and announced == [image] + assert _summary(again)[image] == (1, True, ()) and _summary(again)[f"{ROOT}/solver"] == (1, True, ()) + + +def test_an_image_is_built_from_the_files_the_ledger_hashes_so_a_reused_one_matches_them(bundle_dir, builds): + layout(bundle_dir, { + "agents/solver/Dockerfile": "FROM scratch\nCOPY . /agent\n", + "agents/solver/agent.toml": '[metadata]\ndefault_model = "claude-sonnet-4-6"\n', + "agents/solver/.dockerignore": "notes.md\n", + "agents/solver/lib/run.sh": "echo hi\n", + "agents/solver/__pycache__/run.cpython-312.pyc": "compiled", + "agents/solver/.DS_Store": "finder", + }) + (bundle_dir / "agents/solver/run.sh").symlink_to("lib/run.sh") + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + + _run(bundle_dir) + (bundle_dir / "agents/solver/__pycache__/run.cpython-312.pyc").write_text("recompiled") + (bundle_dir / "agents/solver/.DS_Store").write_text("moved") + again = _summary(_run(bundle_dir)) + + context = [".dockerignore", "Dockerfile", "agent.toml", "lib/run.sh", "run.sh"] + assert [call["build"][1] for call in builds] == [context] + assert builds[0]["put"][2] == context + assert again[f"{ROOT}/solver__agent_image"] == (1, True, ()) + + +def test_a_changed_build_context_rebuilds_the_image_and_rewrites_its_agent(bundle_dir, builds): + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\nCOPY run.sh /\n", "agents/solver/run.sh": "echo hi\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + _run(bundle_dir) + (bundle_dir / "agents/solver/run.sh").write_text("echo bye\n") + + rebuilt = _summary(_run(bundle_dir)) + + assert len(builds) == 2 + assert rebuilt[f"{ROOT}/solver__agent_image"][:2] == (2, False) + assert rebuilt[f"{ROOT}/solver"][:2] == (2, False) + assert A2AAgent.get(f"{ROOT}/solver").docker_image_artifact.version == 2 + + +def test_a_dry_run_predicts_a_rebuild_and_its_agents_rewrite_without_building(bundle_dir, builds): + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\nCOPY run.sh /\n", "agents/solver/run.sh": "echo hi\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + _run(bundle_dir) + (bundle_dir / "agents/solver/run.sh").write_text("echo bye\n") + image = f"{ROOT}/solver__agent_image" + + predicted = _summary(_run(bundle_dir, dry_run=True)) + + assert len(builds) == 1 + assert predicted[image] == (2, False, ("files changed: run.sh",)) + assert predicted[f"{ROOT}/solver"] == (2, False, (f"artifact {image} is written anew (v1 → v2)",)) + assert _summary(_run(bundle_dir)) == predicted + assert len(builds) == 2 + + +def test_an_agent_toml_edit_rebuilds_the_image_too_since_a_dockerfile_can_copy_it(bundle_dir, builds): + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\nCOPY . /agent\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + _run(bundle_dir) + (bundle_dir / "agents/solver/agent.toml").write_text('[metadata]\ndefault_model = "claude-sonnet-4-6"\n') + + edited = _summary(_run(bundle_dir)) + + assert len(builds) == 2 + assert edited[f"{ROOT}/solver__agent_image"][:3] == (2, False, ("files added: agent.toml",)) + assert edited[f"{ROOT}/solver"][:2] == (2, False) + + +def test_a_failed_build_is_a_bundle_problem_naming_the_agent(bundle_dir, builds, monkeypatch): + def fail(dockerfile, context, tag, *, platform): + output = "\n".join(f"#{n} step" for n in range(60)) + raise materialize_module.DockerBuildError(f"docker build of {tag} failed (exit 1):\n{output}\nERROR: failed to solve") + + monkeypatch.setattr(materialize_module, "build_image", fail) + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + + (problem,) = _problems(lambda: _run(bundle_dir)) + lines = problem.splitlines() + assert lines[0].startswith("agents/solver: docker build of ") and lines[0].endswith(" failed (exit 1):") + assert lines[1:] == [*(f"#{n} step" for n in range(21, 60)), "ERROR: failed to solve"] + + +def test_a_failed_push_is_a_bundle_problem_naming_the_agent(bundle_dir, builds, monkeypatch): + def fail(id, **kwargs): + raise RuntimeError("could not start the local registry on port 5000: port is already allocated") + + monkeypatch.setattr(materialize_module.DockerImageArtifact, "put", fail) + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + + assert _problems(lambda: _run(bundle_dir)) == ( + "agents/solver: could not start the local registry on port 5000: port is already allocated",) + + +@_RUN_OR_DRY_RUN +def test_an_image_to_build_without_docker_on_path_is_refused_before_anything_is_written( + bundle_dir, monkeypatch, dry_run, +): + monkeypatch.setattr(materialize_module.shutil, "which", lambda name: None) + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + + assert _problems(lambda: _run(bundle_dir, dry_run)) == ( + "agents/solver: building its image from Dockerfile needs docker, and it isn't on PATH",) + assert not local_store().path.exists() + + +def test_whether_an_image_needs_docker_is_decided_once_another_run_writing_it_is_done(bundle_dir, builds, monkeypatch): + """Another run may be building the image; once its lock is released, the ledger can reuse what it built.""" + events = [] + locked = materialize_module.materializing + + @contextlib.contextmanager + def materializing(bundle, on_wait=None): + with locked(bundle, on_wait): + events.append("locked") + yield + + monkeypatch.setattr(materialize_module, "materializing", materializing) + check = materialize_module._refuse_builds_without_docker + monkeypatch.setattr(materialize_module, "_refuse_builds_without_docker", + lambda plan, ledger: events.append("docker checked") or check(plan, ledger)) + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + + _run(bundle_dir) + + assert events == ["locked", "docker checked"] + + +@_RUN_OR_DRY_RUN +def test_an_image_the_ledger_reuses_needs_no_docker(bundle_dir, builds, monkeypatch, dry_run): + layout(bundle_dir, {"agents/solver/Dockerfile": "FROM scratch\n"}) + _steps(bundle_dir, [{"id": "agent", "type": "deploy_agent", "env_ids": [], "a2a_agent_id": "solver"}]) + _run(bundle_dir) + monkeypatch.setattr(materialize_module.shutil, "which", lambda name: None) + + assert _summary(_run(bundle_dir, dry_run))[f"{ROOT}/solver__agent_image"] == (1, True, ()) + assert len(builds) == 1 + + def _put_image(id): get_artifact_store().put_document(DockerImageArtifact( id=id, description=id, image_name=f"{id}:v1", tar_gz_s3_url=f"file:///{id}.tar.gz")) diff --git a/tst/unit/bundle/run_test.py b/tst/unit/bundle/run_test.py index 64d2ae52..a3328b91 100644 --- a/tst/unit/bundle/run_test.py +++ b/tst/unit/bundle/run_test.py @@ -10,6 +10,7 @@ from click.testing import CliRunner import agent_env.bundle.run as run_module +from agent_env.bundle import materialize as materialize_module from agent_env.artifact.artifacts.docker_image import DockerImageArtifact from agent_env.artifact.artifacts.file import FileArtifact from agent_env.artifact.store import get_artifact_store @@ -19,6 +20,7 @@ from agent_env.bundle.resolve import resolve_bundle from agent_env.cli import cli from agent_env.config.runtime import Config +from agent_env.entity_refs import EntityRef from agent_env.store import Filter from agent_env.store.routing import namespace_routing from agent_env.task import Task @@ -77,6 +79,24 @@ async def execute(self, context): raise asyncio.CancelledError +class _NamesAgent(_Scored): + """Names an agent, so the agent and its image are written, without deploying it.""" + + type = "names_agent_run_test" + entity_refs = (EntityRef.agent("a2a_agent_id"),) + + def __init__(self, a2a_agent_id=None, **base): + super().__init__(**base) + self.a2a_agent_id = a2a_agent_id + + @classmethod + def from_dict(cls, data): + return cls(**cls._base_from_dict(data), a2a_agent_id=data.get("a2a_agent_id")) + + def to_dict(self): + return {**super().to_dict(), "a2a_agent_id": self.a2a_agent_id} + + def _scored(score): return json.dumps([{"id": "check", "type": "scored_run_test", "verifier_id": "v", "score": score}]) @@ -93,7 +113,7 @@ def _scored(score): @pytest.fixture(autouse=True) def registries(monkeypatch): - steps = {**Config().task_step_registry(), **{cls.type: cls for cls in (_Scored, _Failing, _Cancelled)}} + steps = {**Config().task_step_registry(), **{cls.type: cls for cls in (_Scored, _Failing, _Cancelled, _NamesAgent)}} monkeypatch.setattr(Config, "task_step_registry", lambda self: steps) monkeypatch.setattr(_Scored, "seen", []) monkeypatch.setattr(_Scored, "peak", 0) @@ -105,6 +125,25 @@ def bundle_dir(tmp_path, monkeypatch, local_stores): return layout(tmp_path / "triage", LAYOUT) +BUILT_AGENT = { + "agents/solver/Dockerfile": "FROM scratch\n", + "tasks/agent.json": json.dumps([{"id": "agent", "type": "names_agent_run_test", "a2a_agent_id": "solver"}]), +} + + +@pytest.fixture +def docker(monkeypatch): + """Stands in for docker: each build's tag is recorded, and the image written as a docker_image document.""" + builds = [] + monkeypatch.setattr(materialize_module.shutil, "which", lambda name: f"/usr/bin/{name}") + monkeypatch.setattr(materialize_module, "build_image", + lambda dockerfile, context, tag, *, platform: builds.append(tag)) + monkeypatch.setattr(materialize_module.DockerImageArtifact, "put", lambda id, **kwargs: get_artifact_store( + ).put_document(DockerImageArtifact(id=id, description="d", image_name=kwargs["image_name"], + tar_gz_s3_url=f"file:///{id}.tar.gz"))) + return builds + + def _names(runs): return [run.entry.name for run in runs] @@ -256,6 +295,21 @@ def broken(line): assert calls == ["tasks/a.json: v1, unchanged"] +def test_a_run_says_before_it_builds_an_image_and_names_the_image_apart_from_its_agent(bundle_dir, docker): + layout(bundle_dir, BUILT_AGENT) + lines = [] + + run_bundle(bundle_dir, tasks=["agent"], on_progress=lines.append) + + assert lines[:4] == [ + "agents/solver (Dockerfile image): building with docker, which can take minutes", + "agents/solver (Dockerfile image): v1 (new)", + "agents/solver: v1 (new)", + "tasks/agent.json: v1 (new)", + ] + assert len(docker) == 1 + + def test_a_cancelled_run_raises_rather_than_being_kept_as_a_failure(bundle_dir): (bundle_dir / "tasks/c.json").write_text(json.dumps([{"id": "stop", "type": "cancelled_run_test"}])) @@ -468,6 +522,26 @@ def test_the_cli_dry_run_reads_another_bundles_entity_and_says_why_an_agent_pinn ) +def test_the_cli_dry_run_lists_an_image_it_would_build_and_builds_nothing(bundle_dir, docker, quiet_logs): + layout(bundle_dir, BUILT_AGENT) + + result = CliRunner().invoke(cli, ["run", str(bundle_dir), "--task", "agent", "--dry-run"]) + + assert result.exit_code == 0, result.output + assert result.output == ( + DRY_RUN + + "agents/solver (Dockerfile image): v1 (new)\n" + "agents/solver: v1 (new)\n" + "tasks/agent.json: v1 (new)\n" + "\n" + "Would run:\n" + " tasks/agent.json v1\n" + + DRY_RUN + ) + assert docker == [] + assert not local_store().path.exists() + + def test_the_cli_dry_run_prints_a_problem_as_one_line_and_exits_1(bundle_dir, quiet_logs): result = CliRunner().invoke(cli, ["run", str(bundle_dir), "--task", "nope", "--dry-run"]) diff --git a/tst/unit/store/local_backends_composition_test.py b/tst/unit/store/local_backends_composition_test.py index cecbe86d..12966764 100644 --- a/tst/unit/store/local_backends_composition_test.py +++ b/tst/unit/store/local_backends_composition_test.py @@ -229,8 +229,9 @@ def wait(self, timeout=None): @pytest.mark.parametrize("entity_id, registry, repository, tarball", [ (HOSTILE, "localhost:5000", HOSTILE_SEGMENT, - f"artifacts/docker_image/{HOSTILE_SEGMENT}/1/local-my-work-tickets-v2-ee679b8e5d3a-v1.tar.gz"), - ("legacy-image", "fake.registry", "legacy-image", "artifacts/docker_image/legacy-image/1/legacy-image-v1.tar.gz"), + (f"artifacts/docker_image/{HOSTILE_SEGMENT}/1-", "/local-my-work-tickets-v2-ee679b8e5d3a-v1.tar.gz")), + ("legacy-image", "fake.registry", "legacy-image", + ("artifacts/docker_image/legacy-image/1-", "/legacy-image-v1.tar.gz")), ]) def test_a_docker_image_names_its_repository_and_tarball_from_the_encoded_id( local_stores, cli_routing, monkeypatch, entity_id, registry, repository, tarball, @@ -247,7 +248,9 @@ def test_a_docker_image_names_its_repository_and_tarball_from_the_encoded_id( assert images.repositories + local_registry == [repository] assert pushed == [f"{registry}/{repository}:v1"] and art.image_name == pushed[0] - assert local_stores.get_object_store().get_object_key(art.tar_gz_object_url) == tarball + before, after = tarball + key = local_stores.get_object_store().get_object_key(art.tar_gz_object_url) + assert re.fullmatch(f"{re.escape(before)}[0-9a-f]{{8}}{re.escape(after)}", key), key assert gzip.decompress(docker_image.DockerImageArtifact.get(entity_id).load()) == b"image-tar-bytes" @@ -315,6 +318,25 @@ def interrupted(**kwargs): assert {key: fa.load() for key, fa in universe.get_file_artifacts().items()} == {"a.txt": b"A", "b.txt": b"B"} +def test_a_docker_image_put_that_stops_before_its_document_doesnt_block_the_next(local_stores, monkeypatch): + set_image_store(FakeImageStore()) + monkeypatch.setattr(docker_image, "_push_local_image", lambda src, ref, store: None) + monkeypatch.setattr(docker_image.subprocess, "Popen", _DockerSave) + put_tar = docker_image.DockerImageArtifact.put_tar + + def interrupted(*args, **kwargs): + raise KeyboardInterrupt + + monkeypatch.setattr(docker_image.DockerImageArtifact, "put_tar", interrupted) + with pytest.raises(KeyboardInterrupt): + docker_image.DockerImageArtifact.put(id="interrupted", description="d", image_name="src:latest") + monkeypatch.setattr(docker_image.DockerImageArtifact, "put_tar", put_tar) + art = docker_image.DockerImageArtifact.put(id="interrupted", description="d", image_name="src:latest") + + assert art.version == 1 + assert gzip.decompress(art.load()) == b"image-tar-bytes" + + def test_get_many_writes_each_local_id_into_its_own_encoded_directory(local_stores, cli_routing, tmp_path): store = local_stores.get_object_store() for uid, name in (("@local/t/A", "1"), ("@local/t/A/1", "f")):