From f27aba7c2f8ce4e6f6a826d4516bac856f314d37 Mon Sep 17 00:00:00 2001 From: Paulo Date: Mon, 24 Aug 2026 20:11:07 +0200 Subject: [PATCH] ENG-882 -- Workflows declare their sandbox: ensured at install, leased by hash --- backend/druks/agents.py | 10 +- backend/druks/api/server.py | 5 + backend/druks/apps/sandboxes.py | 75 ++++ backend/druks/contrib/ship/app.py | 1 + backend/druks/doctor.py | 53 +++ backend/druks/sandbox/client.py | 36 +- backend/druks/sandbox/declaration.py | 46 +++ backend/druks/sandbox/exceptions.py | 12 + .../app_template/package/workflows.py-tpl | 4 +- backend/druks/settings.py | 2 + backend/druks/setup_env.py | 12 +- backend/druks/user_settings/routes.py | 2 + backend/druks/workflows.py | 20 +- .../druks_field_notes/templates/sandbox.sh | 5 + .../druks_field_notes/workflows.py | 3 +- backend/tests/test_apps.py | 38 ++ backend/tests/test_author_surface.py | 1 + backend/tests/test_declared_sandboxes.py | 327 ++++++++++++++++++ backend/tests/test_doctor.py | 95 +++++ backend/tests/test_sandbox_aws_direct.py | 24 ++ backend/tests/test_setup_env.py | 25 ++ backend/tests/test_warm_host_rotation.py | 10 +- docs/configuration.md | 12 + docs/writing-an-app.md | 25 +- druks.toml.example | 4 + 25 files changed, 835 insertions(+), 12 deletions(-) create mode 100644 backend/druks/apps/sandboxes.py create mode 100644 backend/druks/sandbox/declaration.py create mode 100644 backend/tests/druks-field_notes/druks_field_notes/templates/sandbox.sh create mode 100644 backend/tests/test_declared_sandboxes.py diff --git a/backend/druks/agents.py b/backend/druks/agents.py index 7dbab869..be04bac7 100644 --- a/backend/druks/agents.py +++ b/backend/druks/agents.py @@ -11,6 +11,7 @@ from pydantic import BaseModel, ConfigDict from druks.apps.registry import agents +from druks.apps.sandboxes import resolve_declared_sandbox from druks.database import db_session from druks.durable.activity import set_run_phase from druks.durable.engine import _step_engine, step_session @@ -50,7 +51,14 @@ async def _runner( if host_id: vm = sandbox_client.attach(host_id=host_id) else: - vm = sandbox_client.ephemeral(idempotency_key=f"{workflow_id}:{step}") + image_override = template = None + if workflow.sandbox: + image_override, template = await resolve_declared_sandbox(workflow.sandbox) + vm = sandbox_client.ephemeral( + idempotency_key=f"{workflow_id}:{step}", + image_override=image_override, + template=template, + ) async with vm as box: yield await workflow.get_workspace(box) diff --git a/backend/druks/api/server.py b/backend/druks/api/server.py index c16901f0..b17e5a72 100644 --- a/backend/druks/api/server.py +++ b/backend/druks/api/server.py @@ -20,6 +20,7 @@ from druks.api.subjects import router as subjects_router from druks.apps.loader import iter_apps, load from druks.apps.routes import router as apps_router +from druks.apps.sandboxes import ensure_declared_sandboxes from druks.browser.exceptions import BrowserApiError from druks.browser.routes import router as browser_sessions_router from druks.core.templates import render_page @@ -98,6 +99,10 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: logging.getLogger(__name__).exception( "app %r on_startup failed", registered_app.name ) + try: + await ensure_declared_sandboxes() + except Exception: + logging.getLogger(__name__).exception("declared sandbox ensure failed") try: yield diff --git a/backend/druks/apps/sandboxes.py b/backend/druks/apps/sandboxes.py new file mode 100644 index 00000000..0576af6e --- /dev/null +++ b/backend/druks/apps/sandboxes.py @@ -0,0 +1,75 @@ +import asyncio +from typing import TYPE_CHECKING + +from druks.durable.activity import set_run_phase +from druks.sandbox.client import sandbox_client +from druks.sandbox.declaration import Sandbox, get_content_hash +from druks.sandbox.exceptions import SandboxTemplateNotFound, SandboxTemplateUnavailable +from druks.settings import load_settings + +from . import loader + +if TYPE_CHECKING: + from druks.workflows import Workflow + +_TEMPLATE_POLL_SECONDS = 5 + + +def collect_declared_sandboxes() -> dict[str, tuple[str, bytes, list[type["Workflow"]]]]: + declared: dict[str, tuple[str, bytes, list[type[Workflow]]]] = {} + for app in loader.iter_apps(): + for workflow in app.workflows(): + if sandbox := workflow.sandbox: + base, script = sandbox.resolve() + content_hash = get_content_hash(base, script) + if content_hash in declared: + declared[content_hash][2].append(workflow) + else: + declared[content_hash] = (base, script, [workflow]) + return declared + + +async def ensure_declared_sandboxes() -> dict[str, tuple[str, bytes, list[type["Workflow"]]]]: + if not load_settings().sandbox.service_url: + return {} + + declared = collect_declared_sandboxes() + for requirements_hash, (base_image, script, _) in declared.items(): + await sandbox_client.ensure_template( + base_image=base_image, + script=script, + requirements_hash=requirements_hash, + ) + return declared + + +async def resolve_declared_sandbox(sandbox: Sandbox) -> tuple[str | None, str | None]: + base_image, script = sandbox.resolve() + requirements_hash = get_content_hash(base_image, script) + if pinned_image := load_settings().sandbox.pins.get(requirements_hash): + return pinned_image, None + + try: + template = await sandbox_client.get_template(requirements_hash=requirements_hash) + except SandboxTemplateNotFound as error: + raise SandboxTemplateUnavailable( + f"sandbox template {requirements_hash} is missing. " + "Reinstall the app or run `druks doctor`." + ) from error + + is_building = template.status == "building" + if is_building: + await set_run_phase("sandbox_building") + while template.status == "building": + await asyncio.sleep(_TEMPLATE_POLL_SECONDS) + template = await sandbox_client.get_template(requirements_hash=requirements_hash) + + if template.status == "available": + if is_building: + await set_run_phase("provisioning_vm") + return None, template.id + + raise SandboxTemplateUnavailable( + f"sandbox template {requirements_hash} has status {template.status!r}. " + "Fix its setup and run `druks doctor`." + ) diff --git a/backend/druks/contrib/ship/app.py b/backend/druks/contrib/ship/app.py index c6cd79f4..5d1678b0 100644 --- a/backend/druks/contrib/ship/app.py +++ b/backend/druks/contrib/ship/app.py @@ -26,6 +26,7 @@ # to name it, so the phase that clears provisioning maps to nothing. _PHASE_META: dict[str, SubjectActivity] = { "provisioning_vm": SubjectActivity(label="Building sandbox VM…", kind="infra"), + "sandbox_building": SubjectActivity(label="Building sandbox…", kind="infra"), } diff --git a/backend/druks/doctor.py b/backend/druks/doctor.py index 6238527c..a86010c6 100644 --- a/backend/druks/doctor.py +++ b/backend/druks/doctor.py @@ -21,11 +21,13 @@ from .agents import Agent from .apps.loader import iter_apps from .apps.registry import _ROLES, agents, autodiscover, services, webhooks, workflows +from .apps.sandboxes import ensure_declared_sandboxes from .core.apis.github import get_github_client from .database import create_async_engine_from_url, create_engine_from_url, session_scope from .harnesses.models import HarnessConnection from .harnesses.registry import get_harnesses from .sandbox.client import sandbox_client +from .sandbox.exceptions import SandboxTemplateNotFound from .services import Service, ServiceNotConnectedError from .services.models import ServiceIdentity from .settings import Settings, load_settings @@ -293,6 +295,56 @@ async def check_sandbox_e2e(settings: Settings) -> CheckResult: return CheckResult(name="sandbox_e2e", ok=True, detail=detail) +async def check_declared_sandboxes(settings: Settings) -> CheckResult | list[CheckResult]: + if not settings.sandbox.service_url: + return CheckResult(name="sandbox_templates", ok=True, detail="not configured") + + try: + declared = await ensure_declared_sandboxes() + except Exception as error: # noqa: BLE001 — doctor reports, never raises + return CheckResult( + name="sandbox_templates", + ok=False, + detail=f"could not ensure declared sandboxes: {error}", + ) + + if not declared: + return CheckResult(name="sandbox_templates", ok=True, detail="no declared sandboxes") + + results = [] + for requirements_hash, (_, _, workflow_classes) in declared.items(): + workflow_names = ", ".join(workflow.kind for workflow in workflow_classes) + name = f"sandbox:{requirements_hash[:12]}" + detail = f"{workflow_names}; hash {requirements_hash}" + try: + template = await sandbox_client.get_template(requirements_hash=requirements_hash) + except SandboxTemplateNotFound: + result = CheckResult( + name=name, + ok=False, + detail=f"{detail}; missing", + ) + except Exception as error: # noqa: BLE001 — one lookup failure is one result + result = CheckResult( + name=name, + ok=False, + detail=f"{detail}; lookup failed: {error}", + ) + else: + ok, pending = { + "available": (True, False), + "building": (False, True), + }.get(template.status, (False, False)) + result = CheckResult( + name=name, + ok=ok, + pending=pending, + detail=f"{detail}; {template.status}", + ) + results.append(result) + return results + + async def _sandbox_e2e() -> str: start = time.monotonic() # acquire rolls its own host back on failure; once it yields, we own @@ -482,6 +534,7 @@ async def _run_app_check(app_name: str, check) -> CheckResult: check_drukbox, check_capability_modules, check_apps, + check_declared_sandboxes, ) diff --git a/backend/druks/sandbox/client.py b/backend/druks/sandbox/client.py index ce2ea276..30f34d24 100644 --- a/backend/druks/sandbox/client.py +++ b/backend/druks/sandbox/client.py @@ -4,6 +4,7 @@ from contextlib import asynccontextmanager from datetime import UTC, datetime, timedelta from pathlib import Path +from typing import Any import asyncssh from drukbox_sdk import SandboxAPI, SandboxHost @@ -19,7 +20,7 @@ from druks.settings import load_settings from .constants import SANDBOX_HOST_LEASE_SECONDS -from .exceptions import HostGone, SandboxError, SandboxUnreachable +from .exceptions import HostGone, SandboxError, SandboxTemplateNotFound, SandboxUnreachable from .host import Sandbox from .layout import get_helper_script_path, get_remote_home @@ -53,6 +54,7 @@ async def ephemeral( image_override: str | None = None, provider: str | None = None, sandbox_env: dict[str, str] | None = None, + template: str | None = None, ) -> AsyncIterator[Sandbox]: """One-shot lifecycle: acquire → yield → release. For callers whose sandbox is bound to a single context manager body.""" @@ -64,6 +66,7 @@ async def ephemeral( image_override=image_override, provider=provider, sandbox_env=sandbox_env, + template=template, ) as sandbox: host_id = sandbox.id yield sandbox @@ -79,6 +82,7 @@ async def acquire( image_override: str | None = None, provider: str | None = None, sandbox_env: dict[str, str] | None = None, + template: str | None = None, ) -> AsyncIterator[Sandbox]: """Create a new host (or reuse one matching ``idempotency_key``) and yield it with SSH connected. Closes SSH on exit but does NOT @@ -92,6 +96,10 @@ async def acquire( # Fixed lease: drukbox reaps the host when this lapses, so a run whose # worker dies frees its VM without a druks-side reconciler. expires_at = datetime.now(UTC) + timedelta(seconds=SANDBOX_HOST_LEASE_SECONDS) + create_host_kwargs: dict[str, Any] = {} + if template: + # SDK 0.0.7 rejects template= even when unset, so ordinary leases omit it. + create_host_kwargs["template"] = template try: record = await api.create_host( expires_at=expires_at, @@ -99,6 +107,7 @@ async def acquire( idempotency_key=key, image=image or None, provider=provider, + **create_host_kwargs, ) except (SandboxProvisioningError, SandboxUnavailableError) as exc: # Transient control-plane failures — a 502 the service raises @@ -148,6 +157,29 @@ async def list_hosts(self) -> list[SandboxHost]: finally: await api.aclose() + async def ensure_template( + self, *, base_image: str, script: bytes, requirements_hash: str + ) -> Any: + api: Any = self._api() + try: + return await api.create_template( + base_image=base_image, + script=script, + requirements_hash=requirements_hash, + ) + finally: + await api.aclose() + + async def get_template(self, *, requirements_hash: str) -> Any: + api: Any = self._api() + try: + for template in await api.list_templates(): + if template.requirements_hash == requirements_hash: + return template + finally: + await api.aclose() + raise SandboxTemplateNotFound(f"sandbox template {requirements_hash} does not exist") + @staticmethod async def _best_effort_delete(api: SandboxAPI, host_id: str) -> None: try: @@ -193,6 +225,7 @@ async def provision( image_override: str | None = None, provider: str | None = None, sandbox_env: dict[str, str] | None = None, + template: str | None = None, ) -> Sandbox: """Create a host and return its handle without holding an SSH connection — the handle reconnects lazily when used (its id and lease expiry are readable @@ -202,6 +235,7 @@ async def provision( image_override=image_override, provider=provider, sandbox_env=sandbox_env, + template=template, ) as sandbox: return sandbox diff --git a/backend/druks/sandbox/declaration.py b/backend/druks/sandbox/declaration.py new file mode 100644 index 00000000..45be631f --- /dev/null +++ b/backend/druks/sandbox/declaration.py @@ -0,0 +1,46 @@ +import hashlib +import importlib.util +from dataclasses import dataclass +from pathlib import Path + +from druks.apps import loader +from druks.settings import load_settings + +from .exceptions import SandboxSetupError + + +def get_content_hash(base: str, script: bytes) -> str: + content = base.encode("utf-8") + b"\0" + script + return hashlib.sha256(content).hexdigest() + + +@dataclass(frozen=True) +class Sandbox: + setup: str + + def resolve(self) -> tuple[str, bytes]: + app_name, separator, relative_path = self.setup.partition("/") + if not separator or not relative_path: + raise SandboxSetupError(f"sandbox setup {self.setup!r} must be '/'") + + for app in loader.iter_apps(): + if app.name == app_name: + spec = importlib.util.find_spec(app.package) + if not spec or not spec.submodule_search_locations: + raise SandboxSetupError( + f"sandbox setup {self.setup!r} cannot find app package {app.package!r}" + ) + path = Path(spec.submodule_search_locations[0]) / "templates" / relative_path + try: + script = path.read_bytes() + except OSError as error: + raise SandboxSetupError( + f"sandbox setup {self.setup!r} cannot be read: {error}" + ) from error + return load_settings().sandbox.image, script + + raise SandboxSetupError(f"sandbox setup {self.setup!r} names unknown app {app_name!r}") + + @property + def content_hash(self) -> str: + return get_content_hash(*self.resolve()) diff --git a/backend/druks/sandbox/exceptions.py b/backend/druks/sandbox/exceptions.py index 9da25184..743bb682 100644 --- a/backend/druks/sandbox/exceptions.py +++ b/backend/druks/sandbox/exceptions.py @@ -2,6 +2,18 @@ class SandboxError(Exception): """Base for everything ``druks.sandbox`` raises out of its layer.""" +class SandboxSetupError(SandboxError): + """A declared setup path cannot resolve to package bytes.""" + + +class SandboxTemplateNotFound(SandboxError): + """Drukbox has no template for a requirements hash.""" + + +class SandboxTemplateUnavailable(SandboxError): + """A declared template cannot be leased.""" + + class SandboxUnreachable(SandboxError): """The SSH connection to the VM cannot be (re-)established. diff --git a/backend/druks/scaffolding/app_template/package/workflows.py-tpl b/backend/druks/scaffolding/app_template/package/workflows.py-tpl index eacfc730..77aea167 100644 --- a/backend/druks/scaffolding/app_template/package/workflows.py-tpl +++ b/backend/druks/scaffolding/app_template/package/workflows.py-tpl @@ -1,4 +1,4 @@ -from druks.workflows import FatalError, Gate, Workflow, step # noqa: F401 +from druks.workflows import FatalError, Gate, Sandbox, Workflow, step # noqa: F401 # Durable workflows. Subclass ``Workflow`` and implement exactly one of: # async def run(self, ...) — a single durable operation @@ -6,6 +6,8 @@ from druks.workflows import FatalError, Gate, Workflow, step # noqa: F401 # The parameters ARE the workflow's input: plain typed params, validated at start(). # Set ``every = ""`` to schedule one; raise ``FatalError`` for a clean domain # stop; a ``Gate`` subclass parks the run for human input. +# Set ``sandbox = Sandbox(setup="/sandbox.sh")`` when the workflow needs +# tools beyond the platform base. Ship the shell file under the package's templates/. # # Declare what a workflow's runs are about with ``subject = YourSubject`` — the class, # with a table or without — and druks gives that subject a board and a page. diff --git a/backend/druks/settings.py b/backend/druks/settings.py index a38913fe..deb19e92 100644 --- a/backend/druks/settings.py +++ b/backend/druks/settings.py @@ -154,6 +154,8 @@ class Sandbox(BaseModel): service_token: str = "" # Empty → drukbox decides. image: str = "" + # Declared sandbox content hash → operator-pinned provider artifact. + pins: dict[str, str] = {} # The browser home: browser containers boot on this provider with this image. browser_sandbox_provider: str = "docker" browser_sandbox_image: str = "ghcr.io/czpython/druks-browser:latest" diff --git a/backend/druks/setup_env.py b/backend/druks/setup_env.py index 565b0a53..07d5a3d5 100644 --- a/backend/druks/setup_env.py +++ b/backend/druks/setup_env.py @@ -72,6 +72,7 @@ "service_url", "service_token", "image", + "pins", "browser_login_proxy", "browser_login_tz", "timeout", @@ -269,7 +270,16 @@ def _canonical_config(raw: dict[str, Any]) -> dict[str, Any]: if not isinstance(table, dict): raise ValueError(f"druks.toml: [{table_name}] must be a table") for key in keys: - value = table.setdefault(key, "") + value = table.setdefault(key, {} if key == "pins" else "") + if table_name == "sandbox" and key == "pins": + if not isinstance(value, dict): + raise ValueError("druks.toml: [sandbox.pins] must be a table") + for content_hash, artifact in value.items(): + if not isinstance(artifact, str): + raise ValueError( + f"druks.toml: sandbox.pins.{content_hash} must be a string" + ) + continue if table_name == "sandbox" and key == "timeout": if type(value) not in (str, int, float): raise ValueError("druks.toml: sandbox.timeout must be a number or string") diff --git a/backend/druks/user_settings/routes.py b/backend/druks/user_settings/routes.py index a4155bed..4f0405f3 100644 --- a/backend/druks/user_settings/routes.py +++ b/backend/druks/user_settings/routes.py @@ -6,6 +6,7 @@ from druks.accounts.models import Account from druks.apps.loader import get_app, iter_apps from druks.apps.registry import workflows +from druks.apps.sandboxes import ensure_declared_sandboxes from druks.durable.engine import apply_schedules from druks.harnesses.exceptions import HarnessError from druks.harnesses.models import HarnessConnection @@ -208,4 +209,5 @@ async def update_app_settings(body: AppsSettingsUpdate) -> AppsSettingsResponse: # the just-written overrides off this request's session. await apply_schedules() + await ensure_declared_sandboxes() return await get_app_settings() diff --git a/backend/druks/workflows.py b/backend/druks/workflows.py index 8c17904b..21a1d963 100644 --- a/backend/druks/workflows.py +++ b/backend/druks/workflows.py @@ -29,6 +29,7 @@ from druks.accounts.context import current_account_id from druks.apps.loader import resolve_workflow_app from druks.apps.registry import workflows +from druks.apps.sandboxes import resolve_declared_sandbox from druks.apps.settings import ( coerce_setting_value, validate_setting_override, @@ -53,6 +54,7 @@ from druks.notifications.outbox import notifications_queue, send_notification from druks.sandbox.client import sandbox_client from druks.sandbox.constants import SANDBOX_HOST_ROTATE_BEFORE_SECONDS +from druks.sandbox.declaration import Sandbox from druks.signals import publish from druks.user_settings.models import SettingsOverride, UserSettings from druks.workspaces import Workspace @@ -69,6 +71,7 @@ "Journal", "OperatorReply", "RunResponse", + "Sandbox", "Subject", "SubjectActivity", "SubjectStatus", @@ -82,7 +85,7 @@ ] if TYPE_CHECKING: - from druks.sandbox.host import Sandbox + from druks.sandbox.host import Sandbox as SandboxHost # A human gate can park for days; a long recv TTL still caps zombie parks. GATE_TTL_SECONDS = 14 * 24 * 60 * 60 @@ -683,6 +686,8 @@ class Workflow: # True holds one warm VM across the run's agent calls (released at gate parks); # False gives each call a throwaway VM. steps_reuse_sandbox: ClassVar[bool] = False + # The environment this workflow's agents need. None uses the platform base. + sandbox: ClassVar[Sandbox | None] = None # The Workspace subclass agents run in; an app sets it (default: the bare VM). workspace_class: ClassVar[type[Workspace]] = Workspace # The Journal subclass the run keeps; an app sets it for named projections. @@ -770,7 +775,7 @@ def __init__(self) -> None: self.journal = self.journal_class() # The run's warm VM, provisioned lazily and reaped at segment boundaries; # its lease expiry decides when it must rotate. - self._host: Sandbox | None = None + self._host: SandboxHost | None = None async def announce(self, topic: str, **facts: Any) -> None: # The workflow announcing a domain event in its app's vocabulary @@ -819,12 +824,12 @@ async def get_prompt_context(self, **context: Any) -> dict[str, Any]: # values (expensive derived prose); the call's own kwargs win on collision. return dict(context) - async def get_workspace_kwargs(self, sandbox: "Sandbox") -> dict[str, Any]: + async def get_workspace_kwargs(self, sandbox: "SandboxHost") -> dict[str, Any]: # Extend via super() to add the fields workspace_class needs (an app clones + mints # here). Base: just the VM. return {"sandbox": sandbox} - async def get_workspace(self, sandbox: "Sandbox") -> Workspace: + async def get_workspace(self, sandbox: "SandboxHost") -> Workspace: # What an agent runs in on this run's VM, built per agent call from workspace_class # + the app's kwargs — so short-lived tokens (git) mint fresh each call. return self.workspace_class(**await self.get_workspace_kwargs(sandbox)) @@ -842,8 +847,13 @@ async def _ensure_host(self) -> str | None: # host it lands on (state lives in git), so a bare VM is fine. await self._reap_run() if not self._host: + image_override = template = None + if self.sandbox: + image_override, template = await resolve_declared_sandbox(self.sandbox) self._host = await sandbox_client.provision( - idempotency_key=f"{self._workflow_id}:sandbox" + idempotency_key=f"{self._workflow_id}:sandbox", + image_override=image_override, + template=template, ) return self._host.id diff --git a/backend/tests/druks-field_notes/druks_field_notes/templates/sandbox.sh b/backend/tests/druks-field_notes/druks_field_notes/templates/sandbox.sh new file mode 100644 index 00000000..c4f255a6 --- /dev/null +++ b/backend/tests/druks-field_notes/druks_field_notes/templates/sandbox.sh @@ -0,0 +1,5 @@ +#!/bin/sh +set -eu + +apt-get update +apt-get install -y jq diff --git a/backend/tests/druks-field_notes/druks_field_notes/workflows.py b/backend/tests/druks-field_notes/druks_field_notes/workflows.py index 55d31041..312b3fc0 100644 --- a/backend/tests/druks-field_notes/druks_field_notes/workflows.py +++ b/backend/tests/druks-field_notes/druks_field_notes/workflows.py @@ -1,4 +1,4 @@ -from druks.workflows import Workflow +from druks.workflows import Sandbox, Workflow from druks_field_notes.app import FieldNotes from druks_field_notes.models import Note @@ -9,6 +9,7 @@ class Summarize(Workflow): produces the line, and the run stores it on the note.""" subject = Note + sandbox = Sandbox(setup="field_notes/sandbox.sh") async def run(self) -> None: note = await self.subject diff --git a/backend/tests/test_apps.py b/backend/tests/test_apps.py index 038eb237..f4237bf4 100644 --- a/backend/tests/test_apps.py +++ b/backend/tests/test_apps.py @@ -1,9 +1,12 @@ from types import ModuleType, SimpleNamespace +from unittest.mock import AsyncMock, Mock +import druks.api.server as server import pytest from druks.apps import App, loader from druks.apps.exceptions import AppSubjectContractError from druks.apps.loader import iter_apps, load +from druks.testing import make_settings from druks.workflows import Subject from fastapi import APIRouter, FastAPI from fastapi.testclient import TestClient @@ -114,6 +117,41 @@ def discover(cls) -> list[ModuleType]: assert client.get("/health/ping").status_code == 404 +async def test_server_lifespan_ensures_declared_sandboxes_after_app_startup( + monkeypatch, + tmp_path, +): + calls = [] + on_startup = AsyncMock(side_effect=lambda: calls.append("startup")) + ensure = AsyncMock(side_effect=lambda: calls.append("ensure")) + registered_app = SimpleNamespace(name="empty", on_startup=on_startup) + settings = make_settings( + tmp_path, + identity={"mode": "header", "header": "X-Operator"}, + ) + engine = SimpleNamespace(dispose=AsyncMock()) + monkeypatch.setattr(server, "load_settings", lambda: settings) + monkeypatch.setattr(server, "ensure_data_dirs", Mock()) + monkeypatch.setattr(server, "create_async_engine_from_url", Mock(return_value=engine)) + monkeypatch.setattr(server, "configure_session", Mock()) + monkeypatch.setattr(server, "setup_logging", Mock()) + monkeypatch.setattr(server, "load_mcp_catalog", Mock()) + monkeypatch.setattr(server, "init_dbos", Mock()) + monkeypatch.setattr(server, "launch", AsyncMock()) + monkeypatch.setattr(server, "shutdown", Mock()) + monkeypatch.setattr(server, "close_client", AsyncMock()) + monkeypatch.setattr(server, "iter_apps", lambda: [registered_app]) + monkeypatch.setattr(server, "ensure_declared_sandboxes", ensure) + api = FastAPI() + + async with server.lifespan(api): + calls.append("running") + + assert calls == ["startup", "ensure", "running"] + on_startup.assert_awaited_once_with() + ensure.assert_awaited_once_with() + + def _routes_module(name: str, **routers: APIRouter) -> ModuleType: module = ModuleType(f"{name}.routes") module.__dict__.update(routers) diff --git a/backend/tests/test_author_surface.py b/backend/tests/test_author_surface.py index 69abc835..8292e7a2 100644 --- a/backend/tests/test_author_surface.py +++ b/backend/tests/test_author_surface.py @@ -25,6 +25,7 @@ "Journal", "OperatorReply", "RunResponse", + "Sandbox", "Subject", "SubjectActivity", "SubjectStatus", diff --git a/backend/tests/test_declared_sandboxes.py b/backend/tests/test_declared_sandboxes.py new file mode 100644 index 00000000..91bcdeac --- /dev/null +++ b/backend/tests/test_declared_sandboxes.py @@ -0,0 +1,327 @@ +import hashlib +from contextlib import asynccontextmanager +from datetime import UTC, datetime, timedelta +from types import SimpleNamespace +from unittest.mock import AsyncMock, call + +import druks.agents as agent_module +import druks.workflows as workflow_module +import pytest +from druks.apps import sandboxes +from druks.contrib.ship.app import _PHASE_META +from druks.sandbox import declaration +from druks.sandbox.client import Client +from druks.sandbox.declaration import Sandbox, get_content_hash +from druks.sandbox.exceptions import SandboxTemplateNotFound, SandboxTemplateUnavailable +from druks.user_settings import routes as settings_routes +from druks.user_settings.schemas import AppsSettingsUpdate +from druks.workflows import Workflow + + +def test_sandbox_resolves_package_bytes_and_hashes_base_with_script(monkeypatch, tmp_path): + package = tmp_path / "site_builder" + templates = package / "templates" + templates.mkdir(parents=True) + (package / "__init__.py").write_text("") + (templates / "sandbox.sh").write_bytes(b"#!/bin/sh\ninstall-tool\n") + monkeypatch.syspath_prepend(tmp_path) + monkeypatch.setattr( + declaration.loader, + "iter_apps", + lambda: [SimpleNamespace(name="site_builder", package="site_builder")], + ) + monkeypatch.setattr( + declaration, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(image="base-image")), + ) + sandbox = Sandbox(setup="site_builder/sandbox.sh") + + base, script = sandbox.resolve() + + assert sandbox.setup == "site_builder/sandbox.sh" + assert base == "base-image" + assert script.startswith(b"#!/bin/sh\n") + assert sandbox.content_hash == hashlib.sha256(b"base-image\0" + script).hexdigest() + + +def test_collect_declared_sandboxes_deduplicates_by_content(monkeypatch): + shared = Sandbox(setup="notes/sandbox.sh") + + class First: + kind = "notes.first" + sandbox = shared + + class Second: + kind = "notes.second" + sandbox = shared + + app = SimpleNamespace(workflows=lambda: [First, Second]) + monkeypatch.setattr(sandboxes.loader, "iter_apps", lambda: [app]) + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + + declared = sandboxes.collect_declared_sandboxes() + + requirements_hash = get_content_hash("base", b"setup") + assert declared == {requirements_hash: ("base", b"setup", [First, Second])} + + +def test_ship_maps_the_sandbox_building_phase(): + assert _PHASE_META["sandbox_building"].label == "Building sandbox…" + assert _PHASE_META["sandbox_building"].kind == "infra" + + +async def test_ensure_declared_sandboxes_requests_each_unique_template(monkeypatch): + requirements_hash = get_content_hash("base", b"setup") + ensure_template = AsyncMock() + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(service_url="http://drukbox")), + ) + monkeypatch.setattr( + sandboxes, + "collect_declared_sandboxes", + lambda: {requirements_hash: ("base", b"setup", [SimpleNamespace(kind="notes")])}, + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace(ensure_template=ensure_template), + ) + + declared = await sandboxes.ensure_declared_sandboxes() + + assert declared == {requirements_hash: ("base", b"setup", [SimpleNamespace(kind="notes")])} + ensure_template.assert_awaited_once_with( + base_image="base", + script=b"setup", + requirements_hash=requirements_hash, + ) + + +async def test_app_settings_save_ensures_declared_sandboxes(monkeypatch): + ensure = AsyncMock() + response = SimpleNamespace(apps=[]) + monkeypatch.setattr(settings_routes, "ensure_declared_sandboxes", ensure) + monkeypatch.setattr( + settings_routes, + "get_app_settings", + AsyncMock(return_value=response), + ) + + result = await settings_routes.update_app_settings(AppsSettingsUpdate()) + + assert result is response + ensure.assert_awaited_once_with() + + +async def test_resolve_declared_sandbox_uses_available_template(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + template = SimpleNamespace(id="template-1", status="available") + get_template = AsyncMock(return_value=template) + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(pins={})), + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace(get_template=get_template), + ) + + assert await sandboxes.resolve_declared_sandbox(sandbox) == (None, "template-1") + get_template.assert_awaited_once_with(requirements_hash=get_content_hash("base", b"setup")) + + +async def test_resolve_declared_sandbox_waits_with_visible_phase(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + get_template = AsyncMock( + side_effect=[ + SimpleNamespace(id="template-1", status="building"), + SimpleNamespace(id="template-1", status="available"), + ] + ) + set_run_phase = AsyncMock() + sleep = AsyncMock() + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(pins={})), + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace(get_template=get_template), + ) + monkeypatch.setattr(sandboxes, "set_run_phase", set_run_phase) + monkeypatch.setattr(sandboxes.asyncio, "sleep", sleep) + + assert await sandboxes.resolve_declared_sandbox(sandbox) == (None, "template-1") + assert set_run_phase.await_args_list == [ + call("sandbox_building"), + call("provisioning_vm"), + ] + sleep.assert_awaited_once_with(sandboxes._TEMPLATE_POLL_SECONDS) + assert get_template.await_count == 2 + + +async def test_resolve_declared_sandbox_rejects_missing_template(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(pins={})), + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace(get_template=AsyncMock(side_effect=SandboxTemplateNotFound("missing"))), + ) + + with pytest.raises(SandboxTemplateUnavailable, match="missing.*druks doctor"): + await sandboxes.resolve_declared_sandbox(sandbox) + + +async def test_resolve_declared_sandbox_rejects_failed_template(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(pins={})), + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace( + get_template=AsyncMock(return_value=SimpleNamespace(id="template-1", status="failed")) + ), + ) + + with pytest.raises(SandboxTemplateUnavailable, match="failed.*druks doctor"): + await sandboxes.resolve_declared_sandbox(sandbox) + + +async def test_resolve_declared_sandbox_uses_operator_pin_without_lookup(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + requirements_hash = get_content_hash("base", b"setup") + get_template = AsyncMock() + monkeypatch.setattr(Sandbox, "resolve", lambda self: ("base", b"setup")) + monkeypatch.setattr( + sandboxes, + "load_settings", + lambda: SimpleNamespace(sandbox=SimpleNamespace(pins={requirements_hash: "ami-pinned"})), + ) + monkeypatch.setattr( + sandboxes, + "sandbox_client", + SimpleNamespace(get_template=get_template), + ) + + assert await sandboxes.resolve_declared_sandbox(sandbox) == ("ami-pinned", None) + get_template.assert_not_awaited() + + +async def test_warm_lease_uses_workflow_template(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + host = SimpleNamespace( + id="host-1", + expires_at=datetime.now(UTC) + timedelta(hours=1), + ) + provision = AsyncMock(return_value=host) + resolve = AsyncMock(return_value=(None, "template-1")) + workflow = Workflow.__new__(Workflow) + workflow.steps_reuse_sandbox = True + workflow.sandbox = sandbox + workflow._workflow_id = "run-1" + workflow._host = None + monkeypatch.setattr( + workflow_module, + "sandbox_client", + SimpleNamespace(provision=provision), + ) + monkeypatch.setattr(workflow_module, "resolve_declared_sandbox", resolve) + + assert await workflow._ensure_host() == "host-1" + resolve.assert_awaited_once_with(sandbox) + provision.assert_awaited_once_with( + idempotency_key="run-1:sandbox", + image_override=None, + template="template-1", + ) + + +async def test_ephemeral_lease_uses_workflow_template(monkeypatch): + sandbox = Sandbox(setup="notes/sandbox.sh") + calls = [] + box = SimpleNamespace(id="host-1") + + @asynccontextmanager + async def ephemeral(**kwargs): + calls.append(kwargs) + yield box + + workflow = SimpleNamespace( + sandbox=sandbox, + get_workspace=AsyncMock(return_value="workspace"), + ) + resolve = AsyncMock(return_value=(None, "template-1")) + monkeypatch.setattr( + agent_module, + "sandbox_client", + SimpleNamespace(ephemeral=ephemeral), + ) + monkeypatch.setattr(agent_module, "resolve_declared_sandbox", resolve) + + async with agent_module._runner(workflow, None, "run-1", "summarize") as runner: + assert runner == "workspace" + + resolve.assert_awaited_once_with(sandbox) + assert calls == [ + { + "idempotency_key": "run-1:summarize", + "image_override": None, + "template": "template-1", + } + ] + + +async def test_client_template_primitives_use_sdk_contract(monkeypatch): + created = SimpleNamespace(id="template-1", status="building") + listed = SimpleNamespace( + id="template-1", + status="available", + requirements_hash="requirements-1", + ) + + class FakeAPI: + def __init__(self): + self.create_template = AsyncMock(return_value=created) + self.list_templates = AsyncMock(return_value=[listed]) + self.aclose = AsyncMock() + + api = FakeAPI() + client = Client() + monkeypatch.setattr(Client, "_api", lambda self: api) + + assert ( + await client.ensure_template( + base_image="base", + script=b"setup", + requirements_hash="requirements-1", + ) + is created + ) + assert await client.get_template(requirements_hash="requirements-1") is listed + api.create_template.assert_awaited_once_with( + base_image="base", + script=b"setup", + requirements_hash="requirements-1", + ) + api.list_templates.assert_awaited_once_with() + assert api.aclose.await_count == 2 diff --git a/backend/tests/test_doctor.py b/backend/tests/test_doctor.py index a6df8a51..a45b09dd 100644 --- a/backend/tests/test_doctor.py +++ b/backend/tests/test_doctor.py @@ -1,11 +1,14 @@ from contextlib import asynccontextmanager from datetime import UTC, datetime, timedelta from pathlib import Path +from types import SimpleNamespace +from unittest.mock import AsyncMock import httpx import pytest from druks import doctor from druks.database import db_session +from druks.sandbox.exceptions import SandboxTemplateNotFound from druks.services.models import ServiceIdentity from druks.testing import make_settings @@ -194,6 +197,98 @@ async def test_drukbox_passes_when_unconfigured(tmp_path: Path) -> None: assert "not configured" in result.detail +async def test_declared_sandboxes_pass_when_not_configured(tmp_path: Path) -> None: + result = await doctor.check_declared_sandboxes(make_settings(tmp_path)) + + assert result == doctor.CheckResult( + name="sandbox_templates", + ok=True, + detail="not configured", + ) + + +async def test_declared_sandboxes_pass_when_none_are_declared( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + ensure = AsyncMock(return_value={}) + monkeypatch.setattr(doctor, "ensure_declared_sandboxes", ensure) + settings = make_settings(tmp_path, sandbox={"service_url": "http://drukbox"}) + + result = await doctor.check_declared_sandboxes(settings) + + assert result == doctor.CheckResult( + name="sandbox_templates", + ok=True, + detail="no declared sandboxes", + ) + ensure.assert_awaited_once_with() + + +@pytest.mark.parametrize( + ("status", "ok", "pending"), + [ + ("available", True, False), + ("building", False, True), + ("failed", False, False), + ], +) +async def test_declared_sandboxes_report_template_status( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + status: str, + ok: bool, + pending: bool, +) -> None: + workflow = SimpleNamespace(kind="field_notes.summarize") + ensure = AsyncMock(return_value={"requirements-1": ("base", b"setup", [workflow])}) + lookup = AsyncMock(return_value=SimpleNamespace(status=status)) + monkeypatch.setattr(doctor, "ensure_declared_sandboxes", ensure) + monkeypatch.setattr( + doctor, + "sandbox_client", + SimpleNamespace(get_template=lookup), + ) + settings = make_settings(tmp_path, sandbox={"service_url": "http://drukbox"}) + + results = await doctor.check_declared_sandboxes(settings) + + ensure.assert_awaited_once_with() + lookup.assert_awaited_once_with(requirements_hash="requirements-1") + assert len(results) == 1 + assert results[0].ok is ok + assert results[0].pending is pending + assert "field_notes.summarize" in results[0].detail + assert "requirements-1" in results[0].detail + assert status in results[0].detail + + +async def test_declared_sandboxes_report_missing_template( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + workflow = SimpleNamespace(kind="field_notes.summarize") + monkeypatch.setattr( + doctor, + "ensure_declared_sandboxes", + AsyncMock(return_value={"requirements-1": ("base", b"setup", [workflow])}), + ) + monkeypatch.setattr( + doctor, + "sandbox_client", + SimpleNamespace(get_template=AsyncMock(side_effect=SandboxTemplateNotFound("missing"))), + ) + settings = make_settings(tmp_path, sandbox={"service_url": "http://drukbox"}) + + result = (await doctor.check_declared_sandboxes(settings))[0] + + assert not result.ok + assert not result.pending + assert "field_notes.summarize" in result.detail + assert "requirements-1" in result.detail + assert "missing" in result.detail + + async def test_run_checks_covers_all_check_names(tmp_path: Path) -> None: settings = make_settings(tmp_path) diff --git a/backend/tests/test_sandbox_aws_direct.py b/backend/tests/test_sandbox_aws_direct.py index e3b42f7d..e17b0da7 100644 --- a/backend/tests/test_sandbox_aws_direct.py +++ b/backend/tests/test_sandbox_aws_direct.py @@ -201,3 +201,27 @@ async def test_acquire_persists_private_key_when_returned( key_path = tmp_path / "host-aws" assert key_path.read_text().startswith("-----BEGIN OPENSSH") assert oct(os.stat(key_path).st_mode & 0o777) == "0o600" + + +async def test_acquire_passes_template_to_drukbox( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +): + client, calls = _stub_acquire_settings(monkeypatch, tmp_path, sandbox_image="base") + + async with client.acquire(idempotency_key="op", template="template-1"): + pass + + assert calls[0]["template"] == "template-1" + + +async def test_acquire_omits_unset_template_from_drukbox( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +): + client, calls = _stub_acquire_settings(monkeypatch, tmp_path, sandbox_image="base") + + async with client.acquire(idempotency_key="op"): + pass + + assert "template" not in calls[0] diff --git a/backend/tests/test_setup_env.py b/backend/tests/test_setup_env.py index 320cc948..5e65e8aa 100644 --- a/backend/tests/test_setup_env.py +++ b/backend/tests/test_setup_env.py @@ -195,6 +195,31 @@ def test_non_string_known_setting_is_refused(tmp_path): _run(env_path) +def test_sandbox_pins_are_a_known_string_mapping(tmp_path): + env_path = tmp_path / ".env" + content_hash = "a" * 64 + + assert ( + _run( + env_path, + provider="docker", + set_values=(f"sandbox.pins.{content_hash}=ami-pinned",), + ) + == 0 + ) + toml_path = tmp_path / "druks.toml" + assert _read_toml(toml_path)["sandbox"]["pins"] == {content_hash: "ami-pinned"} + assert "ami-pinned" not in env_path.read_text() + + body = toml_path.read_text().replace( + f'{content_hash} = "ami-pinned"', + f"{content_hash} = 123", + ) + toml_path.write_text(body) + with pytest.raises(ValueError, match=f"sandbox.pins.{content_hash} must be a string"): + _run(env_path) + + def test_set_updates_toml_and_rerender_preserves_the_values(tmp_path): env_path = tmp_path / ".env" _run( diff --git a/backend/tests/test_warm_host_rotation.py b/backend/tests/test_warm_host_rotation.py index c5abf935..63c901df 100644 --- a/backend/tests/test_warm_host_rotation.py +++ b/backend/tests/test_warm_host_rotation.py @@ -19,7 +19,15 @@ def __init__(self, *, lease: timedelta) -> None: self.provisions: list[str] = [] self.released: list[str] = [] - async def provision(self, *, idempotency_key: str) -> _FakeSandbox: + async def provision( + self, + *, + idempotency_key: str, + image_override: str | None, + template: str | None, + ) -> _FakeSandbox: + assert image_override is None + assert template is None self.provisions.append(idempotency_key) host_id = f"host-{len(self.provisions)}" return _FakeSandbox(id=host_id, expires_at=datetime.now(UTC) + self.lease) diff --git a/docs/configuration.md b/docs/configuration.md index b037065e..ccc34427 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -260,12 +260,24 @@ connected. | `sandbox.service_token` | Drukbox API token | | `sandbox.timeout` | Control-plane request timeout; default 180 seconds | | `sandbox.image` | Optional provider image override | +| `sandbox.pins` | Optional declared-sandbox content hashes pinned to provider artifacts | | `sandbox.browser_login_proxy` | Login-window egress proxy; empty keeps the box IP | | `sandbox.browser_login_tz` | Login-window timezone (IANA zone); empty keeps the container default | `DRUKS_SANDBOX_KEYS_DIR` remains a process environment override for the per-host SSH private-key directory. +An app normally ships each workflow's sandbox setup. To pin one declared +sandbox to an existing provider artifact, add its content hash to the same +settings plane: + +```toml +[sandbox.pins] +"" = "" +``` + +The pin bypasses template lookup for that declaration. It is optional. + `[sandbox].browser_login_proxy` sends the browser **login window** through an HTTP proxy. The login then leaves from a different IP than the box. Use it for sign-in flows that refuse a login from the box IP. Only the login window uses the diff --git a/docs/writing-an-app.md b/docs/writing-an-app.md index 051f4bbd..8421ca49 100644 --- a/docs/writing-an-app.md +++ b/docs/writing-an-app.md @@ -130,6 +130,29 @@ decorate its side-effecting operations instead. An agent called directly from the orchestration body gets its own step. An agent called inside `@step` or `run()` shares that enclosing checkpoint. +### Declare the sandbox environment + +A workflow can ship the tools its agents need as a plain shell file: + +```python +from druks.workflows import Sandbox, Workflow + + +class BuildSite(Workflow): + sandbox = Sandbox(setup="site_builder/sandbox.sh") +``` + +Place the file at `templates/sandbox.sh` inside the installed app package (for +example, `site_builder/templates/sandbox.sh`). Druks reads the raw bytes. It does +not render the file or run it during import. Drukbox builds a reusable template +from the platform base and the script. +A run waits with a visible sandbox-building phase when that template is still +building. A workflow with no declaration uses the platform base unchanged. + +The content hash of the base and script identifies the template. An operator can +pin that hash to a provider artifact through `sandbox.pins`; app authors do not +name provider images. + Completed checkpoints are reused on recovery. An interrupted operation can run again, so use provider idempotency keys for writes. Keep decisions in replayable control flow and I/O inside steps. See @@ -1138,7 +1161,7 @@ Import from concern namespaces, not from `druks.durable` or internal modules: | `druks.services` | `Service`, `ServiceConnectError`, `ServiceNotConnectedError`, `OauthClient`, `OauthExchangeError`, `OauthRefreshError` | | `druks.secrets.fields` | `EncryptedJsonField`, `SecretsMapping` | | `druks.agents` | `Agent`, `AgentOutput` | -| `druks.workflows` | `Workflow`, `Gate`, `step`, run/agent response types, lifecycle enums and workflow errors | +| `druks.workflows` | `Workflow`, `Sandbox`, `Gate`, `step`, run/agent response types, lifecycle enums and workflow errors | | `druks.db` | `Base`, `StoredSubject`, `db_session` | | `druks.schemas` | `BaseResponse` | | `druks.signals` | `subscribe` | diff --git a/druks.toml.example b/druks.toml.example index b0fa9436..015bcc08 100644 --- a/druks.toml.example +++ b/druks.toml.example @@ -22,3 +22,7 @@ service_url = "" service_token = "" image = "" timeout = 180 + +# Optional declared-sandbox content hash to provider artifact overrides. +# [sandbox.pins] +# "" = ""