From ee50f8318ee45d63b6ca7f2ba135f756e991367d Mon Sep 17 00:00:00 2001 From: Chris Christiansen Date: Tue, 23 Jun 2026 09:51:32 +0000 Subject: [PATCH] feat(agent): initial implementation of OPAX agent Adds the core logic for the OPAX agent, including the main agent file, a deployment script, and an MCP client. This provides the foundation for the new agent-based architecture. --- agents/core-logic/agent.py | 234 +++++++++++++++++- agents/core-logic/deploy_agent.py | 41 ++- agents/core-logic/opax_mcp_client/__init__.py | 212 ++++++++++++++++ agents/core-logic/root_agent.py | 233 +++++++++++++++++ 4 files changed, 693 insertions(+), 27 deletions(-) create mode 100644 agents/core-logic/opax_mcp_client/__init__.py create mode 100644 agents/core-logic/root_agent.py diff --git a/agents/core-logic/agent.py b/agents/core-logic/agent.py index 1becba2..a81b9f1 100644 --- a/agents/core-logic/agent.py +++ b/agents/core-logic/agent.py @@ -1 +1,233 @@ -404: Not Found \ No newline at end of file +#!/usr/bin/env python3 +""" +agent.py — OSVauco OPAX agent (Jason Vauger). +opax_mcp_client er inlina direkte her for Vertex AI Agent Engine-kompatibilitet. +""" + +import asyncio +import os +import logging +import uuid +import httpx +import google.auth +import google.auth.transport.requests +from typing import Literal + +from google.adk.agents import Agent +from google.adk.tools import FunctionTool +from google.adk.runners import Runner +from google.adk.sessions.in_memory_session_service import InMemorySessionService +from google.genai import types + +logger = logging.getLogger(__name__) + +PROJECT_ID = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5") +LOCATION = os.environ.get("GOOGLE_CLOUD_LOCATION", "global") +ORCHESTRATOR_MODEL = os.environ.get("ORCHESTRATOR_MODEL", "gemini-2.5-flash") +SUBAGENT_MODEL = os.environ.get("SUBAGENT_MODEL", "gemini-2.5-flash") +REASONING_MODEL = os.environ.get("REASONING_MODEL", "gemini-2.5-flash") +HEAVY_ORCHESTRATOR = os.environ.get("HEAVY_ORCHESTRATOR_MODEL", "gemini-2.5-pro") +HEAVY_SUBAGENT = os.environ.get("HEAVY_SUBAGENT_MODEL", "gemini-2.5-pro") +HEAVY_REASONING = os.environ.get("HEAVY_REASONING_MODEL", "gemini-2.5-flash") +BUDGET_LIGHT = float(os.environ.get("BUDGET_A_USD_PER_TASK", "1.0")) +BUDGET_HEAVY = float(os.environ.get("HEAVY_MODE_BUDGET_USD_PER_DAY", "3.0")) +HEAVY_MODE_ALLOWED_USERS = ["opax", "admin"] +Mode = Literal["light", "heavy"] +APP_NAME = "opax" + +try: + from token_logger import log_token_usage, create_bq_table_if_not_exists +except ImportError: + def log_token_usage(*args, **kwargs): pass + def create_bq_table_if_not_exists(): pass + +try: + import sys as _sys, pathlib as _pathlib + + from token_budget import trim_context, route_model, TokenBudgetExceeded, budget_summary +except ImportError as _e: + def trim_context(history, system_prompt="", max_tokens=32000): + return history[-10:] if len(history) > 10 else history + def route_model(message, mode, flash_model, pro_model): + return flash_model if mode != "heavy" else pro_model + class TokenBudgetExceeded(Exception): pass + +create_bq_table_if_not_exists() + +_OPAX_MCP_URL = os.environ.get("MCP_SERVER_URL", "https://opax-mcp-zjbqp3prqq-uc.a.run.app") +_MCP_SECRET = os.environ.get("MCP_SECRET", "") + +def _identity_token() -> str: + try: + resp = httpx.get( + f"http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/default/identity?audience={_OPAX_MCP_URL}&format=full", + headers={"Metadata-Flavor": "Google"}, timeout=5) + if resp.status_code == 200 and resp.text.strip(): + return resp.text.strip() + except Exception: + pass + credentials, _ = google.auth.default() + credentials.refresh(google.auth.transport.requests.Request()) + return credentials.token + +def _mcp_headers() -> dict: + return {"Authorization": f"Bearer {_identity_token()}", "X-MCP-Secret": _MCP_SECRET, "Content-Type": "application/json"} + +def _call_tool(tool: str, params: dict = None) -> dict: + resp = httpx.post(f"{_OPAX_MCP_URL}/tools/call", headers=_mcp_headers(), json={"tool": tool, "params": params or {}}, timeout=30) + resp.raise_for_status() + return resp.json() + +def get_billing_summary() -> dict: + """Hent billing-oversikt for OSVauco (GCP-kostnader, token-forbruk).""" + return _call_tool("get_billing_summary") +def get_billing_credits() -> dict: + """Hent gjenværende GCP-kreditter og burn-rate.""" + return _call_tool("get_billing_credits") +def get_billing_anomalies() -> dict: + """Sjekk for kostnadanomalier og uventede spiker.""" + return _call_tool("get_billing_anomalies") +def get_billing_history() -> dict: + """Hent historisk kostnadsdata (90 dagar).""" + return _call_tool("get_billing_history") +def get_billing_budget() -> dict: + """Hent gjeldande budsjett for OSVauco.""" + return _call_tool("get_billing_budget") +def set_billing_budget(amount: float) -> dict: + """Sett nytt budsjettgrense. amount er beløpet i USD.""" + return _call_tool("set_billing_budget", {"amount": amount}) +def create_invite(email: str, company: str = "", tier: str = "starter") -> dict: + """Opprett invite-link for ny kunde. tier: starter | guard | shield | enterprise.""" + return _call_tool("create_invite", {"email": email, "company": company, "tier": tier}) +def list_customers() -> dict: + """List alle onboardede kunder og systemtilstand.""" + return _call_tool("list_customers") +def send_webhook(message: str, url: str = "", title: str = "OPAX varsel") -> dict: + """Send webhook-varsling til Slack/Teams/Discord.""" + return _call_tool("send_webhook", {"message": message, "url": url, "title": title}) +def send_email(to: str, subject: str, body: str = "") -> dict: + """Send e-post via SendGrid.""" + return _call_tool("send_email", {"to": to, "subject": subject, "body": body}) +def get_notify_channels() -> dict: + """List konfigurerte varslingskanalar.""" + return _call_tool("get_notify_channels") +def get_health() -> dict: + """Sjekk helsestatus for OPAX-plattformen.""" + return _call_tool("get_health") +def get_build_status() -> dict: + """Hent status på siste Cloud Build-kjøring.""" + return _call_tool("get_build_status") +def get_state() -> dict: + """Hent gjeldande systemtilstand og agent-statistikk for OSVauco.""" + return _call_tool("get_state") +def get_telemetry() -> dict: + """Hent token-forbruk og ytingsdata per modul.""" + return _call_tool("get_telemetry") +def run_terminal(command: str) -> dict: + """Kjør whitelist-kommando i OPAX terminal. Gyldige: health, billing, build, logs, help.""" + return _call_tool("run_terminal", {"command": command}) +def list_commits(limit: int = 10) -> dict: + """List siste commits i OSVauco-repoet på Gitea.""" + return _call_tool("list_commits", {"limit": limit}) +def get_file(path: str) -> dict: + """Hent innhald i ein fil frå Gitea-repoet.""" + return _call_tool("get_file", {"path": path}) +def list_open_issues() -> dict: + """List alle opne issues i OSVauco Gitea-repoet.""" + return _call_tool("list_open_issues", {}) +def create_issue(title: str, body: str = "") -> dict: + """Opprett nytt issue i Gitea-repoet.""" + return _call_tool("create_issue", {"title": title, "body": body}) +def push_file(path: str, content: str, message: str = "", sha: str = "") -> dict: + """Push/oppdater ein fil i Gitea.""" + return _call_tool("push_file", {"path": path, "content": content, "message": message, "sha": sha}) + +OPAX_TOOLS = [ + FunctionTool(func=get_billing_summary), FunctionTool(func=get_billing_credits), + FunctionTool(func=get_billing_anomalies), FunctionTool(func=get_billing_history), + FunctionTool(func=get_billing_budget), FunctionTool(func=set_billing_budget), + FunctionTool(func=create_invite), FunctionTool(func=list_customers), + FunctionTool(func=send_webhook), FunctionTool(func=send_email), + FunctionTool(func=get_notify_channels), FunctionTool(func=get_health), + FunctionTool(func=get_build_status), FunctionTool(func=get_state), + FunctionTool(func=get_telemetry), FunctionTool(func=run_terminal), + FunctionTool(func=list_commits), FunctionTool(func=get_file), + FunctionTool(func=list_open_issues), FunctionTool(func=create_issue), + FunctionTool(func=push_file), +] + +def _normalize_mode(mode: str) -> str: + mapping = {"A": "light", "A+": "heavy", "light": "light", "heavy": "heavy"} + if mode not in mapping: + raise ValueError(f"Invalid mode '{mode}'. Must be 'light' or 'heavy'.") + return mapping[mode] + +def authorize_mode(user_id: str, mode: str) -> None: + mode = _normalize_mode(mode) + if mode == "heavy" and user_id not in HEAVY_MODE_ALLOWED_USERS: + raise PermissionError(f"User '{user_id}' is not authorized for heavy mode.") + +def get_models_for_mode(mode: Mode) -> dict: + if mode == "heavy": + return {"orchestrator": HEAVY_ORCHESTRATOR, "subagent": HEAVY_SUBAGENT, "reasoning": HEAVY_REASONING, "budget_usd": BUDGET_HEAVY} + return {"orchestrator": ORCHESTRATOR_MODEL, "subagent": SUBAGENT_MODEL, "reasoning": REASONING_MODEL, "budget_usd": BUDGET_LIGHT} + +_INSTRUCTION_TEMPLATE = ( + "Du er Jason Vauger — OPAX-agenten for Vauco AS. Modus: {mode}. Modell: {model}. " + "Du har tilgang til alle OPAX-tools via opax-mcp. " + "Når brukaren ber om data — KALL ALLTID rett tool før du svarar. " + "Svar på norsk (bokmål) med mindre annet er bedt om. " + "HITL: ikkje kjør terraform/deploy utan godkjenning frå Chris." +) + +def build_agent(mode: str = "light", message: str = "") -> tuple: + mode = _normalize_mode(mode) + models = get_models_for_mode(mode) + actual_model = route_model(message=message, mode=mode, flash_model=ORCHESTRATOR_MODEL, pro_model=models["orchestrator"]) + instruction = _INSTRUCTION_TEMPLATE.format(mode=mode, model=actual_model) + agent = Agent(model=actual_model, name="jason_vauger", description=f"Jason Vauger — OPAX enterprise agent [{mode}]", instruction=instruction, tools=OPAX_TOOLS) + return agent, actual_model, instruction + +root_agent, _, _ = build_agent(mode="light") + +async def _run_async(message: str, user_id: str, session_id: str, mode: str, caller_type: str = "agent") -> str: + mode = _normalize_mode(mode) + agent, actual_model, instruction = build_agent(mode=mode, message=message) + module_name = f"jason/{mode}" + try: + trim_context([types.Content(role="user", parts=[types.Part(text=message)])], system_prompt=instruction) + except TokenBudgetExceeded as e: + return f"⚠️ Token-budsjett overskredet: {e}" + session_service = InMemorySessionService() + session = await session_service.create_session(app_name=APP_NAME, user_id=user_id, session_id=session_id) + runner = Runner(agent=agent, app_name=APP_NAME, session_service=session_service) + new_message = types.Content(role="user", parts=[types.Part(text=message)]) + final_text = "" + input_tokens = 0 + output_tokens = 0 + request_id = str(uuid.uuid4()) + async for event in runner.run_async(user_id=user_id, session_id=session.id, new_message=new_message): + if event.is_final_response() and event.content and event.content.parts: + final_text = event.content.parts[0].text or "" + if hasattr(event, "usage_metadata") and event.usage_metadata: + um = event.usage_metadata + input_tokens += getattr(um, "prompt_token_count", 0) or 0 + output_tokens += getattr(um, "candidates_token_count", 0) or 0 + if input_tokens > 0 or output_tokens > 0: + log_token_usage(agent_name=module_name, model_name=actual_model, input_tokens=input_tokens, output_tokens=output_tokens, request_id=request_id, module_name=module_name, caller_type=caller_type, session_id=session_id) + return final_text + +def run(message: str, user_id: str = "opax", session_id: str = "default", mode: str = "light", caller_type: str = "agent") -> str: + mode = _normalize_mode(mode) + authorize_mode(user_id, mode) + return asyncio.run(_run_async(message=message, user_id=user_id, session_id=session_id, mode=mode, caller_type=caller_type)) + +if __name__ == "__main__": + import sys + logging.basicConfig(level=logging.WARNING) + query = sys.argv[1] if len(sys.argv) > 1 else "Hva er OPAX?" + mode = sys.argv[2] if len(sys.argv) > 2 else "light" + caller_type = sys.argv[3] if len(sys.argv) > 3 else "cli" + print(f"Mode: {mode} | Query: {query}") + print("-" * 60) + print(run(message=query, user_id="opax", mode=mode, caller_type=caller_type)) diff --git a/agents/core-logic/deploy_agent.py b/agents/core-logic/deploy_agent.py index 9ceb465..6438ac0 100644 --- a/agents/core-logic/deploy_agent.py +++ b/agents/core-logic/deploy_agent.py @@ -1,50 +1,39 @@ #!/usr/bin/env python3 -""" -deploy_agent.py — Deploy ADK agent til Vertex AI Agent Engine -Usage: python3 deploy_agent.py [--project PROJECT_ID] [--region REGION] \ - [--display-name DISPLAY_NAME] [--staging-bucket GS_URI] - -Krever: google-cloud-aiplatform>=1.157.0, google-adk>=2.0.0 -""" - import argparse import sys +from pathlib import Path import vertexai -from vertexai import agent_engines # riktig import for SDK >= 1.157.0 +from vertexai import agent_engines - -def deploy(project: str, region: str, display_name: str, staging_bucket: str): +def deploy(project, region, display_name, staging_bucket): print(f"Initialiserer Vertex AI: project={project}, region={region}") vertexai.init(project=project, location=region, staging_bucket=staging_bucket) - # Importer Jason ADK agent - sys.path.insert(0, ".") - import agent as my_agent + base_dir = Path(__file__).resolve().parent + sys.path.insert(0, str(base_dir)) - print(f"Deployer agent '{display_name}' til Vertex AI Agent Engine i {region}...") - remote_agent = agent_engines.create( - my_agent.root_agent, + import agent as _agent_module + root_agent = _agent_module.root_agent + + print(f"Deployer agent '{display_name}'...") + remote = agent_engines.create( + root_agent, requirements=[ "google-cloud-aiplatform[adk,agent_engines]>=1.157.0", "google-adk>=2.2.0", "httpx>=0.27.0", "google-auth>=2.29.0", ], + extra_packages=[str(base_dir / "agent.py")], display_name=display_name, ) - print(f"\n✅ Agent deployet!") - print(f" Resource name : {remote_agent.resource_name}") - print(f" Console : https://console.cloud.google.com/ai/agents?project={project}") - print(f" Region : {region}") - print("\n⚠️ Husk: kjør teardown når du er ferdig for å unngå unnødige kostnader.") - return remote_agent - + print(f"\n✅ Agent deployet!\n Resource name: {remote.resource_name}") if __name__ == "__main__": - parser = argparse.ArgumentParser(description="Deploy Jason ADK agent til Vertex AI") + parser = argparse.ArgumentParser() parser.add_argument("--project", default="propane-will-491900-m5") parser.add_argument("--region", default="us-central1") - parser.add_argument("--display-name", default="jason-vauger-v1") + parser.add_argument("--display-name", default="jason-vauger-v12") parser.add_argument("--staging-bucket", default="gs://propane-will-491900-m5-agent-staging") args = parser.parse_args() deploy(args.project, args.region, args.display_name, args.staging_bucket) diff --git a/agents/core-logic/opax_mcp_client/__init__.py b/agents/core-logic/opax_mcp_client/__init__.py new file mode 100644 index 0000000..c174d99 --- /dev/null +++ b/agents/core-logic/opax_mcp_client/__init__.py @@ -0,0 +1,212 @@ +""" +opax_mcp_client.py — REST-klient for opax-mcp Cloud Run service. +Kopiert inn i core-logic/ slik at Vertex AI Agent Engine finn den ved deploy. + +Auth-lag: + 1. Cloud Run IAM : Authorization: Bearer (automatisk) + 2. Tool-level : X-MCP-Secret: (frå env) + +URL: https://opax-mcp-zjbqp3prqq-uc.a.run.app +""" + +import os +import httpx +import google.auth +import google.auth.transport.requests + +OPAX_MCP_URL = os.environ.get("MCP_SERVER_URL", "https://opax-mcp-zjbqp3prqq-uc.a.run.app") +MCP_SECRET = os.environ.get("MCP_SECRET", "") + + +def _identity_token() -> str: + metadata_url = ( + "http://metadata.google.internal/computeMetadata/v1/instance" + f"/service-accounts/default/identity?audience={OPAX_MCP_URL}&format=full" + ) + try: + resp = httpx.get(metadata_url, headers={"Metadata-Flavor": "Google"}, timeout=5) + if resp.status_code == 200 and resp.text.strip(): + return resp.text.strip() + except Exception: + pass + credentials, _ = google.auth.default() + credentials.refresh(google.auth.transport.requests.Request()) + return credentials.token + + +def _headers() -> dict: + return { + "Authorization": f"Bearer {_identity_token()}", + "X-MCP-Secret": MCP_SECRET, + "Content-Type": "application/json", + } + + +def call_tool(tool: str, params: dict = None) -> dict: + resp = httpx.post( + f"{OPAX_MCP_URL}/tools/call", + headers=_headers(), + json={"tool": tool, "params": params or {}}, + timeout=30, + ) + resp.raise_for_status() + return resp.json() + + +# --------------------------------------------------------------------------- +# Billing +# --------------------------------------------------------------------------- +def get_billing_summary() -> dict: + """Hent billing-oversikt for OSVauco (GCP-kostnader, token-forbruk).""" + return call_tool("get_billing_summary") + +def get_billing_forecast() -> dict: + """Hent token-estimat og kostnadsframskriving.""" + return call_tool("get_billing_forecast") + +def get_billing_credits() -> dict: + """Hent gjenværende GCP-kreditter og burn-rate.""" + return call_tool("get_billing_credits") + +def get_billing_anomalies() -> dict: + """Sjekk for kostnadanomalier og uventede spiker.""" + return call_tool("get_billing_anomalies") + +def get_billing_history() -> dict: + """Hent historisk kostnadsdata (90 dagar).""" + return call_tool("get_billing_history") + +def get_billing_budget() -> dict: + """Hent gjeldande budsjett for OSVauco.""" + return call_tool("get_billing_budget") + +def set_billing_budget(amount: float) -> dict: + """Sett nytt budsjettgrense i NOK/USD. amount er beløpet.""" + return call_tool("set_billing_budget", {"amount": amount}) + + +# --------------------------------------------------------------------------- +# Onboarding +# --------------------------------------------------------------------------- +def create_invite(email: str, company: str = "", tier: str = "starter") -> dict: + """Opprett invite-link for ny kunde. tier: starter | guard | shield | enterprise.""" + return call_tool("create_invite", {"email": email, "company": company, "tier": tier}) + +def list_customers() -> dict: + """List alle onboardede kunder og systemtilstand.""" + return call_tool("list_customers") + + +# --------------------------------------------------------------------------- +# Notify +# --------------------------------------------------------------------------- +def send_webhook(message: str, url: str = "", title: str = "OPAX varsel") -> dict: + """Send webhook-varsling til Slack/Teams/Discord. url er valgfri override.""" + return call_tool("send_webhook", {"message": message, "url": url, "title": title}) + +def send_email(to: str, subject: str, body: str = "") -> dict: + """Send e-post via SendGrid. to er mottaker-adresse.""" + return call_tool("send_email", {"to": to, "subject": subject, "body": body}) + +def send_sms(to: str, message: str, tier: str = "guard") -> dict: + """Send SMS via Twilio. Krev Guard+-tier. to er telefonnummer med landkode.""" + return call_tool("send_sms", {"to": to, "message": message, "tier": tier}) + +def get_notify_channels() -> dict: + """List konfigurerte varslingskanalar (e-post, SMS, webhook).""" + return call_tool("get_notify_channels") + + +# --------------------------------------------------------------------------- +# Agents +# --------------------------------------------------------------------------- +def run_jason(prompt: str) -> dict: + """Kall Jason-agenten med ein prompt via OPAX /run.""" + return call_tool("run_jason", {"prompt": prompt}) + +def run_emma(prompt: str) -> dict: + """Kall Emma-agenten (Gemma lokal) med ein prompt via OPAX /emma.""" + return call_tool("run_emma", {"prompt": prompt}) + + +# --------------------------------------------------------------------------- +# Platform +# --------------------------------------------------------------------------- +def get_health() -> dict: + """Sjekk helsestatus for OPAX-plattformen.""" + return call_tool("get_health") + +def get_build_status() -> dict: + """Hent status på siste Cloud Build-kjøring.""" + return call_tool("get_build_status") + +def get_state() -> dict: + """Hent gjeldande systemtilstand og agent-statistikk for OSVauco.""" + return call_tool("get_state") + +def get_telemetry() -> dict: + """Hent token-forbruk og ytingsdata per modul.""" + return call_tool("get_telemetry") + +def run_terminal(command: str) -> dict: + """Kjør whitelist-kommando i OPAX terminal. Gyldige: health, billing, build, logs, help.""" + return call_tool("run_terminal", {"command": command}) + + +# --------------------------------------------------------------------------- +# Gitea (repo-tools) +# --------------------------------------------------------------------------- +def list_commits(limit: int = 10) -> dict: + """List siste commits i OSVauco-repoet på Gitea.""" + return call_tool("list_commits", {"limit": limit}) + +def get_file(path: str) -> dict: + """Hent innhald i ein fil frå Gitea-repoet. path er relativ til rot.""" + return call_tool("get_file", {"path": path}) + +def list_open_issues() -> dict: + """List alle opne issues i OSVauco Gitea-repoet.""" + return call_tool("list_open_issues", {}) + +def create_issue(title: str, body: str = "") -> dict: + """Opprett nytt issue i Gitea-repoet.""" + return call_tool("create_issue", {"title": title, "body": body}) + +def push_file(path: str, content: str, message: str = "", sha: str = "") -> dict: + """Push/oppdater ein fil i Gitea. sha krevst ved oppdatering av eksisterande fil.""" + return call_tool("push_file", {"path": path, "content": content, "message": message, "sha": sha}) + + +# --------------------------------------------------------------------------- +# ADK FunctionTools — for bruk i agent.py +# --------------------------------------------------------------------------- +def get_all_function_tools() -> list: + """Returner alle tools som ADK FunctionTool-liste for Jason.""" + from google.adk.tools import FunctionTool + return [ + FunctionTool(func=get_billing_summary), + FunctionTool(func=get_billing_forecast), + FunctionTool(func=get_billing_credits), + FunctionTool(func=get_billing_anomalies), + FunctionTool(func=get_billing_history), + FunctionTool(func=get_billing_budget), + FunctionTool(func=set_billing_budget), + FunctionTool(func=create_invite), + FunctionTool(func=list_customers), + FunctionTool(func=send_webhook), + FunctionTool(func=send_email), + FunctionTool(func=send_sms), + FunctionTool(func=get_notify_channels), + FunctionTool(func=run_jason), + FunctionTool(func=run_emma), + FunctionTool(func=get_health), + FunctionTool(func=get_build_status), + FunctionTool(func=get_state), + FunctionTool(func=get_telemetry), + FunctionTool(func=run_terminal), + FunctionTool(func=list_commits), + FunctionTool(func=get_file), + FunctionTool(func=list_open_issues), + FunctionTool(func=create_issue), + FunctionTool(func=push_file), + ] diff --git a/agents/core-logic/root_agent.py b/agents/core-logic/root_agent.py new file mode 100644 index 0000000..a81b9f1 --- /dev/null +++ b/agents/core-logic/root_agent.py @@ -0,0 +1,233 @@ +#!/usr/bin/env python3 +""" +agent.py — OSVauco OPAX agent (Jason Vauger). +opax_mcp_client er inlina direkte her for Vertex AI Agent Engine-kompatibilitet. +""" + +import asyncio +import os +import logging +import uuid +import httpx +import google.auth +import google.auth.transport.requests +from typing import Literal + +from google.adk.agents import Agent +from google.adk.tools import FunctionTool +from google.adk.runners import Runner +from google.adk.sessions.in_memory_session_service import InMemorySessionService +from google.genai import types + +logger = logging.getLogger(__name__) + +PROJECT_ID = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5") +LOCATION = os.environ.get("GOOGLE_CLOUD_LOCATION", "global") +ORCHESTRATOR_MODEL = os.environ.get("ORCHESTRATOR_MODEL", "gemini-2.5-flash") +SUBAGENT_MODEL = os.environ.get("SUBAGENT_MODEL", "gemini-2.5-flash") +REASONING_MODEL = os.environ.get("REASONING_MODEL", "gemini-2.5-flash") +HEAVY_ORCHESTRATOR = os.environ.get("HEAVY_ORCHESTRATOR_MODEL", "gemini-2.5-pro") +HEAVY_SUBAGENT = os.environ.get("HEAVY_SUBAGENT_MODEL", "gemini-2.5-pro") +HEAVY_REASONING = os.environ.get("HEAVY_REASONING_MODEL", "gemini-2.5-flash") +BUDGET_LIGHT = float(os.environ.get("BUDGET_A_USD_PER_TASK", "1.0")) +BUDGET_HEAVY = float(os.environ.get("HEAVY_MODE_BUDGET_USD_PER_DAY", "3.0")) +HEAVY_MODE_ALLOWED_USERS = ["opax", "admin"] +Mode = Literal["light", "heavy"] +APP_NAME = "opax" + +try: + from token_logger import log_token_usage, create_bq_table_if_not_exists +except ImportError: + def log_token_usage(*args, **kwargs): pass + def create_bq_table_if_not_exists(): pass + +try: + import sys as _sys, pathlib as _pathlib + + from token_budget import trim_context, route_model, TokenBudgetExceeded, budget_summary +except ImportError as _e: + def trim_context(history, system_prompt="", max_tokens=32000): + return history[-10:] if len(history) > 10 else history + def route_model(message, mode, flash_model, pro_model): + return flash_model if mode != "heavy" else pro_model + class TokenBudgetExceeded(Exception): pass + +create_bq_table_if_not_exists() + +_OPAX_MCP_URL = os.environ.get("MCP_SERVER_URL", "https://opax-mcp-zjbqp3prqq-uc.a.run.app") +_MCP_SECRET = os.environ.get("MCP_SECRET", "") + +def _identity_token() -> str: + try: + resp = httpx.get( + f"http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/default/identity?audience={_OPAX_MCP_URL}&format=full", + headers={"Metadata-Flavor": "Google"}, timeout=5) + if resp.status_code == 200 and resp.text.strip(): + return resp.text.strip() + except Exception: + pass + credentials, _ = google.auth.default() + credentials.refresh(google.auth.transport.requests.Request()) + return credentials.token + +def _mcp_headers() -> dict: + return {"Authorization": f"Bearer {_identity_token()}", "X-MCP-Secret": _MCP_SECRET, "Content-Type": "application/json"} + +def _call_tool(tool: str, params: dict = None) -> dict: + resp = httpx.post(f"{_OPAX_MCP_URL}/tools/call", headers=_mcp_headers(), json={"tool": tool, "params": params or {}}, timeout=30) + resp.raise_for_status() + return resp.json() + +def get_billing_summary() -> dict: + """Hent billing-oversikt for OSVauco (GCP-kostnader, token-forbruk).""" + return _call_tool("get_billing_summary") +def get_billing_credits() -> dict: + """Hent gjenværende GCP-kreditter og burn-rate.""" + return _call_tool("get_billing_credits") +def get_billing_anomalies() -> dict: + """Sjekk for kostnadanomalier og uventede spiker.""" + return _call_tool("get_billing_anomalies") +def get_billing_history() -> dict: + """Hent historisk kostnadsdata (90 dagar).""" + return _call_tool("get_billing_history") +def get_billing_budget() -> dict: + """Hent gjeldande budsjett for OSVauco.""" + return _call_tool("get_billing_budget") +def set_billing_budget(amount: float) -> dict: + """Sett nytt budsjettgrense. amount er beløpet i USD.""" + return _call_tool("set_billing_budget", {"amount": amount}) +def create_invite(email: str, company: str = "", tier: str = "starter") -> dict: + """Opprett invite-link for ny kunde. tier: starter | guard | shield | enterprise.""" + return _call_tool("create_invite", {"email": email, "company": company, "tier": tier}) +def list_customers() -> dict: + """List alle onboardede kunder og systemtilstand.""" + return _call_tool("list_customers") +def send_webhook(message: str, url: str = "", title: str = "OPAX varsel") -> dict: + """Send webhook-varsling til Slack/Teams/Discord.""" + return _call_tool("send_webhook", {"message": message, "url": url, "title": title}) +def send_email(to: str, subject: str, body: str = "") -> dict: + """Send e-post via SendGrid.""" + return _call_tool("send_email", {"to": to, "subject": subject, "body": body}) +def get_notify_channels() -> dict: + """List konfigurerte varslingskanalar.""" + return _call_tool("get_notify_channels") +def get_health() -> dict: + """Sjekk helsestatus for OPAX-plattformen.""" + return _call_tool("get_health") +def get_build_status() -> dict: + """Hent status på siste Cloud Build-kjøring.""" + return _call_tool("get_build_status") +def get_state() -> dict: + """Hent gjeldande systemtilstand og agent-statistikk for OSVauco.""" + return _call_tool("get_state") +def get_telemetry() -> dict: + """Hent token-forbruk og ytingsdata per modul.""" + return _call_tool("get_telemetry") +def run_terminal(command: str) -> dict: + """Kjør whitelist-kommando i OPAX terminal. Gyldige: health, billing, build, logs, help.""" + return _call_tool("run_terminal", {"command": command}) +def list_commits(limit: int = 10) -> dict: + """List siste commits i OSVauco-repoet på Gitea.""" + return _call_tool("list_commits", {"limit": limit}) +def get_file(path: str) -> dict: + """Hent innhald i ein fil frå Gitea-repoet.""" + return _call_tool("get_file", {"path": path}) +def list_open_issues() -> dict: + """List alle opne issues i OSVauco Gitea-repoet.""" + return _call_tool("list_open_issues", {}) +def create_issue(title: str, body: str = "") -> dict: + """Opprett nytt issue i Gitea-repoet.""" + return _call_tool("create_issue", {"title": title, "body": body}) +def push_file(path: str, content: str, message: str = "", sha: str = "") -> dict: + """Push/oppdater ein fil i Gitea.""" + return _call_tool("push_file", {"path": path, "content": content, "message": message, "sha": sha}) + +OPAX_TOOLS = [ + FunctionTool(func=get_billing_summary), FunctionTool(func=get_billing_credits), + FunctionTool(func=get_billing_anomalies), FunctionTool(func=get_billing_history), + FunctionTool(func=get_billing_budget), FunctionTool(func=set_billing_budget), + FunctionTool(func=create_invite), FunctionTool(func=list_customers), + FunctionTool(func=send_webhook), FunctionTool(func=send_email), + FunctionTool(func=get_notify_channels), FunctionTool(func=get_health), + FunctionTool(func=get_build_status), FunctionTool(func=get_state), + FunctionTool(func=get_telemetry), FunctionTool(func=run_terminal), + FunctionTool(func=list_commits), FunctionTool(func=get_file), + FunctionTool(func=list_open_issues), FunctionTool(func=create_issue), + FunctionTool(func=push_file), +] + +def _normalize_mode(mode: str) -> str: + mapping = {"A": "light", "A+": "heavy", "light": "light", "heavy": "heavy"} + if mode not in mapping: + raise ValueError(f"Invalid mode '{mode}'. Must be 'light' or 'heavy'.") + return mapping[mode] + +def authorize_mode(user_id: str, mode: str) -> None: + mode = _normalize_mode(mode) + if mode == "heavy" and user_id not in HEAVY_MODE_ALLOWED_USERS: + raise PermissionError(f"User '{user_id}' is not authorized for heavy mode.") + +def get_models_for_mode(mode: Mode) -> dict: + if mode == "heavy": + return {"orchestrator": HEAVY_ORCHESTRATOR, "subagent": HEAVY_SUBAGENT, "reasoning": HEAVY_REASONING, "budget_usd": BUDGET_HEAVY} + return {"orchestrator": ORCHESTRATOR_MODEL, "subagent": SUBAGENT_MODEL, "reasoning": REASONING_MODEL, "budget_usd": BUDGET_LIGHT} + +_INSTRUCTION_TEMPLATE = ( + "Du er Jason Vauger — OPAX-agenten for Vauco AS. Modus: {mode}. Modell: {model}. " + "Du har tilgang til alle OPAX-tools via opax-mcp. " + "Når brukaren ber om data — KALL ALLTID rett tool før du svarar. " + "Svar på norsk (bokmål) med mindre annet er bedt om. " + "HITL: ikkje kjør terraform/deploy utan godkjenning frå Chris." +) + +def build_agent(mode: str = "light", message: str = "") -> tuple: + mode = _normalize_mode(mode) + models = get_models_for_mode(mode) + actual_model = route_model(message=message, mode=mode, flash_model=ORCHESTRATOR_MODEL, pro_model=models["orchestrator"]) + instruction = _INSTRUCTION_TEMPLATE.format(mode=mode, model=actual_model) + agent = Agent(model=actual_model, name="jason_vauger", description=f"Jason Vauger — OPAX enterprise agent [{mode}]", instruction=instruction, tools=OPAX_TOOLS) + return agent, actual_model, instruction + +root_agent, _, _ = build_agent(mode="light") + +async def _run_async(message: str, user_id: str, session_id: str, mode: str, caller_type: str = "agent") -> str: + mode = _normalize_mode(mode) + agent, actual_model, instruction = build_agent(mode=mode, message=message) + module_name = f"jason/{mode}" + try: + trim_context([types.Content(role="user", parts=[types.Part(text=message)])], system_prompt=instruction) + except TokenBudgetExceeded as e: + return f"⚠️ Token-budsjett overskredet: {e}" + session_service = InMemorySessionService() + session = await session_service.create_session(app_name=APP_NAME, user_id=user_id, session_id=session_id) + runner = Runner(agent=agent, app_name=APP_NAME, session_service=session_service) + new_message = types.Content(role="user", parts=[types.Part(text=message)]) + final_text = "" + input_tokens = 0 + output_tokens = 0 + request_id = str(uuid.uuid4()) + async for event in runner.run_async(user_id=user_id, session_id=session.id, new_message=new_message): + if event.is_final_response() and event.content and event.content.parts: + final_text = event.content.parts[0].text or "" + if hasattr(event, "usage_metadata") and event.usage_metadata: + um = event.usage_metadata + input_tokens += getattr(um, "prompt_token_count", 0) or 0 + output_tokens += getattr(um, "candidates_token_count", 0) or 0 + if input_tokens > 0 or output_tokens > 0: + log_token_usage(agent_name=module_name, model_name=actual_model, input_tokens=input_tokens, output_tokens=output_tokens, request_id=request_id, module_name=module_name, caller_type=caller_type, session_id=session_id) + return final_text + +def run(message: str, user_id: str = "opax", session_id: str = "default", mode: str = "light", caller_type: str = "agent") -> str: + mode = _normalize_mode(mode) + authorize_mode(user_id, mode) + return asyncio.run(_run_async(message=message, user_id=user_id, session_id=session_id, mode=mode, caller_type=caller_type)) + +if __name__ == "__main__": + import sys + logging.basicConfig(level=logging.WARNING) + query = sys.argv[1] if len(sys.argv) > 1 else "Hva er OPAX?" + mode = sys.argv[2] if len(sys.argv) > 2 else "light" + caller_type = sys.argv[3] if len(sys.argv) > 3 else "cli" + print(f"Mode: {mode} | Query: {query}") + print("-" * 60) + print(run(message=query, user_id="opax", mode=mode, caller_type=caller_type))