claude-code-mcp / app.py
DedeProGames's picture
Redesign monochrome admin with live metrics and inline command logs
c66dc5d verified
Raw History Blame Contribute Delete
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
@app.middleware("http")
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)
@app.get("/", response_class=HTMLResponse)
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)
@app.post("/admin/login")
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)
@app.post("/admin/logout")
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]
@app.get("/admin/api/overview")
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}
@app.get("/admin/api/metrics")
def admin_metrics():
return _resources.snapshot()
@app.get("/admin/api/workspaces/{workspace_id}/live")
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(), {})
@app.get("/admin/api/live")
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"})
@app.get("/admin/api/workspaces/{workspace_id}/jobs/{job_id}/logs")
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))
@app.get("/healthz")
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)