mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-30 05:24:50 +02:00
defluff (frontend + backend): strip em-dashes + shorten docstrings + drop dead UI files (cosmetic only, no schedule code)
This commit is contained in:
@@ -11,12 +11,7 @@ import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# In-flight dedup map for generate-group-meta. Keyed by (session_id, group_id).
|
||||
# When the frontend issues N concurrent requests for the same group (which it
|
||||
# can during heavy streaming), we only fire ONE upstream Anthropic call and
|
||||
# return the same Future to all callers. Eliminates the 429 thundering herd
|
||||
# without changing retry/fallback semantics — each unique (session, group)
|
||||
# still gets its full retry budget, just not multiplied by N callers.
|
||||
# Dedup concurrent generate-group-meta calls; collapses the 429 thundering herd by sharing one upstream Future per (session, group).
|
||||
_group_meta_inflight: dict[tuple[str, str], asyncio.Future] = {}
|
||||
|
||||
@asynccontextmanager
|
||||
@@ -32,7 +27,6 @@ async def agents_lifespan():
|
||||
|
||||
agents = SubApp("agents", agents_lifespan)
|
||||
|
||||
# REST Endpoints
|
||||
|
||||
@agents.router.get("/sessions")
|
||||
async def list_sessions(dashboard_id: str = ""):
|
||||
@@ -71,12 +65,7 @@ async def send_message(session_id: str, body: dict):
|
||||
if not prompt:
|
||||
raise HTTPException(status_code=400, detail="prompt is required")
|
||||
|
||||
# Pre-flight MCP suggestion (Phase 3, Layer N). Runs in parallel with
|
||||
# the agent launch path — if it produces suggestions, they're
|
||||
# surfaced inline in the chat via agent:mcp_suggestions WS event.
|
||||
# Fails open: any error from the classifier is swallowed and the
|
||||
# agent proceeds normally. The classifier is short-circuited for
|
||||
# obviously-local prompts (greetings, shell commands, file paths).
|
||||
# Run MCP-suggestion classifier in parallel with the agent launch; fails open.
|
||||
try:
|
||||
from backend.apps.agents.mcp_preflight import run_preflight
|
||||
from backend.apps.agents.ws_manager import ws_manager as _ws
|
||||
@@ -93,7 +82,6 @@ async def send_message(session_id: str, body: dict):
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Non-blocking — don't gate the agent on the classifier.
|
||||
import asyncio as _asyncio
|
||||
_asyncio.create_task(_emit_preflight())
|
||||
except Exception:
|
||||
@@ -161,11 +149,7 @@ async def generate_group_meta(session_id: str, body: dict):
|
||||
if not group_id or not tool_calls:
|
||||
raise HTTPException(status_code=400, detail="group_id and tool_calls are required")
|
||||
|
||||
# In-flight dedup. If an identical request is already running, await its
|
||||
# result instead of firing another Anthropic call. This is the entire fix
|
||||
# for the 429 storm we were seeing — N concurrent identical requests
|
||||
# collapse to 1 upstream call. Refinement requests bypass dedup since
|
||||
# they may legitimately want fresh results with different inputs.
|
||||
# Dedup: share an in-flight Future across callers; refinement requests bypass since they may want fresh results.
|
||||
is_refinement = body.get("is_refinement", False)
|
||||
key = (session_id, group_id)
|
||||
if not is_refinement:
|
||||
@@ -174,8 +158,7 @@ async def generate_group_meta(session_id: str, body: dict):
|
||||
try:
|
||||
return await existing
|
||||
except Exception:
|
||||
# If the in-flight call failed, fall through and try again
|
||||
# ourselves rather than propagating someone else's error.
|
||||
# In-flight call failed; retry ourselves rather than propagate someone else's error.
|
||||
pass
|
||||
|
||||
future: asyncio.Future = asyncio.get_event_loop().create_future()
|
||||
@@ -197,7 +180,6 @@ async def generate_group_meta(session_id: str, body: dict):
|
||||
future.set_exception(e)
|
||||
raise
|
||||
finally:
|
||||
# Always clear our slot if we own it, so the next request runs fresh.
|
||||
if not is_refinement and _group_meta_inflight.get(key) is future:
|
||||
_group_meta_inflight.pop(key, None)
|
||||
|
||||
@@ -267,12 +249,7 @@ async def resume_session(session_id: str):
|
||||
|
||||
@agents.router.post("/sessions/{session_id}/warm-cache")
|
||||
async def warm_session_cache(session_id: str):
|
||||
"""Fire a max_tokens=1 dummy request through the agent path so
|
||||
Anthropic processes the system+tools prefix and writes the prompt
|
||||
cache. The next real user turn lands a cache hit instead of paying
|
||||
cold-start TTFT. Non-blocking, fire-and-forget on the frontend.
|
||||
Returns 200 even on failure (best-effort).
|
||||
"""
|
||||
"""Fire a max_tokens=1 dummy request to prime the Anthropic prompt cache; best-effort."""
|
||||
try:
|
||||
await agent_manager.warm_prompt_cache(session_id)
|
||||
except Exception:
|
||||
@@ -280,10 +257,6 @@ async def warm_session_cache(session_id: str):
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 9Router / Subscription endpoints
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@agents.router.get("/subscriptions/status")
|
||||
async def subscriptions_status():
|
||||
"""Check if 9Router is running and list connected providers."""
|
||||
@@ -292,8 +265,7 @@ async def subscriptions_status():
|
||||
return {"running": False, "providers": [], "models": []}
|
||||
connections = await get_providers()
|
||||
models = await get_models()
|
||||
# Frontend consumers (OnboardingModal, Settings) read
|
||||
# `data.providers.connections` — preserve that envelope here.
|
||||
# Frontend reads data.providers.connections; preserve the envelope.
|
||||
return {"running": True, "providers": {"connections": connections}, "models": models}
|
||||
|
||||
|
||||
@@ -310,11 +282,7 @@ async def subscriptions_connect(body: dict):
|
||||
if not is_running():
|
||||
raise HTTPException(status_code=503, detail="9Router not available. Please install Node.js.")
|
||||
|
||||
# If reconnecting a primary lane (e.g. gemini-cli), drop its cascade
|
||||
# siblings first. The registry prefers antigravity over gemini-cli
|
||||
# when both are present, so a stale antigravity token would keep
|
||||
# 400ing even after gemini-cli refreshes. Wiping the sibling forces
|
||||
# the registry onto the freshly reconnected lane.
|
||||
# Reconnecting gemini-cli must wipe antigravity; registry prefers AG and a stale AG token would 400 after gemini-cli refreshes.
|
||||
cascade = _PROVIDER_CASCADE_REMOVES.get(provider, [])
|
||||
if cascade:
|
||||
try:
|
||||
@@ -325,7 +293,6 @@ async def subscriptions_connect(body: dict):
|
||||
try:
|
||||
result = await start_oauth(provider)
|
||||
|
||||
# For auth_code flows, store pending state so the callback can exchange
|
||||
if result.get("flow") == "authorization_code" and result.get("state"):
|
||||
from backend.main import _pending_oauth
|
||||
_pending_oauth[result["state"]] = {
|
||||
@@ -399,8 +366,7 @@ async def subscriptions_models():
|
||||
|
||||
@agents.router.post("/probe-model")
|
||||
async def probe_model(body: dict):
|
||||
"""1-token health probe. Returns {ok, latency_ms} or {ok:false, error}
|
||||
or {ok:true, skipped:true} when the route's ambiguous (silent beats wrong)."""
|
||||
"""1-token health probe; returns latency or skipped when the route is ambiguous (silent beats wrong)."""
|
||||
import time as _time
|
||||
short_name = (body or {}).get("model") or ""
|
||||
if not short_name:
|
||||
@@ -459,8 +425,7 @@ async def probe_model(body: dict):
|
||||
except Exception as e:
|
||||
msg = str(e).splitlines()[0] if str(e) else type(e).__name__
|
||||
low = msg.lower()
|
||||
# Suppress transients — chat will retry naturally and probe-time aliasing
|
||||
# 404s often differ from how the chat path resolves the same id.
|
||||
# Suppress transients: chat retries naturally and probe-time alias 404s often differ from chat resolution.
|
||||
if any(s in low for s in (
|
||||
"timeout", "timed out",
|
||||
"connection reset", "connection aborted",
|
||||
@@ -489,7 +454,7 @@ async def list_models():
|
||||
try:
|
||||
conns = await _9r_providers()
|
||||
raw_providers = {c.get("provider", "") for c in conns if c.get("isActive") or c.get("testStatus") == "active"}
|
||||
# 9Router uses "claude"; our models use api="anthropic" — map across.
|
||||
# 9Router uses "claude"; our models use api="anthropic". Map across.
|
||||
_9R_TO_API = {
|
||||
"claude": "anthropic",
|
||||
"codex": "codex",
|
||||
@@ -501,8 +466,7 @@ async def list_models():
|
||||
logger.debug(f"Failed to fetch 9Router providers: {e}")
|
||||
|
||||
def _serialize(models: list[dict]) -> list[dict]:
|
||||
# Native models. Tiers describe the model itself; billing_kind
|
||||
# describes the user's wallet for it. Pricing is shown only for paid.
|
||||
# Tiers describe the model; billing_kind describes the wallet. Pricing shown only for paid.
|
||||
from backend.apps.agents.providers.registry import (
|
||||
COST_PER_1M_TOKENS,
|
||||
compute_tiers,
|
||||
@@ -533,7 +497,7 @@ async def list_models():
|
||||
"reasoning": bool(m.get("reasoning", False)),
|
||||
"input_cost_per_1m": input_cost,
|
||||
"output_cost_per_1m": output_cost,
|
||||
# Strict — subscription doesn't count. Pickerside uses Subscription chip.
|
||||
# Strict free; subscriptions show via the picker's Subscription chip.
|
||||
"is_free": billing_kind == "free",
|
||||
"billing_kind": billing_kind,
|
||||
"tiers": list(tiers),
|
||||
@@ -554,8 +518,7 @@ async def list_models():
|
||||
cc_variants = [m for m in anthropic_models if m.get("route") == "cc"]
|
||||
api_variants = [m for m in anthropic_models if m.get("route") == "api"]
|
||||
|
||||
# Pro mode shows two groups (Pro proxy + Anthropic alternates via cc/api);
|
||||
# own-key mode collapses to one Anthropic group using adaptive routing.
|
||||
# Pro mode splits into Pro proxy + Anthropic alternates; own-key collapses to one adaptive group.
|
||||
notes: list[dict] = []
|
||||
if is_openswarm_pro:
|
||||
result["OpenSwarm Pro"] = _serialize(adaptive)
|
||||
@@ -620,8 +583,7 @@ async def list_models():
|
||||
if visible:
|
||||
result[provider_name] = visible
|
||||
|
||||
# OR catalog fetched straight from openrouter.ai (independent of 9Router
|
||||
# boot state) so picker populates the moment a key lands.
|
||||
# Fetch OpenRouter catalog directly (independent of 9Router) so picker fills the moment a key lands.
|
||||
if has_openrouter_key:
|
||||
try:
|
||||
from backend.apps.agents.providers.registry import fetch_openrouter_models
|
||||
@@ -669,10 +631,7 @@ async def list_models():
|
||||
entries = sorted(by_vendor[vendor], key=lambda x: x["label"].lower())
|
||||
result[f"OpenRouter · {pretty}"] = entries
|
||||
|
||||
# User-configured custom OpenAI-compatible providers (Ollama Cloud, Together, etc).
|
||||
# Each provider becomes its own group in the picker; each model is addressed via
|
||||
# the `custom/<slug>/<model_id>` value, which `_find_builtin_model` synthesises
|
||||
# into a route='api' / api='custom' entry at request time.
|
||||
# Custom OpenAI-compatible providers (Ollama Cloud, Together, etc); addressed via custom/<slug>/<model_id>.
|
||||
from backend.apps.agents.providers.registry import _custom_provider_slug_for_lookup
|
||||
for cp in (getattr(settings, "custom_providers", None) or []):
|
||||
cp_name = (getattr(cp, "name", "") or "").strip()
|
||||
@@ -707,27 +666,14 @@ async def list_models():
|
||||
return {"models": result, "notes": notes}
|
||||
|
||||
|
||||
# Google's two OAuth lanes (gemini-cli and antigravity) share user-facing
|
||||
# meaning (both = "Google subscription") but 9Router treats them as
|
||||
# separate connections with independent token lifecycles. The registry
|
||||
# prefers `ag/` over `gc/` whenever AG is active because AG bypasses the
|
||||
# thoughtSignature validator that breaks multi-step tool turns. That
|
||||
# preference becomes a footgun when AG's token expires silently: the
|
||||
# user reconnects "Google", only gemini-cli refreshes, and every request
|
||||
# still routes through the stale AG token -> 400 Invalid argument.
|
||||
#
|
||||
# Cascade is one-directional. gemini-cli is the primary lane the UI
|
||||
# exposes; operations on it sweep antigravity too. Direct operations on
|
||||
# antigravity (e.g. an explicit AG opt-in/out path) MUST NOT cascade
|
||||
# back to gemini-cli or we'd nuke the user's main Google connection.
|
||||
# gemini-cli and antigravity are two Google OAuth lanes; registry prefers AG, so we cascade-wipe AG when reconnecting gemini-cli to avoid stale-AG 400s. One-directional: AG operations MUST NOT cascade back.
|
||||
_PROVIDER_CASCADE_REMOVES: dict[str, list[str]] = {
|
||||
"gemini-cli": ["antigravity"],
|
||||
}
|
||||
|
||||
|
||||
async def _delete_provider_connections(providers: list[str]) -> int:
|
||||
"""Delete all 9Router connections whose provider is in the given list.
|
||||
Returns the count actually removed. Silent if 9Router is unreachable."""
|
||||
"""Delete 9Router connections in `providers`; returns count removed, silent on 9Router unreachable."""
|
||||
import httpx
|
||||
from backend.apps.nine_router import NINE_ROUTER_API, get_providers
|
||||
try:
|
||||
@@ -748,12 +694,7 @@ async def _delete_provider_connections(providers: list[str]) -> int:
|
||||
|
||||
@agents.router.post("/subscriptions/disconnect")
|
||||
async def subscriptions_disconnect(body: dict):
|
||||
"""Disconnect a subscription provider via 9Router.
|
||||
|
||||
For Google's paired lanes (gemini-cli + antigravity), wipe BOTH so a
|
||||
subsequent reconnect lands on a clean slate instead of resurrecting
|
||||
a stale sibling.
|
||||
"""
|
||||
"""Disconnect a subscription provider via 9Router; cascades-wipe Google's paired lanes."""
|
||||
provider = body.get("provider", "")
|
||||
if not provider:
|
||||
raise HTTPException(status_code=400, detail="provider required")
|
||||
|
||||
@@ -1,23 +1,5 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Stdio MCP server exposing the MCP activation gate.
|
||||
|
||||
Tools:
|
||||
- MCPList: enumerate installed MCP servers (active + available).
|
||||
- MCPSearch(query): rank servers by relevance to a free-form query.
|
||||
- MCPActivate(server_name): activate a server for the rest of the session.
|
||||
|
||||
The activation gate is the dispatch-layer enforcement of the product invariant
|
||||
"all MCP actions only via ToolSearch": the model can only reach an MCP server's
|
||||
tools if the user has approved MCPActivate for that server, which appends to
|
||||
session.active_mcps. _build_mcp_servers in agent_manager.py intersects connected
|
||||
MCPs with that list before handing them to the SDK, so unactivated servers are
|
||||
literally unreachable — the gate cannot be bypassed by ignoring prompt rules.
|
||||
|
||||
HITL: the model's invocation of MCPActivate goes through agent_manager's pre-
|
||||
tool approval hook just like any other tool call — the user is prompted to
|
||||
approve activation in the standard ApprovalBar UI. No separate HITL inside this
|
||||
server.
|
||||
"""
|
||||
"""Stdio MCP server exposing the MCP activation gate (MCPList/MCPSearch/MCPActivate)."""
|
||||
|
||||
import json
|
||||
import os
|
||||
@@ -69,7 +51,7 @@ TOOLS = [
|
||||
"description": (
|
||||
"Request activation of an MCP server for this session. Triggers a "
|
||||
"user approval prompt; on approve the server's tools become callable "
|
||||
"next turn. Always confirm the server name via MCPList/MCPSearch first — "
|
||||
"next turn. Always confirm the server name via MCPList/MCPSearch first; "
|
||||
"invalid names return alternatives instead of activating."
|
||||
),
|
||||
"inputSchema": {
|
||||
@@ -81,7 +63,7 @@ TOOLS = [
|
||||
},
|
||||
"reason": {
|
||||
"type": "string",
|
||||
"description": "Why you need it — shown to the user in the approval prompt.",
|
||||
"description": "Why you need it; shown to the user in the approval prompt.",
|
||||
},
|
||||
},
|
||||
"required": ["server_name"],
|
||||
@@ -133,7 +115,7 @@ def format_servers(servers: list[dict], heading: str = "") -> str:
|
||||
name = s.get("name", "")
|
||||
desc = s.get("description") or f"{name} integration"
|
||||
status = s.get("status", "available")
|
||||
lines.append(f"- `{name}` [{status}] — {desc}")
|
||||
lines.append(f"- `{name}` [{status}]; {desc}")
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
@@ -189,15 +171,17 @@ def handle_tool_call(tool_name: str, arguments: dict) -> dict:
|
||||
"isError": True,
|
||||
}
|
||||
if result.get("status") == "already_active":
|
||||
return {"content": [{"type": "text", "text": f"`{server_name}` is already active for this session — its tools should be callable now."}]}
|
||||
return {"content": [{"type": "text", "text": f"`{server_name}` is already active for this session; its tools should be callable now."}]}
|
||||
if result.get("status") == "activated":
|
||||
return {
|
||||
"content": [{
|
||||
"type": "text",
|
||||
"text": (
|
||||
f"Activated `{server_name}`. Its tools (`mcp__{server_name}__*`) "
|
||||
f"will be callable on the NEXT turn. End this turn now and the user's "
|
||||
f"next message will see the new tools."
|
||||
f"are NOT callable in this turn; the transport snapshot is "
|
||||
f"already locked. This turn will end automatically and a "
|
||||
f"hidden continuation turn will fire with the new tools "
|
||||
f"loaded. Do not attempt any other tool call now."
|
||||
),
|
||||
}],
|
||||
}
|
||||
|
||||
@@ -11,7 +11,7 @@ class AgentConfig(BaseModel):
|
||||
system_prompt: Optional[str] = None
|
||||
allowed_tools: list[str] = Field(default_factory=lambda: ["Read", "Edit", "Write", "Bash", "Glob", "Grep", "AskUserQuestion"])
|
||||
max_turns: Optional[int] = None
|
||||
target_directory: Optional[str] = None # if None, uses repo root
|
||||
target_directory: Optional[str] = None
|
||||
dashboard_id: Optional[str] = None
|
||||
|
||||
class ApprovalRequest(BaseModel):
|
||||
@@ -55,26 +55,15 @@ class Message(BaseModel):
|
||||
forced_tools: Optional[list[str]] = None
|
||||
images: Optional[list[dict]] = None
|
||||
hidden: bool = False
|
||||
# Optional client-generated id used by the frontend to reconcile an
|
||||
# optimistic message bubble (rendered synchronously on send) with the
|
||||
# server-confirmed echo. Plumbed through send_message and round-tripped
|
||||
# back via the agent:message WS event so the frontend can dedupe.
|
||||
# Frontend-generated id for optimistic-bubble dedup against the server echo.
|
||||
client_message_id: Optional[str] = None
|
||||
# Wall-clock duration in milliseconds spent producing this message's
|
||||
# content. For thinking blocks: time from content_block_start →
|
||||
# content_block_stop. Lets the persisted ThinkingBubble show
|
||||
# "Thought for Ns" on reload instead of falling back to the static
|
||||
# "Thoughts" label. Optional for back-compat with messages saved
|
||||
# before this field existed.
|
||||
# Wall-clock ms producing this message's content; for thinking, content_block_start -> stop. Lets reloaded bubbles show "Thought for Ns".
|
||||
elapsed_ms: Optional[int] = None
|
||||
# Approximate output tokens for this message's content. For thinking
|
||||
# blocks we use the same char/3.6 heuristic the live UI uses so the
|
||||
# number frozen on the persisted bubble matches what the user saw
|
||||
# rising during the stream. Pure display, not billing.
|
||||
# Approx output tokens; thinking uses char/3.6 to match the live UI's count. Display only.
|
||||
tokens: Optional[int] = None
|
||||
# tool_count drives the "3 tools used" segment on the thinking pill.
|
||||
# Drives the "N tools used" segment on the thinking pill.
|
||||
tool_count: Optional[int] = None
|
||||
# combined input + output + children tokens for the turn (overloaded name).
|
||||
# Combined input + output + children tokens for the turn (overloaded name).
|
||||
input_tokens: Optional[int] = None
|
||||
|
||||
class MessageBranch(BaseModel):
|
||||
@@ -101,40 +90,22 @@ class AgentSession(BaseModel):
|
||||
allowed_tools: list[str] = Field(default_factory=list)
|
||||
max_turns: Optional[int] = None
|
||||
cwd: Optional[str] = None
|
||||
# Origin remote and branch resolved at session start. Persisted so a
|
||||
# resumed session reattaches to the same project even if the user has
|
||||
# since `cd`'d elsewhere; also surfaced in the session list UI so the
|
||||
# user can tell two sessions apart by repo.
|
||||
# Resolved at session start so resume reattaches to the same repo even after the user cd's elsewhere.
|
||||
repo_url: Optional[str] = None
|
||||
branch: Optional[str] = None
|
||||
created_at: datetime = Field(default_factory=datetime.now)
|
||||
closed_at: Optional[datetime] = None
|
||||
# Wall-clock of the first stream event from the agent SDK. Set once
|
||||
# at the start of the first turn so resumed sessions can show "first
|
||||
# response was at HH:MM" in the session list without rescanning the
|
||||
# message log.
|
||||
# Wall-clock of the first stream event so resumed sessions can show "first response at HH:MM" without rescan.
|
||||
first_response_at: Optional[datetime] = None
|
||||
# Operational log of HITL approval decisions, one entry per request:
|
||||
# {tool, behavior, decision_ms}. Persisted alongside the session so a
|
||||
# reload restores the full approval timeline (which calls were
|
||||
# approved, denied, and how long each took).
|
||||
# HITL approval log: {tool, behavior, decision_ms} per entry.
|
||||
approval_decisions: list[dict] = Field(default_factory=list)
|
||||
cost_usd: float = 0.0
|
||||
tokens: dict[str, int] = Field(default_factory=lambda: {"input": 0, "output": 0})
|
||||
# Total wall-clock ms the agent spent in `status="running"`. Accumulates
|
||||
# across turns; persists across resume. Used by the session-close
|
||||
# report so we can report "agent active time" alongside total session
|
||||
# duration. Off by default so legacy sessions deserialize cleanly.
|
||||
# Total ms in status="running", accumulated across turns/resume; powers session-close "agent active time".
|
||||
agent_active_ms: int = 0
|
||||
# Accumulated wall-clock ms spent on each model. Updated when the
|
||||
# active model changes (model switch) or on close. Surfaced in the
|
||||
# session header so the user can see "Sonnet: 45s · Haiku: 12s"
|
||||
# without scanning turns by hand.
|
||||
# Per-model wall-clock ms; updated on model switch or close.
|
||||
time_per_model: dict[str, int] = Field(default_factory=dict)
|
||||
# Per-tool latency rollup: { tool_name: { count, total_ms, max_ms } }.
|
||||
# Populated as tools complete. Surfaced in the session "tools used"
|
||||
# row so the user can see which tool calls were slow without
|
||||
# opening every turn.
|
||||
# Per-tool latency: { tool_name: { count, total_ms, max_ms } }.
|
||||
tool_latencies: dict[str, dict] = Field(default_factory=dict)
|
||||
browser_domains: list[str] = Field(default_factory=list)
|
||||
messages: list[Message] = Field(default_factory=list)
|
||||
@@ -146,58 +117,20 @@ class AgentSession(BaseModel):
|
||||
browser_id: Optional[str] = None
|
||||
parent_session_id: Optional[str] = None
|
||||
needs_fork: bool = False
|
||||
# Stronger than needs_fork: when True, the next turn drops `resume=`
|
||||
# entirely and replays history into a brand-new sdk_session_id. This
|
||||
# is the only way to make the bundled CLI re-read mcp_servers from
|
||||
# the rebuilt options dict — `fork_session=True` only forks the
|
||||
# conversation tree, it inherits the original transport's MCP server
|
||||
# set. Set after MCPActivate when prior turns exist so the newly
|
||||
# activated server's tools actually reach the model.
|
||||
# Stronger than needs_fork: drop resume= and replay history into a fresh sdk_session_id; fork_session alone won't re-read mcp_servers.
|
||||
needs_fresh_session: bool = False
|
||||
# Set when MCPActivate (or analogous activation) wants the agent to
|
||||
# auto-continue immediately after the current turn ends — without
|
||||
# requiring the user to type another message. The agent loop reads
|
||||
# this at the end of `_run_agent_loop`; if set, it clears it and
|
||||
# dispatches a new hidden turn with `pending_continuation_prompt` as
|
||||
# the prompt. Race-free vs. the original asyncio-task approach.
|
||||
# Auto-continue: agent loop dispatches a hidden turn at end-of-loop using pending_continuation_prompt. Race-free vs background tasks.
|
||||
pending_continuation: bool = False
|
||||
pending_continuation_prompt: Optional[str] = None
|
||||
# Sanitized server names (matching tools_lib._sanitize_server_name) of MCP
|
||||
# servers the model has explicitly activated this session via the
|
||||
# MCPActivate meta-tool. Empty by default — the gate in
|
||||
# _build_mcp_servers intersects connected MCPs with this list, so no
|
||||
# MCP tool is callable until the model searches for and activates a
|
||||
# server. The product invariant is that this is non-bypassable: the
|
||||
# filter lives at the dispatch layer (mcp_servers passed to the SDK),
|
||||
# not the prompt layer.
|
||||
# Sanitized server names model has explicitly activated this session; _build_mcp_servers intersects connected MCPs with this. Non-bypassable; dispatch-layer gate.
|
||||
active_mcps: list[str] = Field(default_factory=list)
|
||||
# Estimated framework preamble tokens (preset + tool defs + MCP descs +
|
||||
# composed prompt). Subtracted from displayed input for honest "this turn"
|
||||
# numbers. Heuristic; clamped >= 0.
|
||||
# Heuristic preamble tokens (preset + tool defs + MCP descs + composed prompt); subtracted from displayed input.
|
||||
framework_overhead_tokens: int = 0
|
||||
# Compaction state. compact_threshold_pct is the live ctx_used ratio
|
||||
# that triggers _maybe_compact at the next turn boundary — turn-based
|
||||
# thresholds break under uneven workloads (one big Bash dump fills
|
||||
# context fast; 30 chitchat turns barely move it). 0.65 = 130K of the
|
||||
# 200K standard tier. compacted_through_msg_id is the last message id
|
||||
# covered by the most recent summary so we don't re-summarize on
|
||||
# every turn.
|
||||
# Live ctx_used ratio triggering _maybe_compact at the next turn boundary; turn-based thresholds break under uneven workloads. 0.65 = 130K of 200K.
|
||||
compact_threshold_pct: float = 0.65
|
||||
compacted_through_msg_id: Optional[str] = None
|
||||
# Pre-send hard guard. Fires later than the compaction threshold —
|
||||
# 0.90 of 200K = 180K — to give the auto-compact path a chance to
|
||||
# bring the request back under the ceiling. If still over after
|
||||
# compaction, LRU-trim the oldest active_mcps. Past this we surface
|
||||
# the friendly context-overflow card instead of letting a 429 hit.
|
||||
# Hard pre-send guard at 0.90 (= 180K); past compaction we LRU-trim active_mcps, then surface the overflow card.
|
||||
context_soft_cap_pct: float = 0.90
|
||||
context_window: int = 200_000
|
||||
# How much the model should "think" before answering. Provider-agnostic
|
||||
# value that gets translated per-API in agent_manager:
|
||||
# off — no thinking
|
||||
# low — minimal thinking (fastest)
|
||||
# medium — balanced
|
||||
# high — extensive thinking (slowest, smartest)
|
||||
# auto — let the model / provider default decide (recommended)
|
||||
# Only applies to models flagged with reasoning: True in the registry.
|
||||
# Existing sessions without this field will default to "auto".
|
||||
# Provider-agnostic thinking level (off/low/medium/high/auto), translated per-API in agent_manager; only affects reasoning-flagged models.
|
||||
thinking_level: Literal["off", "low", "medium", "high", "auto"] = "auto"
|
||||
|
||||
@@ -9,19 +9,7 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ConnectionManager:
|
||||
"""Manages WebSocket connections and bridges HITL approval requests.
|
||||
|
||||
Every outbound event flows through the seq log so reconnecting
|
||||
clients can replay missed events. The send happens *under* the
|
||||
per-session lock yielded by `seq_log.stamp(...)`, which guarantees
|
||||
wire order matches seq order even under concurrent broadcasts.
|
||||
|
||||
A WS disconnect (`disconnect_session`) ONLY removes the socket
|
||||
from the connection registry. It does NOT cancel the underlying
|
||||
agent task. The task lives on `agent_manager.tasks`; only an
|
||||
explicit `agent:stop`, REST `/close`, natural completion, or
|
||||
process shutdown ends a run.
|
||||
"""
|
||||
"""Manages WebSocket connections and HITL approval bridging; events flow through seq_log so reconnects can replay."""
|
||||
|
||||
def __init__(self):
|
||||
self.connections: dict[str, list[WebSocket]] = {}
|
||||
@@ -53,19 +41,7 @@ class ConnectionManager:
|
||||
]
|
||||
|
||||
async def send_to_session(self, session_id: str, event: str, data: dict):
|
||||
"""Broadcast a session event with monotonic sequencing.
|
||||
|
||||
The send to every socket happens inside the seq_log lock so a
|
||||
slow/dead WS doesn't reorder events on the fast ones. If a
|
||||
single send raises (broken pipe, half-open socket), we log and
|
||||
continue — the ring buffer still has the event so the client
|
||||
will replay it on reconnect.
|
||||
|
||||
For terminal status events (completed/stopped/error) we also
|
||||
atomically persist the payload to disk; a client that returns
|
||||
after a process restart can then resolve the spinner via
|
||||
`seq_log.load_terminal(...)` instead of being stuck.
|
||||
"""
|
||||
"""Broadcast a session event with monotonic sequencing; terminal statuses also persist to disk."""
|
||||
async with seq_log.stamp(session_id, event, data) as (seq, payload_str):
|
||||
for ws in list(self.connections.get(session_id, [])):
|
||||
try:
|
||||
@@ -77,38 +53,17 @@ class ConnectionManager:
|
||||
await ws.send_text(payload_str)
|
||||
except Exception:
|
||||
logger.debug("send_to_session: global send failed", exc_info=True)
|
||||
# Persist terminal events under the lock so a concurrent
|
||||
# `agent:status: running` can't race past and overwrite
|
||||
# the disk file with a stale state.
|
||||
# Persist under the lock so a concurrent running status can't race past and overwrite with stale state.
|
||||
if event == "agent:status" and data.get("status") in TERMINAL_STATUSES:
|
||||
seq_log.persist_terminal(session_id, payload_str)
|
||||
|
||||
async def replay_to(
|
||||
self, session_id: str, websocket: WebSocket, last_seq: int
|
||||
) -> dict:
|
||||
"""Replay buffered events with seq > last_seq to one socket.
|
||||
|
||||
Returns a small ack envelope describing what happened so the
|
||||
caller (the WS handler) can send a `server:resume_ack` frame.
|
||||
|
||||
Three cases:
|
||||
1. `events` non-empty: replay them in order; ack carries
|
||||
`from_seq`, `to_seq`.
|
||||
2. No buffer at all (process restarted, session evicted)
|
||||
but a persisted terminal exists: send it; ack signals
|
||||
`terminal_only=True`.
|
||||
3. `last_seq` predates the oldest buffered seq: emit
|
||||
`agent:gap_detected`; client REST-refreshes the session.
|
||||
"""
|
||||
"""Replay buffered events with seq > last_seq; returns ack envelope for the resume handshake."""
|
||||
oldest, newest, events = seq_log.replay(session_id, last_seq)
|
||||
|
||||
# Check for gap FIRST. If the client's last_seq is below the
|
||||
# buffer's oldest seq, we can't deliver everything they
|
||||
# missed — silently replaying only the in-buffer tail would
|
||||
# leave a hole in their state. Tell them to REST-refresh
|
||||
# instead, even if the tail looks safe to send.
|
||||
# Treat last_seq=0 as "fresh client" — they want a full
|
||||
# replay of whatever's in the buffer, not a gap signal.
|
||||
# Gap-check first: if last_seq predates the buffer, signal REST-refresh; last_seq=0 means fresh client (full replay).
|
||||
if last_seq > 0 and oldest is not None and last_seq < oldest - 1:
|
||||
gap_payload = json.dumps({
|
||||
"event": "agent:gap_detected",
|
||||
@@ -162,7 +117,6 @@ class ConnectionManager:
|
||||
"to_seq": newest,
|
||||
}
|
||||
|
||||
# Nothing in memory. Try a persisted terminal event.
|
||||
terminal = seq_log.load_terminal(session_id)
|
||||
if terminal is not None:
|
||||
try:
|
||||
@@ -171,7 +125,6 @@ class ConnectionManager:
|
||||
pass
|
||||
return {"ok": True, "replayed": 1, "terminal_only": True}
|
||||
|
||||
# Nothing missed, nothing to replay. Caller's caught up.
|
||||
return {
|
||||
"ok": True,
|
||||
"replayed": 0,
|
||||
@@ -225,12 +178,7 @@ class ConnectionManager:
|
||||
return out
|
||||
|
||||
async def broadcast_global(self, event: str, data: dict):
|
||||
"""Send a message to all global (dashboard) connections.
|
||||
|
||||
Dashboard-scoped events don't go through the per-session seq
|
||||
log — they're not session-bound and the dashboard WS has its
|
||||
own resume story (full state refetch on reconnect).
|
||||
"""
|
||||
"""Send to all dashboard connections; bypasses seq_log (dashboard resumes via full state refetch)."""
|
||||
payload = json.dumps({"event": event, "data": data})
|
||||
for ws in list(self.global_connections):
|
||||
try:
|
||||
@@ -245,12 +193,7 @@ class ConnectionManager:
|
||||
sensitive_label: str | None = None,
|
||||
sensitive_why: str | None = None,
|
||||
) -> dict:
|
||||
"""Send an approval request and wait for the user's response.
|
||||
|
||||
Returns the approval decision dict. Times out after `timeout`
|
||||
seconds (default 10 minutes) so a forgotten request doesn't
|
||||
permanently park the agent.
|
||||
"""
|
||||
"""Send an approval request and wait for the user's decision; 10-minute timeout prevents permanent park."""
|
||||
future = asyncio.get_event_loop().create_future()
|
||||
self.pending_futures[request_id] = future
|
||||
|
||||
|
||||
@@ -13,16 +13,10 @@ async def health_lifespan():
|
||||
|
||||
health = SubApp("health", health_lifespan)
|
||||
|
||||
######################################
|
||||
# Health Check Endpoints #
|
||||
######################################
|
||||
|
||||
@health.router.get("/check")
|
||||
@typechecked
|
||||
async def check() -> PlainTextResponse:
|
||||
debug("Health check successful")
|
||||
# Use PlainTextResponse instead of JSONResponse for AWS ALB compatibility
|
||||
# ALB health checks can be sensitive to JSON responses and Content-Length headers
|
||||
return PlainTextResponse(
|
||||
content="OK",
|
||||
status_code=status.HTTP_200_OK,
|
||||
|
||||
@@ -55,9 +55,9 @@ BUILTIN_MODES: list[Mode] = [
|
||||
Mode(
|
||||
id="ask",
|
||||
name="Ask",
|
||||
description="Read-only conversation. Browse the codebase, search the web, and discuss ideas — but no edits, shells, or file writes.",
|
||||
description="Read-only conversation. Browse the codebase, search the web, and discuss ideas; but no edits, shells, or file writes.",
|
||||
system_prompt=(
|
||||
"You are in Ask mode — a read-only assistant. Keep responses "
|
||||
"You are in Ask mode; a read-only assistant. Keep responses "
|
||||
"natural and conversational. You CAN read files, search the "
|
||||
"codebase, and search/fetch the web. You CANNOT edit files, run "
|
||||
"shell commands, or otherwise modify anything; if the user asks "
|
||||
@@ -88,21 +88,21 @@ BUILTIN_MODES: list[Mode] = [
|
||||
name="App Builder",
|
||||
description="Create and iterate on reusable App artifacts.",
|
||||
system_prompt=(
|
||||
"You are an App Builder — an AI assistant that creates self-contained "
|
||||
"You are an App Builder; an AI assistant that creates self-contained "
|
||||
"web apps rendered in an iframe preview.\n\n"
|
||||
"Your working directory is a dedicated workspace folder pre-seeded with "
|
||||
"template files. Read the existing files before making changes.\n\n"
|
||||
"## Critical rules\n\n"
|
||||
"- The entry point MUST be named `index.html`. Never rename it or create "
|
||||
"a different HTML file as the main entry point.\n"
|
||||
"- Write files immediately when you have code ready — the user sees a "
|
||||
"- Write files immediately when you have code ready; the user sees a "
|
||||
"live preview that auto-refreshes from these files.\n"
|
||||
"- Always write the complete file content on first creation (do not use "
|
||||
"Edit for partial patches on new files).\n"
|
||||
"- For complex apps, split code into separate files (JS, CSS, etc.) "
|
||||
"and reference them from index.html with relative paths.\n"
|
||||
"- Always update meta.json with a short name and one-sentence description.\n"
|
||||
"- Build beautiful, polished UIs with modern design — dark themes, smooth "
|
||||
"- Build beautiful, polished UIs with modern design; dark themes, smooth "
|
||||
"transitions, proper spacing, and responsive layouts.\n\n"
|
||||
"Read the SKILL.md reference in your workspace for the full technical "
|
||||
"specification of the App platform (available globals, file conventions, "
|
||||
@@ -120,17 +120,17 @@ BUILTIN_MODES: list[Mode] = [
|
||||
name="Skill Builder",
|
||||
description="Create and iterate on skills using AI-assisted vibe coding.",
|
||||
system_prompt=(
|
||||
"You are a Skill Builder — an AI assistant that helps users create, "
|
||||
"You are a Skill Builder; an AI assistant that helps users create, "
|
||||
"refine, and iterate on Claude skills (SKILL.md files).\n\n"
|
||||
"## How Skills Work\n\n"
|
||||
"A skill is a Markdown file that teaches Claude how to perform a specific task. "
|
||||
"Skills have YAML frontmatter with `name` and `description` fields, followed by "
|
||||
"the skill body in Markdown. The description is the primary triggering mechanism — "
|
||||
"the skill body in Markdown. The description is the primary triggering mechanism; "
|
||||
"it tells Claude when to use the skill.\n\n"
|
||||
"## Your Working Directory\n\n"
|
||||
"Your working directory is a dedicated workspace folder for this skill. "
|
||||
"Write your output directly to these files using the Write tool:\n\n"
|
||||
"1. **SKILL.md** — The complete skill file with YAML frontmatter and Markdown body. "
|
||||
"1. **SKILL.md**; The complete skill file with YAML frontmatter and Markdown body. "
|
||||
"Example frontmatter:\n"
|
||||
" ```\n"
|
||||
" ---\n"
|
||||
@@ -138,34 +138,34 @@ BUILTIN_MODES: list[Mode] = [
|
||||
" description: When to trigger and what this skill does.\n"
|
||||
" ---\n"
|
||||
" ```\n\n"
|
||||
"2. **meta.json** — Metadata for the skill builder UI. Always write this file. Example:\n"
|
||||
"2. **meta.json**; Metadata for the skill builder UI. Always write this file. Example:\n"
|
||||
' {"name":"My Skill","description":"A short description","command":"my-skill"}\n\n'
|
||||
"Write these files immediately when you have content ready. The user can see "
|
||||
"a live preview that auto-refreshes from these files. Always write the "
|
||||
"complete file content (do not use Edit for partial patches on first creation).\n\n"
|
||||
"## Skill Creation Process\n\n"
|
||||
"1. **Understand intent** — Ask what the skill should do, when it should trigger, "
|
||||
"1. **Understand intent**; Ask what the skill should do, when it should trigger, "
|
||||
"and what the expected output format is.\n"
|
||||
"2. **Draft the skill** — Write a SKILL.md with clear instructions, examples, "
|
||||
"2. **Draft the skill**; Write a SKILL.md with clear instructions, examples, "
|
||||
"and good progressive disclosure.\n"
|
||||
"3. **Iterate** — Refine based on user feedback. Update the files each time.\n\n"
|
||||
"3. **Iterate**; Refine based on user feedback. Update the files each time.\n\n"
|
||||
"## Skill Writing Best Practices\n\n"
|
||||
"- Keep SKILL.md under 500 lines; use bundled reference files for large content.\n"
|
||||
"- The `description` frontmatter is the primary trigger. Make it slightly \"pushy\" — "
|
||||
"- The `description` frontmatter is the primary trigger. Make it slightly \"pushy\"; "
|
||||
"include both what the skill does AND specific contexts for when to use it.\n"
|
||||
"- Use imperative form in instructions.\n"
|
||||
"- Include examples with input/output pairs when helpful.\n"
|
||||
"- Define output formats explicitly with templates.\n"
|
||||
"- Use theory of mind — explain *why* things matter rather than just MUST directives.\n"
|
||||
"- Use theory of mind; explain *why* things matter rather than just MUST directives.\n"
|
||||
"- Think about edge cases, error handling, and progressive disclosure.\n\n"
|
||||
"## Skill Anatomy\n\n"
|
||||
"```\n"
|
||||
"skill-name/\n"
|
||||
"├── SKILL.md (required) — YAML frontmatter + Markdown instructions\n"
|
||||
"├── SKILL.md (required); YAML frontmatter + Markdown instructions\n"
|
||||
"└── Bundled Resources (optional)\n"
|
||||
" ├── scripts/ — Executable code for repetitive tasks\n"
|
||||
" ├── references/ — Docs loaded into context as needed\n"
|
||||
" └── assets/ — Files used in output\n"
|
||||
" ├── scripts/ ; Executable code for repetitive tasks\n"
|
||||
" ├── references/; Docs loaded into context as needed\n"
|
||||
" └── assets/ ; Files used in output\n"
|
||||
"```\n\n"
|
||||
"Be collaborative and flexible. If the user wants to \"just vibe\", skip the formal "
|
||||
"process and iterate freely. Always write updated files so the preview stays current."
|
||||
|
||||
@@ -14,10 +14,7 @@ from backend.config.paths import MODES_DIR as DATA_DIR
|
||||
@asynccontextmanager
|
||||
async def modes_lifespan():
|
||||
os.makedirs(DATA_DIR, exist_ok=True)
|
||||
# One-time migration: Chat was merged into Ask. Remove a stale built-in
|
||||
# chat.json if it still has its is_builtin=True signature so users don't
|
||||
# see two near-identical modes in the picker. Leave alone if a user has
|
||||
# diverged it (we don't want to wipe customizations).
|
||||
# Migration: Chat merged into Ask; drop a stale built-in chat.json but leave customized copies alone.
|
||||
chat_path = os.path.join(DATA_DIR, "chat.json")
|
||||
if os.path.exists(chat_path):
|
||||
try:
|
||||
|
||||
@@ -127,7 +127,7 @@ class OutputExecute(BaseModel):
|
||||
# running if the backend code touches anything outside the safe
|
||||
# data-shaping allowlist. The UI shows those warnings to the user and
|
||||
# re-submits with force=True after they click "Run Anyway." This is
|
||||
# a UX gate, not a security one — anyone holding the auth token can
|
||||
# a UX gate, not a security one; anyone holding the auth token can
|
||||
# set force=True; the value is providing the user explicit visibility
|
||||
# of what's about to execute.
|
||||
force: bool = False
|
||||
|
||||
@@ -139,7 +139,7 @@ def _inject_token_into_relative_urls(html: str, token: str) -> str:
|
||||
relative `<link href="styles.css">` / `<script src="x.js">`, so without
|
||||
this rewrite the sub-resource fetch lands at the auth middleware with no
|
||||
credentials and gets a 401. Idempotent: skips URLs that already carry a
|
||||
`token=` param. Skips absolute URLs (CDN, data:, etc.) — see prefix list.
|
||||
`token=` param. Skips absolute URLs (CDN, data:, etc.); see prefix list.
|
||||
"""
|
||||
if not token:
|
||||
return html
|
||||
@@ -235,7 +235,7 @@ def load_output(output_id: str) -> Output | None:
|
||||
# descend into. Without this skip-list the workspace endpoint reads
|
||||
# `node_modules/` (300 MB of MUI source, when it's a real dir and not a
|
||||
# symlink), `.venv/` (10k+ Python files from the hardlinked cache),
|
||||
# `__pycache__/`, `dist/`, `.git/`, etc — every 2 seconds while the
|
||||
# `__pycache__/`, `dist/`, `.git/`, etc; every 2 seconds while the
|
||||
# agent is active. Result: backend CPU pegged on JSON-serializing
|
||||
# auto-generated chunks the frontend will then throw away. The frontend
|
||||
# already filters these for display; this skip is the real fix.
|
||||
@@ -266,14 +266,14 @@ _WALK_MAX_FILE_BYTES = 256 * 1024
|
||||
def _walk_directory(folder: str) -> dict[str, str]:
|
||||
"""Walk a directory tree and return {relative_path: content} for all
|
||||
text files the user is actually authoring. Skips build/install
|
||||
directories AND truncates oversize files — both critical for the
|
||||
directories AND truncates oversize files; both critical for the
|
||||
polling endpoint, which is called every 2 s while the agent is
|
||||
writing code and would otherwise serialize hundreds of MB per poll."""
|
||||
files: dict[str, str] = {}
|
||||
if not os.path.isdir(folder):
|
||||
return files
|
||||
for root, dirs, filenames in os.walk(folder):
|
||||
# Mutate `dirs` in place — that's how os.walk skips a subtree.
|
||||
# Mutate `dirs` in place; that's how os.walk skips a subtree.
|
||||
# Doing it here means we never even stat the children, so a
|
||||
# 10k-file `.venv/` costs ~one stat (on the dir itself) instead
|
||||
# of 10k.
|
||||
@@ -288,7 +288,7 @@ def _walk_directory(folder: str) -> dict[str, str]:
|
||||
# mis-parsed.
|
||||
rel_path = os.path.relpath(full_path, folder).replace(os.sep, "/")
|
||||
try:
|
||||
# Stat first — cheap, lets us skip giant files without
|
||||
# Stat first; cheap, lets us skip giant files without
|
||||
# opening + reading them.
|
||||
size = os.path.getsize(full_path)
|
||||
if size > _WALK_MAX_FILE_BYTES:
|
||||
@@ -328,7 +328,7 @@ async def serve_workspace_file(workspace_id: str, filepath: str, _d: str = ""):
|
||||
content = _inject_data_into_html(content, input_json, result_json, backend_url_json)
|
||||
# Iframe sub-resource fetches (<link>, <script src>, <img>) drop the
|
||||
# parent's ?token= query string, so rewrite the HTML to put the token
|
||||
# back on every relative URL — otherwise sub-resources 401.
|
||||
# back on every relative URL; otherwise sub-resources 401.
|
||||
content = _inject_token_into_relative_urls(content, get_auth_token())
|
||||
|
||||
mime, _ = mimetypes.guess_type(filepath)
|
||||
@@ -512,7 +512,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
openswarm-ai/webapp-template snapshot (React + Vite + TS frontend
|
||||
with an optional FastAPI backend) into the workspace, allocates a
|
||||
free FRONTEND_PORT and writes it into both `.env` and
|
||||
`.env.example`. BACKEND_PORT stays NONE — the agent opts in with
|
||||
`.env.example`. BACKEND_PORT stays NONE; the agent opts in with
|
||||
`bash backend_init.sh`. Runtime spawn flips to `bash run.sh` and
|
||||
the preview pane points at `http://localhost:{FRONTEND_PORT}/`.
|
||||
`body.files` is ignored in this mode; the snapshot is the source
|
||||
@@ -524,7 +524,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
# An explicit non-empty `files` payload means the caller has flat-mode
|
||||
# content to write (a saved legacy Output being reseeded). Don't
|
||||
# clobber that with the React template even if template_mode is the
|
||||
# new default — the migration helper has its own path for that.
|
||||
# new default; the migration helper has its own path for that.
|
||||
effective_mode = body.template_mode
|
||||
if body.files:
|
||||
effective_mode = "flat"
|
||||
@@ -533,7 +533,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
# Idempotency guard: re-seeding an existing webapp_template
|
||||
# workspace would clobber the agent's edits (the helper uses
|
||||
# dirs_exist_ok=True + copytree). If `run.sh` already exists,
|
||||
# the workspace was seeded on a previous visit — skip the file
|
||||
# the workspace was seeded on a previous visit; skip the file
|
||||
# copy and only re-derive the frontend port from .env.
|
||||
from backend.apps.outputs.runtime import _find_free_port, _read_env_value
|
||||
already_seeded = os.path.exists(os.path.join(folder, "run.sh"))
|
||||
@@ -546,7 +546,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
else:
|
||||
frontend_port = _find_free_port()
|
||||
seed_webapp_template_workspace(folder, frontend_port)
|
||||
# SKILL.md still goes in workspace root — agent reads it for
|
||||
# SKILL.md still goes in workspace root; agent reads it for
|
||||
# context. Live content (user-editable via Skills page) is
|
||||
# injected into the system prompt regardless.
|
||||
with open(os.path.join(folder, "SKILL.md"), "w") as f:
|
||||
@@ -559,7 +559,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
# the Apps sidebar the moment the user kicks off generation.
|
||||
# Previously the record only landed when the editor's autosave
|
||||
# fired, which itself was gated on `files['index.html']` being
|
||||
# non-empty (a flat-template invariant) — meaning React+Vite
|
||||
# non-empty (a flat-template invariant); meaning React+Vite
|
||||
# apps that navigated-away mid-build had no way back. The record
|
||||
# is a thin pointer (name + workspace_id); the workspace itself
|
||||
# remains the source of truth for the code.
|
||||
@@ -591,7 +591,7 @@ async def seed_workspace(body: WorkspaceSeedRequest):
|
||||
"already_seeded": already_seeded,
|
||||
}
|
||||
|
||||
# Legacy flat path — unchanged.
|
||||
# Legacy flat path; unchanged.
|
||||
if body.files:
|
||||
for rel_path, content in body.files.items():
|
||||
full_path = os.path.normpath(os.path.join(folder, rel_path))
|
||||
@@ -693,7 +693,7 @@ async def runtime_restart(workspace_id: str):
|
||||
from backend.apps.outputs.runtime import manager as runtime_manager
|
||||
# Restart only if something's attached; otherwise this is a no-op
|
||||
# silently (a hard-reload click while the runtime was already torn
|
||||
# down — we'd rather not silently respawn an orphan).
|
||||
# down; we'd rather not silently respawn an orphan).
|
||||
rt = runtime_manager.get(workspace_id)
|
||||
if rt:
|
||||
await runtime_manager.restart(workspace_id, os.path.abspath(folder))
|
||||
@@ -725,7 +725,7 @@ async def write_workspace_file(workspace_id: str, filepath: str, body: dict):
|
||||
folder_norm = os.path.normpath(folder)
|
||||
full_path = os.path.normpath(os.path.join(folder, filepath))
|
||||
# `startswith(folder_norm + os.sep)` (not just folder_norm) so a workspace
|
||||
# `abc-123` can't be tricked into writing into a sibling `abc-1234-evil` —
|
||||
# `abc-123` can't be tricked into writing into a sibling `abc-1234-evil` ,
|
||||
# prefix-string collision rather than path-component containment. Today's
|
||||
# UUID-format ids make the collision unlikely in practice, but the check
|
||||
# is one character and immunizes future id schemes.
|
||||
@@ -931,7 +931,7 @@ async def execute_output(body: OutputExecute):
|
||||
# HITL gate: collect warnings up front. If the caller hasn't opted
|
||||
# in via force=True AND the code touches anything outside the safe
|
||||
# allowlist, return the warnings + the code itself so the UI can
|
||||
# show a preview dialog. No subprocess is spawned on this path —
|
||||
# show a preview dialog. No subprocess is spawned on this path ,
|
||||
# zero-cost when warnings exist, identical-to-before when they
|
||||
# don't.
|
||||
if not body.force:
|
||||
|
||||
@@ -1,16 +1,4 @@
|
||||
"""Per-workspace persistent backend runtime.
|
||||
|
||||
Each App (workspace) has at most one long-running `backend.py` subprocess
|
||||
managed by `AppRuntime`. Lifetime is reference-counted via the module-level
|
||||
`manager` singleton: when the first ViewEditor / DashboardViewCard /
|
||||
TerminalPanel attaches to a workspace, the process is spawned; when the
|
||||
last detaches, it's terminated. Multiple subscribers share the same
|
||||
process and the same in-memory log ring buffer.
|
||||
|
||||
This replaces the old one-shot `execute_backend_code` model for the
|
||||
"backend serves real HTTP endpoints" use case. The one-shot path stays
|
||||
around (see `executor.py`) for legacy `/api/outputs/execute` callers.
|
||||
"""
|
||||
"""Per-workspace persistent backend.py runtime; one AppRuntime per workspace, refcounted by manager singleton."""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
@@ -25,70 +13,28 @@ from typing import Callable, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Recent log lines kept in memory per runtime. Lets a Terminal tab that
|
||||
# opens mid-session replay the context that was already printed instead
|
||||
# of seeing a blank pane. 2000 lines ≈ a few hundred KB at worst —
|
||||
# bounded and predictable.
|
||||
# 2000 lines per runtime; lets a Terminal tab opened mid-session replay context. ~few hundred KB at worst.
|
||||
_LOG_BUFFER_LINES = 2000
|
||||
|
||||
# Seconds to wait after SIGTERM before escalating to SIGKILL. Most
|
||||
# well-behaved Python servers shut down well under a second; this is the
|
||||
# upper bound before we move on so a wedged process can't block a
|
||||
# workspace tear-down forever.
|
||||
# SIGTERM grace; well-behaved servers shut down under a second so 3s is enough.
|
||||
_TERMINATE_GRACE_SECONDS = 3
|
||||
|
||||
# How long we'll wait for Vite (or whatever frontend server bash run.sh
|
||||
# spawns) to bind on FRONTEND_PORT before giving up and reporting the
|
||||
# frontend as "not ready." Covers cold-start `npm install` (~60-90s on
|
||||
# typical hardware for the template's dependency set) plus the Vite
|
||||
# bind itself. After this we keep the runtime running — the user can
|
||||
# check the Terminal pane to see what went wrong — but stop blocking
|
||||
# the preview pane on a port that may never come up.
|
||||
# 180s covers npm install (60-90s on typical hardware) plus the Vite bind.
|
||||
_FRONTEND_BIND_TIMEOUT_SECONDS = 180
|
||||
# Drop from 0.5 → 0.08 because that 500ms window was ENTIRELY user-visible
|
||||
# preview latency — after Vite actually binds we'd wait up to half a second
|
||||
# before noticing and emitting runtime:status to the editor. 80ms TCP
|
||||
# probes are cheap (async open_connection on localhost, no DNS, no
|
||||
# handshake to a real upstream) and shave the perceived cold-start by
|
||||
# roughly half a second. The asyncio.open_connection call has its own
|
||||
# 500ms connect timeout for the failure case so a wedged listener won't
|
||||
# turn this into a tight CPU loop.
|
||||
# 80ms probe: dropping from 500ms was pure user-visible preview latency win; cheap on localhost.
|
||||
_FRONTEND_BIND_POLL_INTERVAL = 0.08
|
||||
|
||||
|
||||
# Process-wide mutex that serializes new-mode workspace boots so only
|
||||
# ONE vite optimizeDeps run is in flight at a time. Acquired in
|
||||
# `AppRuntime.start` (new-mode branch only) BEFORE the run.sh spawn,
|
||||
# released by `_await_frontend_bind` the instant vite emits its
|
||||
# "frontend ready" log line — or by the timeout / failure paths.
|
||||
#
|
||||
# Why a module-level asyncio.Lock and not part of AppRuntimeManager:
|
||||
# the lock has to be acquired BEFORE the runtime is registered in
|
||||
# manager.runtimes (which happens inside manager.attach's own
|
||||
# `_lock`), and we can't hold both locks at once without inviting
|
||||
# deadlock. Lifting to the module keeps the two locks fully
|
||||
# independent — the manager lock guards the runtime dict, this one
|
||||
# guards "is anyone currently mid-MUI-bundle?"
|
||||
# Module-level lock so only ONE vite optimizeDeps runs at a time; must be acquired before manager._lock to avoid deadlock with manager.attach.
|
||||
_vite_boot_lock = asyncio.Lock()
|
||||
|
||||
# Number of idle (zero-attachment) runtimes the manager keeps alive in
|
||||
# its LRU before reaping the oldest. Trades memory for instant
|
||||
# switch-back: clicking a previously-opened App reattaches to an
|
||||
# already-running vite + uvicorn instead of paying the ~1-2s spawn
|
||||
# cost. Bumped beyond 1 because the typical "App Builder" user keeps
|
||||
# 2-3 in-progress apps and ping-pongs between them.
|
||||
# Idle runtimes kept in LRU; trades memory for instant switch-back, beyond 1 because typical users ping-pong 2-3 apps.
|
||||
_MAX_IDLE_RUNTIMES = 3
|
||||
|
||||
# Cap on recent error lines kept per workspace runtime. The agent only
|
||||
# needs a snapshot of "what broke since my last write" — older errors
|
||||
# get dropped. 50 is enough to catch a babel error message + its stack
|
||||
# trace + a couple of related warnings without bloating the context.
|
||||
# Cap on recent error lines the agent gets; 50 is enough for babel error + stack + a few warnings.
|
||||
_RECENT_ERRORS_MAX = 50
|
||||
|
||||
# Regex that matches lines we want to surface back to the agent. Picks
|
||||
# up the common JS/TS/Python build-error formats vite, babel, tsc, and
|
||||
# uvicorn emit. Kept narrow on purpose so routine info logs and
|
||||
# deprecation warnings don't pollute the agent's context.
|
||||
# Narrow regex for build errors (vite, babel, tsc, uvicorn); keeps routine logs out of agent context.
|
||||
import re as _re
|
||||
_ERROR_PATTERNS = _re.compile(
|
||||
r"(?:"
|
||||
@@ -114,10 +60,10 @@ def _suspend_process_tree(proc: Optional[asyncio.subprocess.Process]) -> None:
|
||||
PROCESS GROUP (negative PID) when the child is a session leader,
|
||||
so vite + uvicorn + their npm/python subchildren all pause together.
|
||||
|
||||
No-op on Windows (SIGSTOP has no equivalent — the `OpenProcessToken` +
|
||||
No-op on Windows (SIGSTOP has no equivalent; the `OpenProcessToken` +
|
||||
`NtSuspendProcess` route works but isn't worth the win32 surface
|
||||
here; idle Windows runtimes just stay running, which is the current
|
||||
behavior). Failures here are swallowed — if the process already died
|
||||
behavior). Failures here are swallowed; if the process already died
|
||||
a stop signal is meaningless."""
|
||||
if proc is None or os.name == "nt":
|
||||
return
|
||||
@@ -126,7 +72,7 @@ def _suspend_process_tree(proc: Optional[asyncio.subprocess.Process]) -> None:
|
||||
return
|
||||
os.kill(proc.pid, signal.SIGSTOP)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
# Already-dead or out-of-permission — both safe to ignore.
|
||||
# Already-dead or out-of-permission; both safe to ignore.
|
||||
pass
|
||||
|
||||
|
||||
@@ -272,7 +218,7 @@ def _write_env_value(env_path: str, key: str, value: str) -> None:
|
||||
def _is_new_mode(workspace_path: str) -> bool:
|
||||
"""A workspace is "new-mode" (webapp-template scaffold) if it has a
|
||||
`run.sh` at its root. Old-mode workspaces are flat `index.html`-only
|
||||
apps that pre-date the template swap — they're served by OpenSwarm's
|
||||
apps that pre-date the template swap; they're served by OpenSwarm's
|
||||
own `/api/outputs/workspace/{ws}/serve/...` FastAPI route and have an
|
||||
optional `backend.py` we spawn directly.
|
||||
|
||||
@@ -299,7 +245,7 @@ def _read_env_value(env_path: str, key: str) -> Optional[str]:
|
||||
if k.strip() != key:
|
||||
continue
|
||||
v = v.strip()
|
||||
# Strip an inline `# comment`. Naive — bash semantics are
|
||||
# Strip an inline `# comment`. Naive; bash semantics are
|
||||
# more permissive, but values we write don't contain `#`.
|
||||
if "#" in v:
|
||||
v = v.split("#", 1)[0].rstrip()
|
||||
@@ -350,7 +296,7 @@ class AppRuntime:
|
||||
self.process: Optional[asyncio.subprocess.Process] = None
|
||||
self.log_buffer: deque[LogLine] = deque(maxlen=_LOG_BUFFER_LINES)
|
||||
self._subscribers: set[LogSubscriber] = set()
|
||||
# Recent build/runtime errors scraped from stderr — drained by
|
||||
# Recent build/runtime errors scraped from stderr; drained by
|
||||
# the agent's post-tool hook after Write/Edit so the agent sees
|
||||
# vite/babel/uvicorn errors in its next turn and can self-fix
|
||||
# instead of leaving the user with a red iframe overlay.
|
||||
@@ -402,7 +348,7 @@ class AppRuntime:
|
||||
without waiting for the subprocess to print anything.
|
||||
|
||||
- **Old-mode** (no `run.sh`): spawn `python -u backend.py` if
|
||||
present, with `PORT` env var. This is the legacy path —
|
||||
present, with `PORT` env var. This is the legacy path ,
|
||||
unchanged so flat-index.html apps keep working.
|
||||
|
||||
Returns True if a process is running after this call. False is
|
||||
@@ -432,7 +378,7 @@ class AppRuntime:
|
||||
ok = await self._start_new_mode()
|
||||
if not ok:
|
||||
# Spawn failed before the bind-poll task was
|
||||
# created — release synchronously so we don't
|
||||
# created; release synchronously so we don't
|
||||
# wedge the next workspace.
|
||||
_vite_boot_lock.release()
|
||||
return ok
|
||||
@@ -466,7 +412,7 @@ class AppRuntime:
|
||||
self.frontend_port = new_port
|
||||
_write_env_value(env_path, "FRONTEND_PORT", str(new_port))
|
||||
# BACKEND_PORT may be the literal string "NONE" (frontend-only
|
||||
# app — the common case) or a number once `backend_init.sh` has
|
||||
# app; the common case) or a number once `backend_init.sh` has
|
||||
# run. Only populate self.port when there's a real backend.
|
||||
if bp_raw and bp_raw != "NONE":
|
||||
try:
|
||||
@@ -518,7 +464,7 @@ class AppRuntime:
|
||||
self.process = None
|
||||
return False
|
||||
backend_note = f" + backend on {self.port}" if self.port else ""
|
||||
self._broadcast(LogLine("runtime", f"[runtime] bash run.sh started — frontend on {self.frontend_port}{backend_note} (pid {self.process.pid})"))
|
||||
self._broadcast(LogLine("runtime", f"[runtime] bash run.sh started; frontend on {self.frontend_port}{backend_note} (pid {self.process.pid})"))
|
||||
self._stdout_task = asyncio.create_task(self._pipe_stream(self.process.stdout, "stdout"))
|
||||
self._stderr_task = asyncio.create_task(self._pipe_stream(self.process.stderr, "stderr"))
|
||||
self._wait_task = asyncio.create_task(self._await_exit())
|
||||
@@ -536,7 +482,7 @@ class AppRuntime:
|
||||
`frontend_url` property reads.
|
||||
|
||||
Also responsible for releasing the module-level `_vite_boot_lock`
|
||||
— every exit path (success, process death, hard timeout) MUST
|
||||
; every exit path (success, process death, hard timeout) MUST
|
||||
release exactly once so the next queued workspace can start its
|
||||
own vite spawn. A try/finally on the lock guarantees that even
|
||||
an exception in the poll body doesn't strand the lock holding."""
|
||||
@@ -562,7 +508,7 @@ class AppRuntime:
|
||||
port = self.frontend_port
|
||||
deadline = asyncio.get_event_loop().time() + _FRONTEND_BIND_TIMEOUT_SECONDS
|
||||
while asyncio.get_event_loop().time() < deadline:
|
||||
# Stop polling if the process died — pointless to keep
|
||||
# Stop polling if the process died; pointless to keep
|
||||
# checking a port nothing will bind.
|
||||
if self.process is None or self.process.returncode is not None:
|
||||
return
|
||||
@@ -583,7 +529,7 @@ class AppRuntime:
|
||||
f"[runtime] frontend ready at http://127.0.0.1:{port}/",
|
||||
))
|
||||
# Release the vite-boot mutex the INSTANT vite is
|
||||
# ready — the next queued workspace can start its
|
||||
# ready; the next queued workspace can start its
|
||||
# own bundle now even though we'll keep streaming
|
||||
# logs for this one.
|
||||
_release_boot_lock()
|
||||
@@ -591,12 +537,12 @@ class AppRuntime:
|
||||
except (OSError, asyncio.TimeoutError):
|
||||
pass
|
||||
await asyncio.sleep(_FRONTEND_BIND_POLL_INTERVAL)
|
||||
# Timed out — keep the runtime up (Terminal might show useful
|
||||
# Timed out; keep the runtime up (Terminal might show useful
|
||||
# errors) but surface why the preview never appeared.
|
||||
self._broadcast(LogLine(
|
||||
"runtime",
|
||||
f"[runtime] frontend did NOT bind on port {port} after "
|
||||
f"{_FRONTEND_BIND_TIMEOUT_SECONDS}s — check the Terminal "
|
||||
f"{_FRONTEND_BIND_TIMEOUT_SECONDS}s; check the Terminal "
|
||||
f"for npm/vite errors.",
|
||||
))
|
||||
finally:
|
||||
@@ -613,7 +559,7 @@ class AppRuntime:
|
||||
self.port = _find_free_port()
|
||||
env = self._spawn_env_base()
|
||||
env["PORT"] = str(self.port)
|
||||
env["BACKEND_PORT"] = str(self.port) # alias — both common names work
|
||||
env["BACKEND_PORT"] = str(self.port) # alias; both common names work
|
||||
try:
|
||||
# -u forces unbuffered stdout/stderr so the Terminal pane
|
||||
# sees lines in real time, not whenever Python decides to
|
||||
@@ -648,7 +594,7 @@ class AppRuntime:
|
||||
async with self._lock:
|
||||
if not self.process or self.process.returncode is not None:
|
||||
# Still cancel the bind poller in case stop() races a
|
||||
# never-launched runtime — defensive no-op otherwise.
|
||||
# never-launched runtime; defensive no-op otherwise.
|
||||
if self._frontend_ready_task and not self._frontend_ready_task.done():
|
||||
self._frontend_ready_task.cancel()
|
||||
return
|
||||
@@ -695,7 +641,7 @@ class AppRuntime:
|
||||
|
||||
def _broadcast(self, line: LogLine) -> None:
|
||||
self.log_buffer.append(line)
|
||||
# Snapshot subscribers — they can self-remove during dispatch.
|
||||
# Snapshot subscribers; they can self-remove during dispatch.
|
||||
for cb in list(self._subscribers):
|
||||
try:
|
||||
cb(line)
|
||||
@@ -704,7 +650,7 @@ class AppRuntime:
|
||||
|
||||
def _maybe_capture_error(self, text: str) -> None:
|
||||
"""If a stderr/stdout line matches a known build-error pattern,
|
||||
record it for the next agent-tool drain. Tests every line —
|
||||
record it for the next agent-tool drain. Tests every line ,
|
||||
cheap (single regex search) and only the matching ones land in
|
||||
the buffer."""
|
||||
if _ERROR_PATTERNS.search(text):
|
||||
@@ -739,7 +685,7 @@ class AppRuntimeManager:
|
||||
Reference-counts attachments so we don't kill a backend when one
|
||||
Terminal closes while another is still subscribed. First attach
|
||||
spawns; final detach moves the runtime into an LRU idle pool
|
||||
instead of stopping it immediately — so re-clicking a recent App
|
||||
instead of stopping it immediately; so re-clicking a recent App
|
||||
is instant. The oldest runtime gets reaped once the pool exceeds
|
||||
_MAX_IDLE_RUNTIMES."""
|
||||
|
||||
@@ -755,14 +701,14 @@ class AppRuntimeManager:
|
||||
|
||||
async def attach(self, workspace_id: str, workspace_path: str) -> AppRuntime:
|
||||
revived = False
|
||||
# Defined here so every code path below leaves it bound — the
|
||||
# Defined here so every code path below leaves it bound; the
|
||||
# revive-idle branch used to skip the assignment, leaving the
|
||||
# post-lock `if dead is not None:` check throwing UnboundLocalError.
|
||||
dead: Optional[AppRuntime] = None
|
||||
async with self._lock:
|
||||
rt = self.runtimes.get(workspace_id)
|
||||
if rt is None:
|
||||
# Maybe the runtime is sitting idle in the LRU — revive
|
||||
# Maybe the runtime is sitting idle in the LRU; revive
|
||||
# it without paying the spawn cost again.
|
||||
idle_rt = self._idle_lru.pop(workspace_id, None)
|
||||
if idle_rt is not None and idle_rt.running:
|
||||
@@ -775,7 +721,7 @@ class AppRuntimeManager:
|
||||
_resume_process_tree(rt.process)
|
||||
else:
|
||||
if idle_rt is not None:
|
||||
# Stale idle entry — process died while idling.
|
||||
# Stale idle entry; process died while idling.
|
||||
# Drop and spawn a fresh one below; old one
|
||||
# gets stopped outside the lock.
|
||||
dead = idle_rt
|
||||
@@ -784,7 +730,7 @@ class AppRuntimeManager:
|
||||
else:
|
||||
# Workspace paths shouldn't change for a given id, but if
|
||||
# somehow they did (e.g. the user moved the workspace
|
||||
# folder), trust the latest caller — they have the
|
||||
# folder), trust the latest caller; they have the
|
||||
# current truth.
|
||||
rt.workspace_path = workspace_path
|
||||
self._attached[workspace_id] = self._attached.get(workspace_id, 0) + 1
|
||||
@@ -811,7 +757,7 @@ class AppRuntimeManager:
|
||||
if rt is None:
|
||||
return
|
||||
# If the process is already dead, no point keeping it
|
||||
# around — just clean up. Otherwise move to the LRU AND
|
||||
# around; just clean up. Otherwise move to the LRU AND
|
||||
# SIGSTOP the process tree so it consumes 0% CPU while
|
||||
# idle. The matching SIGCONT lives in attach() above.
|
||||
if not rt.running:
|
||||
@@ -853,7 +799,7 @@ class AppRuntimeManager:
|
||||
"""If `file_path` falls under one of the live workspace
|
||||
runtimes' workspace_path, drain that workspace's recent
|
||||
build/runtime errors. Returns [] if no workspace owns the path
|
||||
or no errors are queued — caller can treat empty as 'all clear'.
|
||||
or no errors are queued; caller can treat empty as 'all clear'.
|
||||
Used by agent_manager's post-tool hook so the agent sees vite /
|
||||
babel / uvicorn errors right after a Write/Edit completes."""
|
||||
if not file_path:
|
||||
@@ -862,7 +808,7 @@ class AppRuntimeManager:
|
||||
abs_path = os.path.abspath(file_path)
|
||||
except Exception:
|
||||
return []
|
||||
# Walk both active and idle runtimes — the user might have
|
||||
# Walk both active and idle runtimes; the user might have
|
||||
# navigated away from the workspace mid-build, but the agent
|
||||
# could still be editing files; the LRU keeps the runtime alive
|
||||
# for ~3 idle slots.
|
||||
|
||||
@@ -1,5 +1 @@
|
||||
"""(Reserved for future use; intentionally empty.)
|
||||
|
||||
The service-sync layer ships opaque payload dicts through `submit()` —
|
||||
no Pydantic shape exposed in the public repo.
|
||||
"""
|
||||
"""Reserved; service-sync ships opaque payload dicts via submit(), no Pydantic shape exposed."""
|
||||
|
||||
@@ -1,8 +1,4 @@
|
||||
"""Centralized credential resolution for LLM API calls.
|
||||
|
||||
Supports multiple providers: Anthropic (native), OpenAI, Gemini,
|
||||
OpenRouter, and user-configured custom providers.
|
||||
"""
|
||||
"""Resolve LLM credentials for the configured provider."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -26,18 +22,14 @@ def _check_9router() -> bool:
|
||||
|
||||
|
||||
def validate_credentials(settings: AppSettings, provider: str = "anthropic") -> None:
|
||||
"""Raise ValueError if credentials are missing for the given provider.
|
||||
|
||||
Allows through if 9Router is running as a fallback.
|
||||
Handles both display names ('Anthropic') and lowercase ('anthropic').
|
||||
"""
|
||||
"""Raise ValueError if the provider has no usable credentials."""
|
||||
p = provider.lower().strip()
|
||||
|
||||
# 9Router-backed providers don't need traditional credentials
|
||||
# 9Router handles its own credentials.
|
||||
if p == "9router":
|
||||
return
|
||||
|
||||
# If 9Router is running, all providers are accessible
|
||||
# 9Router proxies every provider, so if it's up we don't need keys here.
|
||||
if _check_9router():
|
||||
return
|
||||
|
||||
@@ -62,21 +54,20 @@ def validate_credentials(settings: AppSettings, provider: str = "anthropic") ->
|
||||
return
|
||||
raise ValueError("OpenRouter API key not configured. Set it in Settings.")
|
||||
elif p in ("xai", "meta", "deepseek", "mistral", "qwen", "cohere"):
|
||||
# These route through OpenRouter — need either OpenRouter key or 9Router
|
||||
# These providers route through OpenRouter, so its key is required.
|
||||
if getattr(settings, "openrouter_api_key", None):
|
||||
return
|
||||
raise ValueError(f"{provider} requires an OpenRouter API key, or connect a subscription via 9Router.")
|
||||
else:
|
||||
# Custom provider — check if it exists in custom_providers
|
||||
for cp in getattr(settings, "custom_providers", []):
|
||||
if cp.name.lower() == p:
|
||||
return
|
||||
# Unknown provider — allow through (create_provider will handle the error)
|
||||
# Let create_provider raise for unknown providers; not our job here.
|
||||
return
|
||||
|
||||
|
||||
def get_provider_credentials(settings: AppSettings, provider: str) -> dict[str, str]:
|
||||
"""Return credential dict for a specific provider."""
|
||||
"""Return the credential dict for the given provider."""
|
||||
p = provider.lower().strip()
|
||||
validate_credentials(settings, provider)
|
||||
|
||||
@@ -97,45 +88,17 @@ def get_provider_credentials(settings: AppSettings, provider: str) -> dict[str,
|
||||
if p == "openrouter":
|
||||
return {"api_key": getattr(settings, "openrouter_api_key", "") or ""}
|
||||
|
||||
# Custom provider
|
||||
for cp in getattr(settings, "custom_providers", []):
|
||||
if cp.name.lower() == p:
|
||||
# Substitute a placeholder when the user left api_key blank
|
||||
# — local OpenAI-compatible servers (LM Studio, Ollama, etc.)
|
||||
# ignore the Bearer header but downstream callers may insist
|
||||
# on non-empty values.
|
||||
# Local OpenAI-compatible servers (LM Studio, Ollama) ignore the key; placeholder keeps downstream callers happy.
|
||||
key = (cp.api_key or "").strip() or "no-auth-required"
|
||||
return {"api_key": key, "base_url": cp.base_url}
|
||||
|
||||
raise ValueError(f"No credentials for provider: {provider}")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Legacy helpers (kept for backward compat during migration)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def get_agent_sdk_env(settings: AppSettings) -> dict[str, str]:
|
||||
"""Return the env dict for ClaudeAgentOptions based on connection mode.
|
||||
|
||||
DEPRECATED: Use create_provider() from providers.registry instead.
|
||||
"""
|
||||
validate_credentials(settings, "anthropic")
|
||||
|
||||
if getattr(settings, "connection_mode", "own_key") == "openswarm-pro":
|
||||
proxy_url = getattr(settings, "openswarm_proxy_url", None) or OPENSWARM_DEFAULT_PROXY_URL
|
||||
return {
|
||||
"ANTHROPIC_AUTH_TOKEN": getattr(settings, "openswarm_bearer_token", ""),
|
||||
"ANTHROPIC_BASE_URL": proxy_url,
|
||||
}
|
||||
|
||||
return {"ANTHROPIC_API_KEY": settings.anthropic_api_key}
|
||||
|
||||
|
||||
def get_anthropic_client(settings: AppSettings) -> anthropic.AsyncAnthropic:
|
||||
"""Return a configured AsyncAnthropic client based on connection mode.
|
||||
|
||||
Priority: managed mode → 9Router subscription → API key
|
||||
"""
|
||||
"""Return an AsyncAnthropic client for the user's current connection mode."""
|
||||
import anthropic
|
||||
|
||||
if getattr(settings, "connection_mode", "own_key") == "openswarm-pro":
|
||||
@@ -145,11 +108,11 @@ def get_anthropic_client(settings: AppSettings) -> anthropic.AsyncAnthropic:
|
||||
base_url=proxy_url,
|
||||
)
|
||||
|
||||
# Prefer API key when set
|
||||
# Prefer the user's own API key when present.
|
||||
if settings.anthropic_api_key:
|
||||
return anthropic.AsyncAnthropic(api_key=settings.anthropic_api_key)
|
||||
|
||||
# Fall back to 9Router subscription (free for users with Claude/ChatGPT/Gemini subscriptions)
|
||||
# Fall back to 9Router (free for users with Claude/ChatGPT/Gemini subscriptions).
|
||||
if _check_9router():
|
||||
return anthropic.AsyncAnthropic(
|
||||
api_key="9router",
|
||||
@@ -160,17 +123,7 @@ def get_anthropic_client(settings: AppSettings) -> anthropic.AsyncAnthropic:
|
||||
|
||||
|
||||
def get_anthropic_client_for_model(settings: AppSettings, api_model: str) -> anthropic.AsyncAnthropic:
|
||||
"""Return a client configured for the given resolved model id.
|
||||
|
||||
When api_model carries a 9Router prefix (cc/, cx/, gc/, cp-), the client
|
||||
targets 9Router directly — even if connection_mode is openswarm-pro. This
|
||||
is what lets pinned-route models like "sonnet-cc" actually reach the
|
||||
user's own subscription instead of getting sent through the managed proxy
|
||||
with an unrecognizable model id. cp- is the prefix we use when registering
|
||||
user-configured custom OpenAI-compatible providers in 9Router.
|
||||
Otherwise delegates to get_anthropic_client() for the default mode-driven
|
||||
routing.
|
||||
"""
|
||||
"""Route 9Router-prefixed models (cc/, cx/, gc/, cp-) straight to 9Router so user subscriptions reach their own accounts."""
|
||||
import anthropic
|
||||
if isinstance(api_model, str) and (
|
||||
api_model.startswith(("cc/", "cx/", "gc/")) or api_model.startswith("cp-")
|
||||
|
||||
@@ -4,27 +4,27 @@ from typing import Optional, Any, Literal
|
||||
DEFAULT_SYSTEM_PROMPT = (
|
||||
"You are a personal AI assistant running inside OpenSwarm.\n\n"
|
||||
"## Core Behavior\n"
|
||||
"Act, don't ask. When a tool can accomplish the task, call it immediately — "
|
||||
"Act, don't ask. When a tool can accomplish the task, call it immediately; "
|
||||
"do not describe what you would do, do not ask for confirmation, just execute. "
|
||||
"The user expects results, not plans.\n"
|
||||
"If ANY available tool is relevant to the user's request, use it. Never respond "
|
||||
'with "I can do X for you" or "Would you like me to..." — just do it. '
|
||||
'with "I can do X for you" or "Would you like me to..."; just do it. '
|
||||
"A tool call is always better than a text explanation of what the tool would do.\n"
|
||||
"For multi-step tasks, chain tool calls in sequence — don't stop after one step "
|
||||
"For multi-step tasks, chain tool calls in sequence; don't stop after one step "
|
||||
"to ask if you should continue. Complete the entire task, then report the results.\n"
|
||||
"Be adaptable. If one approach fails, try a different tool or strategy instead of "
|
||||
"giving up or repeating the same action. Always stay focused on what the user "
|
||||
"actually wants to accomplish — their intent matters more than the specific method.\n\n"
|
||||
"actually wants to accomplish; their intent matters more than the specific method.\n\n"
|
||||
"## Tool Priority\n"
|
||||
"1. Connected MCP tools — fastest and most reliable. Use ToolSearch to discover "
|
||||
"1. Connected MCP tools; fastest and most reliable. Use ToolSearch to discover "
|
||||
"what integrations are available if you're unsure.\n"
|
||||
"2. WebSearch / WebFetch — for general web lookups when no MCP tool fits.\n"
|
||||
"3. BrowserAgent — last resort, only for visual interaction with websites, "
|
||||
"2. WebSearch / WebFetch; for general web lookups when no MCP tool fits.\n"
|
||||
"3. BrowserAgent; last resort, only for visual interaction with websites, "
|
||||
"filling forms, or tasks no other tool can handle.\n\n"
|
||||
"## Style\n"
|
||||
"Do not narrate routine tool calls — just call the tool.\n"
|
||||
"Do not narrate routine tool calls; just call the tool.\n"
|
||||
"After tool calls complete, present the results directly. Do not recap which "
|
||||
"tools you called or why — the user can see tool calls in the UI.\n"
|
||||
"tools you called or why; the user can see tool calls in the UI.\n"
|
||||
"Keep responses brief and direct. Use plain language.\n"
|
||||
"If you genuinely need clarification on something ambiguous, use the "
|
||||
"AskUserQuestion tool. Never ask questions inline in plain text.\n"
|
||||
@@ -40,61 +40,39 @@ class AppSettings(BaseModel):
|
||||
default_thinking_level: Literal["off", "low", "medium", "high", "auto"] = "auto"
|
||||
zoom_sensitivity: float = 50.0
|
||||
theme: str = "dark"
|
||||
# App Builder workspaces seed a React template that ships with its own
|
||||
# theme toggle ("Light" / "Dark" at the bottom of the sidebar). By
|
||||
# default the template should follow the user's OS appearance; once
|
||||
# the user explicitly toggles it inside any one app the override
|
||||
# persists across every subsequently-built app via this field
|
||||
# (the template fetches /api/settings on mount and PUTs back here on
|
||||
# toggle, so the preference is shared even though each app runs from
|
||||
# its own vite port / localStorage origin).
|
||||
# null = follow system / no override; 'light' or 'dark' = sticky.
|
||||
# Shared across App Builder workspaces (each runs its own vite port / localStorage origin); null = follow system.
|
||||
app_template_theme_override: Optional[Literal["light", "dark"]] = None
|
||||
new_agent_shortcut: str = "Meta+l"
|
||||
anthropic_api_key: Optional[str] = None
|
||||
browser_homepage: str = "https://www.google.com"
|
||||
# Multi-provider API keys
|
||||
openai_api_key: Optional[str] = None
|
||||
google_api_key: Optional[str] = None
|
||||
openrouter_api_key: Optional[str] = None
|
||||
custom_providers: list["CustomProvider"] = Field(default_factory=list)
|
||||
# Dashboard / UI preferences
|
||||
auto_select_mode_on_new_agent: bool = False
|
||||
expand_new_chats_in_dashboard: bool = False
|
||||
auto_reveal_sub_agents: bool = True
|
||||
dev_mode: bool = False
|
||||
allow_experimental_updates: bool = False
|
||||
# Subscription tokens (from CLI tools, alternative to API keys)
|
||||
claude_subscription_token: Optional[str] = None
|
||||
openai_subscription_token: Optional[str] = None
|
||||
gemini_subscription_token: Optional[str] = None
|
||||
# User profile (collected during onboarding)
|
||||
user_name: Optional[str] = None
|
||||
user_email: Optional[str] = None
|
||||
user_use_case: Optional[str] = None
|
||||
user_referral_source: Optional[str] = None
|
||||
# Per-MCP dismissal map for the preflight suggestion modal. Keyed by
|
||||
# the curated ToolDefinition.name (e.g. "Google Workspace"); value is
|
||||
# an ISO timestamp of dismissal. Used by mcp_preflight._build_available_shortlist
|
||||
# to suppress suggestions the user has explicitly waved off.
|
||||
# Suppresses preflight suggestion modal entries the user dismissed; keyed by ToolDefinition.name, value ISO timestamp.
|
||||
dismissed_mcp_suggestions: dict[str, str] = Field(default_factory=dict)
|
||||
# Analytics: opted in by default, user can toggle off
|
||||
analytics_opt_in: bool = True
|
||||
installation_id: Optional[str] = None
|
||||
first_opened_at: Optional[str] = None # ISO timestamp of first app open
|
||||
# OpenSwarm Pro subscription
|
||||
connection_mode: str = "own_key" # "own_key" | "openswarm-pro"
|
||||
first_opened_at: Optional[str] = None
|
||||
connection_mode: str = "own_key"
|
||||
openswarm_bearer_token: Optional[str] = None
|
||||
openswarm_proxy_url: Optional[str] = None # default resolved in credentials.py
|
||||
openswarm_subscription_plan: Optional[str] = None # "hobby"|"pro"|"pro_plus"|"ultra"
|
||||
openswarm_subscription_expires: Optional[str] = None # ISO 8601
|
||||
openswarm_usage_cached: Optional[dict] = None # {count, limit, window_end_at}
|
||||
# Identity (v1.0.29+). Populated after a successful sign-in via the cloud's
|
||||
# /api/auth/signin-activate endpoint (Google OAuth or email magic link).
|
||||
# Stripe checkout also populates these because the cloud's bearer-mint
|
||||
# always returns user info. Distinct from user_email above which was
|
||||
# historically a self-reported onboarding field — the values agree once
|
||||
# sign-in completes (server-validated wins).
|
||||
openswarm_proxy_url: Optional[str] = None
|
||||
openswarm_subscription_plan: Optional[str] = None
|
||||
openswarm_subscription_expires: Optional[str] = None
|
||||
openswarm_usage_cached: Optional[dict] = None
|
||||
# Server-validated identity from /api/auth/signin-activate; user_email above is the self-reported onboarding value.
|
||||
user_id: Optional[str] = None
|
||||
signin_method: Optional[Literal["google", "stripe", "email"]] = None
|
||||
|
||||
|
||||
@@ -37,10 +37,7 @@ async def settings_lifespan():
|
||||
import asyncio as _asyncio
|
||||
|
||||
async def _boot_router_then_sync():
|
||||
"""Start 9Router (if any apikey-routed provider is configured)
|
||||
then push our key-based connections into it. Sequential because
|
||||
sync_* helpers no-op when 9Router isn't running yet — running
|
||||
them post-boot guarantees the connections actually land."""
|
||||
"""Boot 9Router then push key-based connections (sequential: sync helpers no-op pre-boot)."""
|
||||
needs_router = any([
|
||||
getattr(s, "google_api_key", None),
|
||||
getattr(s, "openai_api_key", None),
|
||||
@@ -76,13 +73,7 @@ settings = SubApp("settings", settings_lifespan)
|
||||
|
||||
|
||||
def _migrate_legacy_fields(raw: dict) -> dict:
|
||||
"""Translate deprecated field names/values so they survive into the new schema.
|
||||
|
||||
Pre-launch scaffolding used `connection_mode="managed"` and
|
||||
`openswarm_auth_token`; production names are `"openswarm-pro"` and
|
||||
`openswarm_bearer_token`. Zero known users are affected, but keep the
|
||||
mapping for safety.
|
||||
"""
|
||||
"""Translate deprecated pre-launch field names ('managed', 'openswarm_auth_token') into production schema."""
|
||||
if raw.get("connection_mode") == "managed":
|
||||
raw["connection_mode"] = "openswarm-pro"
|
||||
if "openswarm_auth_token" in raw and "openswarm_bearer_token" not in raw:
|
||||
@@ -102,26 +93,19 @@ def load_settings() -> AppSettings:
|
||||
return AppSettings()
|
||||
|
||||
|
||||
# Single threading.Lock guards every write to SETTINGS_FILE — protects against
|
||||
# corruption from two requests racing through the file system. Async callers
|
||||
# offload the actual write to the default thread pool (run_in_executor), so
|
||||
# the lock works for both sync and thread-pool execution paths.
|
||||
# threading.Lock guards every SETTINGS_FILE write; works for sync paths and async run_in_executor paths.
|
||||
_settings_write_lock = threading.Lock()
|
||||
|
||||
|
||||
def _atomic_write_settings(payload: dict) -> None:
|
||||
"""Internal: serialise payload to SETTINGS_FILE atomically.
|
||||
Always called via save_settings* — don't invoke directly."""
|
||||
"""Atomic SETTINGS_FILE write; call via save_settings*, not directly."""
|
||||
with _settings_write_lock:
|
||||
os.makedirs(DATA_DIR, exist_ok=True)
|
||||
fd, tmp = tempfile.mkstemp(prefix=".settings.", suffix=".tmp", dir=DATA_DIR)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
||||
json.dump(payload, f, indent=2)
|
||||
# On Windows, os.replace can transiently fail with PermissionError
|
||||
# if Defender or another reader holds the destination open. One
|
||||
# retry after a short backoff handles every real-world case
|
||||
# without masking genuine permission bugs.
|
||||
# Windows: Defender can briefly lock the destination; one retry handles every real case.
|
||||
for attempt in range(2):
|
||||
try:
|
||||
os.replace(tmp, SETTINGS_FILE)
|
||||
@@ -139,24 +123,17 @@ def _atomic_write_settings(payload: dict) -> None:
|
||||
|
||||
|
||||
def save_settings(settings_obj: AppSettings) -> None:
|
||||
"""Synchronously persist settings atomically. Thread-safe.
|
||||
Use from sync paths (analytics collector, lifespans). Async callers should
|
||||
prefer save_settings_async to avoid blocking the event loop on Windows
|
||||
where Defender scans can stretch the write to 50-200ms."""
|
||||
"""Sync atomic persist; thread-safe. Async callers should prefer save_settings_async (Defender can stretch writes to 50-200ms)."""
|
||||
_atomic_write_settings(settings_obj.model_dump())
|
||||
|
||||
|
||||
async def save_settings_async(settings_obj: AppSettings) -> None:
|
||||
"""Async-safe atomic save. Runs the file I/O in the default thread pool
|
||||
so the FastAPI event loop stays responsive while the write completes.
|
||||
Shares the threading.Lock with the sync variant for safe interleaving."""
|
||||
"""Async atomic save via thread pool; shares the lock with the sync variant."""
|
||||
payload = settings_obj.model_dump()
|
||||
loop = asyncio.get_running_loop()
|
||||
await loop.run_in_executor(None, _atomic_write_settings, payload)
|
||||
|
||||
|
||||
# Backward-compat alias. Existing sync callers (analytics collector, analytics
|
||||
# lifespan) continue to work; new async callers should use save_settings_async.
|
||||
def _save_settings(settings_obj: AppSettings) -> None:
|
||||
save_settings(settings_obj)
|
||||
|
||||
@@ -172,14 +149,12 @@ async def update_settings(body: AppSettings):
|
||||
|
||||
old = load_settings()
|
||||
|
||||
# Sync the settings state (secrets stripped).
|
||||
secret_keys = {"anthropic_api_key", "openai_api_key", "google_api_key", "openrouter_api_key",
|
||||
"claude_subscription_token", "openai_subscription_token", "gemini_subscription_token",
|
||||
"openswarm_bearer_token", "installation_id"}
|
||||
safe = {k: v for k, v in body.model_dump().items() if k not in secret_keys}
|
||||
_sync(safe)
|
||||
|
||||
# Identify user in service-sync when profile is set/changed
|
||||
if (body.user_email and body.user_email != getattr(old, "user_email", None)) or \
|
||||
(body.user_name and body.user_name != getattr(old, "user_name", None)):
|
||||
from backend.apps.service.client import identify as _identify
|
||||
@@ -227,8 +202,7 @@ async def update_settings(body: AppSettings):
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Boot+sync runs off the request path — ensure_running() can take 5min
|
||||
# on first install (npm pull) and would freeze the event loop.
|
||||
# Off the request path: ensure_running() can take 5min on first install (npm pull) and would freeze the loop.
|
||||
if google_changed or openai_changed or openrouter_changed or custom_providers_changed:
|
||||
async def _boot_and_sync_keys(
|
||||
google_key: str | None,
|
||||
@@ -275,12 +249,7 @@ async def update_settings(body: AppSettings):
|
||||
any_keyed_added,
|
||||
))
|
||||
|
||||
# When openswarm-pro mode or bearer token changes, register a `claude`
|
||||
# apikey connection in 9Router that proxies through our cloud. This
|
||||
# makes the CLI's built-in WebSearch work on non-Claude primaries for
|
||||
# Pro users — the CLI's Anthropic delegation path now has a working
|
||||
# Claude route via 9Router, instead of hitting "no credentials for
|
||||
# provider: claude".
|
||||
# On pro-mode/bearer change, register a `claude` apikey connection in 9Router so CLI WebSearch works on non-Claude primaries.
|
||||
pro_mode_old = getattr(old, "connection_mode", None) == "openswarm-pro"
|
||||
pro_mode_new = getattr(body, "connection_mode", None) == "openswarm-pro"
|
||||
bearer_old = getattr(old, "openswarm_bearer_token", None)
|
||||
@@ -305,25 +274,13 @@ class AppThemeOverridePayload(BaseModel):
|
||||
|
||||
@settings.router.get("/app-theme-override")
|
||||
async def get_app_theme_override():
|
||||
"""Cross-app theme preference for App Builder workspaces.
|
||||
|
||||
Returns the current override (or `null` for follow-system). Apps
|
||||
served from the template fetch this on mount so a toggle inside
|
||||
any one app sticks across every future app the user builds. Each
|
||||
app workspace runs on its own vite port (separate localStorage
|
||||
origin), so the backend is the only place this can live."""
|
||||
"""Cross-app theme preference for App Builder workspaces; backend-held because each app uses its own localStorage origin."""
|
||||
return {"mode": load_settings().app_template_theme_override}
|
||||
|
||||
|
||||
@settings.router.put("/app-theme-override")
|
||||
async def put_app_theme_override(body: AppThemeOverridePayload):
|
||||
"""MERGE the theme override into AppSettings. The general PUT
|
||||
/api/settings endpoint replaces the whole AppSettings object —
|
||||
sending a partial body there would default every secret-bearing
|
||||
field (api keys, subscription tokens), which logs the user out
|
||||
and pops the SignInGate. This dedicated endpoint mutates only
|
||||
`app_template_theme_override` and leaves every other field
|
||||
untouched."""
|
||||
"""MERGE the override; the general PUT /api/settings replaces the whole object and would blank secrets, logging the user out."""
|
||||
current = load_settings()
|
||||
current.app_template_theme_override = body.mode
|
||||
await save_settings_async(current)
|
||||
|
||||
@@ -10,11 +10,7 @@ class Skill(BaseModel):
|
||||
content: str
|
||||
file_path: str = ""
|
||||
command: str = ""
|
||||
# Skills that OpenSwarm ships as part of the platform (e.g. the App
|
||||
# Builder reference) get this flag set. The UI hides the delete
|
||||
# button for them and the DELETE endpoint refuses with 409. Content
|
||||
# is still editable — the whole point is that users can tune how
|
||||
# the platform-internal agents behave.
|
||||
# Platform-shipped skills (e.g. App Builder): UI hides delete and DELETE returns 409, but content stays editable so users can tune them.
|
||||
built_in: bool = False
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user