#!/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 delivery-kanaler per gcp_project. OQ-31: GET /opax/build-status — ekte Cloud Build-status. """ 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 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 ───────────────────────────────────────────────── @app.get("/opax/build-status") async def build_status(): """ OQ-31 — Ekte Cloud Build-status fra GCP. Erstatter mock/manglende endepunkt. Henter siste bygg via Cloud Build API. """ try: from google.cloud.devtools import cloudbuild_v1 client = cloudbuild_v1.CloudBuildClient() request_obj = cloudbuild_v1.ListBuildsRequest( project_id=PROJECT_ID, filter='trigger_id!=""', page_size=5, ) builds = list(client.list_builds(request=request_obj)) if not builds: return JSONResponse({"status": "unknown", "message": "Ingen builds funnet"}) b = builds[0] status_map = { 1: "queued", 2: "working", 3: "success", 4: "failure", 5: "internal_error", 6: "timeout", 7: "cancelled" } status_str = status_map.get(int(b.status), "unknown") duration_s = None if b.start_time and b.finish_time: duration_s = int(b.finish_time.seconds - b.start_time.seconds) return JSONResponse({ "status": status_str, "build_id": b.id, "trigger_id": b.build_trigger_id or "", "branch": (b.substitutions or {}).get("BRANCH_NAME", "main"), "commit": (b.substitutions or {}).get("SHORT_SHA", ""), "duration_s": duration_s, "start_time": str(b.start_time) if b.start_time else None, "finish_time": str(b.finish_time) if b.finish_time else None, "log_url": b.log_url or "", }) 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) # ── OQ-29: Notify channels per kunde ───────────────────────────────────────── @app.get("/notify/channels") async def notify_channels(gcp_project: str): """ OQ-29 — List delivery-kanaler konfigurert for en kunde. Leser fra Firestore: onboarded_customers/{gcp_project}.channels Returformat: { "gcp_project": "xxx", "channels": ["webhook", "email", "sms"] } """ if not db: return JSONResponse( status_code=503, content={"error": "Firestore ikke tilgjengelig"} ) try: doc = db.collection("onboarded_customers").document(gcp_project).get() if not doc.exists: return JSONResponse( status_code=404, content={"gcp_project": gcp_project, "channels": [], "detail": "Kunde ikke funnet"} ) data = doc.to_dict() return JSONResponse({ "gcp_project": gcp_project, "company": data.get("company", ""), "tier": data.get("tier", ""), "channels": data.get("channels", []), "webhook_url": data.get("webhook_url") or None, "email_to": data.get("email_to") or None, }) except Exception as e: return JSONResponse(status_code=500, content={"error": str(e)}) # ── 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"
  • {s['service']}: kr {s.get('total_cost', 0):.2f}
  • " for s in summary[:3]) except Exception: mtd, daily_avg, top3 = 0, 0, "
  • Data ikke tilgjengelig
  • " extra_rows = "".join( f"{k}{v}" for k, v in (extra_payload or {}).items() ) return f"""

    📊 CostGuard — Daglig oppsummering

    {datetime.date.today()}

    {extra_rows}
    Kostnad MTDkr {mtd:.2f}
    Daglig snitt (7d)kr {daily_avg:.2f}

    Topp 3 GCP-tjenester


    CostGuard by Vauco AS · opax.vauco.no

    """ _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"

    Hei {req.contact}, CostGuard overvåker nå {req.gcp_project} (tier: {TIER_LABELS.get(req.tier, req.tier)}).

    ", )) 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"]) @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)