Spaces:
Running
Running
Download app.py from DedeProGames/claude-code-mcp: direct link, hf CLI and curl.
- Browser
- Download file 18 kB
-
https://huggingface.co/spaces/DedeProGames/claude-code-mcp/resolve/main/app.py
- Command line
-
hf download hf://spaces/DedeProGames/claude-code-mcp/app.py
-
curl -L -o app.py https://huggingface.co/spaces/DedeProGames/claude-code-mcp/resolve/main/app.py
18 kB
| """Claude Code-style MCP server with a private workspace administration dashboard. | |
| The browser dashboard is available at / and asks for MCP_AUTH_TOKEN. | |
| The MCP endpoint is /gradio_api/mcp/ and continues to accept the token as | |
| Authorization: Bearer ..., X-MCP-Token: ... or a query parameter. | |
| """ | |
| import asyncio | |
| import hmac | |
| import json | |
| import os | |
| import re | |
| import threading | |
| import time | |
| from pathlib import Path | |
| import gradio as gr | |
| import uvicorn | |
| from fastapi import FastAPI, HTTPException, Query, Request | |
| from fastapi.responses import HTMLResponse, JSONResponse, PlainTextResponse, Response, StreamingResponse | |
| from starlette.concurrency import run_in_threadpool | |
| import workspaces | |
| import tools as workspace_tools | |
| from admin_ui import DASHBOARD_PAGE, LOCKED_PAGE, LOGIN_PAGE | |
| from tools import ALL_TOOLS | |
| from telemetry import ProcessMonitor, ResourceMonitor | |
| TOKEN = os.environ.get("MCP_AUTH_TOKEN", "") | |
| TOKEN_OK = len(TOKEN) >= 16 | |
| PORT = int(os.environ.get("PORT", 7860)) | |
| SESSION_COOKIE = "ccmcp_admin" | |
| SESSION_TTL_SECONDS = 12 * 60 * 60 | |
| WORKSPACE_ID_RE = re.compile(r"^ws-(?:[0-9a-f]{6}|[0-9a-f]{32})$") | |
| JOB_ID_RE = re.compile(r"^bg-?(?:[0-9a-f]{12}|[0-9]+)$") | |
| _resources = ResourceMonitor(workspaces.ROOT) | |
| _process_resources = ProcessMonitor() | |
| _overview_lock = threading.Lock() | |
| _overview_cache = {} | |
| with gr.Blocks(title="Claude Code MCP") as demo: | |
| # Keep MCP tool schemas registered with Gradio. queue=False avoids serializing | |
| # independent chats through the same Gradio event queue. | |
| for fn in ALL_TOOLS: | |
| gr.api(fn, queue=False) | |
| app = FastAPI(title="Claude Code MCP") | |
| def _session_signature(expires: int) -> str: | |
| message = f"workspace-admin:{expires}".encode("ascii") | |
| return hmac.new(TOKEN.encode("utf-8"), message, digestmod="sha256").hexdigest() | |
| def _make_session(expires: int) -> str: | |
| return f"{expires}.{_session_signature(expires)}" | |
| def _mark_session_cookie_partitioned(response: Response) -> None: | |
| """Partition the admin cookie so it works in HF's cross-site Space iframe.""" | |
| cookie_prefix = f"{SESSION_COOKIE}=".encode("ascii") | |
| response.raw_headers = [ | |
| (name, value + b"; Partitioned") | |
| if name.lower() == b"set-cookie" and value.startswith(cookie_prefix) | |
| else (name, value) | |
| for name, value in response.raw_headers | |
| ] | |
| def _valid_admin_session(request: Request) -> bool: | |
| if not TOKEN_OK: | |
| return False | |
| raw = request.cookies.get(SESSION_COOKIE, "") | |
| try: | |
| expires_text, supplied = raw.split(".", 1) | |
| expires = int(expires_text) | |
| except (ValueError, TypeError): | |
| return False | |
| if expires <= int(time.time()): | |
| return False | |
| expected = _session_signature(expires) | |
| return hmac.compare_digest(supplied, expected) | |
| def _admin_headers(response: Response) -> Response: | |
| response.headers["Cache-Control"] = "no-store" | |
| response.headers["X-Content-Type-Options"] = "nosniff" | |
| response.headers["Referrer-Policy"] = "no-referrer" | |
| # Hugging Face serves Space apps inside an iframe on huggingface.co. Do not | |
| # send X-Frame-Options (DENY/SAMEORIGIN would block that embed); use CSP to | |
| # permit only the HF Space page and same-origin frames. | |
| response.headers["Content-Security-Policy"] = ( | |
| "default-src 'none'; style-src 'unsafe-inline'; script-src 'unsafe-inline'; " | |
| "connect-src 'self'; img-src 'self' data:; form-action 'self'; " | |
| "frame-ancestors 'self' https://huggingface.co; base-uri 'none'" | |
| ) | |
| return response | |
| async def require_token_or_admin_session(request: Request, call_next): | |
| path = request.url.path | |
| # Readiness is intentionally public so Hugging Face can check the Space. | |
| # The page itself only displays the login screen until a signed session exists. | |
| if path in ("/", "/healthz"): | |
| response = await call_next(request) | |
| return _admin_headers(response) if path == "/" else response | |
| # Login is the sole public administrative action. Every admin API, including | |
| # logout, requires the short-lived HttpOnly browser session. | |
| if path.startswith("/admin/"): | |
| if path == "/admin/login" and request.method == "POST": | |
| response = await call_next(request) | |
| return _admin_headers(response) | |
| if not TOKEN_OK: | |
| return JSONResponse({"error": "locked: configure MCP_AUTH_TOKEN"}, status_code=503) | |
| if not _valid_admin_session(request): | |
| return _admin_headers(JSONResponse({"error": "admin login required"}, status_code=401)) | |
| response = await call_next(request) | |
| return _admin_headers(response) | |
| # MCP and Gradio API routes retain their shared-token authentication. | |
| if not TOKEN_OK: | |
| return JSONResponse({"error": "locked: set the MCP_AUTH_TOKEN secret (16+ chars)"}, status_code=503) | |
| auth = request.headers.get("authorization", "") | |
| supplied = ( | |
| request.headers.get("x-mcp-token") | |
| or (auth[7:] if auth.lower().startswith("bearer ") else "") | |
| or request.query_params.get("token", "") | |
| ) | |
| if not hmac.compare_digest(supplied.encode("utf-8"), TOKEN.encode("utf-8")): | |
| return JSONResponse({"error": "unauthorized"}, status_code=401) | |
| return await call_next(request) | |
| def admin_home(request: Request): | |
| if not TOKEN_OK: | |
| return HTMLResponse(LOCKED_PAGE, status_code=503) | |
| page = DASHBOARD_PAGE if _valid_admin_session(request) else LOGIN_PAGE | |
| return HTMLResponse(page) | |
| async def admin_login(request: Request): | |
| if not TOKEN_OK: | |
| return _admin_headers(JSONResponse({"error": "MCP_AUTH_TOKEN is not configured"}, status_code=503)) | |
| body = await request.body() | |
| if len(body) > 4096: | |
| return _admin_headers(JSONResponse({"error": "request too large"}, status_code=413)) | |
| try: | |
| payload = json.loads(body or b"{}") | |
| except (ValueError, TypeError): | |
| return _admin_headers(JSONResponse({"error": "invalid JSON body"}, status_code=400)) | |
| supplied = payload.get("token", "") if isinstance(payload, dict) else "" | |
| if not isinstance(supplied, str) or not hmac.compare_digest(supplied.encode("utf-8"), TOKEN.encode("utf-8")): | |
| return _admin_headers(JSONResponse({"error": "Token inválido"}, status_code=401)) | |
| expires = int(time.time()) + SESSION_TTL_SECONDS | |
| response = JSONResponse({"ok": True, "expires_at": expires}) | |
| forwarded_proto = request.headers.get("x-forwarded-proto", "").split(",", 1)[0].strip().lower() | |
| secure = request.url.scheme == "https" or forwarded_proto == "https" | |
| response.set_cookie( | |
| SESSION_COOKIE, | |
| _make_session(expires), | |
| max_age=SESSION_TTL_SECONDS, | |
| httponly=True, | |
| secure=secure, | |
| # SameSite=None is required in the cross-site huggingface.co -> hf.space | |
| # iframe. Partitioned keeps the session scoped to that top-level site. | |
| samesite="none" if secure else "lax", | |
| path="/", | |
| ) | |
| if secure: | |
| _mark_session_cookie_partitioned(response) | |
| return _admin_headers(response) | |
| def admin_logout(request: Request): | |
| response = JSONResponse({"ok": True}) | |
| forwarded_proto = request.headers.get("x-forwarded-proto", "").split(",", 1)[0].strip().lower() | |
| secure = request.url.scheme == "https" or forwarded_proto == "https" | |
| response.delete_cookie( | |
| SESSION_COOKIE, | |
| httponly=True, | |
| secure=secure, | |
| samesite="none" if secure else "lax", | |
| path="/", | |
| ) | |
| if secure: | |
| _mark_session_cookie_partitioned(response) | |
| return _admin_headers(response) | |
| def _load_job_rows(ws: workspaces.Workspace) -> list[dict]: | |
| rows = [] | |
| for meta_path in ws.shells_dir.glob("bg-*.json"): | |
| try: | |
| record = json.loads(meta_path.read_text(encoding="utf-8")) | |
| except (OSError, ValueError, TypeError): | |
| continue | |
| job_id = str(record.get("id", "")) | |
| if not JOB_ID_RE.fullmatch(job_id): | |
| continue | |
| output_name = Path(str(record.get("output_file") or f"{job_id}.out")).name | |
| output_path = ws.shells_dir / output_name | |
| try: | |
| log_bytes = output_path.stat().st_size | |
| except OSError: | |
| log_bytes = 0 | |
| rows.append( | |
| { | |
| "id": job_id, | |
| "command": str(record.get("command", ""))[:4000], | |
| "status": str(record.get("status", "unknown")), | |
| "pid": record.get("pid"), | |
| "started_at": record.get("started_at"), | |
| "finished_at": record.get("finished_at") or record.get("interrupted_at"), | |
| "exit_code": record.get("exit_code"), | |
| "log_bytes": log_bytes, | |
| "dropped_log_bytes": int(record.get("output_truncated_bytes", 0) or 0), | |
| } | |
| ) | |
| rows.sort(key=lambda row: float(row.get("started_at") or 0), reverse=True) | |
| return rows[:100] | |
| def admin_overview(): | |
| rows = [] | |
| active_count = 0 | |
| paused_count = 0 | |
| running_jobs = 0 | |
| for ws in workspaces.list_all(): | |
| # A checkpoint-only workspace can appear after recovery; restore it once | |
| # so its metadata and saved job logs are visible in the dashboard. | |
| if not ws.work.is_dir(): | |
| ws.restore_latest() | |
| try: | |
| workspace_tools._refresh_background_jobs(ws) | |
| except Exception: | |
| # Keep the admin panel readable if one malformed job record is found. | |
| pass | |
| meta = ws.meta() | |
| status = "paused" if meta.get("status") == "paused" else "active" | |
| jobs = _load_job_rows(ws) | |
| workspace_running = sum(1 for job in jobs if job["status"] == "running") | |
| active_count += int(status == "active") | |
| paused_count += int(status == "paused") | |
| running_jobs += workspace_running | |
| rows.append( | |
| { | |
| "id": ws.id, | |
| "title": str(meta.get("title") or ""), | |
| "kind": str(meta.get("kind") or "workspace"), | |
| "status": status, | |
| "parent_id": meta.get("parent_id"), | |
| "created": meta.get("created"), | |
| "last_used": meta.get("last_used"), | |
| "latest_checkpoint": meta.get("latest_checkpoint"), | |
| "work_dir": str(ws.work), | |
| "running_jobs": workspace_running, | |
| "jobs": jobs, | |
| } | |
| ) | |
| try: | |
| bucket_mounted = Path(workspaces.ROOT).resolve().is_relative_to(Path("/data")) | |
| except (OSError, RuntimeError): | |
| bucket_mounted = str(workspaces.ROOT).startswith("/data/") | |
| return { | |
| "summary": { | |
| "total": len(rows), | |
| "active": active_count, | |
| "paused": paused_count, | |
| "running_jobs": running_jobs, | |
| }, | |
| "storage": { | |
| "bucket_mounted": bucket_mounted, | |
| "root": str(workspaces.ROOT), | |
| }, | |
| "workspaces": rows, | |
| } | |
| def _cached_overview() -> dict: | |
| # Share inventory reads across open dashboard tabs; do not scan the Bucket | |
| # separately for every stream on every second. | |
| with _overview_lock: | |
| now = time.monotonic() | |
| if (_overview_cache.get("root") != str(workspaces.ROOT) | |
| or now - _overview_cache.get("at", -10) >= 5): | |
| _overview_cache.update(root=str(workspaces.ROOT), at=now, data=admin_overview()) | |
| return _overview_cache["data"] | |
| def _read_log_tail(ws: workspaces.Workspace, job_id: str, tail: int) -> dict: | |
| record = workspace_tools._load_job(ws, job_id) | |
| if not record: | |
| raise HTTPException(status_code=404, detail="Job not found") | |
| output_name = Path(str(record.get("output_file") or f"{job_id}.out")).name | |
| output_path = ws.shells_dir / output_name | |
| try: | |
| if not output_path.resolve().is_relative_to(ws.shells_dir.resolve()): | |
| raise OSError("log path outside shells directory") | |
| size = output_path.stat().st_size | |
| offset = max(0, size - tail) | |
| with output_path.open("rb") as source: | |
| source.seek(offset) | |
| content = source.read(tail).decode("utf-8", errors="replace") | |
| except OSError: | |
| size, offset, content = 0, 0, "" | |
| return { | |
| "workspace": ws.id, "job": job_id, "status": record.get("status", "unknown"), | |
| "content": content, "total_bytes": size, "offset": offset, | |
| "truncated_before": offset > 0 or bool(record.get("output_truncated_bytes")), | |
| "dropped_log_bytes": int(record.get("output_truncated_bytes", 0) or 0), | |
| } | |
| def _live_workspace(workspace_id: str, metrics: dict, log_signatures: dict) -> dict: | |
| ws = workspaces.Workspace(workspace_id) | |
| if not ws.meta_file.is_file(): | |
| raise HTTPException(status_code=404, detail="Workspace not found") | |
| try: | |
| workspace_tools._refresh_background_jobs(ws) | |
| except Exception: | |
| pass | |
| jobs = _load_job_rows(ws) | |
| history_total = len(jobs) | |
| jobs.sort(key=lambda row: (row["status"] != "running", -float(row.get("started_at") or 0))) | |
| jobs = jobs[:20] | |
| processes = _process_resources.snapshot(jobs, metrics["cpu"]["cores"]) | |
| for job in jobs: | |
| job["resources"] = processes.get(job["id"]) | |
| record = workspace_tools._load_job(ws, job["id"]) or {} | |
| path = ws.shells_dir / Path(str(record.get("output_file") or f'{job["id"]}.out')).name | |
| try: | |
| stat = path.stat() | |
| signature = (stat.st_size, stat.st_mtime_ns) | |
| except OSError: | |
| signature = (0, 0) | |
| if log_signatures.get(job["id"]) != signature: | |
| # Dashboard tail reads never consume the agent's BashOutput cursor. | |
| job["log"] = _read_log_tail(ws, job["id"], 16 * 1024) | |
| log_signatures[job["id"]] = signature | |
| return {"id": ws.id, "jobs": jobs, "history_total": history_total} | |
| def admin_metrics(): | |
| return _resources.snapshot() | |
| def admin_workspace_live(workspace_id: str): | |
| if not WORKSPACE_ID_RE.fullmatch(workspace_id): | |
| raise HTTPException(status_code=404, detail="Workspace not found") | |
| return _live_workspace(workspace_id, _resources.snapshot(), {}) | |
| async def admin_live(request: Request, workspace: str = Query(default="")): | |
| if workspace and not WORKSPACE_ID_RE.fullmatch(workspace): | |
| raise HTTPException(status_code=400, detail="Invalid workspace") | |
| async def events(): | |
| log_signatures = {} | |
| last_inventory_at = -10. | |
| yield ": connected\n\n" | |
| while not await request.is_disconnected(): | |
| if not _valid_admin_session(request): | |
| yield 'event: session-expired\ndata: {}\n\n' | |
| return | |
| def collect(): | |
| nonlocal last_inventory_at | |
| metrics = _resources.snapshot() | |
| payload = {"metrics": metrics} | |
| now = time.monotonic() | |
| if now - last_inventory_at >= 5: | |
| payload["overview"] = _cached_overview() | |
| last_inventory_at = now | |
| if workspace: | |
| try: | |
| payload["workspace"] = _live_workspace(workspace, metrics, log_signatures) | |
| except HTTPException: | |
| payload["workspace"] = {"id": workspace, "jobs": [], "missing": True} | |
| return payload | |
| payload = await run_in_threadpool(collect) | |
| yield "event: snapshot\ndata: " + json.dumps(payload, ensure_ascii=False, separators=(",", ":")) + "\n\n" | |
| await asyncio.sleep(1) | |
| return StreamingResponse(events(), media_type="text/event-stream", | |
| headers={"Cache-Control": "no-store", "X-Accel-Buffering": "no"}) | |
| def admin_job_logs( | |
| workspace_id: str, | |
| job_id: str, | |
| tail: int = Query(default=64 * 1024, ge=1024, le=128 * 1024), | |
| ): | |
| if not WORKSPACE_ID_RE.fullmatch(workspace_id) or not JOB_ID_RE.fullmatch(job_id): | |
| raise HTTPException(status_code=404, detail="Workspace or job not found") | |
| ws = next((item for item in workspaces.list_all() if item.id == workspace_id), None) | |
| if ws is None: | |
| raise HTTPException(status_code=404, detail="Workspace not found") | |
| if not ws.work.is_dir() and not ws.restore_latest(): | |
| raise HTTPException(status_code=404, detail="Workspace data not found") | |
| try: | |
| workspace_tools._refresh_background_jobs(ws) | |
| record = workspace_tools._load_job(ws, job_id) | |
| except workspaces.WorkspaceError as exc: | |
| raise HTTPException(status_code=404, detail="Job not found") from exc | |
| if not record: | |
| raise HTTPException(status_code=404, detail="Job not found") | |
| return _read_log_tail(ws, job_id, int(tail)) | |
| def healthz(): | |
| state = "ready" if TOKEN_OK else "LOCKED: add the MCP_AUTH_TOKEN secret in Settings > Variables and secrets" | |
| storage = ( | |
| "persistent when a Bucket is mounted at /data" | |
| if str(workspaces.ROOT).startswith("/data/") | |
| else "ephemeral; configure a Bucket at /data for persistence" | |
| ) | |
| return PlainTextResponse( | |
| f"Claude Code MCP server: {state}\n" | |
| f"Admin dashboard: / (MCP_AUTH_TOKEN login required)\n" | |
| f"MCP endpoint: /gradio_api/mcp/ (MCP_AUTH_TOKEN required)\n" | |
| f"workspace storage: {storage}\n" | |
| ) | |
| # ssr_mode=False is required on Spaces: HF sets GRADIO_SSR_MODE=True, and mounting then starts | |
| # a Node SSR server on port 7860, so uvicorn below dies with "address already in use". | |
| app = gr.mount_gradio_app(app, demo, path="/", mcp_server=True, ssr_mode=False) | |
| if __name__ == "__main__": | |
| uvicorn.run(app, host="0.0.0.0", port=PORT) | |