mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-28 20:44:50 +02:00
[eric] runner: ephemeral cloud container that boots the backend, runs one workflow, and exits
This commit is contained in:
@@ -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 <resources>/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"]
|
||||
@@ -0,0 +1,10 @@
|
||||
*
|
||||
!backend
|
||||
!openswarm-runner/runner
|
||||
backend/data
|
||||
backend/.venv
|
||||
backend/uv-bin
|
||||
backend/tests
|
||||
**/__pycache__
|
||||
**/*.pyc
|
||||
**/.DS_Store
|
||||
@@ -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": "<already refreshed by the control plane>",
|
||||
"expires_at": "2026-07-31T20:00:00Z" }
|
||||
],
|
||||
"callback": { "url": "https://api.openswarm.com/api/cloud-runs/cr_01J.../report",
|
||||
"token": "<two-party runner token, not a user credential>" },
|
||||
"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.
|
||||
@@ -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
|
||||
@@ -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)
|
||||
@@ -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())
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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"),
|
||||
)
|
||||
@@ -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
|
||||
@@ -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),
|
||||
)
|
||||
@@ -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)) == []
|
||||
@@ -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)
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user