diff --git a/backend/apps/agents/browser/browser_agent.py b/backend/apps/agents/browser/browser_agent.py index b9c2235d..ddc78055 100644 --- a/backend/apps/agents/browser/browser_agent.py +++ b/backend/apps/agents/browser/browser_agent.py @@ -2195,8 +2195,8 @@ def _find_reusable_card(dashboard_id: str, url: str, parent_session_id: str | No if not (dashboard_id and want): return "" try: - from backend.apps.dashboards.dashboards import _load - cards = _load(dashboard_id).layout.browser_cards + from backend.apps.dashboards.dashboards import load as load_dashboard + cards = load_dashboard(dashboard_id).layout.browser_cards except Exception: return "" from backend.apps.agents.agent_manager import agent_manager @@ -2218,10 +2218,10 @@ def _find_reusable_card(dashboard_id: str, url: str, parent_session_id: str | No async def _create_browser_card(dashboard_id: str, url: str, parent_session_id: str | None = None) -> str: """Create a new browser card on the dashboard and return its browser_id.""" - from backend.apps.dashboards.dashboards import _load, _save + from backend.apps.dashboards.dashboards import load as load_dashboard, save as save_dashboard from backend.apps.dashboards.models import BrowserCardPosition, BrowserTab - dashboard = _load(dashboard_id) + dashboard = load_dashboard(dashboard_id) browser_id = f"browser-{uuid4().hex[:8]}" tab_id = f"tab-{uuid4().hex[:8]}" tab = BrowserTab(id=tab_id, url=url or "https://www.google.com", title="") @@ -2238,7 +2238,7 @@ async def _create_browser_card(dashboard_id: str, url: str, parent_session_id: s ) dashboard.layout.browser_cards[browser_id] = card dashboard.updated_at = datetime.now() - _save(dashboard) + save_dashboard(dashboard) await ws_manager.broadcast_global("dashboard:browser_card_added", { "dashboard_id": dashboard_id, diff --git a/backend/apps/agents/manager/prompt/prompt_context.py b/backend/apps/agents/manager/prompt/prompt_context.py index ba243566..7b7c1a86 100644 --- a/backend/apps/agents/manager/prompt/prompt_context.py +++ b/backend/apps/agents/manager/prompt/prompt_context.py @@ -107,7 +107,7 @@ def build_browser_context(dashboard_id: str | None, selected_browser_ids: list[s if not dashboard_id: return None try: - from backend.apps.dashboards.dashboards import _load as load_dashboard + from backend.apps.dashboards.dashboards import load as load_dashboard dashboard = load_dashboard(dashboard_id) except Exception: return None diff --git a/backend/apps/dashboards/dashboards.py b/backend/apps/dashboards/dashboards.py index a6453304..7b0158a7 100644 --- a/backend/apps/dashboards/dashboards.py +++ b/backend/apps/dashboards/dashboards.py @@ -22,7 +22,7 @@ from backend.config.json_store import read_json_or_none, atomic_write_json OLD_LAYOUT_FILE = os.path.join(OLD_LAYOUT_DIR, "layout.json") -def _load_all() -> list[Dashboard]: +def p_load_all() -> list[Dashboard]: result = [] if not os.path.exists(DATA_DIR): return result @@ -40,11 +40,11 @@ def _load_all() -> list[Dashboard]: return result -def _save(dashboard: Dashboard): +def save(dashboard: Dashboard): atomic_write_json(os.path.join(DATA_DIR, f"{dashboard.id}.json"), dashboard.model_dump(mode="json")) -def _load(dashboard_id: str) -> Dashboard: +def load(dashboard_id: str) -> Dashboard: path = os.path.join(DATA_DIR, f"{dashboard_id}.json") data = read_json_or_none(path) if data is None: @@ -52,15 +52,15 @@ def _load(dashboard_id: str) -> Dashboard: return Dashboard(**data) -def _delete(dashboard_id: str): +def p_delete(dashboard_id: str): path = os.path.join(DATA_DIR, f"{dashboard_id}.json") if os.path.exists(path): os.remove(path) -def _migrate_if_needed(): +def p_migrate_if_needed(): """One-time migration: if no dashboards exist, create 'Dashboard 1' from old layout.""" - existing = _load_all() + existing = p_load_all() if existing: return @@ -78,7 +78,7 @@ def _migrate_if_needed(): logger.exception("Failed to read old layout.json, using empty layout") dashboard = Dashboard(name="Dashboard 1", layout=layout) - _save(dashboard) + save(dashboard) logger.info(f"Created default dashboard: {dashboard.id}") if os.path.exists(SESSIONS_DIR): @@ -102,7 +102,7 @@ def _migrate_if_needed(): @asynccontextmanager async def dashboards_lifespan(): os.makedirs(DATA_DIR, exist_ok=True) - _migrate_if_needed() + p_migrate_if_needed() yield @@ -111,7 +111,7 @@ dashboards = SubApp("dashboards", dashboards_lifespan) @dashboards.router.get("/list") async def list_dashboards(): - all_dashboards = _load_all() + all_dashboards = p_load_all() all_dashboards.sort(key=lambda d: d.updated_at or d.created_at, reverse=True) items = [] for d in all_dashboards: @@ -132,7 +132,7 @@ async def list_dashboards(): @dashboards.router.post("/create") async def create_dashboard(body: DashboardCreate): dashboard = Dashboard(name=body.name) - _save(dashboard) + save(dashboard) return dashboard.model_dump(mode="json") @@ -212,7 +212,7 @@ async def seed_orchestration_demo(dashboard_id: str): drags it into a new agent and asks for a PDF report; which delegates back to this seeded agent. """ - _load(dashboard_id) # validate dashboard exists + load(dashboard_id) # validate dashboard exists session_id = uuid4().hex now = datetime.now() @@ -294,7 +294,7 @@ async def seed_orchestration_demo(dashboard_id: str): @dashboards.router.post("/{dashboard_id}/generate-name") async def generate_name(dashboard_id: str): - dashboard = _load(dashboard_id) + dashboard = load(dashboard_id) if not dashboard.auto_named and dashboard.name != "Untitled Dashboard": return {"name": dashboard.name, "auto_named": dashboard.auto_named} @@ -319,7 +319,7 @@ async def generate_name(dashboard_id: str): from backend.apps.settings.credentials import get_anthropic_client_for_model from backend.apps.agents.providers.registry import resolve_aux_model global_settings = load_settings() - aux_model, _aux_base = await resolve_aux_model(global_settings, preferred_tier="haiku") + aux_model, _ = await resolve_aux_model(global_settings, preferred_tier="haiku") client = get_anthropic_client_for_model(global_settings, aux_model) # Mirrors generate_title's hardening: the tasks are inert text to LABEL, never answer, @@ -352,19 +352,19 @@ async def generate_name(dashboard_id: str): dashboard.name = fallback dashboard.auto_named = True dashboard.updated_at = datetime.now() - _save(dashboard) + save(dashboard) return {"name": dashboard.name, "auto_named": True} @dashboards.router.get("/{dashboard_id}") async def get_dashboard(dashboard_id: str): - dashboard = _load(dashboard_id) + dashboard = load(dashboard_id) return dashboard.model_dump(mode="json") @dashboards.router.put("/{dashboard_id}") async def update_dashboard(dashboard_id: str, body: DashboardUpdate): - dashboard = _load(dashboard_id) + dashboard = load(dashboard_id) if body.name is not None: dashboard.name = body.name dashboard.auto_named = False @@ -377,13 +377,13 @@ async def update_dashboard(dashboard_id: str, body: DashboardUpdate): # Only a real screenshot write moves the sort key; layout/rename saves don't reorder. dashboard.preview_updated_at = now dashboard.updated_at = now - _save(dashboard) + save(dashboard) return dashboard.model_dump(mode="json") @dashboards.router.delete("/{dashboard_id}") async def delete_dashboard(dashboard_id: str): - _load(dashboard_id) + load(dashboard_id) if os.path.exists(SESSIONS_DIR): for fname in os.listdir(SESSIONS_DIR): @@ -409,13 +409,13 @@ async def delete_dashboard(dashboard_id: str): except Exception: logger.warning(f"Failed to delete active session {sid} during dashboard deletion") - _delete(dashboard_id) + p_delete(dashboard_id) return {"ok": True} @dashboards.router.post("/{dashboard_id}/duplicate") async def duplicate_dashboard(dashboard_id: str): - source = _load(dashboard_id) + source = load(dashboard_id) source_data = source.model_dump(mode="json") new_id = uuid4().hex now = datetime.now().isoformat() diff --git a/backend/apps/discord_mcp_shim/server.py b/backend/apps/discord_mcp_shim/server.py index a806841f..caf61c98 100644 --- a/backend/apps/discord_mcp_shim/server.py +++ b/backend/apps/discord_mcp_shim/server.py @@ -24,10 +24,10 @@ ALLOWED_GUILDS = set( # -- MCP tool definitions (exposed to the agent) --------------------------- # Names match the original mcp-discord surface so prompts that referenced -# `discord_send` etc. keep working. inputSchema deliberately matches what +# `discordp_send` etc. keep working. inputSchema deliberately matches what # the original package documented. -TOOLS = [ +P_TOOLS = [ { "name": "discord_login", "description": "Verify the Discord bot helper is reachable. Returns the bot's joined guilds.", @@ -214,7 +214,7 @@ TOOLS = [ # -- HTTP plumbing --------------------------------------------------------- -def _call( +def p_call( method: str, path: str, *, @@ -267,17 +267,17 @@ def _call( return 0, f"Request failed: {e!r}" -def _err(text: str) -> dict: +def p_err(text: str) -> dict: return {"content": [{"type": "text", "text": f"Error: {text}"}], "isError": True} -def _ok(payload) -> dict: +def p_ok(payload) -> dict: if isinstance(payload, str): return {"content": [{"type": "text", "text": payload}]} return {"content": [{"type": "text", "text": json.dumps(payload, indent=2, default=str)}]} -def _check_guild(guild_id: str) -> str | None: +def p_check_guild(guild_id: str) -> str | None: """Return an error string if guild_id is outside the user-authorized set, else None. The set is sourced from OPENSWARM_DISCORD_GUILD_IDS env var (CSV) which @@ -297,71 +297,71 @@ def _check_guild(guild_id: str) -> str | None: # -- Tool implementations -------------------------------------------------- -def handle_tool_call(name: str, args: dict) -> dict: +def p_handle_tool_call(name: str, args: dict) -> dict: if name == "discord_login": - status, body = _call("GET", "/users/@me/guilds") + status, body = p_call("GET", "/users/@me/guilds") if status != 200: - return _err(f"Discord proxy unreachable (HTTP {status}): {body}") - return _ok({"connected": True, "guilds": body}) + return p_err(f"Discord proxy unreachable (HTTP {status}): {body}") + return p_ok({"connected": True, "guilds": body}) if name == "discord_get_server_info": gid = str(args.get("guild_id", "")) - if (e := _check_guild(gid)): return _err(e) - status, body = _call("GET", f"/guilds/{gid}") - return _ok(body) if status == 200 else _err(f"HTTP {status}: {body}") + if (e := p_check_guild(gid)): return p_err(e) + status, body = p_call("GET", f"/guilds/{gid}") + return p_ok(body) if status == 200 else p_err(f"HTTP {status}: {body}") if name == "discord_list_channels": gid = str(args.get("guild_id", "")) - if (e := _check_guild(gid)): return _err(e) - status, body = _call("GET", f"/guilds/{gid}/channels") - return _ok(body) if status == 200 else _err(f"HTTP {status}: {body}") + if (e := p_check_guild(gid)): return p_err(e) + status, body = p_call("GET", f"/guilds/{gid}/channels") + return p_ok(body) if status == 200 else p_err(f"HTTP {status}: {body}") if name == "discord_create_text_channel": gid = str(args.get("guild_id", "")) - if (e := _check_guild(gid)): return _err(e) + if (e := p_check_guild(gid)): return p_err(e) payload: dict = {"name": args.get("name", ""), "type": 0} if args.get("parent_id"): payload["parent_id"] = args["parent_id"] if args.get("topic"): payload["topic"] = args["topic"] - status, body = _call("POST", f"/guilds/{gid}/channels", body=payload) - return _ok(body) if status in (200, 201) else _err(f"HTTP {status}: {body}") + status, body = p_call("POST", f"/guilds/{gid}/channels", body=payload) + return p_ok(body) if status in (200, 201) else p_err(f"HTTP {status}: {body}") if name == "discord_create_category": gid = str(args.get("guild_id", "")) - if (e := _check_guild(gid)): return _err(e) - status, body = _call("POST", f"/guilds/{gid}/channels", body={"name": args.get("name", ""), "type": 4}) - return _ok(body) if status in (200, 201) else _err(f"HTTP {status}: {body}") + if (e := p_check_guild(gid)): return p_err(e) + status, body = p_call("POST", f"/guilds/{gid}/channels", body={"name": args.get("name", ""), "type": 4}) + return p_ok(body) if status in (200, 201) else p_err(f"HTTP {status}: {body}") if name == "discord_edit_category": cid = str(args.get("channel_id", "")) payload: dict = {} if args.get("name"): payload["name"] = args["name"] - status, body = _call("PATCH", f"/channels/{cid}", body=payload) - return _ok(body) if status == 200 else _err(f"HTTP {status}: {body}") + status, body = p_call("PATCH", f"/channels/{cid}", body=payload) + return p_ok(body) if status == 200 else p_err(f"HTTP {status}: {body}") if name == "discord_delete_category" or name == "discord_delete_channel": cid = str(args.get("channel_id", "")) - status, body = _call("DELETE", f"/channels/{cid}") - return _ok({"deleted": True}) if status in (200, 204) else _err(f"HTTP {status}: {body}") + status, body = p_call("DELETE", f"/channels/{cid}") + return p_ok({"deleted": True}) if status in (200, 204) else p_err(f"HTTP {status}: {body}") if name == "discord_send": cid = str(args.get("channel_id", "")) content = str(args.get("content", "")) - status, body = _call("POST", f"/channels/{cid}/messages", body={"content": content}) - return _ok(body) if status in (200, 201) else _err(f"HTTP {status}: {body}") + status, body = p_call("POST", f"/channels/{cid}/messages", body={"content": content}) + return p_ok(body) if status in (200, 201) else p_err(f"HTTP {status}: {body}") if name == "discord_read_messages": cid = str(args.get("channel_id", "")) limit = max(1, min(int(args.get("limit", 50) or 50), 100)) - status, body = _call("GET", f"/channels/{cid}/messages", query={"limit": limit}) - return _ok(body) if status == 200 else _err(f"HTTP {status}: {body}") + status, body = p_call("GET", f"/channels/{cid}/messages", query={"limit": limit}) + return p_ok(body) if status == 200 else p_err(f"HTTP {status}: {body}") if name == "discord_add_reaction": cid = str(args.get("channel_id", "")) mid = str(args.get("message_id", "")) emoji = str(args.get("emoji", "")) # Discord's URL needs the emoji urlencoded; passes through. - status, body = _call("PUT", f"/channels/{cid}/messages/{mid}/reactions/{urllib.parse.quote(emoji, safe='')}/@me") - return _ok({"added": emoji}) if status in (200, 204) else _err(f"HTTP {status}: {body}") + status, body = p_call("PUT", f"/channels/{cid}/messages/{mid}/reactions/{urllib.parse.quote(emoji, safe='')}/@me") + return p_ok({"added": emoji}) if status in (200, 204) else p_err(f"HTTP {status}: {body}") if name == "discord_add_multiple_reactions": cid = str(args.get("channel_id", "")) @@ -369,46 +369,46 @@ def handle_tool_call(name: str, args: dict) -> dict: emojis = args.get("emojis", []) or [] results = [] for e in emojis: - status, body = _call("PUT", f"/channels/{cid}/messages/{mid}/reactions/{urllib.parse.quote(str(e), safe='')}/@me") + status, body = p_call("PUT", f"/channels/{cid}/messages/{mid}/reactions/{urllib.parse.quote(str(e), safe='')}/@me") results.append({"emoji": e, "ok": status in (200, 204), "status": status}) - return _ok({"reactions": results}) + return p_ok({"reactions": results}) if name == "discord_get_forum_channels": gid = str(args.get("guild_id", "")) - if (e := _check_guild(gid)): return _err(e) - status, body = _call("GET", f"/guilds/{gid}/channels") - if status != 200: return _err(f"HTTP {status}: {body}") + if (e := p_check_guild(gid)): return p_err(e) + status, body = p_call("GET", f"/guilds/{gid}/channels") + if status != 200: return p_err(f"HTTP {status}: {body}") # Filter to type 15 (forum). Discord channel types reference: # GUILD_FORUM = 15 forums = [ch for ch in (body or []) if isinstance(ch, dict) and ch.get("type") == 15] - return _ok(forums) + return p_ok(forums) if name == "discord_create_forum_post": fid = str(args.get("forum_id", "")) - status, body = _call("POST", f"/channels/{fid}/threads", body={ + status, body = p_call("POST", f"/channels/{fid}/threads", body={ "name": args.get("name", ""), "message": {"content": args.get("content", "")}, }) - return _ok(body) if status in (200, 201) else _err(f"HTTP {status}: {body}") + return p_ok(body) if status in (200, 201) else p_err(f"HTTP {status}: {body}") if name == "discord_get_forum_post": cid = str(args.get("channel_id", "")) mid = str(args.get("message_id", "")) - status, body = _call("GET", f"/channels/{cid}/messages/{mid}") - return _ok(body) if status == 200 else _err(f"HTTP {status}: {body}") + status, body = p_call("GET", f"/channels/{cid}/messages/{mid}") + return p_ok(body) if status == 200 else p_err(f"HTTP {status}: {body}") if name == "discord_reply_to_forum": cid = str(args.get("channel_id", "")) content = str(args.get("content", "")) - status, body = _call("POST", f"/channels/{cid}/messages", body={"content": content}) - return _ok(body) if status in (200, 201) else _err(f"HTTP {status}: {body}") + status, body = p_call("POST", f"/channels/{cid}/messages", body={"content": content}) + return p_ok(body) if status in (200, 201) else p_err(f"HTTP {status}: {body}") - return _err(f"Unknown tool: {name}") + return p_err(f"Unknown tool: {name}") # -- JSON-RPC stdio loop --------------------------------------------------- -def _send(id_, result=None, error=None): +def p_send(id_, result=None, error=None): msg = {"jsonrpc": "2.0", "id": id_} if error is not None: msg["error"] = error @@ -433,7 +433,7 @@ def main(): params = msg.get("params", {}) or {} if method == "initialize": - _send(id_, { + p_send(id_, { "protocolVersion": "2024-11-05", "capabilities": {"tools": {}}, "serverInfo": {"name": "openswarm-discord", "version": "1.0.0"}, @@ -441,18 +441,18 @@ def main(): elif method == "notifications/initialized": pass elif method == "tools/list": - _send(id_, {"tools": TOOLS}) + p_send(id_, {"tools": P_TOOLS}) elif method == "tools/call": name = params.get("name", "") args = params.get("arguments", {}) or {} try: - _send(id_, handle_tool_call(name, args)) + p_send(id_, p_handle_tool_call(name, args)) except Exception as e: - _send(id_, _err(f"shim crashed: {e!r}")) + p_send(id_, p_err(f"shim crashed: {e!r}")) elif method == "ping": - _send(id_, {}) + p_send(id_, {}) elif id_ is not None: - _send(id_, error={"code": -32601, "message": f"Method not found: {method}"}) + p_send(id_, error={"code": -32601, "message": f"Method not found: {method}"}) if __name__ == "__main__": diff --git a/backend/apps/google_workspace_mcp_shim/run.py b/backend/apps/google_workspace_mcp_shim/run.py index 103c5932..a4b48189 100644 --- a/backend/apps/google_workspace_mcp_shim/run.py +++ b/backend/apps/google_workspace_mcp_shim/run.py @@ -24,7 +24,7 @@ from google.oauth2.credentials import Credentials @functools.lru_cache(maxsize=1) -def _patched_get_credentials(): +def p_patched_get_credentials(): refresh_token = os.environ.get("GOOGLE_WORKSPACE_REFRESH_TOKEN") if not refresh_token: raise ValueError("GOOGLE_WORKSPACE_REFRESH_TOKEN env var is required") @@ -40,10 +40,10 @@ def _patched_get_credentials(): ) -gauth.get_credentials = _patched_get_credentials # vulture-ignore: get_credentials +gauth.get_credentials = p_patched_get_credentials # vulture-ignore: get_credentials -from google_workspace_mcp import __main__ as _gw_main # noqa: E402,F401 +from google_workspace_mcp import __main__ as gw_main # noqa: E402,F401 from google_workspace_mcp.app import mcp # noqa: E402 diff --git a/backend/apps/mcp_registry/mcp_registry.py b/backend/apps/mcp_registry/mcp_registry.py index 4a8bdeac..1bfe4370 100644 --- a/backend/apps/mcp_registry/mcp_registry.py +++ b/backend/apps/mcp_registry/mcp_registry.py @@ -12,21 +12,21 @@ from backend.config.Apps import SubApp logger = logging.getLogger(__name__) -REGISTRY_BASE = "https://registry.modelcontextprotocol.io/v0.1" -PAGE_LIMIT = 100 -REFRESH_INTERVAL_S = 3600 +P_REGISTRY_BASE = "https://registry.modelcontextprotocol.io/v0.1" +P_PAGE_LIMIT = 100 +P_REFRESH_INTERVAL_S = 3600 -GITHUB_TOKEN = os.environ.get("GITHUB_TOKEN", "") -GITHUB_BATCH = 4000 if GITHUB_TOKEN else 50 -GITHUB_CONCURRENT = 10 +P_GITHUB_TOKEN = os.environ.get("GITHUB_TOKEN", "") +P_GITHUB_BATCH = 4000 if P_GITHUB_TOKEN else 50 +P_GITHUB_CONCURRENT = 10 -_cache: dict[str, dict] = {} -_cache_updated_at: float = 0 -_refresh_task: Optional[asyncio.Task] = None -_stars_cache: dict[str, int] = {} +P_CACHE: dict[str, dict] = {} +P_CACHE_UPDATED_AT: float = 0 +P_REFRESH_TASK: Optional[asyncio.Task] = None +P_STARS_CACHE: dict[str, int] = {} -def _extract_gh_repo(repo_url: str) -> Optional[str]: +def p_extract_gh_repo(repo_url: str) -> Optional[str]: """Parse 'owner/repo' from a GitHub URL.""" if not repo_url or "github.com" not in repo_url: return None @@ -42,7 +42,7 @@ def _extract_gh_repo(repo_url: str) -> Optional[str]: return None -def _extract_server(entry: dict) -> Optional[dict]: +def p_extract_server(entry: dict) -> Optional[dict]: """Extract a flat server record from a registry entry, keeping only latest versions.""" meta = entry.get("_meta", {}).get("io.modelcontextprotocol.registry/official", {}) if not meta.get("isLatest"): @@ -96,7 +96,7 @@ def _extract_server(entry: dict) -> Optional[dict]: } -async def _fetch_all_servers() -> dict[str, dict]: +async def p_fetch_all_servers() -> dict[str, dict]: """Paginate through the full registry and return a dict keyed by server name.""" servers: dict[str, dict] = {} cursor: Optional[str] = None @@ -104,12 +104,12 @@ async def _fetch_all_servers() -> dict[str, dict]: async with httpx.AsyncClient(timeout=30.0) as client: while True: - params: dict = {"limit": PAGE_LIMIT} + params: dict = {"limit": P_PAGE_LIMIT} if cursor: params["cursor"] = cursor try: - resp = await client.get(f"{REGISTRY_BASE}/servers", params=params) + resp = await client.get(f"{P_REGISTRY_BASE}/servers", params=params) resp.raise_for_status() data = resp.json() except Exception as e: @@ -121,7 +121,7 @@ async def _fetch_all_servers() -> dict[str, dict]: break for entry in entries: - record = _extract_server(entry) + record = p_extract_server(entry) if record: servers[record["name"]] = record @@ -135,16 +135,16 @@ async def _fetch_all_servers() -> dict[str, dict]: return servers -GOOGLE_README_URL = "https://raw.githubusercontent.com/google/mcp/main/README.md" -GOOGLE_ICON_URL = "https://github.com/google.png?size=64" -_ENTRY_RE = re.compile(r"\[\*\*(.+?)\*\*\]\((.+?)\)(?:[,\s]*(.+))?") +P_GOOGLE_README_URL = "https://raw.githubusercontent.com/google/mcp/main/README.md" +P_GOOGLE_ICON_URL = "https://github.com/google.png?size=64" +P_ENTRY_RE = re.compile(r"\[\*\*(.+?)\*\*\]\((.+?)\)(?:[,\s]*(.+))?") -def _slugify(name: str) -> str: +def p_slugify(name: str) -> str: return re.sub(r"[^a-z0-9]+", "-", name.lower()).strip("-") -def _parse_google_readme(text: str) -> dict[str, dict]: +def p_parse_google_readme(text: str) -> dict[str, dict]: servers: dict[str, dict] = {} section: Optional[str] = None @@ -164,7 +164,7 @@ def _parse_google_readme(text: str) -> dict[str, dict]: if section is None: continue - m = _ENTRY_RE.search(stripped) + m = P_ENTRY_RE.search(stripped) if not m: continue @@ -172,7 +172,7 @@ def _parse_google_readme(text: str) -> dict[str, dict]: url = m.group(2).strip() desc_raw = (m.group(3) or "").strip().rstrip(".") - slug = _slugify(title) + slug = p_slugify(title) key = f"google/{slug}" is_github = "github.com" in url or "go.dev" in url @@ -195,7 +195,7 @@ def _parse_google_readme(text: str) -> dict[str, dict]: "repositoryUrl": repo_url, "remoteUrl": "", "remoteType": remote_type, - "iconUrl": GOOGLE_ICON_URL, + "iconUrl": P_GOOGLE_ICON_URL, "environmentVariables": [], "keywords": ["google", section], "license": "Apache-2.0", @@ -206,13 +206,13 @@ def _parse_google_readme(text: str) -> dict[str, dict]: return servers -async def _fetch_google_servers() -> dict[str, dict]: +async def p_fetch_google_servers() -> dict[str, dict]: """Fetch and parse Google's MCP server catalog from their GitHub README.""" try: async with httpx.AsyncClient(timeout=15.0) as client: - resp = await client.get(GOOGLE_README_URL) + resp = await client.get(P_GOOGLE_README_URL) resp.raise_for_status() - servers = _parse_google_readme(resp.text) + servers = p_parse_google_readme(resp.text) logger.info(f"Google MCP catalog: parsed {len(servers)} servers") return servers except Exception as e: @@ -220,40 +220,40 @@ async def _fetch_google_servers() -> dict[str, dict]: return {} -async def _fetch_github_stars(servers: dict[str, dict]): +async def p_fetch_github_stars(servers: dict[str, dict]): """Batch-fetch GitHub star counts for servers with GitHub repos. Uses an in-memory cache so stars accumulate across refresh cycles even when rate-limited (60 req/hr unauthenticated, 5 000 with GITHUB_TOKEN). """ - global _stars_cache + global P_STARS_CACHE needed: list[str] = [] for srv in servers.values(): - gh = _extract_gh_repo(srv.get("repositoryUrl", "")) - if gh and gh not in _stars_cache and gh not in needed: + gh = p_extract_gh_repo(srv.get("repositoryUrl", "")) + if gh and gh not in P_STARS_CACHE and gh not in needed: needed.append(gh) if not needed: - logger.info(f"GitHub stars: all {len(_stars_cache)} repos cached, 0 to fetch") - _apply_stars(servers) + logger.info(f"GitHub stars: all {len(P_STARS_CACHE)} repos cached, 0 to fetch") + p_apply_stars(servers) return - to_fetch = needed[: GITHUB_BATCH] + to_fetch = needed[: P_GITHUB_BATCH] logger.info( f"GitHub stars: fetching {len(to_fetch)} repos " - f"({len(_stars_cache)} cached, {len(needed)} pending)" + f"({len(P_STARS_CACHE)} cached, {len(needed)} pending)" ) headers: dict[str, str] = {"Accept": "application/vnd.github.v3+json"} - if GITHUB_TOKEN: - headers["Authorization"] = f"token {GITHUB_TOKEN}" + if P_GITHUB_TOKEN: + headers["Authorization"] = f"token {P_GITHUB_TOKEN}" - sem = asyncio.Semaphore(GITHUB_CONCURRENT) + sem = asyncio.Semaphore(P_GITHUB_CONCURRENT) rate_limited = False fetched = 0 - async def _fetch_one(client: httpx.AsyncClient, repo: str): + async def p_fetch_one(client: httpx.AsyncClient, repo: str): nonlocal rate_limited, fetched if rate_limited: return @@ -265,56 +265,56 @@ async def _fetch_github_stars(servers: dict[str, dict]): f"https://api.github.com/repos/{repo}", headers=headers ) if resp.status_code == 200: - _stars_cache[repo] = resp.json().get("stargazers_count", 0) + P_STARS_CACHE[repo] = resp.json().get("stargazers_count", 0) fetched += 1 elif resp.status_code in (403, 429): rate_limited = True logger.warning("GitHub API rate-limited, stopping star fetch") elif resp.status_code == 404: - _stars_cache[repo] = 0 + P_STARS_CACHE[repo] = 0 fetched += 1 except Exception as exc: logger.debug(f"GitHub stars fetch failed for {repo}: {exc}") async with httpx.AsyncClient(timeout=15.0) as client: - await asyncio.gather(*[_fetch_one(client, r) for r in to_fetch]) + await asyncio.gather(*[p_fetch_one(client, r) for r in to_fetch]) - logger.info(f"GitHub stars: fetched {fetched} new, {len(_stars_cache)} total cached") - _apply_stars(servers) + logger.info(f"GitHub stars: fetched {fetched} new, {len(P_STARS_CACHE)} total cached") + p_apply_stars(servers) -def _apply_stars(servers: dict[str, dict]): +def p_apply_stars(servers: dict[str, dict]): for srv in servers.values(): - gh = _extract_gh_repo(srv.get("repositoryUrl", "")) - srv["stars"] = _stars_cache.get(gh) if gh else None + gh = p_extract_gh_repo(srv.get("repositoryUrl", "")) + srv["stars"] = P_STARS_CACHE.get(gh) if gh else None -async def _refresh_loop(): +async def p_refresh_loop(): """Background loop that refreshes the cache on startup and then hourly.""" - global _cache, _cache_updated_at + global P_CACHE, P_CACHE_UPDATED_AT while True: try: community, google = await asyncio.gather( - _fetch_all_servers(), - _fetch_google_servers(), + p_fetch_all_servers(), + p_fetch_google_servers(), ) - _cache = {**community, **google} - await _fetch_github_stars(_cache) - _cache_updated_at = time.time() + P_CACHE = {**community, **google} + await p_fetch_github_stars(P_CACHE) + P_CACHE_UPDATED_AT = time.time() except Exception as e: logger.exception(f"MCP registry refresh error: {e}") - await asyncio.sleep(REFRESH_INTERVAL_S) + await asyncio.sleep(P_REFRESH_INTERVAL_S) @asynccontextmanager async def mcp_registry_lifespan(): - global _refresh_task - _refresh_task = asyncio.create_task(_refresh_loop()) + global P_REFRESH_TASK + P_REFRESH_TASK = asyncio.create_task(p_refresh_loop()) yield - if _refresh_task: - _refresh_task.cancel() + if P_REFRESH_TASK: + P_REFRESH_TASK.cancel() try: - await _refresh_task + await P_REFRESH_TASK except asyncio.CancelledError: pass @@ -324,13 +324,13 @@ mcp_registry = SubApp("mcp-registry", mcp_registry_lifespan) @mcp_registry.router.get("/stats") async def registry_stats(): - google = sum(1 for s in _cache.values() if s.get("source") == "google") - community = sum(1 for s in _cache.values() if s.get("source") == "community") + google = sum(1 for s in P_CACHE.values() if s.get("source") == "google") + community = sum(1 for s in P_CACHE.values() if s.get("source") == "community") return { - "total": len(_cache), + "total": len(P_CACHE), "google": google, "community": community, - "lastUpdated": _cache_updated_at, + "lastUpdated": P_CACHE_UPDATED_AT, } @@ -342,7 +342,7 @@ async def registry_search( sort: str = Query("name", description="Sort by: name, stars"), source: str = Query("", description="Filter by source: google, community, or empty for all"), ): - pool = _cache.values() + pool = P_CACHE.values() if source: pool = [s for s in pool if s.get("source") == source] @@ -387,7 +387,7 @@ async def registry_search( @mcp_registry.router.get("/detail/{server_name:path}") async def registry_detail(server_name: str): - srv = _cache.get(server_name) + srv = P_CACHE.get(server_name) if not srv: return {"error": "Server not found"}, 404 return {"server": srv} diff --git a/backend/apps/modes/modes.py b/backend/apps/modes/modes.py index 9cc897fc..499fee84 100644 --- a/backend/apps/modes/modes.py +++ b/backend/apps/modes/modes.py @@ -1,5 +1,6 @@ import os import logging +import json from contextlib import asynccontextmanager from fastapi import HTTPException from backend.config.Apps import SubApp @@ -18,10 +19,9 @@ async def modes_lifespan(): chat_path = os.path.join(DATA_DIR, "chat.json") if os.path.exists(chat_path): try: - import json as _json - with open(chat_path) as _f: - _data = _json.load(_f) - if _data.get("is_builtin") is True and _data.get("id") == "chat": + with open(chat_path) as f: + json_data = json.load(f) + if json_data.get("is_builtin") is True and json_data.get("id") == "chat": os.remove(chat_path) logger.info("Removed deprecated built-in chat.json (merged into ask)") except Exception: @@ -29,14 +29,14 @@ async def modes_lifespan(): for builtin in BUILTIN_MODES: path = os.path.join(DATA_DIR, f"{builtin.id}.json") if not os.path.exists(path): - _save(builtin) + p_save(builtin) yield modes = SubApp("modes", modes_lifespan) -def _load_all() -> list[Mode]: +def p_load_all() -> list[Mode]: result = [] if not os.path.exists(DATA_DIR): return result @@ -52,11 +52,11 @@ def _load_all() -> list[Mode]: return result -def _save(mode: Mode): +def p_save(mode: Mode): atomic_write_json(os.path.join(DATA_DIR, f"{mode.id}.json"), mode.model_dump()) -def _load(mode_id: str) -> Mode: +def p_load(mode_id: str) -> Mode: data = read_json_or_none(os.path.join(DATA_DIR, f"{mode_id}.json")) if data is None: raise HTTPException(status_code=404, detail="Mode not found") @@ -72,12 +72,12 @@ def load_mode(mode_id: str) -> Mode | None: @modes.router.get("/list") async def list_modes(): builtin_defaults = {m.id: m.model_dump() for m in BUILTIN_MODES} - return {"modes": [m.model_dump() for m in _load_all()], "builtin_defaults": builtin_defaults} + return {"modes": [m.model_dump() for m in p_load_all()], "builtin_defaults": builtin_defaults} @modes.router.get("/{mode_id}") async def get_mode(mode_id: str): - return _load(mode_id).model_dump() + return p_load(mode_id).model_dump() @modes.router.post("/create") @@ -93,16 +93,16 @@ async def create_mode(body: ModeCreate): default_folder=body.default_folder, is_builtin=False, ) - _save(mode) + p_save(mode) return {"ok": True, "mode": mode.model_dump()} @modes.router.put("/{mode_id}") async def update_mode(mode_id: str, body: ModeUpdate): - mode = _load(mode_id) + mode = p_load(mode_id) for k, v in body.model_dump(exclude_unset=True).items(): setattr(mode, k, v) - _save(mode) + p_save(mode) return {"ok": True, "mode": mode.model_dump()} @@ -112,13 +112,13 @@ async def reset_mode(mode_id: str): builtin = next((m for m in BUILTIN_MODES if m.id == mode_id), None) if not builtin: raise HTTPException(status_code=400, detail="Only built-in modes can be reset") - _save(builtin) + p_save(builtin) return {"ok": True, "mode": builtin.model_dump()} @modes.router.delete("/{mode_id}") async def delete_mode(mode_id: str): - mode = _load(mode_id) + mode = p_load(mode_id) if mode.is_builtin: raise HTTPException(status_code=403, detail="Cannot delete built-in modes") path = os.path.join(DATA_DIR, f"{mode_id}.json") diff --git a/backend/apps/nine_router/oauth.py b/backend/apps/nine_router/oauth.py index 55b7ba85..b21b0e93 100644 --- a/backend/apps/nine_router/oauth.py +++ b/backend/apps/nine_router/oauth.py @@ -11,7 +11,7 @@ import os import httpx from .process import NINE_ROUTER_API, NINE_ROUTER_PORT, NINE_ROUTER_V1 -from backend.apps.oauth_state import _pending_oauth, _mark_oauth_completed +from backend.apps.oauth_state import PENDING_OAUTH, mark_oauth_completed logger = logging.getLogger(__name__) @@ -22,9 +22,9 @@ logger = logging.getLogger(__name__) # 1455 that serves the same postMessage/BroadcastChannel/localStorage relay so # the frontend's existing popup + msgHandler flow works unchanged. -_CODEX_CALLBACK_PORT = 1455 -_CODEX_CALLBACK_PATH = "/auth/callback" -_CODEX_CALLBACK_HTML = b""" +P_CODEX_CALLBACK_PORT = 1455 +P_CODEX_CALLBACK_PATH = "/auth/callback" +P_CODEX_CALLBACK_HTML = b"""
{desc}
{safe_e}