916 lines
43 KiB
Python
916 lines
43 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
main.py — Cloud Run entrypoint, OSVauco OPAX.
|
|
Modes: light (gemini-2.5-flash) | heavy (gemini-2.5-pro)
|
|
ML-1: telemetri, state store, DAG.
|
|
CG1+CG2: billing endpoints + IAP-beskyttelse.
|
|
CG3: static-mappe serveres fra /static/*.
|
|
AUTH: /auth/login + /auth/callback for klient OAuth2 onboarding.
|
|
CG5a: POST /notify/webhook — push til Slack/Teams/Chat/Discord/etc.
|
|
CG5b: POST /notify/email — SendGrid digest til kunde-e-post.
|
|
CG5c: POST /notify/sms — Twilio SMS, Guard+-tier, spike-varsler.
|
|
CG5-onboard: GET /onboard + POST /onboard/complete — klient-onboarding wizard.
|
|
DOCS: GET /docs/{path} — proxy til privat GitHub-repo via ADC/PAT.
|
|
OQ-29: GET /notify/channels — list konfigurerte delivery-kanaler (env-var-basert). ✅
|
|
OQ-31: GET /opax/build-status — ekte Cloud Build-status via REST API. ✅
|
|
CG4-credits: GET /billing/credits — kreditt-saldo, burn rate og tom-dato.
|
|
Krever GOOGLE_CREDIT_TOTAL_USD i Cloud Run env for runway-beregning.
|
|
TERMINAL: POST /terminal/exec — whitelisted kommandoer: health, billing, build, logs, help. ✅
|
|
VM: GET /vm/ssh-key — henter public key fra Compute Engine metadata. ✅
|
|
"""
|
|
|
|
import os
|
|
import sys
|
|
import time
|
|
import pathlib
|
|
|
|
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "agents", "core-logic"))
|
|
|
|
from fastapi import FastAPI, HTTPException, Request
|
|
from fastapi.responses import FileResponse, JSONResponse, RedirectResponse, Response
|
|
from fastapi.staticfiles import StaticFiles
|
|
from pydantic import BaseModel, Field
|
|
from typing import List, Optional, Dict, Any
|
|
import datetime
|
|
import httpx
|
|
import google.auth
|
|
import google.auth.transport.requests
|
|
from cachetools import cached, TTLCache
|
|
|
|
from agent import run, authorize_mode
|
|
from ml import build_agent_dag, execute_dag, get_store, log_agent_call
|
|
from ml.telemetry import log_dag_execution
|
|
from ml.billing_agent import BillingAgent
|
|
from agents.aws_billing_agent import AWSBillingAgent
|
|
from ml.anomaly_detector import AnomalyDetector
|
|
from auth.token_store import save_token
|
|
|
|
|
|
from google.cloud import firestore
|
|
import firebase_admin
|
|
from firebase_admin import credentials, messaging
|
|
import base64
|
|
import json
|
|
from functools import wraps
|
|
from starlette.middleware.sessions import SessionMiddleware
|
|
from uvicorn.middleware.proxy_headers import ProxyHeadersMiddleware
|
|
from authlib.integrations.starlette_client import OAuth
|
|
import sendgrid
|
|
from sendgrid.helpers.mail import Mail
|
|
|
|
AGENT_ID = "osvauco-opax"
|
|
PROJECT_ID = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5")
|
|
|
|
# Guard+-tiers som får tilgang til SMS-varsler
|
|
SMS_ALLOWED_TIERS = {"guard", "shield", "enterprise"}
|
|
|
|
app = FastAPI(
|
|
title="OSVauco OPAX Agent",
|
|
description="Agent for OSVauco-OPAX platform.",
|
|
version="0.1.0",
|
|
)
|
|
|
|
# Paths exempt from IAP enforcement
|
|
IAP_EXEMPT_PATHS = {"/health", "/healthz", "/readiness", "/liveness", "/onboard"}
|
|
|
|
@app.middleware("http")
|
|
async def require_iap(request: Request, call_next):
|
|
if request.url.path in IAP_EXEMPT_PATHS:
|
|
return await call_next(request)
|
|
if not request.headers.get("x-goog-authenticated-user-email"):
|
|
return Response(status_code=401, content="Unauthorized")
|
|
return await call_next(request)
|
|
|
|
# ── FIREBASE & FIRESTORE INIT ───────────────────────────────────────────────
|
|
db = None
|
|
try:
|
|
firebase_admin.initialize_app()
|
|
db = firestore.Client()
|
|
print("Firestore client initialized successfully.")
|
|
except Exception as e:
|
|
print(f"WARNING: Firestore client failed to initialize: {e}", file=sys.stderr)
|
|
|
|
# ── AUTH & SESSION ────────────────────────────────────────────────────────────
|
|
app.add_middleware(ProxyHeadersMiddleware, trusted_hosts="*")
|
|
app.add_middleware(SessionMiddleware, secret_key=os.environ.get("SESSION_SECRET"))
|
|
|
|
oauth = OAuth()
|
|
oauth.register(
|
|
name='google',
|
|
client_id=os.environ.get("GOOGLE_CLIENT_ID"),
|
|
client_secret=os.environ.get("GOOGLE_CLIENT_SECRET"),
|
|
server_metadata_url='https://accounts.google.com/.well-known/openid-configuration',
|
|
client_kwargs={'scope': 'openid email profile'}
|
|
)
|
|
|
|
|
|
@app.get('/auth/login')
|
|
async def login(request: Request):
|
|
redirect_uri = "https://opax.vauco.no/auth/callback"
|
|
return await oauth.google.authorize_redirect(request, redirect_uri)
|
|
|
|
@app.get('/auth/callback', name='auth')
|
|
async def auth(request: Request):
|
|
token = await oauth.google.authorize_access_token(request)
|
|
user = token.get('userinfo')
|
|
if user:
|
|
request.session['user'] = dict(user)
|
|
return RedirectResponse(url='/static/billing-dashboard.html')
|
|
|
|
|
|
@app.get("/auth/me")
|
|
async def me(request: Request):
|
|
user = request.session.get("user")
|
|
if not user:
|
|
raise HTTPException(status_code=401, detail="Not authenticated")
|
|
return JSONResponse(user)
|
|
|
|
|
|
@app.get('/auth/logout')
|
|
async def logout(request: Request):
|
|
request.session.pop('user', None)
|
|
return RedirectResponse(url='/static/billing-dashboard.html')
|
|
|
|
|
|
ALLOWED_EMAILS = [email.strip() for email in os.environ.get("ALLOWED_EMAILS", "").split(",") if email.strip()]
|
|
ALERT_EMAIL = os.environ.get("ALERT_EMAIL")
|
|
|
|
|
|
def require_auth(func):
|
|
"""Decorator to protect endpoints that require authentication."""
|
|
@wraps(func)
|
|
async def wrapper(request: Request, *args, **kwargs):
|
|
user = request.session.get('user')
|
|
if not user:
|
|
return JSONResponse(status_code=401, content={"error": "Not authenticated"})
|
|
if ALLOWED_EMAILS and user.get('email') not in ALLOWED_EMAILS:
|
|
return JSONResponse(status_code=403, content={"error": "Email not allowed"})
|
|
return await func(request, *args, **kwargs)
|
|
return wrapper
|
|
|
|
|
|
# ── STATIC FILES ──────────────────────────────────────────────────────────────
|
|
_static_dir = pathlib.Path(__file__).parent / "static"
|
|
if _static_dir.is_dir():
|
|
app.mount("/static", StaticFiles(directory=str(_static_dir), html=True), name="static")
|
|
|
|
|
|
# ── DOCS PROXY (MD-filer fra privat GitHub-repo) ──────────────────────────────
|
|
_GH_REPO = "vauco-saas/OSVauco"
|
|
_GH_BRANCH = os.environ.get("DOCS_BRANCH", "main")
|
|
|
|
@app.get("/docs/{path:path}")
|
|
async def proxy_doc(path: str):
|
|
token = os.environ.get("GITHUB_PAT")
|
|
url = f"https://raw.githubusercontent.com/{_GH_REPO}/{_GH_BRANCH}/{path}"
|
|
headers = {"Authorization": f"token {token}"} if token else {}
|
|
try:
|
|
async with httpx.AsyncClient(timeout=10.0) as client:
|
|
resp = await client.get(url, headers=headers)
|
|
except httpx.TimeoutException:
|
|
raise HTTPException(status_code=504, detail="GitHub timeout")
|
|
if resp.status_code == 404:
|
|
raise HTTPException(status_code=404, detail=f"Fil ikke funnet: {path}")
|
|
if resp.status_code == 401:
|
|
raise HTTPException(status_code=503, detail="GitHub auth feilet — sjekk GITHUB_PAT i Secret Manager")
|
|
if resp.status_code != 200:
|
|
raise HTTPException(status_code=resp.status_code, detail="GitHub feil")
|
|
return Response(content=resp.text, media_type="text/plain; charset=utf-8")
|
|
|
|
|
|
# ── OQ-31: Cloud Build status (REST API) ─────────────────────────────────────
|
|
@app.get("/opax/build-status")
|
|
async def build_status():
|
|
try:
|
|
creds, project = google.auth.default()
|
|
creds.refresh(google.auth.transport.requests.Request())
|
|
url = f"https://cloudbuild.googleapis.com/v1/projects/{project}/builds?pageSize=5"
|
|
r = httpx.get(url, headers={"Authorization": f"Bearer {creds.token}"})
|
|
builds = r.json().get("builds", [])
|
|
return {"builds": [{"id": b["id"], "status": b["status"],
|
|
"branch": b.get("substitutions", {}).get("BRANCH_NAME", ""),
|
|
"createTime": b["createTime"]} for b in builds]}
|
|
except Exception as e:
|
|
import logging
|
|
logging.getLogger(__name__).error(f"[/opax/build-status] Failed: {e}", exc_info=True)
|
|
return JSONResponse({"status": "error", "message": str(e)}, status_code=200)
|
|
|
|
|
|
# ── TERMINAL: POST /terminal/exec ─────────────────────────────────────────────
|
|
TERMINAL_WHITELIST = {"health", "billing", "build", "logs", "help"}
|
|
|
|
class TerminalExecRequest(BaseModel):
|
|
cmd: str
|
|
|
|
@app.post("/terminal/exec")
|
|
async def terminal_exec(req: TerminalExecRequest):
|
|
"""
|
|
Whitelisted terminal-kommandoer:
|
|
health → GET /health internt
|
|
billing → GET /billing/summary
|
|
build → GET /opax/build-status
|
|
logs → siste 20 linjer fra Cloud Logging
|
|
help → returner kommandoliste
|
|
"""
|
|
cmd = req.cmd.strip().lower().split()[0] if req.cmd.strip() else ""
|
|
if cmd not in TERMINAL_WHITELIST:
|
|
return JSONResponse(
|
|
{"output": f"Ukjent kommando: '{req.cmd}'. Lov: {', '.join(sorted(TERMINAL_WHITELIST))}", "exit_code": 1}
|
|
)
|
|
|
|
try:
|
|
if cmd == "help":
|
|
lines = [
|
|
"Tilgjengelige kommandoer:",
|
|
" health — sjekk /health",
|
|
" billing — hent billing summary",
|
|
" build — hent siste build-status",
|
|
" logs — siste 20 linjer fra Cloud Logging",
|
|
" help — vis denne listen",
|
|
]
|
|
return {"output": "\n".join(lines), "exit_code": 0}
|
|
|
|
if cmd == "health":
|
|
return {"output": json.dumps({"status": "ok"}, ensure_ascii=False), "exit_code": 0}
|
|
|
|
if cmd == "billing":
|
|
try:
|
|
summary = BillingAgent().get_summary()
|
|
return {"output": json.dumps(summary, ensure_ascii=False, indent=2, default=str), "exit_code": 0}
|
|
except Exception as e:
|
|
return {"output": f"billing feilet: {e}", "exit_code": 1}
|
|
|
|
if cmd == "build":
|
|
try:
|
|
creds, project = google.auth.default()
|
|
creds.refresh(google.auth.transport.requests.Request())
|
|
url = f"https://cloudbuild.googleapis.com/v1/projects/{project}/builds?pageSize=5"
|
|
r = httpx.get(url, headers={"Authorization": f"Bearer {creds.token}"}, timeout=10)
|
|
builds = r.json().get("builds", [])
|
|
result = [{"id": b["id"], "status": b["status"],
|
|
"branch": b.get("substitutions", {}).get("BRANCH_NAME", ""),
|
|
"createTime": b["createTime"]} for b in builds]
|
|
return {"output": json.dumps(result, ensure_ascii=False, indent=2, default=str), "exit_code": 0}
|
|
except Exception as e:
|
|
return {"output": f"build feilet: {e}", "exit_code": 1}
|
|
|
|
if cmd == "logs":
|
|
try:
|
|
creds, project = google.auth.default()
|
|
creds.refresh(google.auth.transport.requests.Request())
|
|
log_filter = (
|
|
'resource.type="cloud_run_revision" '
|
|
'resource.labels.service_name="osvauco-agent" '
|
|
'severity>=DEFAULT'
|
|
)
|
|
body = {
|
|
"resourceNames": [f"projects/{project}"],
|
|
"filter": log_filter,
|
|
"orderBy": "timestamp desc",
|
|
"pageSize": 20,
|
|
}
|
|
r = httpx.post(
|
|
"https://logging.googleapis.com/v2/entries:list",
|
|
headers={"Authorization": f"Bearer {creds.token}"},
|
|
json=body,
|
|
timeout=10,
|
|
)
|
|
entries = r.json().get("entries", [])
|
|
if not entries:
|
|
return {"output": "(ingen logglinjer funnet)", "exit_code": 0}
|
|
lines = []
|
|
for e in reversed(entries):
|
|
ts = e.get("timestamp", "")[:19].replace("T", " ")
|
|
msg = e.get("textPayload") or json.dumps(e.get("jsonPayload", {}), ensure_ascii=False)
|
|
sev = e.get("severity", "")
|
|
lines.append(f"[{ts}] {sev:<8} {msg}")
|
|
return {"output": "\n".join(lines), "exit_code": 0}
|
|
except Exception as e:
|
|
return {"output": f"logs feilet: {e}", "exit_code": 1}
|
|
|
|
except Exception as e:
|
|
return {"output": f"Intern feil: {e}", "exit_code": 1}
|
|
|
|
return {"output": "", "exit_code": 0}
|
|
|
|
|
|
# ── VM: GET /vm/ssh-key ────────────────────────────────────────────────────────
|
|
@app.get("/vm/ssh-key")
|
|
async def vm_ssh_key(instance: str = "osvauco-dev-vm"):
|
|
"""
|
|
Henter public SSH-nøkkel fra Compute Engine instance metadata.
|
|
Bruker google.auth ADC — krever compute.instances.get på service account.
|
|
"""
|
|
try:
|
|
creds, project = google.auth.default()
|
|
creds.refresh(google.auth.transport.requests.Request())
|
|
zone = "us-central1-b"
|
|
url = (
|
|
f"https://compute.googleapis.com/compute/v1/projects/{project}"
|
|
f"/zones/{zone}/instances/{instance}"
|
|
)
|
|
async with httpx.AsyncClient(timeout=10.0) as client:
|
|
r = await client.get(url, headers={"Authorization": f"Bearer {creds.token}"})
|
|
if r.status_code == 404:
|
|
raise HTTPException(status_code=404, detail=f"Instans ikke funnet: {instance}")
|
|
if r.status_code != 200:
|
|
raise HTTPException(status_code=r.status_code, detail=f"Compute API feil: {r.text[:200]}")
|
|
|
|
data = r.json()
|
|
# SSH-nøkler ligger i metadata items med key "ssh-keys"
|
|
metadata_items = data.get("metadata", {}).get("items", [])
|
|
ssh_keys_raw = next(
|
|
(item["value"] for item in metadata_items if item["key"] == "ssh-keys"),
|
|
None
|
|
)
|
|
if not ssh_keys_raw:
|
|
return {"public_key": None, "instance": instance, "detail": "Ingen SSH-nøkkel funnet i metadata"}
|
|
|
|
# Format: "user:ssh-rsa AAAA..." — returner kun public key-delen
|
|
public_key = ssh_keys_raw.strip()
|
|
if ":" in public_key:
|
|
public_key = public_key.split(":", 1)[1].strip()
|
|
|
|
return {"public_key": public_key, "instance": instance}
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
raise HTTPException(status_code=500, detail=f"vm/ssh-key feilet: {e}")
|
|
|
|
|
|
# ── OQ-29: GET /notify/channels — konfigurerte delivery-kanaler ───────────────
|
|
@app.get("/notify/channels")
|
|
async def get_notify_channels():
|
|
return {
|
|
"channels": [
|
|
{"type": "email", "configured": bool(os.getenv("SENDGRID_API_KEY")), "target": os.getenv("NOTIFY_EMAIL_TO", None)},
|
|
{"type": "sms", "configured": bool(os.getenv("TWILIO_ACCOUNT_SID")), "target": os.getenv("TWILIO_FROM_NUMBER", None)},
|
|
{"type": "webhook", "configured": bool(os.getenv("NOTIFY_WEBHOOK_URL")), "url": os.getenv("NOTIFY_WEBHOOK_URL", None)},
|
|
]
|
|
}
|
|
|
|
|
|
# ── MODELLER ──────────────────────────────────────────────────────────────────
|
|
class RunRequest(BaseModel):
|
|
message: str
|
|
user_id: str = "opax"
|
|
session_id: str = "default"
|
|
mode: str = "light"
|
|
|
|
|
|
class DagRequest(BaseModel):
|
|
messages: List[str] = Field(...)
|
|
user_id: str = "opax"
|
|
session_id: str = "default"
|
|
mode: str = "light"
|
|
scheduler: str = Field("threads")
|
|
|
|
class PushSubscription(BaseModel):
|
|
token: str
|
|
budget_nok: float
|
|
|
|
class BudgetWebhookPayload(BaseModel):
|
|
message: dict
|
|
subscription: str
|
|
|
|
class BudgetUpdateRequest(BaseModel):
|
|
budget: float
|
|
|
|
|
|
# ── NOTIFICATION MODELS (CG5a + CG5b + CG5c) ──────────────────────────────────
|
|
class WebhookNotifyRequest(BaseModel):
|
|
url: str = Field(..., description="Webhook-URL (Slack/Teams/Google Chat/Discord/custom)")
|
|
event: str = Field("custom", description="Event-type: spike | digest | budget | anomaly | custom")
|
|
title: str = Field("CostGuard varsel", description="Tittel på meldingen")
|
|
body: str = Field(..., description="Meldingstekst")
|
|
payload: Optional[Dict[str, Any]] = Field(None, description="Valgfri rådata å inkludere")
|
|
|
|
class EmailNotifyRequest(BaseModel):
|
|
to: Optional[str] = Field(None, description="Mottaker-e-post. Faller tilbake til ALERT_EMAIL hvis tom.")
|
|
subject: Optional[str] = Field(None, description="Emne. Auto-generert fra event hvis tom.")
|
|
event: str = Field("digest", description="Event-type: spike | digest | budget | anomaly | custom")
|
|
body_html: Optional[str] = Field(None, description="HTML-innhold. Auto-generert fra billing-data hvis tom.")
|
|
payload: Optional[Dict[str, Any]] = Field(None, description="Valgfri rådata å inkludere i e-post")
|
|
|
|
class SmsNotifyRequest(BaseModel):
|
|
to: str = Field(..., description="Mottakers mobilnummer i E.164-format", pattern=r"^\+[1-9]\d{7,14}$")
|
|
body: str = Field(..., description="Meldingstekst. Maks 160 tegn anbefalt.", max_length=320)
|
|
event: str = Field("spike", description="Event-type: spike | budget | anomaly | custom")
|
|
tier: str = Field(..., description="Kundens tier. Må være guard, shield eller enterprise.")
|
|
|
|
|
|
# ── ONBOARDING MODEL (CG5-onboard) ────────────────────────────────────────────
|
|
class OnboardCompleteRequest(BaseModel):
|
|
company: str
|
|
contact: str
|
|
email: str
|
|
gcp_project: str
|
|
tier: str
|
|
channels: List[str]
|
|
webhook_url: Optional[str] = None
|
|
email_to: Optional[str] = None
|
|
|
|
|
|
# ── NOTIFICATION SERVICE (CG5a + CG5b + CG5c) ─────────────────────────────────
|
|
class NotificationService:
|
|
EVENT_META = {
|
|
"spike": {"emoji": "🚨", "label": "Kostnadsspike oppdaget"},
|
|
"digest": {"emoji": "📊", "label": "Daglig kostnadsoppsummering"},
|
|
"budget": {"emoji": "⚠️", "label": "Budsjettgrense nærmer seg"},
|
|
"anomaly": {"emoji": "🔍", "label": "Anomali oppdaget"},
|
|
"custom": {"emoji": "📢", "label": "CostGuard varsel"},
|
|
}
|
|
|
|
def _event_meta(self, event: str) -> dict:
|
|
return self.EVENT_META.get(event, self.EVENT_META["custom"])
|
|
|
|
def _build_slack_payload(self, req: WebhookNotifyRequest) -> dict:
|
|
meta = self._event_meta(req.event)
|
|
return {
|
|
"text": f"{meta['emoji']} *{req.title}*\n{req.body}",
|
|
"attachments": [
|
|
{
|
|
"color": "#FF6B35" if req.event in ("spike", "budget") else "#4A90D9",
|
|
"fields": [
|
|
{"title": k, "value": str(v), "short": True}
|
|
for k, v in (req.payload or {}).items()
|
|
]
|
|
}
|
|
] if req.payload else []
|
|
}
|
|
|
|
async def send_webhook(self, req: WebhookNotifyRequest) -> dict:
|
|
slack_payload = self._build_slack_payload(req)
|
|
try:
|
|
async with httpx.AsyncClient(timeout=10.0) as client:
|
|
resp = await client.post(req.url, json=slack_payload, headers={"Content-Type": "application/json"})
|
|
if resp.status_code >= 400:
|
|
return {"status": "error", "http_status": resp.status_code, "detail": resp.text[:200]}
|
|
return {"status": "ok", "http_status": resp.status_code, "event": req.event}
|
|
except httpx.TimeoutException:
|
|
return {"status": "error", "detail": "Webhook timeout (>10s)"}
|
|
except Exception as e:
|
|
return {"status": "error", "detail": str(e)}
|
|
|
|
async def send_email(self, req: EmailNotifyRequest) -> dict:
|
|
api_key = os.environ.get("SENDGRID_API_KEY")
|
|
to_email = req.to or ALERT_EMAIL
|
|
if not api_key:
|
|
return {"status": "not_configured", "detail": "SENDGRID_API_KEY mangler"}
|
|
if not to_email:
|
|
return {"status": "not_configured", "detail": "Ingen mottaker-e-post"}
|
|
meta = self._event_meta(req.event)
|
|
subject = req.subject or f"{meta['emoji']} CostGuard: {meta['label']} — {datetime.date.today()}"
|
|
html_content = req.body_html if req.body_html else await self._build_digest_html(req.payload)
|
|
try:
|
|
message = Mail(from_email="costguard@osvauco.no", to_emails=to_email, subject=subject, html_content=html_content)
|
|
sg = sendgrid.SendGridAPIClient(api_key)
|
|
response = sg.send(message)
|
|
msg_id = response.headers.get("X-Message-Id", "")
|
|
return {"status": "ok", "message_id": msg_id, "to": to_email, "event": req.event}
|
|
except Exception as e:
|
|
return {"status": "error", "detail": str(e)}
|
|
|
|
async def send_sms(self, req: SmsNotifyRequest) -> dict:
|
|
if req.tier.lower() not in SMS_ALLOWED_TIERS:
|
|
return {"status": "not_allowed", "detail": f"SMS krever Guard+-tier. Aktuell tier: '{req.tier}'."}
|
|
account_sid = os.environ.get("TWILIO_ACCOUNT_SID")
|
|
auth_token = os.environ.get("TWILIO_AUTH_TOKEN")
|
|
from_number = os.environ.get("TWILIO_FROM_NUMBER")
|
|
if not all([account_sid, auth_token, from_number]):
|
|
return {"status": "not_configured", "detail": "Twilio env-vars mangler"}
|
|
meta = self._event_meta(req.event)
|
|
sms_body = req.body if req.body.startswith(meta["emoji"]) else f"{meta['emoji']} CostGuard: {req.body}"
|
|
twilio_url = f"https://api.twilio.com/2010-04-01/Accounts/{account_sid}/Messages.json"
|
|
try:
|
|
async with httpx.AsyncClient(timeout=15.0) as client:
|
|
resp = await client.post(twilio_url, data={"From": from_number, "To": req.to, "Body": sms_body}, auth=(account_sid, auth_token))
|
|
data = resp.json()
|
|
if resp.status_code >= 400:
|
|
return {"status": "error", "http_status": resp.status_code, "detail": data.get("message", "")}
|
|
return {"status": "ok", "sid": data.get("sid"), "to": req.to, "sms_status": data.get("status"), "event": req.event}
|
|
except httpx.TimeoutException:
|
|
return {"status": "error", "detail": "Twilio timeout (>15s)"}
|
|
except Exception as e:
|
|
return {"status": "error", "detail": str(e)}
|
|
|
|
async def _build_digest_html(self, extra_payload: Optional[dict] = None) -> str:
|
|
try:
|
|
billing = BillingAgent()
|
|
summary = billing.get_summary().get("summary", [])
|
|
forecast = billing.get_forecast()
|
|
mtd = forecast.get("month_to_date_cost", 0)
|
|
daily_avg = forecast.get("daily_average_last_7_days", 0)
|
|
top3 = "".join(f"<li>{s['service']}: kr {s.get('total_cost', 0):.2f}</li>" for s in summary[:3])
|
|
except Exception:
|
|
mtd, daily_avg, top3 = 0, 0, "<li>Data ikke tilgjengelig</li>"
|
|
extra_rows = "".join(
|
|
f"<tr><td style='padding:4px 8px;color:#666'>{k}</td><td style='padding:4px 8px'><strong>{v}</strong></td></tr>"
|
|
for k, v in (extra_payload or {}).items()
|
|
)
|
|
return f"""
|
|
<div style='font-family:sans-serif;max-width:600px;margin:0 auto'>
|
|
<div style='background:#1a1a2e;padding:20px;border-radius:8px 8px 0 0'>
|
|
<h2 style='color:#fff;margin:0'>📊 CostGuard — Daglig oppsummering</h2>
|
|
<p style='color:#aaa;margin:4px 0 0'>{datetime.date.today()}</p>
|
|
</div>
|
|
<div style='background:#f9f9f9;padding:20px;border-radius:0 0 8px 8px'>
|
|
<table style='width:100%;border-collapse:collapse'>
|
|
<tr><td style='padding:4px 8px;color:#666'>Kostnad MTD</td><td style='padding:4px 8px'><strong>kr {mtd:.2f}</strong></td></tr>
|
|
<tr><td style='padding:4px 8px;color:#666'>Daglig snitt (7d)</td><td style='padding:4px 8px'><strong>kr {daily_avg:.2f}</strong></td></tr>
|
|
{extra_rows}
|
|
</table>
|
|
<h4 style='margin-top:16px'>Topp 3 GCP-tjenester</h4>
|
|
<ul>{top3}</ul>
|
|
<hr style='border:none;border-top:1px solid #ddd;margin:16px 0'>
|
|
<p style='color:#999;font-size:12px'>CostGuard by Vauco AS · opax.vauco.no</p>
|
|
</div>
|
|
</div>
|
|
"""
|
|
|
|
|
|
_notifier = NotificationService()
|
|
|
|
|
|
# ── NOTIFY ENDPOINTS (CG5a + CG5b + CG5c) ─────────────────────────────────────
|
|
async def _notify_webhook(request: Request, req: WebhookNotifyRequest):
|
|
result = await _notifier.send_webhook(req)
|
|
return JSONResponse(status_code=200 if result["status"] == "ok" else 502, content=result)
|
|
app.add_api_route("/notify/webhook", endpoint=require_auth(_notify_webhook), methods=["POST"])
|
|
|
|
async def _notify_email(request: Request, req: EmailNotifyRequest):
|
|
result = await _notifier.send_email(req)
|
|
if result["status"] == "not_configured":
|
|
return JSONResponse(status_code=503, content=result)
|
|
return JSONResponse(status_code=200 if result["status"] == "ok" else 500, content=result)
|
|
app.add_api_route("/notify/email", endpoint=require_auth(_notify_email), methods=["POST"])
|
|
|
|
async def _notify_sms(request: Request, req: SmsNotifyRequest):
|
|
result = await _notifier.send_sms(req)
|
|
if result["status"] == "not_allowed":
|
|
return JSONResponse(status_code=402, content=result)
|
|
if result["status"] == "not_configured":
|
|
return JSONResponse(status_code=503, content=result)
|
|
return JSONResponse(status_code=200 if result["status"] == "ok" else 500, content=result)
|
|
app.add_api_route("/notify/sms", endpoint=require_auth(_notify_sms), methods=["POST"])
|
|
|
|
|
|
# ── ONBOARDING ENDPOINTS (CG5-onboard) ───────────────────────────────────────
|
|
@app.get("/onboard")
|
|
def onboard_wizard():
|
|
return FileResponse("static/onboard.html")
|
|
|
|
|
|
@app.post("/onboard/complete")
|
|
async def onboard_complete(req: OnboardCompleteRequest):
|
|
if db:
|
|
try:
|
|
db.collection("onboarded_customers").document(req.gcp_project).set({
|
|
"company": req.company,
|
|
"contact": req.contact,
|
|
"email": req.email,
|
|
"gcp_project": req.gcp_project,
|
|
"tier": req.tier,
|
|
"channels": req.channels,
|
|
"webhook_url": req.webhook_url,
|
|
"email_to": req.email_to,
|
|
"onboarded_at": firestore.SERVER_TIMESTAMP,
|
|
})
|
|
except Exception as e:
|
|
print(f"[Onboard] Firestore-feil: {e}", file=sys.stderr)
|
|
|
|
notify_results: Dict[str, Any] = {}
|
|
if req.webhook_url and "webhook" in req.channels:
|
|
notify_results["webhook"] = await _notifier.send_webhook(WebhookNotifyRequest(
|
|
url=req.webhook_url, event="custom",
|
|
title=f"🎉 {req.company} er koblet til CostGuard!",
|
|
body=f"Hei {req.contact}! CostGuard overvåker nå *{req.gcp_project}*.",
|
|
payload={"bedrift": req.company, "tier": req.tier, "prosjekt": req.gcp_project, "dato": datetime.date.today().isoformat()},
|
|
))
|
|
if req.email_to and "email" in req.channels:
|
|
TIER_LABELS = {"starter": "Starter ($499/mnd)", "guard": "Guard ($999/mnd)", "shield": "Shield ($1.999/mnd)", "enterprise": "Enterprise ($3.500+/mnd)"}
|
|
notify_results["email"] = await _notifier.send_email(EmailNotifyRequest(
|
|
to=req.email_to, event="custom",
|
|
subject=f"✅ Velkommen til CostGuard — {req.company}",
|
|
body_html=f"<p>Hei {req.contact}, CostGuard overvåker nå {req.gcp_project} (tier: {TIER_LABELS.get(req.tier, req.tier)}).</p>",
|
|
))
|
|
internal_wh = os.environ.get("VAUCO_INTERNAL_WEBHOOK")
|
|
if internal_wh:
|
|
try:
|
|
await _notifier.send_webhook(WebhookNotifyRequest(
|
|
url=internal_wh, event="custom",
|
|
title="🆕 Ny CostGuard-kunde onboardet",
|
|
body=f"{req.company} ({req.gcp_project}) — tier: {req.tier}",
|
|
payload={"kontakt": req.contact, "e-post": req.email, "kanaler": ", ".join(req.channels)},
|
|
))
|
|
except Exception as e:
|
|
print(f"[Onboard] Intern Slack-feil: {e}", file=sys.stderr)
|
|
|
|
return {"status": "ok", "customer": req.gcp_project, "company": req.company, "channels": req.channels, "notify": notify_results}
|
|
|
|
|
|
# ── ROOT LANDING PAGE ───────────────────────────────────────────────────
|
|
@app.get("/")
|
|
def root():
|
|
return FileResponse("static/opax.html")
|
|
|
|
|
|
# ── HEALTH ────────────────────────────────────────────────────────────────────
|
|
@app.get("/health")
|
|
def health():
|
|
return {"status": "ok"}
|
|
|
|
|
|
@app.get("/manifest.json", include_in_schema=False)
|
|
def manifest():
|
|
return FileResponse("static/manifest.json")
|
|
|
|
|
|
@app.get("/sw.js", include_in_schema=False)
|
|
def service_worker():
|
|
return FileResponse("static/sw.js")
|
|
|
|
|
|
# ── ADMIN ─────────────────────────────────────────────────────────────────────
|
|
@app.get('/admin')
|
|
@require_auth
|
|
async def admin_panel(request: Request):
|
|
return FileResponse("static/admin.html")
|
|
|
|
@app.post('/admin/create-customer')
|
|
@require_auth
|
|
async def create_customer(request: Request):
|
|
import subprocess
|
|
data = await request.json()
|
|
customer_name = data.get("customer_name", "").strip()
|
|
project_id = data.get("project_id", "").strip()
|
|
billing_account_id = data.get("billing_account_id", "").strip()
|
|
alert_email = data.get("alert_email", "").strip()
|
|
container_image = data.get("container_image", "").strip()
|
|
region = data.get("region", "europe-north1").strip()
|
|
if not all([customer_name, project_id, billing_account_id, alert_email, container_image]):
|
|
raise HTTPException(status_code=400, detail="Alle felt er påkrevd")
|
|
import re
|
|
if not re.match(r'^[a-z0-9_-]+$', customer_name):
|
|
raise HTTPException(status_code=400, detail="customer_name kan kun inneholde a-z, 0-9, - og _")
|
|
customer_dir = f"infrastructure/terraform/customers/{customer_name}"
|
|
template_dir = "infrastructure/terraform/customers/_template"
|
|
import shutil
|
|
if os.path.exists(customer_dir):
|
|
raise HTTPException(status_code=409, detail=f"Kunde {customer_name} eksisterer allerede")
|
|
shutil.copytree(template_dir, customer_dir)
|
|
tfvars_content = f'''customer_id = "{customer_name}"\nproject_id = "{project_id}"\nregion = "{region}"\nbilling_account_id = "{billing_account_id}"\nalert_email = "{alert_email}"\nbilling_viewer_emails = ["chris.christiansen@vauco.no", "jason.vauger@vauco.no"]\ncontainer_image = "{container_image}"\n'''
|
|
with open(f"{customer_dir}/terraform.tfvars", "w") as f:
|
|
f.write(tfvars_content)
|
|
try:
|
|
init = subprocess.run(["terraform", f"-chdir={customer_dir}", "init", "-no-color"], capture_output=True, text=True, timeout=120)
|
|
if init.returncode != 0:
|
|
shutil.rmtree(customer_dir)
|
|
raise HTTPException(status_code=500, detail=f"terraform init feilet: {init.stderr[-500:]}")
|
|
apply = subprocess.run(["terraform", f"-chdir={customer_dir}", "apply", "-auto-approve", "-no-color"], capture_output=True, text=True, timeout=600)
|
|
if apply.returncode != 0:
|
|
raise HTTPException(status_code=500, detail=f"terraform apply feilet: {apply.stderr[-500:]}")
|
|
return {"status": "ok", "customer": customer_name, "project_id": project_id}
|
|
except subprocess.TimeoutExpired:
|
|
raise HTTPException(status_code=504, detail="Terraform tok for lang tid (>10 min)")
|
|
|
|
|
|
# ── BILLING ENDPOINTS (CG1 + CG2) ────────────────────────────────────────────
|
|
async def get_budget(request: Request):
|
|
if not db:
|
|
return JSONResponse(status_code=500, content={"error": "Firestore is not configured"})
|
|
try:
|
|
doc = db.collection("settings").document("budget").get()
|
|
return JSONResponse(content={"budget": doc.to_dict().get("limit", 500) if doc.exists else 500})
|
|
except Exception as e:
|
|
return JSONResponse(status_code=500, content={"error": str(e)})
|
|
|
|
async def set_budget(request: Request, payload: BudgetUpdateRequest):
|
|
if not db:
|
|
return JSONResponse(status_code=500, content={"error": "Firestore is not configured"})
|
|
try:
|
|
db.collection("settings").document("budget").set({"limit": payload.budget})
|
|
return JSONResponse(content={"status": "ok", "budget": payload.budget})
|
|
except Exception as e:
|
|
return JSONResponse(status_code=500, content={"error": str(e)})
|
|
|
|
app.add_api_route("/billing/budget", endpoint=require_auth(get_budget), methods=["GET"])
|
|
app.add_api_route("/billing/budget", endpoint=require_auth(set_budget), methods=["POST"])
|
|
|
|
|
|
@app.post("/billing/email-report")
|
|
async def trigger_email_report():
|
|
result = await _notifier.send_email(EmailNotifyRequest(event="digest"))
|
|
if result["status"] == "not_configured":
|
|
return JSONResponse(status_code=200, content={"status": "not_configured", "error": result["detail"]})
|
|
return JSONResponse(status_code=200 if result["status"] == "ok" else 500, content=result)
|
|
|
|
|
|
@app.post("/billing/snapshot")
|
|
async def create_daily_snapshot(request: Request):
|
|
await trigger_email_report()
|
|
return JSONResponse(content={"status": "ok", "snapshot_id": str(datetime.date.today())})
|
|
|
|
|
|
async def get_history(request: Request):
|
|
if not db:
|
|
return JSONResponse(status_code=500, content={"error": "Firestore is not configured"})
|
|
try:
|
|
end_date = datetime.date.today()
|
|
start_date = end_date - datetime.timedelta(days=90)
|
|
docs = db.collection("daily_snapshots") \
|
|
.where("created_at", ">=", start_date.isoformat()) \
|
|
.order_by("created_at", direction=firestore.Query.DESCENDING) \
|
|
.limit(90).stream()
|
|
history = [{"date": doc.id, **doc.to_dict()} for doc in docs]
|
|
history.reverse()
|
|
return JSONResponse(content={"history": history})
|
|
except Exception as e:
|
|
return JSONResponse(status_code=500, content={"error": str(e)})
|
|
app.add_api_route("/billing/history", endpoint=require_auth(get_history), methods=["GET"])
|
|
|
|
|
|
async def authenticated_billing_summary(request: Request):
|
|
try:
|
|
return BillingAgent().get_summary()
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
app.add_api_route("/billing/summary", endpoint=require_auth(authenticated_billing_summary), methods=["GET"])
|
|
|
|
|
|
async def authenticated_billing_forecast(request: Request):
|
|
try:
|
|
return BillingAgent().get_forecast()
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
app.add_api_route("/billing/forecast", endpoint=require_auth(authenticated_billing_forecast), methods=["GET"])
|
|
|
|
|
|
async def authenticated_billing_anomalies(request: Request):
|
|
try:
|
|
return AnomalyDetector().detect_anomalies()
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
app.add_api_route("/billing/anomalies", endpoint=require_auth(authenticated_billing_anomalies), methods=["GET"])
|
|
|
|
|
|
# ── CG4-credits: GET /billing/credits ────────────────────────────────────────
|
|
async def authenticated_billing_credits(request: Request, days: int = 90):
|
|
try:
|
|
return JSONResponse(BillingAgent().get_credits_status(days=days))
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
|
|
app.add_api_route("/billing/credits", endpoint=require_auth(authenticated_billing_credits), methods=["GET"])
|
|
|
|
|
|
@app.get("/billing-dashboard")
|
|
def billing_dashboard_view(request: Request):
|
|
return RedirectResponse(url='/static/billing-dashboard.html')
|
|
|
|
@app.post("/billing/subscribe")
|
|
def subscribe_for_push(sub: PushSubscription):
|
|
if not db:
|
|
raise HTTPException(status_code=500, detail="Firestore is not configured")
|
|
try:
|
|
db.collection("push_subscribers").document(sub.token).set({"budget_nok": sub.budget_nok, "subscribed_at": firestore.SERVER_TIMESTAMP})
|
|
return {"status": "ok"}
|
|
except Exception as e:
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
@app.post("/billing/budget-webhook")
|
|
async def budget_webhook(payload: BudgetWebhookPayload):
|
|
if not db:
|
|
return {"status": "error", "detail": "Firestore not configured"}
|
|
try:
|
|
data = base64.b64decode(payload.message.get("data", "")).decode("utf-8")
|
|
data_json = json.loads(data)
|
|
cost_amount = data_json.get("costAmount", 0)
|
|
budget_amount = data_json.get("budgetAmount", 0)
|
|
if budget_amount > 0 and (cost_amount / budget_amount) > 0.8:
|
|
percent_used = round((cost_amount / budget_amount) * 100)
|
|
tokens = [s.id for s in db.collection("push_subscribers").stream()]
|
|
if tokens:
|
|
notification = messaging.Notification(
|
|
title="⚠️ CostGuard Varsel",
|
|
body=f"Du har brukt {percent_used}% av budsjett (kr {int(cost_amount)} av kr {int(budget_amount)})"
|
|
)
|
|
messaging.send_multicast(messaging.MulticastMessage(tokens=tokens, notification=notification))
|
|
return {"status": "processed"}
|
|
except Exception as e:
|
|
return {"status": "error", "detail": str(e)}
|
|
|
|
|
|
@app.get("/billing/live")
|
|
@cached(TTLCache(maxsize=1, ttl=3600))
|
|
async def billing_live(request: Request):
|
|
user = request.session.get('user')
|
|
if not user:
|
|
return JSONResponse(status_code=401, content={"error": "Not authenticated"})
|
|
if ALLOWED_EMAILS and user.get('email') not in ALLOWED_EMAILS:
|
|
return JSONResponse(status_code=403, content={"error": "Email not allowed"})
|
|
bq_forecast = BillingAgent().get_forecast()
|
|
bq_forecast["data_source"] = "BigQuery Fallback"
|
|
return bq_forecast
|
|
|
|
|
|
# ── AWS BILLING ───────────────────────────────────────────────────────────────
|
|
async def authenticated_aws_billing_summary(request: Request):
|
|
try:
|
|
return AWSBillingAgent().get_summary()
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
app.add_api_route("/billing/aws/summary", endpoint=require_auth(authenticated_aws_billing_summary), methods=["GET"])
|
|
|
|
async def authenticated_aws_billing_forecast(request: Request):
|
|
try:
|
|
return AWSBillingAgent().get_forecast()
|
|
except Exception as exc:
|
|
return JSONResponse(status_code=500, content={"error": str(exc)})
|
|
app.add_api_route("/billing/aws/forecast", endpoint=require_auth(authenticated_aws_billing_forecast), methods=["GET"])
|
|
|
|
|
|
# ── AGENT ENDPOINTS ───────────────────────────────────────────────────────────
|
|
@app.post("/run")
|
|
def run_agent(req: RunRequest):
|
|
try:
|
|
authorize_mode(req.user_id, req.mode)
|
|
except ValueError as e:
|
|
raise HTTPException(status_code=400, detail=str(e))
|
|
except PermissionError as e:
|
|
raise HTTPException(status_code=403, detail=str(e))
|
|
store = get_store()
|
|
start = time.monotonic()
|
|
error_msg = None
|
|
response = None
|
|
try:
|
|
response = run(message=req.message, user_id=req.user_id, session_id=req.session_id, mode=req.mode)
|
|
except Exception as e:
|
|
error_msg = str(e)
|
|
raise HTTPException(status_code=500, detail=error_msg)
|
|
finally:
|
|
duration = round(time.monotonic() - start, 3)
|
|
success = error_msg is None
|
|
model = "gemini-2.5-flash" if req.mode in ("light", "A") else "gemini-2.5-pro"
|
|
log_agent_call(agent_id=AGENT_ID, input_payload={"message": req.message, "mode": req.mode}, output=response, model_used=model, mode=req.mode, duration_s=duration, success=success, error=error_msg)
|
|
store.push(AGENT_ID, "last_duration_s", duration)
|
|
store.push(AGENT_ID, "last_mode", req.mode)
|
|
store.push(AGENT_ID, "last_success", success)
|
|
if response:
|
|
store.push(AGENT_ID, "last_result_preview", response[:200])
|
|
return {"response": response}
|
|
|
|
|
|
@app.post("/run/dag")
|
|
def run_dag(req: DagRequest):
|
|
try:
|
|
authorize_mode(req.user_id, req.mode)
|
|
except ValueError as e:
|
|
raise HTTPException(status_code=400, detail=str(e))
|
|
except PermissionError as e:
|
|
raise HTTPException(status_code=403, detail=str(e))
|
|
def make_agent_fn(msg: str, idx: int):
|
|
def _agent_fn(payload: dict):
|
|
return run(message=msg, user_id=payload["user_id"], session_id=f"{payload['session_id']}-dag-{idx}", mode=payload["mode"])
|
|
_agent_fn.__name__ = f"opax-dag-{idx}"
|
|
return _agent_fn
|
|
payload = {"user_id": req.user_id, "session_id": req.session_id, "mode": req.mode}
|
|
agent_fns = [make_agent_fn(msg, i) for i, msg in enumerate(req.messages)]
|
|
agent_ids = [f"opax-dag-{i}" for i in range(len(req.messages))]
|
|
start = time.monotonic()
|
|
tasks = build_agent_dag(agent_fns, payload, agent_ids)
|
|
results = execute_dag(tasks, scheduler=req.scheduler)
|
|
total_dur = round(time.monotonic() - start, 3)
|
|
store = get_store()
|
|
log_dag_execution(dag_id=req.session_id, agent_results=results, total_duration_s=total_dur)
|
|
store.push(AGENT_ID, "last_dag_total_duration_s", total_dur)
|
|
store.push(AGENT_ID, "last_dag_agent_count", len(agent_ids))
|
|
store.push(AGENT_ID, "last_dag_success_count", sum(1 for r in results if r["success"]))
|
|
return {
|
|
"total_duration_s": total_dur,
|
|
"results": [{"index": i, "message": req.messages[i], "response": r["result"], "success": r["success"], "duration_s": r["duration_s"], "error": r["error"]} for i, r in enumerate(results)],
|
|
}
|
|
|
|
|
|
@app.get("/state")
|
|
def get_state():
|
|
return get_store().snapshot()
|
|
|
|
@app.get("/state/agents")
|
|
def list_agents():
|
|
return {"agents": get_store().list_agents()}
|
|
|
|
@app.get("/state/aggregate/{key}")
|
|
def aggregate_key(key: str):
|
|
return get_store().aggregate(key)
|
|
|
|
@app.get("/telemetry/history")
|
|
def telemetry_history(limit: int = 50):
|
|
return {"history": get_store().history(limit=limit)}
|
|
|
|
|
|
if __name__ == "__main__":
|
|
import uvicorn
|
|
port = int(os.environ.get("PORT", 8080))
|
|
uvicorn.run(app, host="0.0.0.0", port=port)
|