From 2cc6bc017f4023e5a29e2a16c3c89116cdf52e58 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Fri, 31 Jul 2026 12:44:35 -0700 Subject: [PATCH] [eric] runner: ephemeral cloud container that boots the backend, runs one workflow, and exits --- openswarm-runner/Dockerfile | 93 ++++++++ openswarm-runner/Dockerfile.dockerignore | 10 + openswarm-runner/README.md | 68 ++++++ openswarm-runner/fly.toml | 51 ++++ openswarm-runner/runner/__init__.py | 0 openswarm-runner/runner/backend_process.py | 111 +++++++++ openswarm-runner/runner/main.py | 199 ++++++++++++++++ openswarm-runner/runner/report.py | 69 ++++++ openswarm-runner/runner/run_spec.py | 139 +++++++++++ openswarm-runner/runner/seed/__init__.py | 0 openswarm-runner/runner/seed/data_root.py | 81 +++++++ .../runner/seed/router_credentials.py | 128 ++++++++++ openswarm-runner/runner/workflow_run.py | 224 ++++++++++++++++++ .../tests/test_router_credentials.py | 111 +++++++++ openswarm-runner/tests/test_run_spec.py | 84 +++++++ openswarm-runner/tests/test_workflow_run.py | 36 +++ 16 files changed, 1404 insertions(+) create mode 100644 openswarm-runner/Dockerfile create mode 100644 openswarm-runner/Dockerfile.dockerignore create mode 100644 openswarm-runner/README.md create mode 100644 openswarm-runner/fly.toml create mode 100644 openswarm-runner/runner/__init__.py create mode 100644 openswarm-runner/runner/backend_process.py create mode 100644 openswarm-runner/runner/main.py create mode 100644 openswarm-runner/runner/report.py create mode 100644 openswarm-runner/runner/run_spec.py create mode 100644 openswarm-runner/runner/seed/__init__.py create mode 100644 openswarm-runner/runner/seed/data_root.py create mode 100644 openswarm-runner/runner/seed/router_credentials.py create mode 100644 openswarm-runner/runner/workflow_run.py create mode 100644 openswarm-runner/tests/test_router_credentials.py create mode 100644 openswarm-runner/tests/test_run_spec.py create mode 100644 openswarm-runner/tests/test_workflow_run.py diff --git a/openswarm-runner/Dockerfile b/openswarm-runner/Dockerfile new file mode 100644 index 00000000..46519520 --- /dev/null +++ b/openswarm-runner/Dockerfile @@ -0,0 +1,93 @@ +# syntax=docker/dockerfile:1 +# +# One OpenSwarm workflow run, then exit. Build context is the REPO ROOT, not this +# directory, because the image needs backend/ and requirements.lock: +# +# docker build --platform linux/amd64 -f openswarm-runner/Dockerfile -t openswarm-runner . +# +# Layout the image commits to (all three are load-bearing, backend code resolves +# them with zero changes when OPENSWARM_PACKAGED=1): +# /app/backend the FastAPI orchestrator +# /app/router 9router's standalone server, found by p_find_9router_dir() +# /app/python-env UV_PYTHON target probed by tools_lib/mcp_config.py + +ARG PYTHON_VERSION=3.13 +ARG NODE_VERSION=20 +ARG ROUTER_VERSION=0.3.60 +ARG UV_VERSION=0.11.8 + +FROM node:${NODE_VERSION}-bookworm-slim AS node + +# 9router 0.3.60 is pure JavaScript; --ignore-scripts skips a postinstall that only rebuilds a native addon the standalone server never loads. +FROM node AS router +ARG ROUTER_VERSION +WORKDIR /stage +RUN printf '{"name":"router-stage","version":"0.0.0","private":true}\n' > package.json \ + && npm install "9router@${ROUTER_VERSION}" --no-save --no-audit --no-fund --silent --ignore-scripts \ + && test -f node_modules/9router/app/server.js \ + && test -z "$(find node_modules/9router -name '*.node' -print -quit)" + +FROM python:${PYTHON_VERSION}-slim-bookworm AS uv +ARG UV_VERSION +ARG TARGETARCH +RUN set -eux; \ + apt-get update && apt-get install -y --no-install-recommends curl ca-certificates; \ + case "${TARGETARCH}" in \ + amd64) triple=x86_64-unknown-linux-gnu; sha=56dd1b66701ecb62fe896abb919444e4b83c5e8645cca953e6ddd496ff8a0feb ;; \ + arm64) triple=aarch64-unknown-linux-gnu; sha=eee8dd658d20e5ac85fec9c2326b6cbc9d83a1eef09ef07433e58698ac849591 ;; \ + *) echo "unsupported TARGETARCH ${TARGETARCH}" >&2; exit 1 ;; \ + esac; \ + curl -fsSL -o /tmp/uv.tar.gz "https://github.com/astral-sh/uv/releases/download/${UV_VERSION}/uv-${triple}.tar.gz"; \ + echo "${sha} /tmp/uv.tar.gz" | sha256sum -c -; \ + mkdir -p /stage; \ + tar -xzf /tmp/uv.tar.gz -C /stage --strip-components=1 + +# Wheels only: the runtime image ships no compiler, so a source build here is a build-time failure rather than a 3am surprise. +FROM python:${PYTHON_VERSION}-slim-bookworm AS pydeps +COPY backend/requirements.lock /tmp/requirements.lock +RUN pip install --no-cache-dir --require-hashes --only-binary=:all: \ + --prefix=/opt/pydeps -r /tmp/requirements.lock + +FROM python:${PYTHON_VERSION}-slim-bookworm +ARG NODE_VERSION + +RUN set -eux; \ + apt-get update; \ + apt-get install -y --no-install-recommends git ca-certificates; \ + rm -rf /var/lib/apt/lists/* + +COPY --from=node /usr/local/bin/node /usr/local/bin/node +COPY --from=pydeps /opt/pydeps /usr/local +COPY --from=router /stage/node_modules/9router/app /app/router + +COPY backend /app/backend +COPY openswarm-runner/runner /app/runner + +# After backend/, never before: mcp_config.resolve_command probes uv-bin last, and the repo's own copy is Mach-O. +COPY --from=uv /stage/uv /app/backend/uv-bin/uv +COPY --from=uv /stage/uvx /app/backend/uv-bin/uvx + +# mcp_config.py points uv at /python-env so an MCP server never downloads its own interpreter. +RUN set -eux; \ + mkdir -p /app/python-env/bin; \ + ln -s /usr/local/bin/python3 /app/python-env/bin/python3; \ + find /app/backend -name '__pycache__' -type d -prune -exec rm -rf {} +; \ + useradd --create-home --uid 10001 --shell /usr/sbin/nologin runner; \ + mkdir -p /data; \ + chown -R runner:runner /app /data + +USER runner +WORKDIR /app +ENV HOME=/home/runner \ + PYTHONPATH=/app \ + PYTHONUNBUFFERED=1 \ + PYTHONDONTWRITEBYTECODE=1 \ + OPENSWARM_HEADLESS=1 \ + OPENSWARM_PACKAGED=1 \ + OPENSWARM_DATA_ROOT=/data/openswarm \ + OPENSWARM_HOST=127.0.0.1 \ + OPENSWARM_PORT=8324 \ + DATA_DIR=/data/9router \ + NODE_ENV=production + +ENTRYPOINT ["python3", "-m", "runner.main"] diff --git a/openswarm-runner/Dockerfile.dockerignore b/openswarm-runner/Dockerfile.dockerignore new file mode 100644 index 00000000..a6e1638a --- /dev/null +++ b/openswarm-runner/Dockerfile.dockerignore @@ -0,0 +1,10 @@ +* +!backend +!openswarm-runner/runner +backend/data +backend/.venv +backend/uv-bin +backend/tests +**/__pycache__ +**/*.pyc +**/.DS_Store diff --git a/openswarm-runner/README.md b/openswarm-runner/README.md new file mode 100644 index 00000000..cb99f1d6 --- /dev/null +++ b/openswarm-runner/README.md @@ -0,0 +1,68 @@ +# openswarm-runner + +One ephemeral Linux container that executes ONE OpenSwarm workflow run and exits. +One Fly Firecracker machine per run, no state kept. + +## Build + +The build context is the **repo root**, not this directory (the image needs `backend/` +and `backend/requirements.lock`): + +```bash +docker build --platform linux/amd64 -f openswarm-runner/Dockerfile -t openswarm-runner . +``` + +## Run + +The container is told everything it needs by one JSON run spec in `OPENSWARM_RUN_SPEC` +(or a path in `OPENSWARM_RUN_SPEC_FILE`). See `runner/run_spec.py` for the typed shape. + +```json +{ + "run_id": "cr_01J...", + "workflow": { "id": "wf_1", "title": "Daily digest", "model": "opus-5", + "steps": [{ "text": "summarize my inbox" }] }, + "credentials": [ + { "provider": "claude", "auth_type": "oauth", + "access_token": "", + "expires_at": "2026-07-31T20:00:00Z" } + ], + "callback": { "url": "https://api.openswarm.com/api/cloud-runs/cr_01J.../report", + "token": "" }, + "max_run_seconds": 1800 +} +``` + +Exit codes: `0` ok, `1` runner crash, `2` bad spec, `3` credential expired on arrival, +`4` backend never came up, `5` workflow failed, `6` wall-clock cap hit. + +## The credential rule + +**A `providerConnections[]` entry this runner writes never contains a `refreshToken`.** +9Router's refresh dispatcher bails on `if (!b || !b.refreshToken) return null`, so +omitting the field is what makes the container incapable of rotating the user's grant. +If it ever rotated, the user's laptop would be left replaying a dead token and the +provider would revoke the whole grant family. + +Two independent walls enforce it, and a third makes a leak require deleting the code +that builds the entry: + +1. `ProviderCredential` forbids extra fields, so a spec carrying `refreshToken` fails + validation before the backend boots. +2. `assert_no_refresh_token` re-reads the assembled db payload just before the write. +3. `router_connection` assembles the entry from a fixed key list, never a passthrough. + +All three live in `runner/seed/router_credentials.py`. + +An access token that arrives expired fails the run (exit 3). The runner never refreshes. + +## Test + +```bash +PYTHONPATH=.:openswarm-runner backend/.venv/bin/python3 -m pytest openswarm-runner/tests -q +``` + +## Deploy + +Not deployed. `fly.toml` is written but never applied; read its header first, the app +has to be created onto its own isolated private network by hand before any deploy. diff --git a/openswarm-runner/fly.toml b/openswarm-runner/fly.toml new file mode 100644 index 00000000..ff2ed454 --- /dev/null +++ b/openswarm-runner/fly.toml @@ -0,0 +1,51 @@ +# openswarm-runner: one ephemeral Firecracker machine per workflow run. Boots the +# backend headless, runs the workflow, reports, exits. Nothing here is long-lived. +# +# THIS APP MUST NOT SHARE THE TRUSTED 6PN MESH with openswarm-cloud / openswarm-edge. +# The agent inside has Bash and executes user prose, so from in here +# `curl http://openswarm-cloud.internal:8080` must resolve to nothing. Fly decides an +# app's private network AT CREATE TIME and fly.toml cannot express it, so the app is +# created once, by hand, onto its own isolated network: +# +# fly apps create openswarm-runner --org openswarm --network openswarm-runner-isolated +# fly deploy . --config openswarm-runner/fly.toml --dockerfile openswarm-runner/Dockerfile +# +# (deploy runs from the REPO ROOT: the image needs backend/ in its build context.) +# Verify the isolation after the first deploy, do not assume it: +# fly ssh console -a openswarm-runner -C "getent hosts openswarm-cloud.internal" # must fail +# +# There is deliberately no [http_service] and no [[services]]: the runner takes no +# inbound traffic and gets no public IP. It reaches the control plane outbound over +# the public internet with the callback token in the run spec, which is why the two +# do not need a shared private network in the first place. +# +# Machines are created per run by the control plane (Machines API, auto_destroy=true, +# run spec passed as OPENSWARM_RUN_SPEC). This file is the app-level shape they inherit. + +app = 'openswarm-runner' +primary_region = 'iad' +kill_signal = 'SIGTERM' +kill_timeout = '30s' + +[build] + dockerfile = 'Dockerfile' + +[env] + # Hard wall-clock cap, enforced twice inside the container: the poll loop stops the + # run at this mark, and an independent thread kills the process 90s later. A run + # spec asking for more is clamped down to this, never up. + RUNNER_MAX_RUN_SECONDS = '1800' + OPENSWARM_HEADLESS = '1' + OPENSWARM_PACKAGED = '1' + OPENSWARM_DATA_ROOT = '/data/openswarm' + OPENSWARM_HOST = '127.0.0.1' + OPENSWARM_PORT = '8324' + DATA_DIR = '/data/9router' + +# No [[mounts]]: a run's state is garbage the moment it ends, and an ephemeral rootfs +# means one run cannot leave a credential lying around for the next tenant to find. + +[[vm]] + cpu_kind = 'shared' + cpus = 2 + memory_mb = 4096 diff --git a/openswarm-runner/runner/__init__.py b/openswarm-runner/runner/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/openswarm-runner/runner/backend_process.py b/openswarm-runner/runner/backend_process.py new file mode 100644 index 00000000..b4114933 --- /dev/null +++ b/openswarm-runner/runner/backend_process.py @@ -0,0 +1,111 @@ +"""Boot the OpenSwarm backend inside the container and wait until it answers. + +Spawned the same way the desktop shell spawns it (`python -m uvicorn backend.main:app` +on loopback) so the cloud path and the laptop path are the same code on the same +socket. Loopback, not 0.0.0.0: every caller of this API lives in this container, and +the agent running inside it has Bash, so there is no reason to publish the port onto +the machine's private network. +""" + +import os +import subprocess +import time +from typing import Dict, List, Optional + +import httpx +from pydantic import BaseModel, ConfigDict, InstanceOf +from typeguard import typechecked + +HOST = "127.0.0.1" +HEALTH_PATH = "/api/health/check" +# 9Router is a Next.js standalone server: it binds `process.env.HOSTNAME || '0.0.0.0'`, and Docker sets HOSTNAME to the container id, so left alone it listens on the container's eth0 address and every probe of 127.0.0.1:20128 gets ECONNREFUSED. Set on the backend's env because the backend is what spawns node. +ROUTER_BIND_HOSTNAME = "127.0.0.1" +AUTH_TOKEN_FILENAME = "auth.token" +# The backend imports the whole app graph before it binds; on a cold Fly machine that has been measured in tens of seconds, so the budget is generous rather than tight. +BOOT_TIMEOUT_SECONDS = 120.0 +SHUTDOWN_GRACE_SECONDS = 10.0 + + +class BackendUnavailable(RuntimeError): + """The backend never came up, or died while we were using it.""" + + +class BackendProcess(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + process: InstanceOf[subprocess.Popen] + base_url: str + token: str + + @typechecked + def headers(self) -> Dict[str, str]: + return {"Authorization": f"Bearer {self.token}"} + + @typechecked + def is_alive(self) -> bool: + return self.process.poll() is None + + +@typechecked +def p_command(port: int) -> List[str]: + return ["python3", "-m", "uvicorn", "backend.main:app", "--host", HOST, "--port", str(port)] + + +@typechecked +def p_read_auth_token(data_root: str) -> str: + """The backend mints this before it binds, so by the time health passes the file exists.""" + path = os.path.join(data_root, AUTH_TOKEN_FILENAME) + try: + with open(path, "r", encoding="utf-8") as handle: + return handle.read().strip() + except OSError as exc: + raise BackendUnavailable(f"backend is up but its auth token is unreadable at {path}: {exc}") from exc + + +@typechecked +def start_backend(app_root: str, data_root: str, port: int, deadline: float) -> BackendProcess: + """Spawn the backend and block until it answers health, or raise BackendUnavailable.""" + environment = dict(os.environ) + environment["OPENSWARM_DATA_ROOT"] = data_root + environment["OPENSWARM_HEADLESS"] = "1" + environment["OPENSWARM_PORT"] = str(port) + environment["OPENSWARM_HOST"] = HOST + environment["HOSTNAME"] = ROUTER_BIND_HOSTNAME + environment["PYTHONPATH"] = app_root + + process = subprocess.Popen(p_command(port), cwd=app_root, env=environment) + base_url = f"http://{HOST}:{port}" + budget = min(time.monotonic() + BOOT_TIMEOUT_SECONDS, deadline) + + with httpx.Client(timeout=2.0) as client: + while time.monotonic() < budget: + if process.poll() is not None: + raise BackendUnavailable(f"backend exited during startup with code {process.returncode}") + try: + healthy = client.get(f"{base_url}{HEALTH_PATH}").status_code == 200 + except httpx.HTTPError: + healthy = False + if healthy: + try: + token = p_read_auth_token(data_root) + except BackendUnavailable: + stop_backend(process) + raise + return BackendProcess(process=process, base_url=base_url, token=token) + time.sleep(0.25) + + stop_backend(process) + raise BackendUnavailable(f"backend did not answer {HEALTH_PATH} within its startup budget") + + +@typechecked +def stop_backend(process: Optional[subprocess.Popen]) -> None: + """SIGTERM then SIGKILL. The machine is about to die anyway; this just stops the logs mid-sentence.""" + if process is None or process.poll() is not None: + return + process.terminate() + try: + process.wait(timeout=SHUTDOWN_GRACE_SECONDS) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=SHUTDOWN_GRACE_SECONDS) diff --git a/openswarm-runner/runner/main.py b/openswarm-runner/runner/main.py new file mode 100644 index 00000000..411359fb --- /dev/null +++ b/openswarm-runner/runner/main.py @@ -0,0 +1,199 @@ +"""One workflow run, one container, one exit code. + +Boots the backend headless, executes the workflow the control plane asked for, +reports the result, and dies. Nothing here is meant to survive the run. +""" + +import logging +import os +import subprocess +import sys +import threading +import time +from datetime import datetime, timezone +from typing import Optional + +from pydantic import BaseModel, ConfigDict +from typeguard import typechecked + +from runner.backend_process import BackendProcess, BackendUnavailable, start_backend, stop_backend +from runner.report import RunReport, send_report +from runner.run_spec import CallbackTarget, InvalidRunSpec, RunSpec, load_run_spec +from runner.seed.data_root import seed_data_root +from runner.seed.router_credentials import write_router_db +from runner.workflow_run import RunOutcome, RunProgress, WorkflowRunFailed, execute_workflow + +EXIT_OK = 0 +EXIT_INTERNAL = 1 +EXIT_BAD_SPEC = 2 +EXIT_CREDENTIAL_EXPIRED = 3 +EXIT_BACKEND_UNAVAILABLE = 4 +EXIT_WORKFLOW_FAILED = 5 +EXIT_DEADLINE = 6 + +DEFAULT_APP_ROOT = "/app" +DEFAULT_DATA_ROOT = "/data/openswarm" +DEFAULT_ROUTER_DATA_DIR = "/data/9router" +DEFAULT_PORT = 8324 +# Slack between the soft deadline (stop the run, report it) and the hard one (kill the process). +REPORT_GRACE_SECONDS = 90.0 +# Ceiling the control plane cannot raise. A cap a caller can override is not a cap. +MAX_RUN_SECONDS_ENV = "RUNNER_MAX_RUN_SECONDS" +DEFAULT_MAX_RUN_SECONDS = 1800 + +logger = logging.getLogger("runner") + + +class Heartbeat(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + run_id: str + interval_seconds: float + callback: Optional[CallbackTarget] = None + last_sent: float = 0.0 + + @typechecked + def maybe_send(self, progress: RunProgress) -> None: + now = time.monotonic() + if now - self.last_sent < self.interval_seconds: + return + self.last_sent = now + send_report(self.callback, RunReport( + run_id=self.run_id, + phase="heartbeat", + status=progress.status, + active_step_idx=progress.active_step_idx, + last_tool_label=progress.last_tool_label, + )) + + +@typechecked +def effective_max_run_seconds(spec: RunSpec) -> int: + """The shorter of what the job asked for and what this machine's config allows.""" + try: + ceiling = int(os.environ.get(MAX_RUN_SECONDS_ENV, "") or DEFAULT_MAX_RUN_SECONDS) + except ValueError: + ceiling = DEFAULT_MAX_RUN_SECONDS + return max(60, min(spec.max_run_seconds, ceiling)) + + +@typechecked +def arm_hard_stop(seconds: float) -> None: + """Independent backstop on machine-seconds; fires even if the graceful path is wedged.""" + def p_fire() -> None: + time.sleep(seconds) + logger.error("hard wall-clock stop after %.0fs, killing the run", seconds) + os._exit(EXIT_DEADLINE) + + threading.Thread(target=p_fire, daemon=True, name="hard-stop").start() + + +@typechecked +def p_fail(spec: Optional[RunSpec], status: str, message: str, code: int) -> int: + logger.error("%s: %s", status, message) + send_report( + spec.callback if spec else None, + RunReport( + run_id=spec.run_id if spec else "unknown", + phase="finished", + status=status, + exit_code=code, + error=message, + ), + ) + return code + + +@typechecked +def p_exit_code_for(outcome: RunOutcome) -> int: + if outcome.status in ("success", "ran_late"): + return EXIT_OK + if outcome.status == "timed_out": + return EXIT_DEADLINE + return EXIT_WORKFLOW_FAILED + + +@typechecked +def p_run(spec: RunSpec, deadline: float) -> int: + now = datetime.now(timezone.utc) + expired = spec.expired_credentials(now) + if expired: + names = ", ".join(credential.provider for credential in expired) + return p_fail( + spec, + "failure", + f"access token for {names} is expired or about to expire; the runner never refreshes, " + "so the control plane must re-issue it", + EXIT_CREDENTIAL_EXPIRED, + ) + + app_root = os.environ.get("OPENSWARM_APP_ROOT", DEFAULT_APP_ROOT) + data_root = os.environ.get("OPENSWARM_DATA_ROOT", DEFAULT_DATA_ROOT) + router_data_dir = os.environ.get("DATA_DIR", DEFAULT_ROUTER_DATA_DIR) + port = int(os.environ.get("OPENSWARM_PORT", str(DEFAULT_PORT))) + + write_router_db(router_data_dir, spec.credentials, now) + seed_data_root(data_root, spec) + logger.info("seeded data root %s and router db in %s", data_root, router_data_dir) + + backend: Optional[BackendProcess] = None + process: Optional[subprocess.Popen] = None + try: + backend = start_backend(app_root, data_root, port, deadline) + process = backend.process + logger.info("backend healthy at %s", backend.base_url) + + send_report(spec.callback, RunReport(run_id=spec.run_id, phase="started", status="running")) + heartbeat = Heartbeat( + run_id=spec.run_id, + interval_seconds=float(spec.callback.heartbeat_seconds) if spec.callback else 30.0, + callback=spec.callback, + ) + outcome = execute_workflow(backend, spec.workflow.id, deadline, heartbeat.maybe_send) + except BackendUnavailable as exc: + return p_fail(spec, "failure", str(exc), EXIT_BACKEND_UNAVAILABLE) + except WorkflowRunFailed as exc: + return p_fail(spec, "failure", str(exc), EXIT_WORKFLOW_FAILED) + finally: + stop_backend(process) + + code = p_exit_code_for(outcome) + logger.info("run %s finished as %s (exit %d)", spec.run_id, outcome.status, code) + send_report(spec.callback, RunReport( + run_id=spec.run_id, + phase="finished", + status=outcome.status, + exit_code=code, + error=outcome.error, + cost_usd=outcome.cost_usd, + answer=outcome.answer, + transcript=outcome.transcript, + system_notices=outcome.system_notices, + )) + return code + + +@typechecked +def main() -> int: + logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname).1s %(name)s: %(message)s", + datefmt="%H:%M:%S", + ) + try: + spec = load_run_spec() + except InvalidRunSpec as exc: + return p_fail(None, "failure", str(exc), EXIT_BAD_SPEC) + + budget = effective_max_run_seconds(spec) + arm_hard_stop(budget + REPORT_GRACE_SECONDS) + deadline = time.monotonic() + budget + try: + return p_run(spec, deadline) + except Exception as exc: + logger.exception("runner crashed") + return p_fail(spec, "failure", f"runner crashed: {exc}", EXIT_INTERNAL) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/openswarm-runner/runner/report.py b/openswarm-runner/runner/report.py new file mode 100644 index 00000000..bb5c3de1 --- /dev/null +++ b/openswarm-runner/runner/report.py @@ -0,0 +1,69 @@ +"""Tell the control plane what happened. The terminal report is the run's only receipt.""" + +import logging +import time +from typing import List, Literal, Optional + +import httpx +from pydantic import BaseModel, ConfigDict, Field +from typeguard import typechecked + +from runner.run_spec import CallbackTarget + +logger = logging.getLogger(__name__) + +TERMINAL_ATTEMPTS = 5 +TERMINAL_BACKOFF_SECONDS = 2.0 +REQUEST_TIMEOUT_SECONDS = 15.0 + + +class RunReport(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + run_id: str + phase: Literal["started", "heartbeat", "finished"] + status: str + exit_code: Optional[int] = None + error: Optional[str] = None + cost_usd: float = 0.0 + active_step_idx: Optional[int] = None + last_tool_label: Optional[str] = None + answer: str = "" + transcript: str = "" + # The backend calls a run "success" even when the provider rejected the token; these are how the control plane sees that. + system_notices: List[str] = Field(default_factory=list) + + +@typechecked +def p_post_once(callback: CallbackTarget, report: RunReport) -> bool: + try: + with httpx.Client(timeout=REQUEST_TIMEOUT_SECONDS) as client: + response = client.post( + callback.url, + headers={"Authorization": f"Bearer {callback.token}"}, + json=report.model_dump(mode="json"), + ) + if response.status_code < 300: + return True + logger.warning("report %s rejected with HTTP %s", report.phase, response.status_code) + return False + except httpx.HTTPError as exc: + logger.warning("report %s failed to send: %s", report.phase, exc) + return False + + +@typechecked +def send_report(callback: Optional[CallbackTarget], report: RunReport) -> bool: + """Post a report. Terminal reports retry; heartbeats get one shot and are never retried.""" + if callback is None: + logger.info("no callback configured; %s report kept local: %s", report.phase, report.status) + return True + if report.phase != "finished": + return p_post_once(callback, report) + for attempt in range(TERMINAL_ATTEMPTS): + if p_post_once(callback, report): + return True + if attempt + 1 < TERMINAL_ATTEMPTS: + time.sleep(TERMINAL_BACKOFF_SECONDS * (attempt + 1)) + logger.error("terminal report for run %s never landed after %d attempts", report.run_id, TERMINAL_ATTEMPTS) + return False diff --git a/openswarm-runner/runner/run_spec.py b/openswarm-runner/runner/run_spec.py new file mode 100644 index 00000000..8a63318f --- /dev/null +++ b/openswarm-runner/runner/run_spec.py @@ -0,0 +1,139 @@ +"""Typed description of the single workflow run this container exists to execute. + +The control plane hands the container exactly one of these (JSON in +OPENSWARM_RUN_SPEC, or a path in OPENSWARM_RUN_SPEC_FILE) and nothing else. Every +field is validated before the backend boots, so a malformed job dies in under a +second instead of burning a machine-minute discovering it. +""" + +import json +import os +from datetime import datetime, timedelta, timezone +from typing import List, Literal, Optional + +from pydantic import BaseModel, ConfigDict, Field, model_validator +from typeguard import typechecked + +from backend.apps.workflows.models import Workflow + +SPEC_ENV = "OPENSWARM_RUN_SPEC" +SPEC_FILE_ENV = "OPENSWARM_RUN_SPEC_FILE" + +# Headroom the access token must still have on arrival. The control plane refreshes right before dispatch; anything thinner than this means its clock or its queue is broken, and we must not paper over that by refreshing ourselves. +MIN_TOKEN_LIFETIME = timedelta(minutes=2) + + +class InvalidRunSpec(ValueError): + """The control plane handed us something we refuse to run.""" + + +class ProviderCredential(BaseModel): + """One already-refreshed provider credential, spendable but not rotatable. + + `extra="forbid"` is the first of two walls keeping a refresh token out of this + container: a payload carrying `refreshToken` fails validation here and the run + dies loudly. See runner.router_credentials for the second wall and the why. + """ + + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + provider: str = Field(min_length=1) + auth_type: Literal["oauth", "api_key"] + label: str = "OpenSwarm cloud run" + access_token: Optional[str] = None + api_key: Optional[str] = None + expires_at: Optional[datetime] = None + scope: Optional[str] = None + + @model_validator(mode="after") + def p_require_matching_secret(self) -> "ProviderCredential": + if self.auth_type == "oauth": + if not self.access_token: + raise ValueError(f"credential for {self.provider!r} is oauth but carries no access_token") + if self.expires_at is None: + raise ValueError(f"credential for {self.provider!r} is oauth but carries no expires_at") + if self.api_key: + raise ValueError(f"credential for {self.provider!r} carries both an access_token and an api_key") + else: + if not self.api_key: + raise ValueError(f"credential for {self.provider!r} is api_key but carries no api_key") + if self.access_token: + raise ValueError(f"credential for {self.provider!r} carries both an access_token and an api_key") + return self + + @typechecked + def remaining_lifetime(self, now: datetime) -> Optional[timedelta]: + """How long this credential is still good for; None when it cannot expire.""" + if self.expires_at is None: + return None + return self.expires_at.astimezone(timezone.utc) - now.astimezone(timezone.utc) + + +class CallbackTarget(BaseModel): + """Where the run reports back. The token is a dedicated two-party secret, never a user credential.""" + + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + url: str = Field(min_length=1) + token: str = Field(min_length=1) + heartbeat_seconds: int = Field(default=30, ge=5, le=300) + + +class RunSpec(BaseModel): + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + run_id: str = Field(min_length=1) + workflow: Workflow + credentials: List[ProviderCredential] = Field(min_length=1) + callback: Optional[CallbackTarget] = None + # Hard wall-clock ceiling. Fly bills by machine-second, so an agent that wedges must cost a bounded amount. + max_run_seconds: int = Field(default=1800, ge=60, le=7200) + + @typechecked + def expired_credentials(self, now: datetime) -> List[ProviderCredential]: + """Credentials too close to expiry to spend. The runner cannot refresh, so this is fatal, not a retry.""" + stale: List[ProviderCredential] = [] + for credential in self.credentials: + remaining = credential.remaining_lifetime(now) + if remaining is not None and remaining < MIN_TOKEN_LIFETIME: + stale.append(credential) + return stale + + @typechecked + def workflow_for_disk(self) -> Workflow: + """The workflow as this container should see it: one run, never a schedule. + + A cloud-executed workflow arrives with its schedule still configured. Left + enabled, the container's own scheduler would fire it a second time inside + the box, so the timer is stripped here rather than trusted to stay off. + """ + copy = self.workflow.model_copy(deep=True) + copy.schedule.enabled = False + copy.deleted_at = None + copy.draft_steps = None + copy.next_run_at = None + return copy + + +@typechecked +def load_run_spec() -> RunSpec: + """Parse the run spec from the environment, or raise InvalidRunSpec with a legible reason.""" + raw = os.environ.get(SPEC_ENV, "").strip() + source = SPEC_ENV + if not raw: + path = os.environ.get(SPEC_FILE_ENV, "").strip() + if not path: + raise InvalidRunSpec(f"no run spec: set {SPEC_ENV} to JSON or {SPEC_FILE_ENV} to a file path") + source = f"{SPEC_FILE_ENV}={path}" + try: + with open(path, "r", encoding="utf-8") as handle: + raw = handle.read() + except OSError as exc: + raise InvalidRunSpec(f"cannot read run spec from {source}: {exc}") from exc + + try: + return RunSpec.model_validate_json(raw) + except json.JSONDecodeError as exc: + raise InvalidRunSpec(f"run spec from {source} is not valid JSON: {exc}") from exc + except ValueError as exc: + raise InvalidRunSpec(f"run spec from {source} is not a valid RunSpec: {exc}") from exc diff --git a/openswarm-runner/runner/seed/__init__.py b/openswarm-runner/runner/seed/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/openswarm-runner/runner/seed/data_root.py b/openswarm-runner/runner/seed/data_root.py new file mode 100644 index 00000000..413bf53a --- /dev/null +++ b/openswarm-runner/runner/seed/data_root.py @@ -0,0 +1,81 @@ +"""Lay down the backend's data dir before it boots, so the run is ready on the first tick. + +Everything here is written pre-boot on purpose: the workflow store and the settings +store both load from disk once at startup, so seeding files is cheaper and more +deterministic than replaying create/PATCH calls over HTTP (no aux LLM naming call, +no schedule normalization, no chance of the container inventing a second workflow). +""" + +import json +import os +import tempfile +from typing import Any, Dict + +from typeguard import typechecked + +from backend.apps.settings.models import AppSettings +from runner.run_spec import RunSpec + +# 9Router provider id -> the AppSettings field the backend reads a raw key from. +API_KEY_SETTINGS_FIELD: Dict[str, str] = { + "anthropic": "anthropic_api_key", + "openai": "openai_api_key", + "gemini": "google_api_key", + "google": "google_api_key", + "openrouter": "openrouter_api_key", +} + +# One synthetic id for every cloud run. Without it each ephemeral container mints a fresh uuid and analytics sees a brand-new "install" per run. +CLOUD_RUNNER_INSTALLATION_ID = "openswarm-cloud-runner" + + +@typechecked +def p_write_json(path: str, payload: Any) -> None: + """Atomic, owner-only write; these files hold API keys.""" + directory = os.path.dirname(path) + os.makedirs(directory, mode=0o700, exist_ok=True) + handle, temp_path = tempfile.mkstemp(dir=directory, prefix=".seed-", suffix=".json") + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump(payload, stream, indent=2) + os.chmod(temp_path, 0o600) + os.replace(temp_path, path) + except BaseException: + if os.path.exists(temp_path): + os.unlink(temp_path) + raise + + +@typechecked +def settings_for_run(spec: RunSpec) -> AppSettings: + """The AppSettings a cloud run needs: this workflow's model, this run's keys, no telemetry.""" + settings = AppSettings() + settings.default_model = spec.workflow.model + settings.connection_mode = "own_key" + settings.analytics_opt_in = False + settings.installation_id = CLOUD_RUNNER_INSTALLATION_ID + for credential in spec.credentials: + if credential.auth_type != "api_key": + continue + field = API_KEY_SETTINGS_FIELD.get(credential.provider) + if field is None: + raise ValueError( + f"no settings field for api_key provider {credential.provider!r}; " + f"supported: {', '.join(sorted(set(API_KEY_SETTINGS_FIELD)))}" + ) + setattr(settings, field, credential.api_key) + return settings + + +@typechecked +def seed_data_root(data_root: str, spec: RunSpec) -> None: + """Write the workflow record and the settings file the backend will read at boot.""" + workflow = spec.workflow_for_disk() + p_write_json( + os.path.join(data_root, "workflows", f"{workflow.id}.json"), + workflow.model_dump(mode="json"), + ) + p_write_json( + os.path.join(data_root, "settings", "settings.json"), + settings_for_run(spec).model_dump(mode="json"), + ) diff --git a/openswarm-runner/runner/seed/router_credentials.py b/openswarm-runner/runner/seed/router_credentials.py new file mode 100644 index 00000000..5d2d719f --- /dev/null +++ b/openswarm-runner/runner/seed/router_credentials.py @@ -0,0 +1,128 @@ +"""Write 9Router's credential db for one cloud run, without a refresh token. Ever. + +This container must be structurally incapable of rotating the user's OAuth grant. +9Router's refresh dispatcher bails on `if (!b || !b.refreshToken) return null`, so a +providerConnections entry with no such field is one it can spend and never rotate. +If the runner did rotate, the user's laptop would be left replaying a dead token and +the provider would revoke their entire grant family. That is the whole safety +property of this file, not a style preference. + +Two independent walls hold it up: + 1. ProviderCredential forbids extra fields, so a payload carrying `refreshToken` + never parses into the process at all. + 2. assert_no_refresh_token re-reads the assembled payload just before the write + and refuses anything whose key names a refresh token, however it got there. +""" + +import json +import os +import tempfile +from datetime import datetime, timezone +from typing import Any, Dict, List +from uuid import uuid4 + +from typeguard import typechecked + +from runner.run_spec import ProviderCredential + +DB_FILENAME = "db.json" + +# Normalized key fragment that must never appear anywhere in the db we write. +FORBIDDEN_KEY_FRAGMENT = "refreshtoken" + + +class RefreshTokenLeak(RuntimeError): + """A refresh token reached the credential writer. Fail the run rather than write it.""" + + +@typechecked +def p_normalize_key(key: str) -> str: + return "".join(char for char in key.lower() if char.isalnum()) + + +@typechecked +def p_iso(moment: datetime) -> str: + """9Router timestamps are ISO-8601 UTC with a Z suffix; match it exactly.""" + return moment.astimezone(timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z") + + +@typechecked +def assert_no_refresh_token(payload: Any, path: str = "$") -> None: + """Raise if any key anywhere under `payload` names a refresh token.""" + if isinstance(payload, dict): + for key, value in payload.items(): + if FORBIDDEN_KEY_FRAGMENT in p_normalize_key(str(key)): + raise RefreshTokenLeak( + f"refusing to write a 9Router db containing a refresh token at {path}.{key}" + ) + assert_no_refresh_token(value, f"{path}.{key}") + elif isinstance(payload, list): + for index, value in enumerate(payload): + assert_no_refresh_token(value, f"{path}[{index}]") + + +@typechecked +def router_connection(credential: ProviderCredential, now: datetime) -> Dict[str, Any]: + """Build one providerConnections entry from an allow-list of keys, never a passthrough.""" + if credential.auth_type != "oauth": + raise ValueError(f"credential for {credential.provider!r} is not an oauth connection") + entry: Dict[str, Any] = { + "id": str(uuid4()), + "provider": credential.provider, + "authType": "oauth", + "name": credential.label, + "priority": 1, + "isActive": True, + "createdAt": p_iso(now), + "updatedAt": p_iso(now), + "accessToken": credential.access_token, + "testStatus": "active", + } + if credential.expires_at is not None: + entry["expiresAt"] = p_iso(credential.expires_at) + if credential.scope: + entry["scope"] = credential.scope + return entry + + +@typechecked +def router_db_payload(credentials: List[ProviderCredential], now: datetime) -> Dict[str, Any]: + """A complete 9Router db seeded with this run's subscription connections and nothing else.""" + return { + "providerConnections": [ + router_connection(credential, now) + for credential in credentials + if credential.auth_type == "oauth" + ], + "providerNodes": [], + "proxyPools": [], + "modelAliases": {}, + "mitmAlias": {}, + "combos": [], + "apiKeys": [], + "customModels": [], + "pricing": {}, + "settings": {}, + } + + +@typechecked +def write_router_db(data_dir: str, credentials: List[ProviderCredential], now: datetime) -> str: + """Write $DATA_DIR/db.json owner-only and return its path.""" + payload = router_db_payload(credentials, now) + assert_no_refresh_token(payload) + + os.makedirs(data_dir, mode=0o700, exist_ok=True) + os.chmod(data_dir, 0o700) + path = os.path.join(data_dir, DB_FILENAME) + handle, temp_path = tempfile.mkstemp(dir=data_dir, prefix=".db-", suffix=".json") + try: + with os.fdopen(handle, "w", encoding="utf-8") as stream: + json.dump(payload, stream, indent=2) + os.chmod(temp_path, 0o600) + os.replace(temp_path, path) + except BaseException: + if os.path.exists(temp_path): + os.unlink(temp_path) + raise + return path diff --git a/openswarm-runner/runner/workflow_run.py b/openswarm-runner/runner/workflow_run.py new file mode 100644 index 00000000..7e71f8b5 --- /dev/null +++ b/openswarm-runner/runner/workflow_run.py @@ -0,0 +1,224 @@ +"""Drive one workflow through the backend's own HTTP surface and collect its result. + +Deliberately no shortcuts into agent_manager: the cloud run fires the same route the +Run button fires, so the MCP gate, action filtering, provider routing and history all +behave exactly as they do on a laptop. +""" + +import json +import time +from typing import Any, Callable, Dict, List, Optional + +import httpx +from pydantic import BaseModel, ConfigDict, Field +from typeguard import typechecked + +from runner.backend_process import BackendProcess + +TERMINAL_STATUSES = ("success", "failure", "ran_late", "skipped") +POLL_INTERVAL_SECONDS = 1.0 +TRANSCRIPT_MAX_CHARS = 14000 + + +class WorkflowRunFailed(RuntimeError): + """The backend refused to start the run at all.""" + + +class RunProgress(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + run_id: str + status: str + active_step_idx: Optional[int] = None + last_tool_label: Optional[str] = None + + +class RunOutcome(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + run_id: str + status: str + error: Optional[str] = None + cost_usd: float = 0.0 + session_id: Optional[str] = None + transcript: str = "" + answer: str = "" + system_notices: List[str] = Field(default_factory=list) + + +@typechecked +def p_block_text(block: Dict[str, Any]) -> str: + kind = block.get("type") + if kind == "text": + return str(block.get("text") or "") + if kind == "tool_use": + return f"[tool {block.get('name')}] {json.dumps(block.get('input') or {})[:300]}" + if kind == "tool_result": + inner = block.get("content") + return f"[result] {inner if isinstance(inner, str) else json.dumps(inner)[:300]}" + return "" + + +@typechecked +def p_message_text(message: Dict[str, Any]) -> str: + content = message.get("content") + if isinstance(content, str): + return content + if isinstance(content, list): + parts = [p_block_text(block) for block in content if isinstance(block, dict)] + return "\n".join(part for part in parts if part) + return "" + + +@typechecked +def render_transcript(messages: List[Dict[str, Any]]) -> str: + """Role-tagged flatten, tail-biased so the end of a long run always survives the cap.""" + lines: List[str] = [] + for message in messages: + if message.get("hidden"): + continue + text = p_message_text(message).strip() + if text: + lines.append(f"{str(message.get('role') or '?').upper()}: {text}") + joined = "\n\n".join(lines) + if len(joined) > TRANSCRIPT_MAX_CHARS: + return "...(earlier turns trimmed)...\n\n" + joined[-TRANSCRIPT_MAX_CHARS:] + return joined + + +@typechecked +def final_answer(messages: List[Dict[str, Any]]) -> str: + """Last visible assistant text: the thing a user actually asked the workflow for.""" + for message in reversed(messages): + if message.get("hidden") or message.get("role") != "assistant": + continue + text = p_message_text(message).strip() + if text: + return text + return "" + + +@typechecked +def system_notices(messages: List[Dict[str, Any]]) -> List[str]: + """Every system-role bubble in the session. + + The backend appends a system message only when something went wrong (a dead + provider token, a run error, a blocked tool), and it does NOT fail the run for + those, so a workflow whose credential was rejected still comes back "success". + Keyed on the typed role, not on the prose, and reported rather than judged: the + control plane decides what a notice means for billing and retries. + """ + notices: List[str] = [] + for message in messages: + if message.get("role") != "system" or message.get("hidden"): + continue + text = p_message_text(message).strip() + if text: + notices.append(text) + return notices + + +@typechecked +def p_get_json(client: httpx.Client, backend: BackendProcess, path: str) -> Dict[str, Any]: + response = client.get(f"{backend.base_url}{path}", headers=backend.headers()) + response.raise_for_status() + payload = response.json() + return payload if isinstance(payload, dict) else {} + + +@typechecked +def trigger_run(client: httpx.Client, backend: BackendProcess, workflow_id: str) -> str: + response = client.post( + f"{backend.base_url}/api/workflows/{workflow_id}/run", + headers=backend.headers(), + json={}, + ) + response.raise_for_status() + body = response.json() + run_id = str(body.get("run_id") or "") + if not run_id: + raise WorkflowRunFailed( + f"backend accepted the trigger but never created a run for workflow {workflow_id}" + ) + if body.get("status") == "failure": + raise WorkflowRunFailed(str(body.get("error") or "run failed immediately")) + return run_id + + +@typechecked +def p_find_run(client: httpx.Client, backend: BackendProcess, workflow_id: str, run_id: str) -> Dict[str, Any]: + body = p_get_json(client, backend, f"/api/workflows/{workflow_id}/runs?limit=50") + for record in body.get("runs") or []: + if isinstance(record, dict) and record.get("id") == run_id: + return record + return {} + + +@typechecked +def p_stop_run(client: httpx.Client, backend: BackendProcess, run_id: str) -> None: + try: + client.post(f"{backend.base_url}/api/workflows/runs/{run_id}/stop", headers=backend.headers()) + except httpx.HTTPError: + pass + + +@typechecked +def p_collect_session(client: httpx.Client, backend: BackendProcess, session_id: str) -> List[Dict[str, Any]]: + try: + body = p_get_json(client, backend, f"/api/agents/sessions/{session_id}") + except httpx.HTTPError: + return [] + messages = body.get("messages") + return [m for m in messages if isinstance(m, dict)] if isinstance(messages, list) else [] + + +@typechecked +def execute_workflow( + backend: BackendProcess, + workflow_id: str, + deadline: float, + on_progress: Optional[Callable[[RunProgress], None]] = None, +) -> RunOutcome: + """Fire the workflow, poll it to a terminal state, and pull the transcript back. + + Blowing the deadline stops the run and reports `timed_out`; the caller still gets + whatever the agent produced before the wall came down. + """ + with httpx.Client(timeout=30.0) as client: + run_id = trigger_run(client, backend, workflow_id) + record: Dict[str, Any] = {} + timed_out = False + + while True: + record = p_find_run(client, backend, workflow_id, run_id) or record + status = str(record.get("status") or "running") + if on_progress is not None: + on_progress(RunProgress( + run_id=run_id, + status=status, + active_step_idx=record.get("active_step_idx"), + last_tool_label=record.get("last_tool_label"), + )) + if status in TERMINAL_STATUSES: + break + if not backend.is_alive(): + raise WorkflowRunFailed("backend died while the workflow was running") + if time.monotonic() >= deadline: + timed_out = True + p_stop_run(client, backend, run_id) + record = p_find_run(client, backend, workflow_id, run_id) or record + break + time.sleep(POLL_INTERVAL_SECONDS) + + session_id = record.get("session_id") + messages = p_collect_session(client, backend, str(session_id)) if session_id else [] + return RunOutcome( + run_id=run_id, + status="timed_out" if timed_out else str(record.get("status") or "failure"), + error=("wall-clock cap reached before the workflow finished" if timed_out else record.get("error")), + cost_usd=float(record.get("cost_usd") or 0.0), + session_id=str(session_id) if session_id else None, + transcript=render_transcript(messages), + answer=final_answer(messages), + system_notices=system_notices(messages), + ) diff --git a/openswarm-runner/tests/test_router_credentials.py b/openswarm-runner/tests/test_router_credentials.py new file mode 100644 index 00000000..2dcfb232 --- /dev/null +++ b/openswarm-runner/tests/test_router_credentials.py @@ -0,0 +1,111 @@ +"""The runner must be structurally unable to rotate a user's OAuth grant. + +Every test here exists to make one class of bug unwritable: a refresh token reaching +9Router's db.json. Delete either wall in runner/seed/router_credentials.py and these go red. +""" + +import json +import os +import stat +from datetime import datetime, timedelta, timezone + +import pytest + +from runner.run_spec import InvalidRunSpec, ProviderCredential, RunSpec, load_run_spec +from runner.seed.router_credentials import ( + RefreshTokenLeak, + assert_no_refresh_token, + router_db_payload, + write_router_db, +) + +NOW = datetime(2026, 7, 31, 12, 0, 0, tzinfo=timezone.utc) +ACCESS_TOKEN = "at-test-value-not-a-real-token" +REFRESH_TOKEN = "rt-test-value-not-a-real-token" + + +def spec_json(credential: dict) -> str: + return json.dumps({ + "run_id": "run-1", + "workflow": {"id": "wf-1", "title": "Test", "steps": [{"text": "say hi"}]}, + "credentials": [credential], + }) + + +def oauth_credential() -> ProviderCredential: + return ProviderCredential( + provider="claude", + auth_type="oauth", + access_token=ACCESS_TOKEN, + expires_at=NOW + timedelta(hours=8), + ) + + +@pytest.mark.parametrize("key", ["refreshToken", "refresh_token", "Refresh-Token", "oauthRefreshToken"]) +def test_a_spec_carrying_a_refresh_token_never_parses(key: str, tmp_path, monkeypatch) -> None: + payload = { + "provider": "claude", + "auth_type": "oauth", + "access_token": ACCESS_TOKEN, + "expires_at": (NOW + timedelta(hours=8)).isoformat(), + key: REFRESH_TOKEN, + } + monkeypatch.setenv("OPENSWARM_RUN_SPEC", spec_json(payload)) + with pytest.raises(InvalidRunSpec) as caught: + load_run_spec() + assert key in str(caught.value) + assert not list(tmp_path.iterdir()), "a rejected spec must not leave anything on disk" + + +@pytest.mark.parametrize("key", ["refreshToken", "refresh_token", "Refresh-Token", "oauthRefreshToken"]) +def test_the_writer_guard_rejects_a_poisoned_payload(key: str) -> None: + payload = {"providerConnections": [{"provider": "claude", key: REFRESH_TOKEN}]} + with pytest.raises(RefreshTokenLeak): + assert_no_refresh_token(payload) + + +def test_the_guard_passes_a_clean_payload() -> None: + assert_no_refresh_token(router_db_payload([oauth_credential()], NOW)) + + +def test_written_db_carries_the_access_token_and_no_refresh_token(tmp_path) -> None: + path = write_router_db(str(tmp_path / "9router"), [oauth_credential()], NOW) + raw = open(path, "r", encoding="utf-8").read() + + # Without this the "no refresh token" assertion below would also pass on an empty file. + assert ACCESS_TOKEN in raw + connection = json.loads(raw)["providerConnections"][0] + assert connection["provider"] == "claude" + assert connection["isActive"] is True + assert connection["expiresAt"] == "2026-07-31T20:00:00.000Z" + + assert "refresh" not in raw.lower() + assert not any("refresh" in key.lower() for key in connection) + + +def test_the_db_and_its_directory_are_owner_only(tmp_path) -> None: + path = write_router_db(str(tmp_path / "9router"), [oauth_credential()], NOW) + assert stat.S_IMODE(os.stat(path).st_mode) == 0o600 + assert stat.S_IMODE(os.stat(os.path.dirname(path)).st_mode) == 0o700 + + +def test_an_api_key_credential_never_reaches_the_router_db(tmp_path) -> None: + credential = ProviderCredential(provider="anthropic", auth_type="api_key", api_key="sk-test-not-real") + path = write_router_db(str(tmp_path / "9router"), [credential], NOW) + assert json.loads(open(path, encoding="utf-8").read())["providerConnections"] == [] + + +def test_an_oauth_credential_without_an_access_token_is_rejected() -> None: + with pytest.raises(ValueError, match="no access_token"): + ProviderCredential(provider="claude", auth_type="oauth", expires_at=NOW) + + +def test_an_expired_access_token_is_fatal_not_refreshable() -> None: + spec = RunSpec.model_validate_json(spec_json({ + "provider": "claude", + "auth_type": "oauth", + "access_token": ACCESS_TOKEN, + "expires_at": (NOW + timedelta(seconds=30)).isoformat(), + })) + assert [credential.provider for credential in spec.expired_credentials(NOW)] == ["claude"] + assert spec.expired_credentials(NOW - timedelta(hours=1)) == [] diff --git a/openswarm-runner/tests/test_run_spec.py b/openswarm-runner/tests/test_run_spec.py new file mode 100644 index 00000000..0fba9641 --- /dev/null +++ b/openswarm-runner/tests/test_run_spec.py @@ -0,0 +1,84 @@ +"""The run spec is the only thing the control plane can say to this container.""" + +import json +import os +import stat + +import pytest + +from runner.run_spec import InvalidRunSpec, RunSpec, load_run_spec +from runner.seed.data_root import seed_data_root, settings_for_run + +VALID_CREDENTIAL = {"provider": "anthropic", "auth_type": "api_key", "api_key": "sk-test-not-real"} + + +def spec_body(**overrides) -> dict: + body = { + "run_id": "run-1", + "workflow": { + "id": "wf-1", + "title": "Daily digest", + "model": "opus-5", + "steps": [{"text": "summarize the inbox"}], + "schedule": {"enabled": True, "repeat_unit": "day", "hour": 9}, + }, + "credentials": [VALID_CREDENTIAL], + } + body.update(overrides) + return body + + +def test_a_missing_spec_names_both_env_vars(monkeypatch) -> None: + monkeypatch.delenv("OPENSWARM_RUN_SPEC", raising=False) + monkeypatch.delenv("OPENSWARM_RUN_SPEC_FILE", raising=False) + with pytest.raises(InvalidRunSpec, match="OPENSWARM_RUN_SPEC_FILE"): + load_run_spec() + + +def test_a_spec_file_is_accepted(tmp_path, monkeypatch) -> None: + path = tmp_path / "spec.json" + path.write_text(json.dumps(spec_body()), encoding="utf-8") + monkeypatch.delenv("OPENSWARM_RUN_SPEC", raising=False) + monkeypatch.setenv("OPENSWARM_RUN_SPEC_FILE", str(path)) + assert load_run_spec().workflow.title == "Daily digest" + + +def test_unknown_top_level_fields_are_rejected(monkeypatch) -> None: + monkeypatch.setenv("OPENSWARM_RUN_SPEC", json.dumps(spec_body(surprise="hello"))) + with pytest.raises(InvalidRunSpec, match="surprise"): + load_run_spec() + + +def test_a_run_needs_at_least_one_credential(monkeypatch) -> None: + monkeypatch.setenv("OPENSWARM_RUN_SPEC", json.dumps(spec_body(credentials=[]))) + with pytest.raises(InvalidRunSpec): + load_run_spec() + + +def test_the_container_never_inherits_the_schedule() -> None: + spec = RunSpec.model_validate(spec_body()) + assert spec.workflow.schedule.enabled is True + assert spec.workflow_for_disk().schedule.enabled is False + + +def test_seeding_writes_the_workflow_and_owner_only_settings(tmp_path) -> None: + spec = RunSpec.model_validate(spec_body()) + seed_data_root(str(tmp_path), spec) + + workflow_path = tmp_path / "workflows" / "wf-1.json" + settings_path = tmp_path / "settings" / "settings.json" + assert json.loads(workflow_path.read_text())["schedule"]["enabled"] is False + assert stat.S_IMODE(os.stat(settings_path).st_mode) == 0o600 + + settings = json.loads(settings_path.read_text()) + assert settings["anthropic_api_key"] == "sk-test-not-real" + assert settings["default_model"] == "opus-5" + assert settings["analytics_opt_in"] is False + + +def test_an_unmappable_api_key_provider_fails_loudly() -> None: + spec = RunSpec.model_validate(spec_body( + credentials=[{"provider": "wat", "auth_type": "api_key", "api_key": "x"}] + )) + with pytest.raises(ValueError, match="no settings field"): + settings_for_run(spec) diff --git a/openswarm-runner/tests/test_workflow_run.py b/openswarm-runner/tests/test_workflow_run.py new file mode 100644 index 00000000..612a23f9 --- /dev/null +++ b/openswarm-runner/tests/test_workflow_run.py @@ -0,0 +1,36 @@ +"""Reading a finished session correctly, including the failures the backend calls success.""" + +from runner.workflow_run import final_answer, render_transcript, system_notices + +# Shape taken verbatim from a real container run whose provider token was rejected. +REJECTED_TOKEN_SESSION = [ + {"role": "user", "content": "Reply with exactly the word PONG and nothing else."}, + {"role": "system", "content": "Provider authentication expired. Open Settings, Models and reconnect, then send your message again."}, +] + +ANSWERED_SESSION = [ + {"role": "user", "content": "ping"}, + {"role": "assistant", "content": [{"type": "tool_use", "name": "Bash", "input": {"command": "echo hi"}}]}, + {"role": "assistant", "content": [{"type": "text", "text": "PONG"}]}, + {"role": "assistant", "content": [{"type": "text", "text": "draft"}], "hidden": True}, +] + + +def test_a_rejected_credential_surfaces_as_a_system_notice() -> None: + assert system_notices(REJECTED_TOKEN_SESSION) == [REJECTED_TOKEN_SESSION[1]["content"]] + assert final_answer(REJECTED_TOKEN_SESSION) == "" + + +def test_a_healthy_run_raises_no_notices() -> None: + assert system_notices(ANSWERED_SESSION) == [] + + +def test_the_answer_is_the_last_visible_assistant_text() -> None: + assert final_answer(ANSWERED_SESSION) == "PONG" + + +def test_the_transcript_keeps_tool_calls_and_drops_hidden_turns() -> None: + transcript = render_transcript(ANSWERED_SESSION) + assert "[tool Bash]" in transcript + assert "PONG" in transcript + assert "draft" not in transcript