287 lines
17 KiB
Python
287 lines
17 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
|
|
import google.oauth2.id_token
|
|
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.3.0") # BUMPED
|
|
|
|
# ── Service discovery ───────────────────────────────────────────────────────
|
|
OSVAUCO_AGENT_URL = os.environ.get("OSVAUCO_AGENT_URL", "") # f.eks. https://osvauco-agent-....run.app
|
|
MCP_SECRET = os.environ.get("MCP_SECRET", "")
|
|
INTERNAL_API_KEY = os.environ.get("INTERNAL_API_KEY", "")
|
|
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")
|
|
|
|
# Emma-vm Ollama — direkte tilkobling
|
|
OLLAMA_BASE_URL = os.environ.get("OLLAMA_BASE_URL", "http://34.13.238.133:11434")
|
|
EMMA_MODEL = os.environ.get("EMMA_MODEL", "gemma3:27b")
|
|
EMMA_FAST_MODEL = os.environ.get("EMMA_FAST_MODEL", "gemma3:4b")
|
|
EMMA_LIGHT_MODEL = os.environ.get("EMMA_LIGHT_MODEL", "qwen2.5:3b")
|
|
|
|
logger.info(f"OSVAUCO_AGENT_URL: {OSVAUCO_AGENT_URL}")
|
|
logger.info(f"OLLAMA_BASE_URL: {OLLAMA_BASE_URL}")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Auth (kun for innkommende kall til opax-mcp, f.eks. fra Gemini TUI)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _verify_auth(request: Request):
|
|
"""Sjekker at innkommende kall til DENNE tjenesten (opax-mcp) er autentisert."""
|
|
if not MCP_SECRET:
|
|
logger.warning("MCP_SECRET is not set, skipping auth verification.")
|
|
return
|
|
auth = request.headers.get("Authorization", "")
|
|
if auth.startswith("Bearer ") and auth[7:] == MCP_SECRET:
|
|
return
|
|
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")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Backend Agent helpers (for kall VIDERE til osvauco-agent)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _agent_headers() -> dict:
|
|
"""Headers for maskin-til-maskin kall videre til osvauco-agent via X-Internal-Key."""
|
|
return {
|
|
"X-Internal-Key": INTERNAL_API_KEY, # Bruker den nye delte nøkkelen
|
|
"Content-Type": "application/json"
|
|
}
|
|
|
|
async def _agent_get(path: str) -> Any:
|
|
"""GET-kall til osvauco-agent."""
|
|
url = f"{OSVAUCO_AGENT_URL}{path}"
|
|
async with httpx.AsyncClient(timeout=30) as c:
|
|
r = await c.get(url, headers=_agent_headers())
|
|
r.raise_for_status()
|
|
return r.json()
|
|
|
|
async def _agent_post(path: str, body: dict) -> Any:
|
|
"""POST-kall til osvauco-agent."""
|
|
url = f"{OSVAUCO_AGENT_URL}{path}"
|
|
async with httpx.AsyncClient(timeout=45) as c:
|
|
r = await c.post(url, json=body, headers=_agent_headers())
|
|
r.raise_for_status()
|
|
return r.json()
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Ollama helpers — direkte mot emma-gpu-vm
|
|
# ---------------------------------------------------------------------------
|
|
|
|
async def _ollama_chat(model: str, prompt: str, system: str = "") -> dict:
|
|
# ... (beholdt uendret)
|
|
messages = []
|
|
if system:
|
|
messages.append({"role": "system", "content": system})
|
|
messages.append({"role": "user", "content": prompt})
|
|
payload = { "model": model, "messages": messages, "stream": False }
|
|
try:
|
|
async with httpx.AsyncClient(timeout=120) as c:
|
|
r = await c.post(f"{OLLAMA_BASE_URL}/api/chat", json=payload)
|
|
r.raise_for_status()
|
|
data = r.json()
|
|
return { "model": model, "response": data.get("message", {}).get("content", ""), "done": data.get("done", True), "total_duration_ms": round(data.get("total_duration", 0) / 1e6) }
|
|
except Exception as e:
|
|
logger.error(f"[OLLAMA ERROR] model={model} {type(e).__name__}: {e}")
|
|
raise
|
|
|
|
async def _ollama_models() -> list:
|
|
async with httpx.AsyncClient(timeout=10) as c:
|
|
r = await c.get(f"{OLLAMA_BASE_URL}/api/tags")
|
|
r.raise_for_status()
|
|
return r.json().get("models", [])
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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:
|
|
# ... (beholdt uendret)
|
|
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:
|
|
# ... (beholdt uendret)
|
|
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:
|
|
# ... (beholdt uendret)
|
|
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 (refaktorert til å bruke _agent_get/_agent_post)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
async def get_billing_summary(p): return await _agent_get("/billing/summary")
|
|
async def get_billing_forecast(p): return await _agent_get("/billing/forecast")
|
|
async def get_billing_credits(p): return await _agent_get("/billing/credits")
|
|
async def get_billing_anomalies(p): return await _agent_get("/billing/anomalies")
|
|
async def get_billing_history(p): return await _agent_get("/billing/history")
|
|
async def get_billing_budget(p): return await _agent_get("/billing/budget")
|
|
async def get_billing_tokens_by_module(p): return await _agent_get("/telemetry/history")
|
|
async def set_billing_budget(p): return await _agent_post("/billing/budget", {"budget": p.get("amount", 500)})
|
|
|
|
async def create_invite(p):
|
|
return await _agent_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 _agent_get("/state")
|
|
async def send_webhook(p):
|
|
return await _agent_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 _agent_post("/notify/email", {"to": p.get("to"), "subject": p.get("subject"), "event": p.get("event", "digest")})
|
|
async def send_sms(p): return await _agent_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 _agent_get("/notify/channels")
|
|
async def get_health(p): return await _agent_get("/health")
|
|
async def get_build_status(p): return await _agent_get("/opax/build-status")
|
|
async def get_state(p): return await _agent_get("/state")
|
|
async def get_telemetry(p): return await _agent_get("/telemetry/history")
|
|
async def run_terminal(p): return await _agent_post("/terminal/exec", {"cmd": p.get("command", p.get("cmd", "help"))})
|
|
async def tui_command(p): return await _agent_post("/tui-command", p)
|
|
|
|
# --- Uendrede funksjoner ---
|
|
async def run_jason(p):
|
|
"""Kaller /run på osvauco-agent, som nå har sin egen JASON_BACKEND-logikk."""
|
|
return await _agent_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 _ollama_chat(EMMA_MODEL, p.get("prompt", p.get("message", "")), p.get("system", "Du er Emma Vauger..."))
|
|
async def run_emma_fast(p): return await _ollama_chat(EMMA_FAST_MODEL, p.get("prompt", p.get("message", "")), p.get("system", "Du er en rask og konsis AI-assistent..."))
|
|
async def run_qwen(p): return await _ollama_chat(EMMA_LIGHT_MODEL, p.get("prompt", p.get("message", "")))
|
|
async def list_emma_models(p): return await _ollama_models()
|
|
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):
|
|
# ... (beholdt uendret)
|
|
repo, path = p.get("repo", GITEA_REPO), p.get("path")
|
|
content = base64.b64encode(p.get("content", "").encode()).decode()
|
|
body = {"message": p.get("message", f"chore: update {path} via opax-mcp"), "content": content, "branch": p.get("branch", "main")}
|
|
if not p.get("sha"):
|
|
try:
|
|
existing = await _gitea_get(f"/repos/{repo}/contents/{path}?ref={p.get('branch','main')}")
|
|
body["sha"] = existing["sha"]
|
|
except Exception: pass
|
|
else: body["sha"] = p["sha"]
|
|
if "sha" in body: 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 = {
|
|
# ... (alle verktøy beholdt, peker nå på de refaktorerte funksjonene)
|
|
"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"}}}),
|
|
"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"},"message":{"type":"string"}}}),
|
|
"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, "Emma Vauger (gemma3:27b) — primær lokal AI", {"type":"object","properties":{"prompt":{"type":"string"}}}),
|
|
"run_emma_fast": (run_emma_fast, "Emma Fast (gemma3:4b) — rask versjon", {"type":"object","properties":{"prompt":{"type":"string"}}}),
|
|
"run_qwen": (run_qwen, "Qwen 2.5 3B — lett hjelper", {"type":"object","properties":{"prompt":{"type":"string"}}}),
|
|
"list_emma_models": (list_emma_models, "List alle Ollama-modeller", {}),
|
|
"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", {}),
|
|
"tui_command": (tui_command, "Kjør shell-kommando via TUI-bro", {"type":"object","properties":{"command":{"type":"string"}}}),
|
|
"run_terminal": (run_terminal, "Kjør terminalkommando på VM", {"type":"object","properties":{"command":{"type":"string"}}}),
|
|
"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"}}}),
|
|
"push_file": (push_file, "Opprett eller oppdater fil i Gitea-repo", {"type":"object","properties":{"path":{"type":"string"},"content":{"type":"string"}},"required":["path","content"]}),
|
|
}
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# JSON-RPC Endpoint
|
|
# ---------------------------------------------------------------------------
|
|
|
|
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}}
|
|
|
|
@app.post("/")
|
|
async def mcp_handler(request: Request):
|
|
# ... (beholdt uendret)
|
|
_verify_auth(request)
|
|
try: body = await request.json()
|
|
except Exception: return JSONResponse(_jsonrpc_err(None, -32700, "Parse error"), status_code=400)
|
|
method, req_id, params = body.get("method", ""), body.get("id"), body.get("params", {})
|
|
if method == "initialize":
|
|
return JSONResponse(_jsonrpc_ok(req_id, {"protocolVersion": "2024-11-05", "capabilities": {"tools": {}},"serverInfo": {"name": "opax-mcp", "version": "3.3.0"}})) # BUMPED
|
|
if method == "tools/list":
|
|
return JSONResponse(_jsonrpc_ok(req_id, {"tools": [
|
|
{"name": name, "description": desc, "inputSchema": schema if schema else {"type": "object", "properties": {}}}
|
|
for name, (_, desc, schema) in TOOLS.items()
|
|
]}))
|
|
if method == "tools/call":
|
|
tool_name, tool_args = params.get("name") or params.get("tool"), 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)
|
|
return JSONResponse(_jsonrpc_ok(req_id, {"content": [{"type": "text", "text": json.dumps(result, ensure_ascii=False)}]}))
|
|
except Exception as e:
|
|
logger.error(f"[MCP HANDLER ERROR] tool={tool_name} {type(e).__name__}: {e}", exc_info=True)
|
|
return JSONResponse(_jsonrpc_err(req_id, -32000, str(e)))
|
|
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.3.0", "ollama": OLLAMA_BASE_URL} # BUMPED
|