mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-30 13:34:50 +02:00
[eric] runner: a cloud run's files come home, and one too big to carry says so instead of vanishing
This commit is contained in:
@@ -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__":
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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 == []
|
||||
@@ -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))
|
||||
|
||||
Reference in New Issue
Block a user