OSVauco/opax-mcp/server.py
Chris Christiansen 5d156642b2
Some checks are pending
Check Python Version Consistency / Check Python Version (push) Waiting to run
fix: log exceptions in mcp_handler
2026-06-29 01:20:00 +00:00

309 lines
14 KiB
Python

"""opax-mcp — MCP Streamable HTTP server for OPAX/Vauco
Transport: MCP Streamable HTTP (JSON-RPC 2.0) on POST /
Auth: Authorization: Bearer <secret> OR X-MCP-Secret: <secret> OR api-key: <secret>
"""
import os
import json
import uuid
import httpx
import base64
import google.auth
import google.auth.transport.requests
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import JSONResponse
from typing import Any, Optional
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
app = FastAPI(title="opax-mcp", version="3.0.2")
OPAX_BASE_URL = os.environ.get("OPAX_BASE_URL", "https://opax.vauco.no")
OPAX_IAP_CLIENT_ID = os.environ.get("OPAX_IAP_CLIENT_ID", "")
MCP_SECRET = os.environ.get("MCP_SECRET", "")
GITEA_URL = os.environ.get("GITEA_URL", "http://34.59.131.162:3000")
GITEA_TOKEN = os.environ.get("GITEA_TOKEN", "")
GITEA_REPO = os.environ.get("GITEA_REPO", "chris/OSVauco")
GOOGLE_CLOUD_PROJECT = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5")
logger.info(f"OPAX_IAP_CLIENT_ID: {OPAX_IAP_CLIENT_ID[:10]}...")
# ---------------------------------------------------------------------------
# Auth
# ---------------------------------------------------------------------------
def _verify_auth(request: Request):
if not MCP_SECRET:
return
# Authorization: Bearer <token>
auth = request.headers.get("Authorization", "")
if auth.startswith("Bearer ") and auth[7:] == MCP_SECRET:
return
# X-MCP-Secret or api-key
if request.headers.get("X-MCP-Secret") == MCP_SECRET:
return
if request.headers.get("api-key") == MCP_SECRET:
return
raise HTTPException(status_code=401, detail="Unauthorized")
# ---------------------------------------------------------------------------
# OPAX helpers
# ---------------------------------------------------------------------------
def _opax_identity_token() -> str:
"""Hent identity token for opax.vauco.no (Cloud Run IAP/IAM)."""
try:
metadata_url = (
"http://metadata.google.internal/computeMetadata/v1/instance"
f"/service-accounts/default/identity?audience={OPAX_IAP_CLIENT_ID}&format=full"
)
resp = httpx.get(metadata_url, headers={"Metadata-Flavor": "Google"}, timeout=5)
resp.raise_for_status()
token = resp.text.strip()
if not token:
raise ValueError("Empty token returned from metadata server")
return token
except Exception as e:
logger.error(f"[IAP TOKEN ERROR] {type(e).__name__}: {e}")
raise
def _opax_headers() -> dict:
return {"Authorization": f"Bearer {_opax_identity_token()}", "Content-Type": "application/json"}
async def _opax_get(path: str) -> Any:
async with httpx.AsyncClient(timeout=30) as c:
r = await c.get(f"{OPAX_BASE_URL}{path}", headers=_opax_headers())
r.raise_for_status()
return r.json()
async def _opax_post(path: str, body: dict) -> Any:
async with httpx.AsyncClient(timeout=30) as c:
r = await c.post(f"{OPAX_BASE_URL}{path}", json=body, headers=_opax_headers())
r.raise_for_status()
return r.json()
# ---------------------------------------------------------------------------
# Gitea helpers
# ---------------------------------------------------------------------------
def _gitea_headers() -> dict:
return {"Authorization": f"token {GITEA_TOKEN}", "Content-Type": "application/json", "Accept": "application/json"}
async def _gitea_get(path: str) -> Any:
async with httpx.AsyncClient(timeout=30) as c:
r = await c.get(f"{GITEA_URL}/api/v1{path}", headers=_gitea_headers())
r.raise_for_status()
return r.json()
async def _gitea_post(path: str, body: dict) -> Any:
async with httpx.AsyncClient(timeout=30) as c:
r = await c.post(f"{GITEA_URL}/api/v1{path}", json=body, headers=_gitea_headers())
r.raise_for_status()
return r.json()
async def _gitea_put(path: str, body: dict) -> Any:
async with httpx.AsyncClient(timeout=30) as c:
r = await c.put(f"{GITEA_URL}/api/v1{path}", json=body, headers=_gitea_headers())
r.raise_for_status()
return r.json()
# ---------------------------------------------------------------------------
# Tool implementations
# ---------------------------------------------------------------------------
async def get_billing_summary(p): return await _opax_get("/billing/summary")
async def get_billing_forecast(p): return await _opax_get("/billing/forecast")
async def get_billing_credits(p): return await _opax_get("/billing/credits")
async def get_billing_anomalies(p): return await _opax_get("/billing/anomalies")
async def get_billing_history(p): return await _opax_get("/billing/history")
async def get_billing_budget(p): return await _opax_get("/billing/budget")
async def get_billing_tokens_by_module(p): return await _opax_get("/telemetry/history")
async def set_billing_budget(p):
return await _opax_post("/billing/budget", {"budget": p.get("amount", 500)})
async def create_invite(p):
return await _opax_post("/onboard/invite", {
"company": p.get("company", p.get("name", "")),
"email": p.get("email"),
"tier": p.get("tier", "starter"),
})
async def list_customers(p): return await _opax_get("/state")
async def send_webhook(p):
return await _opax_post("/notify/webhook", {
"url": p.get("url", ""), "event": p.get("event", "custom"),
"title": p.get("title", "OPAX varsel"), "body": p.get("message", p.get("body", "")),
})
async def send_email(p):
return await _opax_post("/notify/email", {"to": p.get("to"), "subject": p.get("subject"), "event": p.get("event", "digest")})
async def send_sms(p):
return await _opax_post("/notify/sms", {"to": p.get("to"), "body": p.get("message", p.get("body", "")), "event": p.get("event", "spike"), "tier": p.get("tier", "guard")})
async def get_notify_channels(p): return await _opax_get("/notify/channels")
async def run_jason(p):
return await _opax_post("/run", {"message": p.get("prompt", p.get("message", "")), "user_id": p.get("user_id", "opax"), "session_id": p.get("session_id", "mcp"), "mode": p.get("mode", "light")})
async def run_emma(p):
return await _opax_post("/emma", {"message": p.get("prompt", p.get("message", "")), "session_id": p.get("session_id", "mcp")})
async def get_health(p): return await _opax_get("/health")
async def get_build_status(p): return await _opax_get("/opax/build-status")
async def get_state(p): return await _opax_get("/state")
async def get_telemetry(p): return await _opax_get("/telemetry/history")
async def run_terminal(p):
return await _opax_post("/terminal/exec", {"cmd": p.get("command", p.get("cmd", "help"))})
async def list_commits(p):
return await _gitea_get(f"/repos/{p.get('repo', GITEA_REPO)}/commits?limit={p.get('limit', 10)}")
async def get_file(p):
return await _gitea_get(f"/repos/{p.get('repo', GITEA_REPO)}/contents/{p.get('path', '')}?ref={p.get('ref', 'main')}")
async def list_open_issues(p):
return await _gitea_get(f"/repos/{p.get('repo', GITEA_REPO)}/issues?state=open&limit=20")
async def create_issue(p):
return await _gitea_post(f"/repos/{p.get('repo', GITEA_REPO)}/issues", {"title": p.get("title"), "body": p.get("body", "")})
async def push_file(p):
repo = p.get("repo", GITEA_REPO)
path = p.get("path")
content = base64.b64encode(p.get("content", "").encode()).decode()
body = {"message": p.get("message", f"chore: oppdater {path} via opax-mcp"), "content": content, "branch": p.get("branch", "main")}
if p.get("sha"):
body["sha"] = p["sha"]
return await _gitea_put(f"/repos/{repo}/contents/{path}", body)
return await _gitea_post(f"/repos/{repo}/contents/{path}", body)
# ---------------------------------------------------------------------------
# Tool registry + MCP schema
# ---------------------------------------------------------------------------
TOOLS = {
"get_billing_summary": (get_billing_summary, "Hent billing-sammendrag for OPAX", {}),
"get_billing_forecast": (get_billing_forecast, "Hent billing-prognose", {}),
"get_billing_credits": (get_billing_credits, "Hent gjenværende kreditter", {}),
"get_billing_anomalies": (get_billing_anomalies, "Hent billing-anomalier", {}),
"get_billing_history": (get_billing_history, "Hent billing-historikk", {}),
"get_billing_budget": (get_billing_budget, "Hent nåværende budsjett", {}),
"get_billing_tokens_by_module":(get_billing_tokens_by_module,"Hent token-forbruk per modul", {}),
"set_billing_budget": (set_billing_budget, "Sett månedlig budsjett", {"type":"object","properties":{"amount":{"type":"number","description":"Budsjett i NOK"}},"required":[]}),
"create_invite": (create_invite, "Inviter ny kunde til OPAX", {"type":"object","properties":{"company":{"type":"string"},"email":{"type":"string"},"tier":{"type":"string"}},"required":["email"]}),
"list_customers": (list_customers, "List alle kunder", {}),
"send_webhook": (send_webhook, "Send webhook-varsling", {"type":"object","properties":{"url":{"type":"string"},"title":{"type":"string"},"message":{"type":"string"}},"required":["url"]}),
"send_email": (send_email, "Send e-post", {"type":"object","properties":{"to":{"type":"string"},"subject":{"type":"string"}},"required":["to","subject"]}),
"send_sms": (send_sms, "Send SMS", {"type":"object","properties":{"to":{"type":"string"},"message":{"type":"string"}},"required":["to","message"]}),
"get_notify_channels": (get_notify_channels, "Hent varslingkanaler", {}),
"run_jason": (run_jason, "Kjør Jason-agenten med en prompt", {"type":"object","properties":{"prompt":{"type":"string"},"mode":{"type":"string"}},"required":["prompt"]}),
"run_emma": (run_emma, "Kjør Emma-agenten med en prompt", {"type":"object","properties":{"prompt":{"type":"string"}},"required":["prompt"]}),
"get_health": (get_health, "Hent helsestatus for OPAX", {}),
"get_build_status": (get_build_status, "Hent siste build-status", {}),
"get_state": (get_state, "Hent platform-tilstand", {}),
"get_telemetry": (get_telemetry, "Hent telemetri-historikk", {}),
"run_terminal": (run_terminal, "Kjør terminalkommando på VM", {"type":"object","properties":{"command":{"type":"string"}},"required":["command"]}),
"list_commits": (list_commits, "List siste commits i Gitea-repo", {}),
"get_file": (get_file, "Hent fil fra Gitea-repo", {"type":"object","properties":{"path":{"type":"string"}},"required":["path"]}),
"list_open_issues": (list_open_issues, "List åpne issues i Gitea-repo", {}),
"create_issue": (create_issue, "Opprett nytt issue i Gitea-repo", {"type":"object","properties":{"title":{"type":"string"},"body":{"type":"string"}},"required":["title"]}),
"push_file": (push_file, "Opprett eller oppdater fil i Gitea-repo", {"type":"object","properties":{"path":{"type":"string"},"content":{"type":"string"},"message":{"type":"string"}},"required":["path","content"]}),
}
def _tool_list_result():
return [
{
"name": name,
"description": desc,
"inputSchema": schema if schema else {"type": "object", "properties": {}},
}
for name, (_, desc, schema) in TOOLS.items()
]
def _jsonrpc_ok(req_id, result):
return {"jsonrpc": "2.0", "id": req_id, "result": result}
def _jsonrpc_err(req_id, code, message):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": code, "message": message}}
# ---------------------------------------------------------------------------
# MCP Streamable HTTP endpoint — POST /
# ---------------------------------------------------------------------------
@app.post("/")
async def mcp_handler(request: Request):
_verify_auth(request)
try:
body = await request.json()
except Exception:
return JSONResponse(_jsonrpc_err(None, -32700, "Parse error"), status_code=400)
method = body.get("method", "")
req_id = body.get("id")
params = body.get("params", {})
# initialize
if method == "initialize":
return JSONResponse(_jsonrpc_ok(req_id, {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {"name": "opax-mcp", "version": "3.0.0"},
}))
# tools/list
if method == "tools/list":
return JSONResponse(_jsonrpc_ok(req_id, {"tools": _tool_list_result()}))
# tools/call
if method == "tools/call":
tool_name = params.get("name") or params.get("tool")
tool_args = params.get("arguments", params.get("params", {}))
entry = TOOLS.get(tool_name)
if not entry:
return JSONResponse(_jsonrpc_err(req_id, -32601, f"Unknown tool: {tool_name}"))
handler, _, _ = entry
try:
result = await handler(tool_args)
except Exception as e:
logger.error(f"[MCP HANDLER ERROR] {type(e).__name__}: {e}")
return JSONResponse(_jsonrpc_err(req_id, -32000, str(e)))
return JSONResponse(_jsonrpc_ok(req_id, {
"content": [{"type": "text", "text": json.dumps(result, ensure_ascii=False)}]
}))
# notifications (fire-and-forget, no response needed)
if method.startswith("notifications/"):
return JSONResponse(status_code=202, content={})
return JSONResponse(_jsonrpc_err(req_id, -32601, f"Method not found: {method}"), status_code=404)
# ---------------------------------------------------------------------------
# Health (keepalive)
# ---------------------------------------------------------------------------
@app.get("/health")
async def health():
return {"status": "ok", "service": "opax-mcp", "version": "3.0.0"}