OSVauco/agents/core-logic/app.py

202 lines
6.1 KiB
Python

# agents/core-logic/app.py
# OSVauco-NMTMD-GCOS — FastAPI HTTP entrypoint for Cloud Run
import os
import sys
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# Allow imports from repo root (e.g. ml.billing_agent)
sys.path.insert(0, os.path.join(os.path.dirname(__file__), '..', '..'))
try:
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.genai.types import Content, Part
from agent import root_agent
from agents.recommendations_engine import get_recommendations
except ImportError as e:
logger.error(f"Failed to import ADK dependencies: {e}")
raise
session_service = InMemorySessionService()
APP_NAME = os.environ.get("CLOUD_RUN_SERVICE", "gcp-orchestrator")
PROJECT_ID = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5")
BQ_BILLING_DATASET = os.environ.get("BQ_BILLING_DATASET", "billing_data")
@asynccontextmanager
async def lifespan(app: FastAPI):
logger.info(f"OSVauco agent '{APP_NAME}' starting up")
yield
logger.info(f"OSVauco agent '{APP_NAME}' shutting down")
app = FastAPI(
title="OSVauco GCP Agent",
description="ADK-based multi-agent orchestrator on Cloud Run",
version="1.0.0",
lifespan=lifespan,
)
class RunRequest(BaseModel):
user_id: str
session_id: str
message: str
class RunResponse(BaseModel):
user_id: str
session_id: str
response: str
async def _ensure_session(user_id: str, session_id: str):
"""Await get_session; create if missing or raises."""
try:
session = await session_service.get_session(
app_name=APP_NAME,
user_id=user_id,
session_id=session_id,
)
if session is not None:
return session
except Exception:
pass
session = await session_service.create_session(
app_name=APP_NAME,
user_id=user_id,
session_id=session_id,
)
logger.info(f"Created new session: {session_id} for user: {user_id}")
return session
@app.get("/health")
async def health():
return JSONResponse({"status": "ok", "service": APP_NAME})
@app.post("/run", response_model=RunResponse)
async def run(req: RunRequest):
try:
await _ensure_session(req.user_id, req.session_id)
runner = Runner(
agent=root_agent,
app_name=APP_NAME,
session_service=session_service,
)
user_content = Content(
role="user",
parts=[Part(text=req.message)],
)
final_response = ""
async for event in runner.run_async(
user_id=req.user_id,
session_id=req.session_id,
new_message=user_content,
):
if event.is_final_response() and event.content:
for part in event.content.parts:
if part.text:
final_response += part.text
logger.info(f"[{req.user_id}/{req.session_id}] Response length: {len(final_response)}")
return RunResponse(
user_id=req.user_id,
session_id=req.session_id,
response=final_response,
)
except Exception as e:
logger.error(f"Agent run failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@app.get("/billing/tokens/summary")
async def billing_tokens_summary():
"""
CG3e — Aggreger LLM token-bruk og estimert kostnad per agent siste 30 dager.
Returnerer JSON-liste sortert etter total_cost DESC.
"""
try:
from google.cloud import bigquery
client = bigquery.Client(project=PROJECT_ID)
query = f"""
SELECT
agent_name,
model_name,
SUM(input_tokens) AS total_input_tokens,
SUM(output_tokens) AS total_output_tokens,
SUM(total_tokens) AS total_tokens,
SUM(estimated_cost_usd) AS total_cost_usd
FROM
`{PROJECT_ID}.{BQ_BILLING_DATASET}.llm_token_usage`
WHERE
timestamp >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 30 DAY)
GROUP BY
agent_name, model_name
ORDER BY
total_cost_usd DESC
"""
results = client.query(query).result()
rows = [
{
"agent_name": row.agent_name,
"model_name": row.model_name,
"total_input_tokens": row.total_input_tokens,
"total_output_tokens": row.total_output_tokens,
"total_tokens": row.total_tokens,
"total_cost_usd": round(float(row.total_cost_usd), 6),
}
for row in results
]
return JSONResponse({"period_days": 30, "rows": rows})
except Exception as e:
logger.error(f"[/billing/tokens/summary] Failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@app.get("/billing/recommendations")
async def billing_recommendations(budget: float = 500.0):
"""
CG4 — Kjorer anbefalings- og anomali-motoren.
Returnerer en liste med anbefalinger og estimert besparelse.
"""
try:
recommendations = get_recommendations(budget)
return JSONResponse(recommendations)
except Exception as e:
logger.error(f"[/billing/recommendations] Failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@app.get("/billing/by-service")
async def billing_by_service(days: int = 30):
"""
CG5 — Returnerer kostnad per tjeneste gruppert med SKU-detaljer.
Brukes av dashboard for grouped drill-down visning.
"""
try:
from ml.billing_agent import BillingAgent
agent = BillingAgent()
return JSONResponse(agent.get_service_totals(days))
except Exception as e:
logger.error(f"[/billing/by-service] Failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))