From 9cf70af5872a7dbc85878c147b41e5f5eb0ad787 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Sat, 1 Aug 2026 01:19:58 -0700 Subject: [PATCH] [eric] runner: a cloud run's files come home, and one too big to carry says so instead of vanishing --- openswarm-runner/runner/main.py | 52 ++++- openswarm-runner/runner/report.py | 69 ------- openswarm-runner/runner/results/__init__.py | 0 .../runner/results/deliverables.py | 193 ++++++++++++++++++ openswarm-runner/runner/results/report.py | 168 +++++++++++++++ openswarm-runner/runner/run_spec.py | 62 ++++++ openswarm-runner/runner/seed/data_root.py | 53 ++++- openswarm-runner/tests/test_deliverables.py | 174 ++++++++++++++++ openswarm-runner/tests/test_run_spec.py | 20 +- 9 files changed, 707 insertions(+), 84 deletions(-) delete mode 100644 openswarm-runner/runner/report.py create mode 100644 openswarm-runner/runner/results/__init__.py create mode 100644 openswarm-runner/runner/results/deliverables.py create mode 100644 openswarm-runner/runner/results/report.py create mode 100644 openswarm-runner/tests/test_deliverables.py diff --git a/openswarm-runner/runner/main.py b/openswarm-runner/runner/main.py index 6c7759cc..ac6b177f 100644 --- a/openswarm-runner/runner/main.py +++ b/openswarm-runner/runner/main.py @@ -18,10 +18,12 @@ from typeguard import typechecked from runner.boot.backend_process import BackendProcess, BackendUnavailable, start_backend, stop_backend from runner.boot.renderer_process import RendererProcess, RendererUnavailable, start_renderer, stop_renderer -from runner.report import RunReport, send_report +from runner.results.deliverables import collect +from runner.results.report import RunReport, deliver_files, send_report from runner.run_spec import CLOUD_RUN_DASHBOARD_ID, 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.seed.skills import write_skills from runner.workflow_run import RunOutcome, RunProgress, WorkflowRunFailed, execute_workflow EXIT_OK = 0 @@ -37,9 +39,13 @@ DEFAULT_APP_ROOT = "/app" DEFAULT_FRONTEND_DIR = "/app/frontend" DEFAULT_DATA_ROOT = "/data/openswarm" DEFAULT_ROUTER_DATA_DIR = "/data/9router" +# The agent's own folder, and the only place on this machine whose contents come home. +DEFAULT_RUN_WORKSPACE = "/data/workspace" 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 +# Has to cover the file upload as well as the report's retries, so it is minutes, not seconds; the +# control plane's own kill sits further out again (dispatch.ts MACHINE_GRACE_MS). +REPORT_GRACE_SECONDS = 240.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 @@ -92,8 +98,25 @@ def arm_hard_stop(seconds: float) -> None: @typechecked -def p_fail(spec: Optional[RunSpec], status: str, message: str, code: int) -> int: +def p_fail( + spec: Optional[RunSpec], + status: str, + message: str, + code: int, + workspace: Optional[str] = None, +) -> int: + """Report a failure, handing over anything the run managed to make first. + + `workspace` is passed only where the agent actually ran: a workflow that died on step 3 may + have written a perfectly good report on step 1, and losing it because a later step threw is + exactly the "the file died with the machine" problem this whole path exists to fix. + """ logger.error("%s: %s", status, message) + files = ( + deliver_files(spec.callback, workspace, collect(workspace)) + if spec is not None and workspace is not None + else [] + ) send_report( spec.callback if spec else None, RunReport( @@ -102,6 +125,7 @@ def p_fail(spec: Optional[RunSpec], status: str, message: str, code: int) -> int status=status, exit_code=code, error=message, + files=files, ), ) return code @@ -133,10 +157,12 @@ def p_run(spec: RunSpec, deadline: float) -> int: 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) + workspace = os.environ.get("OPENSWARM_RUN_WORKSPACE", DEFAULT_RUN_WORKSPACE) 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) + seed_data_root(data_root, workspace, spec) + write_skills(os.path.expanduser("~"), spec.skills) logger.info("seeded data root %s and router db in %s", data_root, router_data_dir) backend: Optional[BackendProcess] = None @@ -173,11 +199,15 @@ def p_run(spec: RunSpec, deadline: float) -> int: # confident wrong answer, which is worse than no answer. return p_fail(spec, "failure", str(exc), EXIT_RENDERER_UNAVAILABLE) except WorkflowRunFailed as exc: - return p_fail(spec, "failure", str(exc), EXIT_WORKFLOW_FAILED) + return p_fail(spec, "failure", str(exc), EXIT_WORKFLOW_FAILED, workspace) finally: stop_renderer(renderer) stop_backend(process) + # Files before the terminal report, always: the callback token is refused the moment the run + # is closed, so this is the only order in which both the files and the receipt can land. + files = deliver_files(spec.callback, workspace, collect(workspace)) + 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( @@ -190,6 +220,7 @@ def p_run(spec: RunSpec, deadline: float) -> int: answer=outcome.answer, transcript=outcome.transcript, system_notices=outcome.system_notices, + files=files, )) return code @@ -213,7 +244,16 @@ def main() -> int: return p_run(spec, deadline) except Exception as exc: logger.exception("runner crashed") - return p_fail(spec, "failure", f"runner crashed: {exc}", EXIT_INTERNAL) + # Re-uploading a file the successful path already sent is harmless: the control plane keys + # a run's files on their path, so a second delivery overwrites one row rather than billing + # the budget twice. + return p_fail( + spec, + "failure", + f"runner crashed: {exc}", + EXIT_INTERNAL, + os.environ.get("OPENSWARM_RUN_WORKSPACE", DEFAULT_RUN_WORKSPACE), + ) if __name__ == "__main__": diff --git a/openswarm-runner/runner/report.py b/openswarm-runner/runner/report.py deleted file mode 100644 index bb5c3de1..00000000 --- a/openswarm-runner/runner/report.py +++ /dev/null @@ -1,69 +0,0 @@ -"""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/results/__init__.py b/openswarm-runner/runner/results/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/openswarm-runner/runner/results/deliverables.py b/openswarm-runner/runner/results/deliverables.py new file mode 100644 index 00000000..5aaba20f --- /dev/null +++ b/openswarm-runner/runner/results/deliverables.py @@ -0,0 +1,193 @@ +"""What the run made, and what of it is allowed to come home. + +A cloud run's machine is destroyed the moment it exits, so a file it wrote is gone +unless something carries it out. This module is the "what": it walks the one +directory a run is given as its working folder and decides, per file, deliver or +refuse. The "how" (handing bytes to the control plane) lives in runner.results.report. + +Refusing loudly is the whole point of the caps. A user who asked for a video and +got a 20MB fragment of one is worse off than a user who was told the video was too +big, so nothing here ever truncates a file: it either arrives whole or it arrives +as a sentence explaining why it did not. +""" + +import hashlib +import logging +import os +from typing import List, Optional, Tuple + +from pydantic import BaseModel, ConfigDict, Field +from typeguard import typechecked + +logger = logging.getLogger(__name__) + +# Per-file ceiling. Deliverables are reports, spreadsheets, charts and small archives; a run that +# produces something bigger is doing a different job than this pipe was built for. +MAX_FILE_BYTES = 20 * 1024 * 1024 +# Per-run ceiling, enforced in walk order so the first files still arrive when a later one blows it. +MAX_TOTAL_BYTES = 50 * 1024 * 1024 +MAX_FILES = 40 +# Longest path we will accept, so a deep tree cannot produce a name no filesystem will take back. +MAX_RELATIVE_PATH_CHARS = 180 + +# Machinery, not deliverables. Everything here is either regenerable (dependencies, caches, +# compiled bytecode) or the run's own plumbing, and shipping it would blow the file budget on +# things nobody asked for. +EXCLUDED_DIRS = frozenset({ + ".git", + ".claude", + "node_modules", + "__pycache__", + ".venv", + "venv", + ".pytest_cache", + ".ruff_cache", + ".mypy_cache", + ".cache", + ".npm", +}) +EXCLUDED_NAMES = frozenset({".DS_Store", ".gitignore", ".gitkeep"}) + +READ_CHUNK_BYTES = 1024 * 1024 + + +class Deliverable(BaseModel): + """One file that fits, addressed by its path relative to the run's workspace.""" + + model_config = ConfigDict(validate_assignment=True) + + path: str + size_bytes: int + sha256: str + + +class Refused(BaseModel): + """One file that does not come home, and the sentence the user gets instead.""" + + model_config = ConfigDict(validate_assignment=True) + + path: str + size_bytes: int + reason: str + + +class Harvest(BaseModel): + model_config = ConfigDict(validate_assignment=True) + + files: List[Deliverable] = Field(default_factory=list) + refused: List[Refused] = Field(default_factory=list) + + @typechecked + def total_bytes(self) -> int: + return sum(item.size_bytes for item in self.files) + + +@typechecked +def human_bytes(count: int) -> str: + if count < 1024: + return f"{count} B" + if count < 1024 * 1024: + return f"{count / 1024:.0f} KB" + if count < 1024 * 1024 * 1024: + return f"{count / (1024 * 1024):.1f} MB" + return f"{count / (1024 * 1024 * 1024):.1f} GB" + + +@typechecked +def p_digest(path: str) -> Optional[str]: + """sha256, streamed. None when the file went away mid-walk, which is not an error.""" + digest = hashlib.sha256() + try: + with open(path, "rb") as handle: + while True: + chunk = handle.read(READ_CHUNK_BYTES) + if not chunk: + break + digest.update(chunk) + except OSError as exc: + logger.warning("could not read %s while harvesting: %s", path, exc) + return None + return digest.hexdigest() + + +@typechecked +def p_walk(workspace: str) -> List[Tuple[str, int]]: + """Every candidate file under the workspace as (relative path, size), sorted for determinism.""" + found: List[Tuple[str, int]] = [] + for directory, subdirs, filenames in os.walk(workspace): + subdirs[:] = sorted(name for name in subdirs if name not in EXCLUDED_DIRS) + for filename in sorted(filenames): + if filename in EXCLUDED_NAMES: + continue + absolute = os.path.join(directory, filename) + # Symlinks are not followed: a run that linked to /etc/passwd must not exfiltrate it. + if os.path.islink(absolute) or not os.path.isfile(absolute): + continue + try: + size = os.path.getsize(absolute) + except OSError: + continue + found.append((os.path.relpath(absolute, workspace), size)) + return found + + +@typechecked +def collect(workspace: str) -> Harvest: + """Decide, for every file in the run's workspace, whether it comes home.""" + harvest = Harvest() + if not os.path.isdir(workspace): + return harvest + + running_total = 0 + for relative, size in p_walk(workspace): + if len(relative) > MAX_RELATIVE_PATH_CHARS: + harvest.refused.append(Refused( + path=relative[:MAX_RELATIVE_PATH_CHARS] + "...", + size_bytes=size, + reason="its path is too long to save anywhere", + )) + continue + if size == 0: + continue + if size > MAX_FILE_BYTES: + harvest.refused.append(Refused( + path=relative, + size_bytes=size, + reason=( + f"it is {human_bytes(size)} and a single cloud-run file cannot exceed " + f"{human_bytes(MAX_FILE_BYTES)}" + ), + )) + continue + if len(harvest.files) >= MAX_FILES: + harvest.refused.append(Refused( + path=relative, + size_bytes=size, + reason=f"this run already produced the maximum of {MAX_FILES} files", + )) + continue + if running_total + size > MAX_TOTAL_BYTES: + harvest.refused.append(Refused( + path=relative, + size_bytes=size, + reason=( + f"the run's files already total {human_bytes(running_total)} and the limit " + f"is {human_bytes(MAX_TOTAL_BYTES)}" + ), + )) + continue + digest = p_digest(os.path.join(workspace, relative)) + if digest is None: + harvest.refused.append(Refused(path=relative, size_bytes=size, reason="it could not be read")) + continue + harvest.files.append(Deliverable(path=relative, size_bytes=size, sha256=digest)) + running_total += size + + logger.info( + "harvested %d file(s) totalling %s from %s, refused %d", + len(harvest.files), + human_bytes(running_total), + workspace, + len(harvest.refused), + ) + return harvest diff --git a/openswarm-runner/runner/results/report.py b/openswarm-runner/runner/results/report.py new file mode 100644 index 00000000..ca769138 --- /dev/null +++ b/openswarm-runner/runner/results/report.py @@ -0,0 +1,168 @@ +"""Tell the control plane what happened, and hand it whatever the run made. + +The terminal report is the run's only receipt. Files go up BEFORE it, on purpose: +the per-run callback token stops working the instant the run reaches a terminal +state, so "upload, then close" is the only order in which both can succeed, and it +means the window for writing files to a run closes exactly when the run does. +""" + +import logging +import os +import time +from typing import List, Literal, Optional +from urllib.parse import quote + +import httpx +from pydantic import BaseModel, ConfigDict, Field +from typeguard import typechecked + +from runner.results.deliverables import Harvest, human_bytes +from runner.run_spec import CallbackTarget + +logger = logging.getLogger(__name__) + +TERMINAL_ATTEMPTS = 5 +TERMINAL_BACKOFF_SECONDS = 2.0 +REQUEST_TIMEOUT_SECONDS = 15.0 +# One file, one shot, generous: 20MB over a cold uplink is slower than any report. +UPLOAD_TIMEOUT_SECONDS = 120.0 +FILE_PATH_HEADER = "X-Openswarm-File-Path" +FILE_SHA256_HEADER = "X-Openswarm-File-Sha256" + + +class ReportedFile(BaseModel): + """One file the run produced, delivered or not, always named.""" + + model_config = ConfigDict(validate_assignment=True) + + path: str + size_bytes: int + delivered: bool + # Present only when delivered is False. Written for a human, because it is shown to one. + reason: Optional[str] = None + + +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) + # Every file the run made, including the ones that were too big to send. A deliverable that + # silently vanished is the failure this list exists to make impossible. + files: List[ReportedFile] = 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 + + +@typechecked +def p_upload_one(client: httpx.Client, callback: CallbackTarget, workspace: str, relative: str, sha256: str) -> Optional[str]: + """Push one file. Returns None on success, or the sentence explaining the failure.""" + try: + with open(os.path.join(workspace, relative), "rb") as handle: + response = client.post( + str(callback.artifacts_url), + headers={ + "Authorization": f"Bearer {callback.token}", + "Content-Type": "application/octet-stream", + FILE_PATH_HEADER: quote(relative, safe="/"), + FILE_SHA256_HEADER: sha256, + }, + content=handle.read(), + ) + except OSError as exc: + return f"it could not be read back off disk ({exc.strerror or exc})" + except httpx.HTTPError as exc: + return f"the upload did not complete ({type(exc).__name__})" + if response.status_code < 300: + return None + # The control plane refuses with prose it wrote for the user; keep its words rather than ours. + detail = (response.text or "").strip() + try: + body = response.json() + if isinstance(body, dict) and isinstance(body.get("error"), str): + detail = body["error"] + except ValueError: + pass + return detail[:300] if detail else f"the storage service answered HTTP {response.status_code}" + + +@typechecked +def deliver_files(callback: Optional[CallbackTarget], workspace: str, harvest: Harvest) -> List[ReportedFile]: + """Hand every deliverable to the control plane and report honestly on each one. + + Never raises and never fails the run: a workflow whose answer is good and whose + attachment did not make it should still deliver the answer, with the miss stated. + """ + reported = [ + ReportedFile(path=item.path, size_bytes=item.size_bytes, delivered=False, reason=item.reason) + for item in harvest.refused + ] + if not harvest.files: + return reported + if callback is None or not callback.artifacts_url: + for item in harvest.files: + reported.append(ReportedFile( + path=item.path, + size_bytes=item.size_bytes, + delivered=False, + reason="this run had nowhere to send files, so it kept them on the machine", + )) + return reported + + with httpx.Client(timeout=UPLOAD_TIMEOUT_SECONDS) as client: + for item in harvest.files: + failure = p_upload_one(client, callback, workspace, item.path, item.sha256) + if failure is None: + logger.info("delivered %s (%s)", item.path, human_bytes(item.size_bytes)) + else: + logger.warning("could not deliver %s: %s", item.path, failure) + reported.append(ReportedFile( + path=item.path, + size_bytes=item.size_bytes, + delivered=failure is None, + reason=failure, + )) + return reported diff --git a/openswarm-runner/runner/run_spec.py b/openswarm-runner/runner/run_spec.py index 735b97c0..db6487e7 100644 --- a/openswarm-runner/runner/run_spec.py +++ b/openswarm-runner/runner/run_spec.py @@ -27,6 +27,11 @@ CLOUD_RUN_DASHBOARD_ID = "cloud-run" # 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) +# A skill is prose plus the odd small script. These bounds keep a run spec from becoming a file +# transfer, and are enforced here so an oversized one dies before a machine is billed for it. +MAX_SKILLS = 60 +MAX_SKILL_FILE_CHARS = 200_000 + class InvalidRunSpec(ValueError): """The control plane handed us something we refuse to run.""" @@ -82,6 +87,59 @@ class CallbackTarget(BaseModel): url: str = Field(min_length=1) token: str = Field(min_length=1) heartbeat_seconds: int = Field(default=30, ge=5, le=300) + # Where the run's files go. Named outright rather than derived from `url`, so a control plane + # that cannot accept files says so by leaving it out instead of being guessed at. + artifacts_url: Optional[str] = None + + +class SkillFile(BaseModel): + """One file inside a skill folder. Text only: a skill is prose plus small scripts.""" + + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + path: str = Field(min_length=1, max_length=180) + text: str = Field(max_length=MAX_SKILL_FILE_CHARS) + + @model_validator(mode="after") + def p_reject_escaping_path(self) -> "SkillFile": + parts = self.path.split("/") + if self.path.startswith("/") or "\\" in self.path or any(p in ("", ".", "..") for p in parts): + raise ValueError(f"skill file path {self.path!r} is not a plain relative path") + return self + + +class SkillPayload(BaseModel): + """One of the user's skills, carried up so the agent has the same know-how it has at home. + + Skills are the user's own writing, not credentials. Nothing in here is a token, and the + control plane never learns anything from it that it could spend. + """ + + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + # Also the folder name under ~/.claude/skills, so it has to survive being a directory. + id: str = Field(min_length=1, max_length=80, pattern=r"^[A-Za-z0-9][A-Za-z0-9_-]*$") + files: List[SkillFile] = Field(min_length=1) + + @model_validator(mode="after") + def p_require_skill_md(self) -> "SkillPayload": + if not any(file.path == "SKILL.md" for file in self.files): + raise ValueError(f"skill {self.id!r} has no SKILL.md, so nothing would ever load it") + return self + + +class McpServerNote(BaseModel): + """A server the user has connected at home, named so the agent can say it cannot reach it. + + Deliberately carries NO transport and NO secret. The user's MCP credentials (Slack session + cookies, Notion and GitHub access tokens, Google refresh tokens) stay on their laptop, so this + is a list of names and nothing else. Its whole job is to stop the agent quietly answering + "update my Notion" from general knowledge because it never knew Notion existed. + """ + + model_config = ConfigDict(validate_assignment=True, extra="forbid") + + name: str = Field(min_length=1, max_length=120) class RunSpec(BaseModel): @@ -91,6 +149,10 @@ class RunSpec(BaseModel): workflow: Workflow credentials: List[ProviderCredential] = Field(min_length=1) callback: Optional[CallbackTarget] = None + # The user's skills, shipped so a cloud run is as capable as the same workflow at home. + skills: List[SkillPayload] = Field(default_factory=list, max_length=MAX_SKILLS) + # Names only. See McpServerNote for why there is no config here. + unavailable_mcp_servers: List[McpServerNote] = Field(default_factory=list, max_length=100) # 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) # Boot Electron under a virtual display so browser steps work. On by default: parity is the diff --git a/openswarm-runner/runner/seed/data_root.py b/openswarm-runner/runner/seed/data_root.py index 9d8e5979..7741ff7c 100644 --- a/openswarm-runner/runner/seed/data_root.py +++ b/openswarm-runner/runner/seed/data_root.py @@ -29,6 +29,17 @@ API_KEY_SETTINGS_FIELD: Dict[str, str] = { # 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" +# Told to the agent in as many words, because it cannot find this out any other way and the +# consequence of not knowing is a report written to a folder that is deleted minutes later. +DELIVERY_NOTE = ( + "You are running in the OpenSwarm cloud on a throwaway machine. Your working directory is " + "the ONLY place that survives: every file you save there is delivered back to the user, and " + "everything else on this machine is destroyed the moment this run ends. So when a task asks " + "for a document, spreadsheet, image or archive, write it to a plainly named file in your " + "working directory rather than only pasting it into your reply. Do not write deliverables to " + "/tmp or to your home directory; they will not come back." +) + @typechecked def p_write_json(path: str, payload: Any) -> None: @@ -48,13 +59,43 @@ def p_write_json(path: str, payload: Any) -> None: @typechecked -def settings_for_run(spec: RunSpec) -> AppSettings: - """The AppSettings a cloud run needs: this workflow's model, this run's keys, no telemetry.""" +def unavailable_apps_note(spec: RunSpec) -> str: + """Name the user's connected apps this run cannot reach, so silence is not mistaken for absence. + + Their MCP credentials never leave the laptop, so the servers are not here and never will be + mid-run. Without this sentence the agent has no way to know the app exists, and "update my + Notion" comes back as a confident paragraph about Notion rather than an admission. + """ + names = [server.name for server in spec.unavailable_mcp_servers] + if not names: + return "" + return ( + "These apps are connected on the user's own computer but NOT reachable from this cloud " + f"run, because their sign-in details stay on that computer: {', '.join(sorted(names))}. " + "If a task needs one of them, say plainly that it cannot be done from a cloud run and " + "that it has to run on their machine. Never guess at, invent, or describe from memory " + "what one of those apps contains." + ) + + +@typechecked +def settings_for_run(spec: RunSpec, workspace: str) -> AppSettings: + """The AppSettings a cloud run needs: this workflow's model, this run's keys, no telemetry. + + `default_folder` is what makes the run's files findable afterwards. Left unset, the agent + falls back to $HOME and the launcher reroutes it into a per-session scratch directory whose + name nothing outside the backend can predict, so the harvest would have nowhere to look. + """ 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 + settings.default_folder = workspace + additions = [DELIVERY_NOTE, unavailable_apps_note(spec)] + settings.default_system_prompt = "\n\n".join( + part for part in [settings.default_system_prompt or "", *additions] if part + ).strip() for credential in spec.credentials: if credential.auth_type != "api_key": continue @@ -69,13 +110,17 @@ def settings_for_run(spec: RunSpec) -> AppSettings: @typechecked -def seed_data_root(data_root: str, spec: RunSpec) -> None: +def seed_data_root(data_root: str, workspace: str, spec: RunSpec) -> None: """Write the workflow, settings and dashboard records the backend will read at boot. The dashboard exists so the Electron window has somewhere to land and browser cards have somewhere to render. Writing it here rather than letting the backend's first-boot migration invent one keeps its id knowable before anything has started. + + The workspace sits OUTSIDE the data root deliberately: it is the agent's own folder, and a + Glob or Grep run inside it should not sweep up the settings file its API keys live in. """ + os.makedirs(workspace, mode=0o700, exist_ok=True) workflow = spec.workflow_for_disk() p_write_json( os.path.join(data_root, "workflows", f"{workflow.id}.json"), @@ -83,7 +128,7 @@ def seed_data_root(data_root: str, spec: RunSpec) -> None: ) p_write_json( os.path.join(data_root, "settings", "settings.json"), - settings_for_run(spec).model_dump(mode="json"), + settings_for_run(spec, workspace).model_dump(mode="json"), ) dashboard = Dashboard(id=CLOUD_RUN_DASHBOARD_ID, name=spec.workflow.title or "Cloud run") p_write_json( diff --git a/openswarm-runner/tests/test_deliverables.py b/openswarm-runner/tests/test_deliverables.py new file mode 100644 index 00000000..becaf220 --- /dev/null +++ b/openswarm-runner/tests/test_deliverables.py @@ -0,0 +1,174 @@ +"""What the run made, what comes home, and what is refused out loud instead of silently. + +The caps are the point. A truncated file is worse than a refused one, and a file that +vanishes with no sentence attached is the failure the whole list exists to prevent. +""" + +import os + +import pytest + +from runner.results.deliverables import ( + MAX_FILE_BYTES, + MAX_FILES, + MAX_TOTAL_BYTES, + collect, + human_bytes, +) +from runner.results.report import deliver_files +from runner.run_spec import CallbackTarget + + +def write(root, relative: str, payload: bytes) -> str: + path = os.path.join(str(root), relative) + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "wb") as handle: + handle.write(payload) + return path + + +def test_a_missing_workspace_is_an_empty_harvest_not_a_crash(tmp_path) -> None: + assert collect(str(tmp_path / "never-made")).files == [] + + +def test_ordinary_files_are_collected_with_their_digest(tmp_path) -> None: + write(tmp_path, "report.md", b"# Digest\n") + write(tmp_path, "data/rows.csv", b"a,b\n1,2\n") + + # Walk order, and it is fixed: this directory's own files first, then subdirectories in name + # order, so a run's file list does not shuffle between two identical runs. + harvest = collect(str(tmp_path)) + assert [f.path for f in harvest.files] == ["report.md", "data/rows.csv"] + assert harvest.files[0].size_bytes == len(b"# Digest\n") + # A digest travels with every file, so a corrupted upload is detectable rather than assumed fine. + assert len(harvest.files[0].sha256) == 64 + assert harvest.refused == [] + + +def test_machinery_is_not_a_deliverable(tmp_path) -> None: + write(tmp_path, "report.md", b"keep me") + write(tmp_path, ".git/config", b"[core]") + write(tmp_path, "node_modules/left-pad/index.js", b"module.exports=1") + write(tmp_path, "__pycache__/x.pyc", b"\x00") + write(tmp_path, ".claude/worktrees/probe/README", b"scratch") + + assert [f.path for f in collect(str(tmp_path)).files] == ["report.md"] + + +def test_an_empty_file_is_neither_delivered_nor_complained_about(tmp_path) -> None: + write(tmp_path, "touched.txt", b"") + harvest = collect(str(tmp_path)) + assert harvest.files == [] + assert harvest.refused == [] + + +def test_a_file_over_the_cap_is_refused_whole_and_says_why(tmp_path) -> None: + write(tmp_path, "render.mp4", b"x" * (MAX_FILE_BYTES + 1)) + write(tmp_path, "notes.md", b"still fine") + + harvest = collect(str(tmp_path)) + assert [f.path for f in harvest.files] == ["notes.md"] + assert len(harvest.refused) == 1 + assert harvest.refused[0].path == "render.mp4" + assert "cannot exceed" in harvest.refused[0].reason + # Never a fragment: an over-sized file is absent, not shortened. + assert all(f.path != "render.mp4" for f in harvest.files) + + +def test_the_run_total_stops_collecting_but_keeps_what_already_fit(tmp_path) -> None: + chunk = b"x" * MAX_FILE_BYTES + for index in range(MAX_TOTAL_BYTES // MAX_FILE_BYTES + 1): + write(tmp_path, f"blob-{index}.bin", chunk) + + harvest = collect(str(tmp_path)) + assert harvest.total_bytes() <= MAX_TOTAL_BYTES + assert len(harvest.files) >= 1 + assert harvest.refused, "the file that blew the budget must be named, not dropped" + assert "limit is" in harvest.refused[0].reason + + +def test_too_many_files_refuses_the_extras_by_name(tmp_path) -> None: + for index in range(MAX_FILES + 3): + write(tmp_path, f"note-{index:03d}.txt", b"hi") + + harvest = collect(str(tmp_path)) + assert len(harvest.files) == MAX_FILES + assert len(harvest.refused) == 3 + assert all("maximum of" in item.reason for item in harvest.refused) + + +def test_a_symlink_out_of_the_workspace_is_never_followed(tmp_path) -> None: + secret = tmp_path / "outside" / "id_rsa" + os.makedirs(secret.parent, exist_ok=True) + secret.write_text("PRIVATE KEY") + workspace = tmp_path / "ws" + os.makedirs(workspace, exist_ok=True) + os.symlink(str(secret), str(workspace / "borrowed.pem")) + + assert collect(str(workspace)).files == [] + + +def test_with_nowhere_to_send_files_the_run_says_so_per_file(tmp_path) -> None: + write(tmp_path, "report.md", b"the answer") + reported = deliver_files(None, str(tmp_path), collect(str(tmp_path))) + + assert len(reported) == 1 + assert reported[0].delivered is False + assert "nowhere to send files" in (reported[0].reason or "") + + +def test_a_control_plane_with_no_file_route_is_reported_not_guessed(tmp_path) -> None: + write(tmp_path, "report.md", b"the answer") + callback = CallbackTarget(url="https://cloud.test/report", token="two-party") + assert callback.artifacts_url is None + + reported = deliver_files(callback, str(tmp_path), collect(str(tmp_path))) + assert reported[0].delivered is False + + +def test_refusals_reach_the_report_even_when_nothing_was_delivered(tmp_path) -> None: + write(tmp_path, "render.mp4", b"x" * (MAX_FILE_BYTES + 1)) + reported = deliver_files(None, str(tmp_path), collect(str(tmp_path))) + + assert len(reported) == 1 + assert reported[0].path == "render.mp4" + assert reported[0].delivered is False + assert "cannot exceed" in (reported[0].reason or "") + + +@pytest.mark.parametrize( + "count,expected", + [(512, "512 B"), (2048, "2 KB"), (5 * 1024 * 1024, "5.0 MB"), (3 * 1024**3, "3.0 GB")], +) +def test_sizes_are_written_the_way_a_person_reads_them(count: int, expected: str) -> None: + assert human_bytes(count) == expected + + +def test_a_failed_run_still_hands_over_what_it_managed_to_make(tmp_path, monkeypatch) -> None: + """A workflow that dies on step 3 may have written a perfectly good report on step 1.""" + from runner import main + from runner.run_spec import RunSpec + + write(tmp_path, "partial-report.md", b"# What I got through\n") + spec = RunSpec.model_validate({ + "run_id": "cr-fail", + "workflow": {"id": "wf-1", "steps": [{"id": "s1", "text": "go"}]}, + "credentials": [{"provider": "anthropic", "auth_type": "api_key", "api_key": "sk-test"}], + }) + + sent: list = [] + monkeypatch.setattr(main, "send_report", lambda callback, report: sent.append(report) or True) + code = main.p_fail(spec, "failure", "step 3 blew up", main.EXIT_WORKFLOW_FAILED, str(tmp_path)) + + assert code == main.EXIT_WORKFLOW_FAILED + assert [f.path for f in sent[0].files] == ["partial-report.md"] + + +def test_a_failure_with_no_workspace_reports_no_files_rather_than_guessing(monkeypatch) -> None: + from runner import main + + sent: list = [] + monkeypatch.setattr(main, "send_report", lambda callback, report: sent.append(report) or True) + main.p_fail(None, "failure", "bad spec", main.EXIT_BAD_SPEC) + + assert sent[0].files == [] diff --git a/openswarm-runner/tests/test_run_spec.py b/openswarm-runner/tests/test_run_spec.py index 0fba9641..1733f5a5 100644 --- a/openswarm-runner/tests/test_run_spec.py +++ b/openswarm-runner/tests/test_run_spec.py @@ -63,10 +63,11 @@ def test_the_container_never_inherits_the_schedule() -> None: 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) + workspace = str(tmp_path / "workspace") + seed_data_root(str(tmp_path / "data"), workspace, spec) - workflow_path = tmp_path / "workflows" / "wf-1.json" - settings_path = tmp_path / "settings" / "settings.json" + workflow_path = tmp_path / "data" / "workflows" / "wf-1.json" + settings_path = tmp_path / "data" / "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 @@ -74,11 +75,20 @@ def test_seeding_writes_the_workflow_and_owner_only_settings(tmp_path) -> None: assert settings["anthropic_api_key"] == "sk-test-not-real" assert settings["default_model"] == "opus-5" assert settings["analytics_opt_in"] is False + # The agent's folder is the deliverable folder, and it exists before the backend boots. + assert settings["default_folder"] == workspace + assert os.path.isdir(workspace) -def test_an_unmappable_api_key_provider_fails_loudly() -> None: +def test_the_agent_is_told_its_files_only_survive_from_the_workspace(tmp_path) -> None: + spec = RunSpec.model_validate(spec_body()) + prompt = settings_for_run(spec, str(tmp_path)).default_system_prompt or "" + assert "delivered back to the user" in prompt + + +def test_an_unmappable_api_key_provider_fails_loudly(tmp_path) -> 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) + settings_for_run(spec, str(tmp_path))