OSVauco/opax-mcp/server.py

1824 lines
80 KiB
Python

"""opax-mcp — MCP Streamable HTTP server for OPAX/Vauco
Transport: MCP Streamable HTTP (JSON-RPC 2.0) on POST /
Auth: Authorization: Bearer <secret> OR X-MCP-Secret: <secret> OR api-key: <secret>
"""
import os
import json
import uuid
import httpx
import base64
from urllib.parse import quote
import difflib
import google.auth
import google.auth.transport.requests
from google.cloud import kms_v1
import google.oauth2.id_token
from googleapiclient.discovery import build
from google.cloud import secretmanager
from google.oauth2 import service_account
from fastapi import FastAPI, Request, HTTPException, Form
from fastapi.responses import JSONResponse, StreamingResponse, HTMLResponse
from fastapi.middleware.cors import CORSMiddleware
import asyncio
import secrets
# --- OPAX Imports ---
from persistence.firestore_store import FirestoreMemoryStore
from persistence.memory_store import MemoryStore
from contracts.common import (
CallerContext,
CallEmmaRequest,
CallEmmaResponse,
AuditEvent as EmmaAuditEvent,
)
from policy.caller_context import derive_caller_context
from persistence.memory_scope import calculate_effective_memory_scope
# --- Tool Imports (with fault tolerance) ---
try:
from tyr.tools.scan_tyr_surface import scan_tyr_surface
except Exception as e:
print(f"Failed to load scan_tyr_surface: {e}")
scan_tyr_surface = None
try:
from tyr.tools.eval_tyr_identity import eval_tyr_identity
except Exception as e:
print(f"Failed to load eval_tyr_identity: {e}")
eval_tyr_identity = None
try:
from tyr.tools.get_tyr_forecast import get_tyr_forecast
except Exception as e:
print(f"Failed to load get_tyr_forecast: {e}")
get_tyr_forecast = None
try:
from tyr.tools.get_tyr_user_risk import get_tyr_user_risk
except Exception as e:
print(f"Failed to load get_tyr_user_risk: {e}")
get_tyr_user_risk = None
try:
from modules.factory.tools.provision_new_mcp_module import provision_new_mcp_module
except Exception as e:
print(f"Failed to load provision_new_mcp_module: {e}")
provision_new_mcp_module = None
from emma_adapter import CanonicalEmma
from gitea_handler import (
handle_list_repo_files,
list_allowed_namespace_repositories,
handle_get_file_content,
resolve_branch_to_commit_sha,
download_repo_archive,
validate_repo_for_write,
validate_branch_for_write,
validate_path_for_write,
)
from capability_bridge import build_capability_system_context
from deployment_policy import get_deployment_target
from deployment_source import normalize_gitea_archive
from email.mime.text import MIMEText
from datetime import datetime, timezone, timedelta
from typing import Any, Optional, Dict, List, Literal
import logging
import re
from pydantic import BaseModel, Field
import hashlib
MAX_PROPOSE_CONTENT_BYTES = 1_048_576
_REPO_ID_RE = re.compile(r"^[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+$")
# A simple regex to validate the KMS key name format
_KMS_KEY_NAME_RE = re.compile(
r"^projects/[^/]+/locations/[^/]+/keyRings/[^/]+/cryptoKeys/[^/]+$"
)
# --- A2H2A Pydantic Models ---
class EncryptedPayloadV1(BaseModel):
"""Versioned Firestore-safe envelope for KMS-encrypted plan content."""
payload_version: Literal["1"] = "1"
kms_key_name: str
ciphertext: str
GITEA_CHANGE_PLAN_STATUS_PENDING = "PENDING"
GITEA_CHANGE_PLAN_STATUS_APPROVED = "APPROVED"
GITEA_CHANGE_PLAN_STATUS_APPLYING = "APPLYING"
GITEA_CHANGE_PLAN_STATUS_APPLIED = "APPLIED"
GITEA_CHANGE_PLAN_STATUS_REJECTED = "REJECTED"
GITEA_CHANGE_PLAN_STATUS_EXPIRED = "EXPIRED"
class GiteaChangePlan(BaseModel):
"""Immutable admin-approved one-file Gitea change plan."""
plan_id: str = Field(
default_factory=lambda: f"gitea-change-{uuid.uuid4().hex}"
)
status: str = GITEA_CHANGE_PLAN_STATUS_PENDING
created_at: datetime = Field(default_factory=datetime.utcnow)
expires_at: datetime = Field(
default_factory=lambda: datetime.utcnow() + timedelta(minutes=15)
)
approved_at: Optional[datetime] = None
approved_by: Optional[str] = None
applied_at: Optional[datetime] = None
repo: str
branch: str
path: str
base_sha: str
content_hash: str
commit_message: str
unified_diff: str
encrypted_payload: Optional[EncryptedPayloadV1] = None
approval_subject_hash: Optional[str] = None
result_commit_sha: Optional[str] = None
existing_file_sha: Optional[str] = None
apply_idempotency_key: Optional[str] = None
apply_status: Optional[str] = None
apply_error_message: Optional[str] = None
def is_gitea_change_plan_expired(plan: GiteaChangePlan) -> bool:
"""Checks if a Gitea change plan has expired."""
now_utc = datetime.now(timezone.utc)
expires_at = plan.expires_at
if expires_at.tzinfo is None:
expires_at = expires_at.replace(tzinfo=timezone.utc)
return now_utc >= expires_at
_ALLOWED_GITEA_PLAN_TRANSITIONS = {
GITEA_CHANGE_PLAN_STATUS_PENDING: {
GITEA_CHANGE_PLAN_STATUS_APPROVED,
GITEA_CHANGE_PLAN_STATUS_REJECTED,
GITEA_CHANGE_PLAN_STATUS_EXPIRED,
},
GITEA_CHANGE_PLAN_STATUS_APPROVED: {
GITEA_CHANGE_PLAN_STATUS_APPLYING,
GITEA_CHANGE_PLAN_STATUS_EXPIRED,
},
GITEA_CHANGE_PLAN_STATUS_APPLYING: {GITEA_CHANGE_PLAN_STATUS_APPLIED},
}
def is_valid_gitea_change_plan_status_transition(
current_status: str,
next_status: str,
) -> bool:
"""Checks if a Gitea change plan status transition is allowed."""
return next_status in _ALLOWED_GITEA_PLAN_TRANSITIONS.get(current_status, set())
async def create_gitea_change_plan(plan: GiteaChangePlan) -> dict:
"""Persist a pending one-file Gitea change plan in Firestore."""
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
plan_data = plan.model_dump(mode="json")
plan_data["status"] = "PENDING"
plan_ref = db.collection("gitea_change_plans").document(plan.plan_id)
await plan_ref.set(plan_data)
logger.info(
"Created Gitea change plan plan_id=%s repo=%s branch=%s path=%s",
plan.plan_id,
plan.repo,
plan.branch,
plan.path,
)
return plan_data
async def _encrypt_payload_v1(plaintext_content: str) -> EncryptedPayloadV1:
"""Encrypt plaintext UTF-8 content with the configured Cloud KMS key."""
key_name = os.environ.get("GITEA_PLAN_KMS_KEY_NAME")
if not key_name or not _KMS_KEY_NAME_RE.fullmatch(key_name):
raise RuntimeError("KMS key is not configured or has an invalid format.")
plaintext_bytes = plaintext_content.encode("utf-8")
async with kms_v1.KeyManagementServiceAsyncClient() as client:
response = await client.encrypt(
request={"name": key_name, "plaintext": plaintext_bytes}
)
ciphertext = getattr(response, "ciphertext", None)
if not isinstance(ciphertext, bytes) or not ciphertext:
raise RuntimeError("KMS encryption returned an invalid empty ciphertext.")
return EncryptedPayloadV1(
kms_key_name=key_name,
ciphertext=base64.b64encode(ciphertext).decode("ascii"),
)
def _calculate_approval_subject_hash(plan: GiteaChangePlan) -> str:
"""Hash the immutable, approval-bound representation of a Gitea plan."""
import binascii
payload = plan.encrypted_payload
if payload is None:
raise ValueError(
"Cannot calculate approval hash without an encrypted payload."
)
try:
ciphertext_bytes = base64.b64decode(payload.ciphertext, validate=True)
except (ValueError, TypeError, binascii.Error) as exc:
raise ValueError("Encrypted payload ciphertext is invalid.") from exc
subject = {
"plan_id": plan.plan_id,
"repo": plan.repo,
"branch": plan.branch,
"path": plan.path,
"base_sha": plan.base_sha,
"existing_file_sha": plan.existing_file_sha,
"content_hash": plan.content_hash,
"commit_message": plan.commit_message,
"unified_diff": plan.unified_diff,
"payload_version": payload.payload_version,
"kms_key_name": payload.kms_key_name,
"ciphertext_hash": hashlib.sha256(ciphertext_bytes).hexdigest(),
}
canonical_json = json.dumps(
subject,
sort_keys=True,
separators=(",", ":"),
).encode("utf-8")
return hashlib.sha256(canonical_json).hexdigest()
async def get_gitea_change_plan(
plan_id: str,
) -> Optional[GiteaChangePlan]:
"""Return a persisted Gitea change plan, or None when it does not exist."""
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
plan_ref = db.collection("gitea_change_plans").document(plan_id)
document = await plan_ref.get()
if not document.exists:
return None
return GiteaChangePlan(**document.to_dict())
async def transition_gitea_change_plan_status(
plan_id: str,
expected_status: str,
next_status: str,
updates: Optional[dict] = None,
) -> GiteaChangePlan:
if not is_valid_gitea_change_plan_status_transition(expected_status, next_status):
raise ValueError("Invalid Gitea change plan status transition")
if updates and "status" in updates:
raise ValueError("Gitea change plan updates cannot include status")
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
plan_ref = db.collection("gitea_change_plans").document(plan_id)
transaction = db.transaction()
@firestore.async_transactional
async def transactional_update(transaction):
snapshot = await plan_ref.get(transaction=transaction)
if not snapshot.exists:
raise ValueError("Gitea change plan not found")
if snapshot.get("status") != expected_status:
raise ValueError("Gitea change plan status changed")
update_data = dict(updates or {})
update_data["status"] = next_status
transaction.update(plan_ref, update_data)
await transactional_update(transaction)
updated_plan = await get_gitea_change_plan(plan_id)
if updated_plan is None:
raise ValueError("Gitea change plan not found")
return updated_plan
def _validate_admin_gitea_path(path: str) -> None:
"""Validate a repository-relative admin file path."""
if not isinstance(path, str) or not path or path.isspace():
raise ValueError("Path cannot be empty.")
if path.startswith("/") or "\\" in path or "\x00" in path:
raise ValueError("Path contains unsupported characters.")
segments = path.split("/")
if any(segment in {".", ".."} for segment in segments):
raise ValueError("Path traversal is not allowed.")
lowered_path = path.lower()
sensitive_markers = (
".env",
".pem",
".key",
"credentials",
"service_account",
"private_key",
)
if any(marker in lowered_path for marker in sensitive_markers):
raise ValueError("Path targets a sensitive file.")
def _validate_gitea_repo_format(repo: str) -> None:
"""Validate the canonical owner/repository identifier format."""
if not isinstance(repo, str) or not _REPO_ID_RE.fullmatch(repo):
raise ValueError("Repository must use the owner/repository format.")
async def _get_gitea_file_details(
repo: str,
path: str,
ref: str,
) -> tuple[Optional[str], Optional[str]]:
"""Return decoded file content and file SHA; return (None, None) on 404."""
api_path = f"/repos/{repo}/contents/{quote(path, safe='/')}?ref={ref}"
try:
payload = await _gitea_get(api_path)
except httpx.HTTPStatusError as exc:
if exc.response.status_code == 404:
return None, None
logger.warning(
"Gitea file detail lookup failed repo=%s path=%s ref=%s status=%s",
repo,
path,
ref,
exc.response.status_code,
)
raise ValueError("Gitea API error while reading file details.") from exc
try:
encoded_content = payload["content"]
content = base64.b64decode(encoded_content).decode("utf-8")
except (KeyError, TypeError, ValueError, UnicodeDecodeError) as exc:
raise ValueError("Gitea returned invalid file content.") from exc
return content, payload.get("sha")
async def propose_gitea_change(p: dict) -> dict:
"""Create and persist a read-only, immutable one-file Gitea change plan."""
repo = p.get("repo")
branch = p.get("branch")
path = p.get("path")
new_content = p.get("new_content")
base_sha = p.get("base_sha")
commit_message = p.get("commit_message")
if not all((repo, branch, path, base_sha, commit_message)):
raise ValueError(
"Missing required fields: repo, branch, path, base_sha, commit_message."
)
if new_content is None:
raise ValueError("Missing required field: new_content.")
if not isinstance(new_content, str):
raise ValueError("new_content must be a string.")
if len(new_content.encode("utf-8")) > MAX_PROPOSE_CONTENT_BYTES:
raise ValueError("new_content exceeds the maximum allowed size.")
_validate_gitea_repo_format(repo)
if repo != GITEA_REPO:
raise ValueError(
"Repository change proposal is not allowed for this repository."
)
_validate_admin_gitea_path(path)
validate_path_for_write(path)
validate_branch_for_write(branch)
live_sha = await resolve_branch_to_commit_sha(
branch_name=branch,
repo_id=repo,
gitea_url=GITEA_URL,
)
if live_sha != base_sha:
raise ValueError(
"Branch head does not match the supplied base_sha. Refresh and retry."
)
old_content, existing_file_sha = await _get_gitea_file_details(
repo=repo,
path=path,
ref=branch,
)
old_lines = (
old_content.splitlines(keepends=True)
if old_content is not None
else []
)
new_lines = new_content.splitlines(keepends=True)
unified_diff = "".join(
difflib.unified_diff(
old_lines,
new_lines,
fromfile=f"a/{path}",
tofile=f"b/{path}",
)
)
if not unified_diff:
raise ValueError("Proposed content produces no file change.")
content_hash = hashlib.sha256(new_content.encode("utf-8")).hexdigest()
encrypted_payload = await _encrypt_payload_v1(new_content)
plan = GiteaChangePlan(
repo=repo,
branch=branch,
path=path,
base_sha=base_sha,
content_hash=content_hash,
commit_message=commit_message,
unified_diff=unified_diff,
existing_file_sha=existing_file_sha,
encrypted_payload=encrypted_payload,
)
plan.approval_subject_hash = _calculate_approval_subject_hash(plan)
await create_gitea_change_plan(plan)
return {
"plan_id": plan.plan_id,
"status": "PENDING",
"expires_at": plan.expires_at.isoformat(),
"unified_diff": unified_diff,
"content_hash": content_hash,
"approval_subject_hash": plan.approval_subject_hash,
"repo": repo,
"branch": branch,
"path": path,
"base_sha": base_sha,
"commit_message": commit_message,
}
class TicketSource(BaseModel):
reporter: str
trigger: str
affected_service: str
project_id: str
region: str
class TicketContext(BaseModel):
summary: str
evidence_logs: List[str]
class ProposedAction(BaseModel):
action_type: str
runbook_reference: str
execution_tool: str
parameters: Dict
parameter_hash: Optional[str] = None
rollback_plan: str
class Governance(BaseModel):
approval_status: str = "PENDING"
authorized_approver: str
requires_mfa: bool = True
timeout_minutes: int = 15
approval_token_hash: Optional[str] = None
csrf_token_hash: Optional[str] = None
approved_at: Optional[datetime] = None
rejected_at: Optional[datetime] = None
verified_approver_email: Optional[str] = None
token_used_at: Optional[datetime] = None
class AuditEvent(BaseModel):
event_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
ticket_id: str
timestamp: datetime = Field(default_factory=datetime.utcnow)
event_type: str # e.g., PROPOSAL_CREATED, TICKET_APPROVED, TICKET_REJECTED
actor: str
details: Dict[str, Any]
class A2H2ATicket(BaseModel):
ticket_id: str
timestamp: datetime
severity: str # CRITICAL, HIGH, MEDIUM, LOW
category: str
source: TicketSource
context: TicketContext
proposed_action: ProposedAction
governance: Governance
def calculate_parameter_hash(self) -> str:
"""Beregn SHA-256 hash av parameters i kanonisk form."""
param_json = json.dumps(self.proposed_action.parameters, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(param_json.encode()).hexdigest()
def is_expired(self) -> bool:
"""Sjekk om ticket er utløpt"""
from datetime import datetime, timedelta
expiry = self.timestamp + timedelta(minutes=self.governance.timeout_minutes)
return datetime.utcnow() > expiry.replace(tzinfo=None)
class ConfirmationRequired(BaseModel):
message: str
preview: dict
requires_confirmation: bool = True
def require_confirmation(confirm: bool, preview: dict, message: str = "This action requires confirmation"):
if not confirm:
raise HTTPException(
status_code=400,
detail=ConfirmationRequired(message=message, preview=preview).model_dump()
)
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
app = FastAPI(title="opax-mcp", version="3.6.0") # BUMPED
app.add_middleware(
CORSMiddleware,
allow_origins=["https://opax-mcp-357036551735.us-central1.run.app"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# --- Metrics ---
A2H2A_METRICS = {
"tickets_created": 0,
"tickets_approved": 0,
"tickets_rejected": 0,
"tickets_expired": 0,
}
@app.get("/metrics")
async def get_metrics():
return A2H2A_METRICS
@app.get("/")
async def root_get():
return {"status": "healthy", "service": "OPAX MCP Server", "version": "3.6.0"}
async def sse_event_stream(request: Request):
"""Yields server-sent events for MCP protocol."""
# MCP requires an immediate connect event
yield f"id: {uuid.uuid4()}\nevent: connect\ndata: {json.dumps({'protocolVersion': '2024-11-05'})}\n\n"
while True:
# Send keepalive pings every 15s
yield f"id: {uuid.uuid4()}\nevent: ping\ndata: {json.dumps({'time': datetime.now(timezone.utc).isoformat()})}\n\n"
await asyncio.sleep(15)
@app.get("/sse")
async def sse_endpoint(request: Request):
"""Handles Server-Sent Events for external clients like Perplexity."""
return StreamingResponse(sse_event_stream(request), media_type="text/event-stream")
# --- A2H2A API Endpoints ---
# --- A2H2A Secure Design Placeholder ---
#
# The following logic is a placeholder for a secure A2H2A (Agent-to-Human-to-Agent)
# approval workflow. The key principles are:
#
# 1. NOTIFICATION, NOT ACTION: Google Chat messages are for notification only. They
# do not carry any identity or authorization guarantee.
# 2. VERIFIED IDENTITY: Human approval must happen through a protected surface that
# verifies the user's identity (e.g., a web app protected by Google IAP).
# 3. ATOMIC, IMMUTABLE ACTIONS: The ticket stored in the backend (Firestore) is the
# source of truth. The proposed action, its parameters, and the server-generated
# parameter hash are immutable.
# 4. ONE-TIME USE TOKENS: The approval UI should be accessed via a link containing a
# cryptographically random, single-use token that is invalidated after the first
# approval or rejection, or after the ticket's TTL expires.
# 5. DECOUPLED EXECUTION: The approval action only changes the ticket's state to
# 'APPROVED'. A separate, secure worker process will poll for approved tickets and
# execute them, ensuring a clean separation of concerns. This worker validates the
# ticket's integrity (non-expired, parameter hash match) before execution.
#
# This implementation only covers the safe recording of a proposed ticket.
# Execution is explicitly disabled.
# A strict allow-list of tools that can be proposed in an A2H2A ticket.
# `execution_enabled: False` means the tool can be proposed but not executed.
A2H2A_TOOL_ALLOW_LIST = {
"example_tool_1": {"execution_enabled": False, "description": "A sample tool that does something harmless."},
"example_tool_2": {"execution_enabled": False, "description": "Another sample tool."},
# Add real, vetted tools here with execution_enabled: True ONLY after extensive security review.
}
async def _create_audit_event(db, event: AuditEvent):
"""Stores an audit event in Firestore."""
event_ref = db.collection('a2h2a_audit_events').document(event.event_id)
await event_ref.set(event.model_dump(mode='json'))
logger.info(f"Created audit event {event.event_type} for ticket {event.ticket_id}")
@app.post("/api/v1/a2h2a/tickets")
async def create_a2h2a_ticket(ticket: A2H2ATicket, request: Request, test_mode: bool = False):
"""
Opprett ny A2H2A ticket og send til Google Chat for godkjenning.
Dette er kun en prototype for sikker registrering, ingen verktøy blir kjørt.
"""
try:
if not test_mode:
await _verify_auth(request)
# Security check: Reject if caller provides a parameter_hash.
if ticket.proposed_action.parameter_hash is not None:
raise HTTPException(status_code=400, detail="Do not provide a parameter_hash. The server will generate it.")
# Security check: Ensure the proposed tool is in the allow-list.
tool_name = ticket.proposed_action.execution_tool
if tool_name not in A2H2A_TOOL_ALLOW_LIST:
raise HTTPException(status_code=400, detail=f"Tool '{tool_name}' is not on the approved A2H2A allow-list.")
# Generate parameter hash server-side for integrity.
ticket.proposed_action.parameter_hash = ticket.calculate_parameter_hash()
# Generate a cryptographically secure, single-use approval token.
# The raw token is returned once in the review URL and never persisted.
# Only its SHA-256 hash is stored server-side.
raw_token = secrets.token_hex(32)
token_hash = hashlib.sha256(raw_token.encode()).hexdigest()
ticket.governance.approval_token_hash = token_hash
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
ticket_data = ticket.model_dump(mode='json')
ticket_data['created_at'] = datetime.utcnow().isoformat()
ticket_data['governance']['approval_status'] = 'PENDING'
doc_ref = db.collection('a2h2a_tickets').document(ticket.ticket_id)
await doc_ref.set(ticket_data)
# --- Create Initial Audit Event ---
audit_event = AuditEvent(
ticket_id=ticket.ticket_id,
event_type="PROPOSAL_CREATED",
actor="opax-mcp-server", # Use a fixed, trusted server identity
details={
"severity": ticket.severity,
"execution_tool": ticket.proposed_action.execution_tool,
"parameter_hash": ticket.proposed_action.parameter_hash,
"authorized_approver": ticket.governance.authorized_approver,
"timeout_minutes": ticket.governance.timeout_minutes,
}
)
await _create_audit_event(db, audit_event)
# The approval URL now contains the raw token for one-time use.
approval_ui_url = f"https://opax.vauco.no/a2h2a/review/{ticket.ticket_id}?token={raw_token}"
google_chat_webhook = os.getenv('A2H2A_GOOGLE_CHAT_WEBHOOK')
if google_chat_webhook and not test_mode:
# The raw_token is NOT included in the chat payload, only in the approval_ui_url
chat_payload = {
"cards": [
{
"header": {
"title": f"🚨 A2H2A Proposal: {ticket.severity}",
"subtitle": "OSVxCC Security (Prototype - No Execution)",
"image": {"imageUrl": "https://www.gstatic.com/images/branding/googlelogo/2x/googlelogo_color_92x30dp.png"}
},
"sections": [
{
"widgets": [
{
"textParagraph": {
"text": f"<b>Category:</b> {ticket.category}<br>" f"<b>Service:</b> {ticket.source.affected_service}<br>" f"<b>Tool:</b> {ticket.proposed_action.execution_tool}<br>" f"<b>Summary:</b> {ticket.context.summary}"
}
},
{
"buttons": [
{
"textButton": {
"text": "REVIEW IN APPROVAL UI",
"onClick": {"openLink": {"url": approval_ui_url}}
}
}
]
}
]
}
]
}
]
}
import httpx
async with httpx.AsyncClient() as client:
await client.post(google_chat_webhook, json=chat_payload)
response_payload = {"status": "success", "ticket_id": ticket.ticket_id, "parameter_hash": ticket.proposed_action.parameter_hash}
if test_mode:
response_payload['raw_token'] = raw_token
logger.info(f"A2H2A ticket created successfully: {ticket.ticket_id}", extra={"ticket_id": ticket.ticket_id, "event_type": "TICKET_CREATED"})
A2H2A_METRICS["tickets_created"] += 1
return response_payload
except Exception as e:
logger.error(f"Feil ved opprettelse av A2H2A ticket: {e}", extra={"ticket_id": ticket.ticket_id if 'ticket' in locals() else 'N/A', "event_type": "TICKET_CREATION_FAILED"}, exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
# --- A2H2A Approval UI (Prototype) ---
# These endpoints implement server-side validation of a single-use token
# and IAP-verified identity. They only change ticket state; no tool execution occurs.
@app.get("/a2h2a/review/{ticket_id}", response_class=HTMLResponse)
async def review_a2h2a_ticket(ticket_id: str, token: str, request: Request):
"""
Displays a review page for an A2H2A ticket.
This endpoint is expected to be protected by IAP.
"""
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
ticket_ref = db.collection('a2h2a_tickets').document(ticket_id)
ticket_doc = await ticket_ref.get()
if not ticket_doc.exists:
logger.warning("A2H2A review attempted for non-existent ticket.", extra={"ticket_id": ticket_id, "event_type": "REVIEW_TICKET_NOT_FOUND"})
return HTMLResponse(content="<h1>404: Ticket not found</h1>", status_code=404)
ticket = A2H2ATicket(**ticket_doc.to_dict())
# --- Security Validation ---
if ticket.governance.approval_status != 'PENDING':
logger.warning(f"A2H2A review attempted for already actioned ticket {ticket_id}.", extra={"ticket_id": ticket_id, "status": ticket.governance.approval_status, "event_type": "REVIEW_TICKET_NOT_PENDING"})
return HTMLResponse(content=f"<h1>400: Ticket Not Pending</h1><p>This ticket has already been actioned. Its current status is: <b>{ticket.governance.approval_status}</b>.</p>", status_code=400)
if ticket.is_expired():
logger.warning(f"A2H2A review attempted for expired ticket {ticket_id}.", extra={"ticket_id": ticket_id, "event_type": "REVIEW_TICKET_EXPIRED"})
A2H2A_METRICS["tickets_expired"] += 1
return HTMLResponse(content="<h1>400: Ticket Expired</h1>", status_code=400)
provided_token_hash = hashlib.sha256(token.encode()).hexdigest()
if not secrets.compare_digest(provided_token_hash, ticket.governance.approval_token_hash):
logger.error(f"Invalid approval token for ticket {ticket_id}.", extra={"ticket_id": ticket_id, "event_type": "REVIEW_INVALID_TOKEN"})
return HTMLResponse(content="<h1>403: Invalid Token</h1>", status_code=403)
user_email = request.headers.get("X-Goog-Authenticated-User-Email", "").replace("accounts.google.com:", "")
if not user_email or user_email != ticket.governance.authorized_approver:
logger.error(f"Unauthorized approver for ticket {ticket_id}.", extra={"ticket_id": ticket_id, "user_email": user_email, "authorized_approver": ticket.governance.authorized_approver, "event_type": "REVIEW_UNAUTHORIZED_APPROVER"})
return HTMLResponse(content=f"<h1>403: Unauthorized</h1><p>You (<b>{user_email}</b>) are not the authorized approver (<b>{ticket.governance.authorized_approver}</b>) for this ticket.</p>", status_code=403)
expiry_time = ticket.timestamp + timedelta(minutes=ticket.governance.timeout_minutes)
# --- Generate and store CSRF token ---
csrf_token = secrets.token_hex(16)
csrf_token_hash = hashlib.sha256(csrf_token.encode()).hexdigest()
await ticket_ref.update({"governance.csrf_token_hash": csrf_token_hash})
# --- Render HTML Page ---
html_content = f"""
<html>
<head>
<title>A2H2A Ticket Review</title>
<style> body {{ font-family: sans-serif; }} </style>
</head>
<body>
<h1>A2H2A Ticket Review: {ticket.severity}</h1>
<p><b>Authorized Approver:</b> {ticket.governance.authorized_approver}</p>
<p><b>Expires At:</b> {expiry_time.isoformat()}Z</p>
<hr>
<p><b>Ticket ID:</b> {ticket.ticket_id}</p>
<p><b>Tool:</b> {ticket.proposed_action.execution_tool}</p>
<p><b>Summary:</b> {ticket.context.summary}</p>
<p><b>Parameters Hash:</b> {ticket.proposed_action.parameter_hash}</p>
<hr>
<form action="/a2h2a/approve" method="post" style="display: inline-block;">
<input type="hidden" name="ticket_id" value="{ticket.ticket_id}">
<input type="hidden" name="token" value="{token}">
<input type="hidden" name="csrf_token" value="{csrf_token}">
<button type="submit" style="background-color: #28a745; color: white; padding: 10px; border: none; border-radius: 5px; cursor: pointer;">Approve</button>
</form>
<form action="/a2h2a/reject" method="post" style="display: inline-block;">
<input type="hidden" name="ticket_id" value="{ticket.ticket_id}">
<input type="hidden" name="token" value="{token}">
<input type="hidden" name="csrf_token" value="{csrf_token}">
<button type="submit" style="background-color: #dc3545; color: white; padding: 10px; border: none; border-radius: 5px; cursor: pointer;">Reject</button>
</form>
</body>
</html>
"""
return HTMLResponse(content=html_content)
async def _process_approval_action(ticket_id: str, token: str, user_email: str, action: str, csrf_token: str):
"""Helper to handle the logic for approving or rejecting a ticket."""
from google.cloud import firestore
db = firestore.AsyncClient(project=GOOGLE_CLOUD_PROJECT)
ticket_ref = db.collection('a2h2a_tickets').document(ticket_id)
ticket_doc = await ticket_ref.get()
if not ticket_doc.exists:
raise HTTPException(status_code=404, detail="Ticket not found")
ticket = A2H2ATicket(**ticket_doc.to_dict())
ticket_dict = ticket_doc.to_dict()
# --- Full Security Re-validation ---
if ticket.governance.approval_status != 'PENDING':
raise HTTPException(status_code=400, detail=f"Ticket has already been actioned (Status: {ticket.governance.approval_status})")
if ticket.is_expired():
A2H2A_METRICS["tickets_expired"] += 1
raise HTTPException(status_code=400, detail="Ticket has expired")
provided_token_hash = hashlib.sha256(token.encode()).hexdigest()
if not secrets.compare_digest(provided_token_hash, ticket.governance.approval_token_hash):
raise HTTPException(status_code=403, detail="Invalid token")
# --- CSRF Token validation ---
stored_csrf_hash = ticket_dict.get("governance", {}).get("csrf_token_hash")
if not stored_csrf_hash:
logger.error(f"CSRF token missing or already used for ticket {ticket_id}.", extra={"ticket_id": ticket_id, "user_email": user_email, "event_type": "APPROVAL_CSRF_TOKEN_MISSING"})
raise HTTPException(status_code=403, detail="CSRF token missing or already used.")
provided_csrf_hash = hashlib.sha256(csrf_token.encode()).hexdigest()
if not secrets.compare_digest(provided_csrf_hash, stored_csrf_hash):
logger.error(f"Invalid CSRF token for ticket {ticket_id}.", extra={"ticket_id": ticket_id, "user_email": user_email, "event_type": "APPROVAL_INVALID_CSRF_TOKEN"})
raise HTTPException(status_code=403, detail="Invalid CSRF token")
if user_email != ticket.governance.authorized_approver:
logger.error(f"Unauthorized approver in approval action for ticket {ticket_id}.", extra={"ticket_id": ticket_id, "user_email": user_email, "authorized_approver": ticket.governance.authorized_approver, "event_type": "APPROVAL_UNAUTHORIZED_APPROVER"})
raise HTTPException(status_code=403, detail="Unauthorized")
now = datetime.utcnow()
update_data = {
"governance.approval_status": action,
"governance.verified_approver_email": user_email,
"governance.token_used_at": now.isoformat(),
f"governance.{action.lower()}_at": now.isoformat(),
"governance.csrf_token_hash": None # Invalidate CSRF token
}
await ticket_ref.update(update_data)
logger.info(f"A2H2A ticket {ticket_id} has been {action.lower()}.")
if action == "APPROVED":
A2H2A_METRICS["tickets_approved"] += 1
elif action == "REJECTED":
A2H2A_METRICS["tickets_rejected"] += 1
audit_event = AuditEvent(
ticket_id=ticket.ticket_id,
event_type=f"TICKET_{action}",
actor=user_email,
details={"message": f"Ticket was {action.lower()} by {user_email}."}
)
await _create_audit_event(db, audit_event)
return {"status": action.lower(), "ticket_id": ticket_id}
@app.post("/a2h2a/approve")
async def approve_a2h2a_ticket(request: Request, ticket_id: str = Form(...), token: str = Form(...), csrf_token: str = Form(...)):
user_email = request.headers.get("X-Goog-Authenticated-User-Email", "").replace("accounts.google.com:", "")
if not user_email: raise HTTPException(status_code=403, detail="IAP header missing or invalid.")
return await _process_approval_action(ticket_id, token, user_email, "APPROVED", csrf_token)
@app.post("/a2h2a/reject")
async def reject_a2h2a_ticket(request: Request, ticket_id: str = Form(...), token: str = Form(...), csrf_token: str = Form(...)):
user_email = request.headers.get("X-Goog-Authenticated-User-Email", "").replace("accounts.google.com:", "")
if not user_email: raise HTTPException(status_code=403, detail="IAP header missing or invalid.")
return await _process_approval_action(ticket_id, token, user_email, "REJECTED", csrf_token)
# ── Service discovery ───────────────────────────────────────────────────────
OSVAUCO_AGENT_URL = os.environ.get("OSVAUCO_AGENT_URL", "") # f.eks. https://osvauco-agent-....run.app
MCP_SECRET = os.environ.get("MCP_SECRET", "")
INTERNAL_API_KEY = os.environ.get("INTERNAL_API_KEY", "")
GITEA_URL = os.environ.get("GITEA_URL", "http://34.59.131.162:3000")
GITEA_TOKEN = os.environ.get("GITEA_TOKEN", "")
GITEA_REPO = os.environ.get("GITEA_REPO", "chris/OSVauco")
GOOGLE_CLOUD_PROJECT = os.environ.get("GOOGLE_CLOUD_PROJECT", "propane-will-491900-m5")
WORKSPACE_ADMIN_USER = "chris.christiansen@vauco.no"
# Emma-vm Ollama — direkte tilkobling
OLLAMA_BASE_URL = os.environ.get("OLLAMA_BASE_URL", "http://34.13.238.133:11434")
EMMA_MODEL = os.environ.get("EMMA_MODEL", "gemma3:27b")
EMMA_FAST_MODEL = os.environ.get("EMMA_FAST_MODEL", "gemma3:4b")
EMMA_LIGHT_MODEL = os.environ.get("EMMA_LIGHT_MODEL", "qwen2.5:3b")
logger.info(f"OSVAUCO_AGENT_URL: {OSVAUCO_AGENT_URL}")
logger.info(f"OLLAMA_BASE_URL: {OLLAMA_BASE_URL}")
# ---------------------------------------------------------------------------
# Auth (kun for innkommende kall til opax-mcp, f.eks. fra Gemini TUI)
# ---------------------------------------------------------------------------
async def _verify_auth(request: Request) -> None:
secret = os.getenv("MCP_SECRET")
if not secret:
logger.error("MCP_SECRET is not configured on the server.")
raise HTTPException(status_code=500, detail="MCP_SECRET not configured")
# Perplexity sends 'api-key', others might send 'X-MCP-Secret' or 'Authorization: Bearer ...'
token = request.headers.get("api-key")
if not token:
token = request.headers.get("X-MCP-Secret")
if not token:
auth_header = request.headers.get("Authorization", "")
if auth_header.startswith("Bearer "):
token = auth_header[7:]
if not token or not secrets.compare_digest(token.encode(), secret.encode()):
raise HTTPException(status_code=401, detail="Invalid or missing API key")
# ---------------------------------------------------------------------------
# Backend Agent helpers (for kall VIDERE til osvauco-agent)
# ---------------------------------------------------------------------------
def _get_identity_token(audience: str) -> str:
"""Fetches a Google-signed identity token for the given audience."""
# OPAX_LOCAL_DEV_MODE bypass is for local testing where metadata server is not available.
if os.environ.get("OPAX_LOCAL_DEV_MODE") == "true":
logger.warning("OPAX_LOCAL_DEV_MODE is enabled, skipping identity token fetch.")
return ""
try:
auth_req = google.auth.transport.requests.Request()
token = google.oauth2.id_token.fetch_id_token(auth_req, audience)
return token
except Exception as e:
logger.error(f"Failed to fetch identity token for audience {audience}: {e}", exc_info=True)
raise RuntimeError(f"Could not obtain identity token for service-to-service authentication.") from e
async def _agent_headers() -> dict:
"""Headers for machine-to-machine calls to osvauco-agent."""
token = await asyncio.to_thread(_get_identity_token, OSVAUCO_AGENT_URL)
headers = {
"X-Internal-Key": INTERNAL_API_KEY,
"Content-Type": "application/json",
}
if token:
headers["Authorization"] = f"Bearer {token}"
return headers
async def _agent_get(path: str) -> Any:
"""GET-kall til osvauco-agent."""
url = f"{OSVAUCO_AGENT_URL}{path}"
async with httpx.AsyncClient(timeout=30) as c:
# Headers are now generated asynchronously for each request.
request_headers = await _agent_headers()
r = await c.get(url, headers=request_headers)
r.raise_for_status()
return r.json()
async def _agent_post(path: str, body: dict) -> Any:
"""POST-kall til osvauco-agent."""
url = f"{OSVAUCO_AGENT_URL}{path}"
async with httpx.AsyncClient(timeout=45) as c:
# Headers are now generated asynchronously for each request.
request_headers = await _agent_headers()
r = await c.post(url, json=body, headers=request_headers)
r.raise_for_status()
return r.json()
# ---------------------------------------------------------------------------
# Google Workspace Helpers
# ---------------------------------------------------------------------------
def _get_workspace_service(api: str, version: str, scopes: list):
"""
Creates an authenticated Google API service object using a service account
key from Secret Manager and impersonating the admin user.
"""
try:
client = secretmanager.SecretManagerServiceClient()
secret_name = "workspace-sa-key"
version_name = f"projects/{GOOGLE_CLOUD_PROJECT}/secrets/{secret_name}/versions/latest"
response = client.access_secret_version(request={"name": version_name})
secret_payload = response.payload.data.decode("UTF-8")
secret_info = json.loads(secret_payload)
creds = service_account.Credentials.from_service_account_info(
secret_info, scopes=scopes
)
delegated_creds = creds.with_subject(WORKSPACE_ADMIN_USER)
return build(api, version, credentials=delegated_creds, cache_discovery=False)
except Exception as e:
logger.error(f"Failed to get Workspace service: {e}", exc_info=True)
raise
async def create_email_alias(p: dict) -> dict:
"""Creates a new email alias for a Google Workspace user."""
user_key = p.get("user_key")
alias_email = p.get("alias")
if not user_key or not alias_email:
raise ValueError("Missing required parameters: 'user_key' and 'alias'")
service = _get_workspace_service('admin', 'directory_v1', ['https://www.googleapis.com/auth/admin.directory.user'])
alias_body = {'alias': alias_email}
logger.info(f"Creating alias '{alias_email}' for user '{user_key}'...")
result = service.users().aliases().insert(userKey=user_key, body=alias_body).execute()
logger.info(f"Successfully created alias.")
return {"status": "success", "alias": result}
async def list_user_aliases(p: dict) -> dict:
"""Lists all aliases for a Google Workspace user."""
user_key = p.get("user_key")
if not user_key:
raise ValueError("Missing required parameter: 'user_key'")
service = _get_workspace_service('admin', 'directory_v1', ['https://www.googleapis.com/auth/admin.directory.user'])
logger.info(f"Listing aliases for user '{user_key}'...")
result = service.users().aliases().list(userKey=user_key).execute()
logger.info(f"Successfully listed aliases.")
return {"status": "success", "aliases": result.get("aliases", [])}
async def delete_email_alias(p: dict) -> dict:
"""Deletes an email alias from a Google Workspace user."""
user_key = p.get("user_key")
alias_email = p.get("alias")
if not user_key or not alias_email:
raise ValueError("Missing required parameters: 'user_key' and 'alias'")
service = _get_workspace_service('admin', 'directory_v1', ['https://www.googleapis.com/auth/admin.directory.user'])
logger.info(f"Deleting alias '{alias_email}' for user '{user_key}'...")
service.users().aliases().delete(userKey=user_key, alias=alias_email).execute()
logger.info(f"Successfully deleted alias.")
return {"status": "success", "detail": f"Alias '{alias_email}' deleted."}
async def send_email_as(p: dict) -> dict:
"""Sends an email as a user or their alias, on their behalf."""
from_address, to_address, subject, body = p.get("from"), p.get("to"), p.get("subject"), p.get("body")
if not all([from_address, to_address, subject, body]):
raise ValueError("Missing required parameters: 'from', 'to', 'subject', 'body'")
service = _get_workspace_service('gmail', 'v1', ['https://www.googleapis.com/auth/gmail.send'])
message = MIMEText(body)
message['to'], message['from'], message['subject'] = to_address, from_address, subject
encoded_message = base64.urlsafe_b64encode(message.as_bytes()).decode()
create_message = {'raw': encoded_message}
logger.info(f"Sending email from '{from_address}' to '{to_address}'...")
send_message = service.users().messages().send(userId='me', body=create_message).execute()
logger.info(f"Successfully sent email with ID: {send_message['id']}")
return {"status": "success", "message_id": send_message['id']}
async def get_workspace_user(p: dict) -> dict:
"""Gets detailed information about a single user in Google Workspace."""
user_key = p.get("user_key")
if not user_key: raise ValueError("Missing required parameter: 'user_key'")
service = _get_workspace_service('admin', 'directory_v1', ['https://www.googleapis.com/auth/admin.directory.user'])
logger.info(f"Getting info for user '{user_key}'...")
user = service.users().get(userKey=user_key).execute()
logger.info(f"Successfully retrieved user info.")
return {"name": user.get("name", {}).get("fullName"), "email": user.get("primaryEmail"), "aliases": user.get("aliases", []), "suspended": user.get("suspended"), "lastLoginTime": user.get("lastLoginTime")}
async def list_workspace_users(p: dict) -> dict:
"""Lists all users in the Google Workspace domain."""
service = _get_workspace_service('admin', 'directory_v1', ['https://www.googleapis.com/auth/admin.directory.user'])
logger.info("Listing all workspace users...")
result = service.users().list(domain='vauco.no', maxResults=50, orderBy='email').execute()
users = result.get('users', [])
logger.info(f"Found {len(users)} users.")
return {"users": [{"email": user.get("primaryEmail"), "name": user.get("name", {}).get("fullName")} for user in users]}
async def get_emails(p: dict) -> dict:
"""Searches a user's mailbox and retrieves a list of emails."""
user_key, query, max_results = p.get('user_key'), p.get('query', ''), p.get('max_results', 10)
if not user_key: raise ValueError("Missing required parameter: 'user_key'")
service = _get_workspace_service('gmail', 'v1', ['https://www.googleapis.com/auth/gmail.modify'])
logger.info(f"Searching emails for user '{user_key}' with query '{query}'...")
list_result = service.users().messages().list(userId=user_key, q=query, maxResults=max_results).execute()
messages = list_result.get('messages', [])
if not messages: return {"messages": []}
email_details = []
for msg in messages:
detail = service.users().messages().get(userId=user_key, id=msg['id'], format='metadata', metadataHeaders=['subject', 'from']).execute()
headers = {h['name']: h['value'] for h in detail['payload']['headers']}
email_details.append({"id": msg['id'], "snippet": detail.get('snippet'), "subject": headers.get('Subject'), "from": headers.get('From')})
logger.info(f"Successfully retrieved {len(email_details)} emails.")
return {"messages": email_details}
async def list_calendar_events(p: dict) -> dict:
"""Lists events from a user's calendar."""
user_key, days_ahead = p.get('user_key'), p.get('days_ahead', 7)
if not user_key: raise ValueError("Missing required parameter: 'user_key'")
service = _get_workspace_service('calendar', 'v3', ['https://www.googleapis.com/auth/calendar'])
logger.info(f"Fetching calendar events for '{user_key}' for the next {days_ahead} days...")
now, time_min = datetime.now(timezone.utc), (datetime.now(timezone.utc)).isoformat()
time_max = (now + timedelta(days=days_ahead)).isoformat()
events_result = service.events().list(calendarId=user_key, timeMin=time_min, timeMax=time_max, singleEvents=True, orderBy='startTime').execute()
events = events_result.get('items', [])
formatted_events = [{"summary": event.get('summary'), "start": event.get('start', {}).get('dateTime', event.get('start', {}).get('date')), "end": event.get('end', {}).get('dateTime', event.get('end', {}).get('date')), "organizer": event.get('organizer', {}).get('email')} for event in events]
logger.info(f"Found {len(formatted_events)} events.")
return {"events": formatted_events}
async def trigger_build(p: dict) -> dict:
"""Triggers a Cloud Build manually by its trigger ID."""
from google.cloud.devtools import cloudbuild_v1
project_id = GOOGLE_CLOUD_PROJECT
trigger_id = p.get("trigger_id")
branch = p.get("branch", "main")
substitutions = p.get("substitutions", {})
if not trigger_id:
raise ValueError("Missing required parameter: 'trigger_id'")
try:
client = cloudbuild_v1.CloudBuildClient()
source = cloudbuild_v1.RepoSource(branch_name=branch, substitutions=substitutions)
response = client.run_build_trigger(
project_id=project_id,
trigger_id=trigger_id,
source=source,
)
build_id = response.metadata.build.id
logger.info(f"Successfully triggered Cloud Build. Build ID: {build_id}")
return {"status": "success", "build_id": build_id}
except Exception as e:
logger.error(f"Failed to trigger Cloud Build: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=f"Failed to trigger Cloud Build: {e}")
# ---------------------------------------------------------------------------
# Ollama helpers — direkte mot emma-gpu-vm
# ---------------------------------------------------------------------------
async def _ollama_chat(model: str, prompt: str, system: str = "", history: Optional[List[Dict[str, str]]] = None) -> dict:
# ... (beholdt uendret)
messages = []
if system:
messages.append({"role": "system", "content": system})
if history:
messages.extend(history)
messages.append({"role": "user", "content": prompt})
payload = {
"model": model,
"messages": messages,
"stream": False
}
try:
async with httpx.AsyncClient(timeout=120) as c:
r = await c.post(f"{OLLAMA_BASE_URL}/api/chat", json=payload)
r.raise_for_status()
data = r.json()
return { "model": model, "response": data.get("message", {}).get("content", ""), "done": data.get("done", True), "total_duration_ms": round(data.get("total_duration", 0) / 1e6) }
except Exception as e:
logger.error(f"[OLLAMA ERROR] model={model} {type(e).__name__}: {e}")
raise
canonical_emma = CanonicalEmma(
chat_function=_ollama_chat,
model="gemma3:4b"
)
async def _ollama_models() -> list:
async with httpx.AsyncClient(timeout=10) as c:
r = await c.get(f"{OLLAMA_BASE_URL}/api/tags")
r.raise_for_status()
return r.json().get("models", [])
# ---------------------------------------------------------------------------
# Gitea helpers
# ---------------------------------------------------------------------------
def _gitea_headers() -> dict:
return {"Authorization": f"token {GITEA_TOKEN}", "Content-Type": "application/json", "Accept": "application/json"}
async def _gitea_get(path: str) -> Any:
# ... (beholdt uendret)
async with httpx.AsyncClient(timeout=30) as c:
r = await c.get(f"{GITEA_URL}/api/v1{path}", headers=_gitea_headers())
r.raise_for_status()
return r.json()
async def _gitea_post(path: str, body: dict) -> Any:
# ... (beholdt uendret)
async with httpx.AsyncClient(timeout=30) as c:
r = await c.post(f"{GITEA_URL}/api/v1{path}", json=body, headers=_gitea_headers())
r.raise_for_status()
return r.json()
async def _gitea_put(path: str, body: dict) -> Any:
# ... (beholdt uendret)
async with httpx.AsyncClient(timeout=30) as c:
r = await c.put(f"{GITEA_URL}/api/v1{path}", json=body, headers=_gitea_headers())
r.raise_for_status()
return r.json()
# ---------------------------------------------------------------------------
# Phase 2C: Local-only, identity-gated helpers for call_emma
# ---------------------------------------------------------------------------
def is_call_emma_local_mode_enabled() -> bool:
"""Checks if the server is configured to run call_emma in local test mode."""
return (
os.getenv("ENABLE_CALL_EMMA_TOOL", "").lower() == "true"
and bool(os.getenv("FIRESTORE_EMULATOR_HOST"))
)
def get_emma_store() -> MemoryStore:
"""Factory to get the appropriate memory store for call_emma."""
if not is_call_emma_local_mode_enabled():
raise RuntimeError(
"call_emma is enabled only for local Firestore emulator mode."
)
project_id = os.getenv("GOOGLE_CLOUD_PROJECT", "opax-phase2c-local")
return FirestoreMemoryStore(project_id=project_id)
def get_call_emma_caller_context() -> CallerContext:
"""
Gets the caller context for call_emma, using a server-injected test identity.
This is for local/test environments only.
"""
if not is_call_emma_local_mode_enabled():
raise RuntimeError(
"Test caller context is available only in local emulator mode."
)
test_caller_id = os.getenv("OPAX_EMMA_TEST_CALLER_ID")
if not test_caller_id:
raise ValueError(
"Configuration error: OPAX_EMMA_TEST_CALLER_ID must be set "
"for local emulator tests."
)
return derive_caller_context(test_caller_id)
async def _local_deterministic_emma_chat(model: str, prompt: str, system: str = "") -> dict:
"""A local, deterministic, non-network chat function for testing."""
return {
"model": "local-test-model",
"response": "This is a deterministic local reply.",
"done": True,
"total_duration_ms": 1,
}
# ---------------------------------------------------------------------------
# Tool implementations (refaktorert til å bruke _agent_get/_agent_post)
# ---------------------------------------------------------------------------
async def get_billing_summary(p): return await _agent_get("/billing/summary")
async def get_billing_forecast(p): return await _agent_get("/billing/forecast")
async def get_billing_credits(p): return await _agent_get("/billing/credits")
async def get_billing_anomalies(p): return await _agent_get("/billing/anomalies")
async def get_billing_history(p): return await _agent_get("/billing/history")
async def get_billing_budget(p): return await _agent_get("/billing/budget")
async def get_billing_tokens_by_module(p): return await _agent_get("/telemetry/history")
async def set_billing_budget(p): return await _agent_post("/billing/budget", {"budget": p.get("amount", 500)})
async def create_invite(p):
return await _agent_post("/onboard/invite", {
"company": p.get("company", p.get("name", "")), "email": p.get("email"), "tier": p.get("tier", "starter"),
})
async def list_customers(p): return await _agent_get("/state")
async def send_webhook(p):
return await _agent_post("/notify/webhook", {
"url": p.get("url", ""), "event": p.get("event", "custom"),
"title": p.get("title", "OPAX varsel"), "body": p.get("message", p.get("body", "")),
})
async def send_email(p): return await _agent_post("/notify/email", {"to": p.get("to"), "subject": p.get("subject"), "event": p.get("event", "digest")})
async def send_sms(p): return await _agent_post("/notify/sms", {"to": p.get("to"), "body": p.get("message", p.get("body", "")), "event": p.get("event", "spike"), "tier": p.get("tier", "guard")})
async def get_notify_channels(p): return await _agent_get("/notify/channels")
async def get_health(p): return await _agent_get("/health")
async def get_build_status(p): return await _agent_get("/opax/build-status")
async def get_state(p): return await _agent_get("/state")
async def get_telemetry(p): return await _agent_get("/telemetry/history")
async def run_terminal(p): return await _agent_post("/terminal/exec", {"cmd": p.get("command", p.get("cmd", "help"))})
async def tui_command(p): return await _agent_post("/tui-command", p)
# --- Uendrede funksjoner ---
async def run_jason(p):
"""Kaller /run på osvauco-agent, som nå har sin egen JASON_BACKEND-logikk."""
return await _agent_post("/run", {"message": p.get("prompt", p.get("message", "")), "user_id": p.get("user_id", "opax"), "session_id": p.get("session_id", "mcp"), "mode": p.get("mode", "light")})
def _normalize_emma_history(value):
if not isinstance(value, list):
return []
normalized = []
for entry in value:
if not isinstance(entry, dict):
continue
role = entry.get("role")
content = entry.get("content")
if role not in {"user", "assistant"}:
continue
if not isinstance(content, str):
continue
content = content.strip()
if not content:
continue
normalized.append({"role": role, "content": content})
return normalized[-40:]
async def run_emma(p: dict) -> dict:
"""Kaller den kanoniske Emma-agenten med en prompt."""
capability_system_context = build_capability_system_context()
return await canonical_emma.run(
prompt=p.get("prompt", p.get("message", "")),
history=_normalize_emma_history(p.get("history")),
system_context=capability_system_context,
)
async def run_emma_fast(p): return await _ollama_chat(EMMA_FAST_MODEL, p.get("prompt", p.get("message", "")), p.get("system", "Du er en rask og konsis AI-assistent..."))
async def run_qwen(p): return await _ollama_chat(EMMA_LIGHT_MODEL, p.get("prompt", p.get("message", "")))
async def list_emma_models(p): return await _ollama_models()
async def list_gitea_repositories(p: dict) -> dict:
"""Wrapper for the read-only, metadata-only repository catalog."""
if p:
raise ValueError("list_gitea_repositories does not accept any arguments.")
return await list_allowed_namespace_repositories()
async def list_commits(p): return await _gitea_get(f"/repos/{p.get('repo', GITEA_REPO)}/commits?limit={p.get('limit', 10)}")
async def get_file(p: dict) -> dict:
"""Wrapper for the hardened get_file_content handler."""
# Note: server_repo_id (GITEA_REPO) is validated inside the handler now.
return await handle_get_file_content(p, GITEA_REPO)
async def list_open_issues(p): return await _gitea_get(f"/repos/{p.get('repo', GITEA_REPO)}/issues?state=open&limit=20")
async def list_repo_files(p: dict) -> dict:
"""Wrapper for handle_list_repo_files that enforces server-side constraints."""
if set(p.keys()) - {"path"}:
raise ValueError("Unsupported list_repo_files input field.")
return await handle_list_repo_files({"path": p.get("path")}, GITEA_REPO)
async def create_issue(p): return await _gitea_post(f"/repos/{p.get('repo', GITEA_REPO)}/issues", {"title": p.get("title"), "body": p.get("body", "")})
async def push_file(p):
# ... (beholdt uendret)
repo, path = p.get("repo", GITEA_REPO), p.get("path")
content = base64.b64encode(p.get("content", "").encode()).decode()
body = {"message": p.get("message", f"chore: update {path} via opax-mcp"), "content": content, "branch": p.get("branch", "main")}
if not p.get("sha"):
try:
existing = await _gitea_get(f"/repos/{repo}/contents/{path}?ref={p.get('branch','main')}")
body["sha"] = existing["sha"]
except Exception: pass
else: body["sha"] = p["sha"]
if "sha" in body: return await _gitea_put(f"/repos/{repo}/contents/{path}", body)
return await _gitea_post(f"/repos/{repo}/contents/{path}", body)
_SHA_RE = re.compile(r"^[0-9a-f]{40}$")
async def fetch_and_normalize_source(p: dict) -> dict:
"""
Resolves a service and Git ref, fetches the source archive, normalizes it,
and returns deterministic build metadata.
"""
service_key = p.get("service_key")
ref = p.get("ref")
if not isinstance(service_key, str) or not service_key:
raise ValueError("Missing or invalid 'service_key'")
if not isinstance(ref, str) or not ref:
raise ValueError("Missing or invalid 'ref'")
target_policy = get_deployment_target(service_key)
repo_id = target_policy["repository"]
required_paths = target_policy["required_source_paths"]
requested_ref_lower = ref.lower()
if _SHA_RE.fullmatch(requested_ref_lower):
resolved_commit_sha = requested_ref_lower
else:
resolved_commit_sha = await resolve_branch_to_commit_sha(
branch_name=ref,
repo_id=repo_id,
gitea_url=GITEA_URL,
)
archive_bytes = await download_repo_archive(
commit_sha=resolved_commit_sha,
repo_id=repo_id,
gitea_url=GITEA_URL,
)
_normalized_bytes, manifest = normalize_gitea_archive(
archive_bytes=archive_bytes,
required_paths=required_paths,
)
return {
"service_key": service_key,
"repository": repo_id,
"requested_ref": ref,
"resolved_commit_sha": resolved_commit_sha,
"sha256": manifest["sha256"],
"source_bytes": manifest["source_bytes"],
"wrapper_directory_stripped": manifest["wrapper_directory_stripped"],
}
# ---------------------------------------------------------------------------
# Memory Bank & Deployment Tools
# ---------------------------------------------------------------------------
PROJECT_DIR = os.path.abspath("project")
PROD_BUILD_TRIGGER_ID = "cf42cbc7-a5e8-4724-a569-dbaac1b46a0b"
def _resolve_memory_path(file_name: str) -> str:
"""Resolves a file path and ensures it is within the project/ directory."""
if not file_name or '..' in file_name or file_name.startswith('/'):
raise ValueError("Invalid file name.")
abs_path = os.path.abspath(os.path.join(PROJECT_DIR, file_name))
if not abs_path.startswith(PROJECT_DIR):
raise ValueError("Security violation: Path is outside the allowed memory bank directory.")
return abs_path
async def read_memory_bank(p: dict) -> dict:
"""Reads content from a file in the memory bank (project/ directory)."""
file_name = p.get("file_name")
if not file_name:
raise ValueError("Missing required parameter: 'file_name'")
file_path = _resolve_memory_path(file_name)
try:
with open(file_path, 'r', encoding='utf-8') as f:
content = f.read()
return {"file_path": file_name, "content": content}
except FileNotFoundError:
raise ValueError(f"File not found: {file_name}")
except Exception as e:
logger.error(f"Error reading memory bank file {file_name}: {e}", exc_info=True)
raise
async def write_memory_bank(p: dict) -> dict:
"""Writes or overwrites content to a file in the memory bank (project/ directory)."""
file_name = p.get("file_name")
content = p.get("content")
if not file_name or content is None:
raise ValueError("Missing required parameters: 'file_name' and 'content'")
file_path = _resolve_memory_path(file_name)
try:
os.makedirs(PROJECT_DIR, exist_ok=True)
with open(file_path, 'w', encoding='utf-8') as f:
f.write(content)
logger.info(f"Successfully wrote to memory bank file: {file_name}")
return {"status": "success", "file_path": file_name}
except Exception as e:
logger.error(f"Error writing to memory bank file {file_name}: {e}", exc_info=True)
raise
async def build_and_deploy_service(p: dict) -> dict:
"""
Triggers the production build and deployment pipeline for the main service.
This action requires manual approval in the Google Cloud Console.
"""
branch_to_deploy = p.get("branch", "main")
logger.info(f"Triggering production deployment for branch: {branch_to_deploy}")
result = await trigger_build({
"trigger_id": PROD_BUILD_TRIGGER_ID,
"branch": branch_to_deploy,
"substitutions": p.get("substitutions", {})
})
return result
async def handle_call_emma(
arguments: Dict[str, Any],
caller: CallerContext,
store: MemoryStore,
) -> CallEmmaResponse:
"""
Handles the call_emma tool, managing conversation state, memory, and auditing
for a single turn of conversation with the canonical Emma agent.
"""
correlation_id = str(uuid.uuid4())
conversation_id = None
status = "failure"
error_type = None
model_name = None
history_message_count = 0
memory_count = 0
proposed_action_count = 0
try:
request = CallEmmaRequest(**arguments)
conversation_id = request.conversation_id
if conversation_id:
conversation = await store.get_conversation(conversation_id, caller.owner_id)
if not conversation:
raise ValueError(f"Conversation {conversation_id} not found or access denied.")
else:
conversation = await store.create_conversation(
owner_id=caller.owner_id,
workspace_id=caller.workspace_id,
created_by=caller.caller_id
)
conversation_id = conversation.conversation_id
user_message = {
"message_id": str(uuid.uuid4()),
"conversation_id": conversation_id,
"owner_id": caller.owner_id,
"workspace_id": caller.workspace_id,
"role": "user",
"content": request.prompt,
"created_at": datetime.now(timezone.utc),
"caller_id": caller.caller_id,
"correlation_id": correlation_id,
}
await store.append_message(conversation_id, user_message)
history = await store.list_messages(conversation_id, limit=20)
history_message_count = len(history)
# Per rule #10, recall_memory is not called as it requires an embedding.
recalled_memories = []
memory_count = 0
if is_call_emma_local_mode_enabled():
local_emma = CanonicalEmma(
chat_function=_local_deterministic_emma_chat,
model="local-test-model"
)
emma_for_call = local_emma
else:
emma_for_call = canonical_emma
emma_response = await emma_for_call.run(
prompt=request.prompt,
history=history
)
model_name = emma_response.get("model")
assistant_message_content = emma_response.get("response", "")
assistant_message = {
"message_id": str(uuid.uuid4()),
"conversation_id": conversation_id,
"owner_id": caller.owner_id,
"workspace_id": caller.workspace_id,
"role": "assistant",
"content": assistant_message_content,
"created_at": datetime.now(timezone.utc),
"caller_id": "agent:emma",
"correlation_id": correlation_id,
"model": model_name,
}
await store.append_message(conversation_id, assistant_message)
status = "success"
# Proposed actions are not implemented yet
proposed_actions = []
proposed_action_count = len(proposed_actions)
return CallEmmaResponse(
reply_text=assistant_message_content,
model=model_name or "unknown",
conversation_id=conversation_id,
correlation_id=correlation_id,
proposed_actions=proposed_actions,
)
except Exception as e:
error_type = type(e).__name__
logger.error(
"[handle_call_emma ERROR] correlation_id=%s error_type=%s",
correlation_id,
error_type,
)
raise RuntimeError(
"An internal error occurred in handle_call_emma. "
f"Correlation ID: {correlation_id}"
) from None
finally:
try:
event_type = "TOOL_CALL" if status == "success" else "TOOL_CALL_FAILED"
audit_details = {
"tool_name": "call_emma",
"conversation_id": conversation_id,
"correlation_id": correlation_id,
"status": status,
"model": model_name,
"history_message_count": history_message_count,
"memory_count": memory_count,
"proposed_action_count": proposed_action_count,
"error_type": error_type,
}
audit_event = EmmaAuditEvent(
event_type=event_type,
caller_id=caller.caller_id,
details={k: v for k, v in audit_details.items() if v is not None},
workspace_id=caller.workspace_id,
)
await store.record_audit_event(audit_event)
except Exception as audit_e:
logger.error(
"[handle_call_emma AUDIT FAILED] correlation_id=%s audit_error_type=%s",
correlation_id,
type(audit_e).__name__,
)
# ---------------------------------------------------------------------------
# Tool registry + MCP schema
# ---------------------------------------------------------------------------
TOOLS = {
# Billing
"get_billing_summary": (get_billing_summary, "Hent billing-sammendrag for OPAX", {}),
"get_billing_forecast": (
get_billing_forecast,
"Get billing forecast",
{},
),
"get_billing_credits": (
get_billing_credits,
"Get available billing credits",
{},
),
"get_billing_anomalies": (
get_billing_anomalies,
"Get billing anomalies",
{},
),
"get_billing_history": (
get_billing_history,
"Get billing history",
{},
),
"get_billing_budget": (
get_billing_budget,
"Get the configured billing budget",
{},
),
"get_telemetry": (
get_billing_tokens_by_module,
"Get token and telemetry history by module",
{},
),
"set_billing_budget": (set_billing_budget, "Sett månedlig budsjett", {"type":"object","properties":{"amount":{"type":"number"}}}),
# Onboarding
"create_invite": (create_invite, "Inviter ny kunde til OPAX", {"type":"object","properties":{"company":{"type":"string"},"email":{"type":"string"},"tier":{"type":"string"}},"required":["email"]}),
# Notifications
"send_webhook": (send_webhook, "Send webhook-varsling", {"type":"object","properties":{"url":{"type":"string"},"message":{"type":"string"}}}),
"send_email": (send_email, "Send e-post via ekstern tjeneste", {"type":"object","properties":{"to":{"type":"string"},"subject":{"type":"string"}},"required":["to","subject"]}),
"send_sms": (send_sms, "Send SMS", {"type":"object","properties":{"to":{"type":"string"},"message":{"type":"string"}},"required":["to","message"]}),
"get_notify_channels": (
get_notify_channels,
"List configured notification channels",
{},
),
# AI Agents
"run_jason": (run_jason, "Kjør Jason-agenten med en prompt", {"type":"object","properties":{"prompt":{"type":"string"},"mode":{"type":"string"}},"required":["prompt"]}),
"run_emma": (run_emma, "Emma Vauger (gemma3:4b) — primær lokal AI", {"type": "object", "properties": {"prompt": {"type": "string"}, "history": {"type": "array", "items": {"type": "object", "properties": {"role": {"type": "string", "enum": ["user", "assistant"]}, "content": {"type": "string"}}}}}}),
# Local models — read-only discovery
"list_emma_models": (
list_emma_models,
"List locally available OSVx model names",
{},
),
# System & Ops
"get_health": (get_health, "Hent helsestatus for OPAX", {}),
"get_build_status": (get_build_status, "Hent siste build-status", {}),
"trigger_build": (trigger_build, "Trigger en Cloud Build manuelt", {"type":"object","properties":{"repo":{"type":"string"},"branch":{"type":"string"},"config":{"type":"string"}}}),
"get_state": (get_state, "Hent platform-tilstand", {}),
"list_customers": (list_customers, "List alle kunder (alias for get_state)", {}),
"run_terminal": (run_terminal, "Kjør terminalkommando på VM", {"type":"object","properties":{"command":{"type":"string"}}}),
# Gitea / VCS
"fetch_and_normalize_source": (
fetch_and_normalize_source,
"Fetches and normalizes a repository source archive, returning build metadata.",
{
"type": "object",
"properties": {
"service_key": {"type": "string"},
"ref": {"type": "string"},
},
"required": ["service_key", "ref"],
},
),
"list_gitea_repositories": (list_gitea_repositories, "Lists repositories in the approved 'chris' namespace. Listing does not grant read, build, or deploy authority.", {}),
"list_commits": (list_commits, "List siste commits i Gitea-repo", {}),
"get_file": (get_file, "Hent fil fra Gitea-repo", {"type":"object","properties":{"path":{"type":"string"}, "ref":{"type":"string"}, "repo":{"type":"string"}},"required":["path", "ref"]}),
"list_repo_files": (list_repo_files, "Lists files and directories in a Gitea repository path.", {"type": "object", "properties": {"path": {"type": "string"}}}),
"propose_gitea_change": (
propose_gitea_change,
"Create a pending, one-file Gitea change plan for review; does not apply or commit changes.",
{
"type": "object",
"properties": {
"repo": {"type": "string"},
"branch": {"type": "string"},
"path": {"type": "string"},
"new_content": {"type": "string"},
"base_sha": {"type": "string"},
"commit_message": {"type": "string"}
},
"required": ["repo", "branch", "path", "new_content", "base_sha", "commit_message"]
},
),
# Google Workspace
"create_email_alias": (create_email_alias, "Opprett et nytt e-postalias", {"type":"object","properties":{"user_key":{"type":"string"},"alias":{"type":"string"}},"required":["user_key","alias"]}),
"list_user_aliases": (list_user_aliases, "List en brukers e-postaliaser", {"type":"object","properties":{"user_key":{"type":"string"}},"required":["user_key"]}),
"delete_email_alias": (delete_email_alias, "Slett et e-postalias", {"type":"object","properties":{"user_key":{"type":"string"},"alias":{"type":"string"}},"required":["user_key","alias"]}),
"send_email_as": (send_email_as, "Send en e-post på vegne av en bruker/alias", {"type":"object","properties":{"from":{"type":"string"},"to":{"type":"string"},"subject":{"type":"string"},"body":{"type":"string"}},"required":["from","to","subject","body"]}),
"get_workspace_user": (get_workspace_user, "Hent detaljer for en Workspace-bruker", {"type":"object","properties":{"user_key":{"type":"string"}},"required":["user_key"]}),
"list_workspace_users": (list_workspace_users, "List alle brukere i Workspace-domenet", {}),
"get_emails": (get_emails, "Søk i en brukers e-post", {"type":"object","properties":{"user_key":{"type":"string"},"query":{"type":"string"}},"required":["user_key"]}),
"list_calendar_events": (list_calendar_events, "List en brukers kalenderhendelser", {"type":"object","properties":{"user_key":{"type":"string"},"days_ahead":{"type":"integer"}},"required":["user_key"]}),
# Memory Bank
"read_memory_bank": (read_memory_bank, "Reads a file from the project memory bank.", {"type": "object", "properties": {"file_name": {"type": "string"}}, "required": ["file_name"]}),
"write_memory_bank": (write_memory_bank, "Writes a file to the project memory bank.", {"type": "object", "properties": {"file_name": {"type": "string"}, "content": {"type": "string"}}, "required": ["file_name", "content"]}),
"build_and_deploy_service": (build_and_deploy_service, "Triggers the production build & deployment pipeline.", {"type": "object", "properties": {"branch": {"type": "string"}}, "required": []}),
}
if is_call_emma_local_mode_enabled():
TOOLS["call_emma"] = (handle_call_emma, "Calls the canonical Emma agent with managed persistence (EMULATOR ONLY)", {"type": "object", "properties": {"prompt": {"type": "string"}}, "required": ["prompt"]})
if provision_new_mcp_module:
TOOLS["provision_new_mcp_module"] = (provision_new_mcp_module, "Provisions a new MCP module", {"type":"object","properties":{"module_name":{"type":"string"},"tool_name":{"type":"string"},"tool_spec":{"type":"object"}},"required":["module_name","tool_name","tool_spec"]})
# ---------------------------------------------------------------------------
# JSON-RPC Endpoint
# ---------------------------------------------------------------------------
def _jsonrpc_ok(req_id, result):
return {"jsonrpc": "2.0", "id": req_id, "result": result}
def _jsonrpc_err(req_id, code, message):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": code, "message": message}}
@app.post("/")
async def mcp_handler(request: Request):
await _verify_auth(request)
try:
body = await request.json()
except Exception: return JSONResponse(_jsonrpc_err(None, -32700, "Parse error"), status_code=400)
method, req_id, params = body.get("method", ""), body.get("id"), body.get("params", {})
if method == "initialize":
return JSONResponse(_jsonrpc_ok(req_id, {"protocolVersion": "2024-11-05", "capabilities": {"tools": {}},"serverInfo": {"name": "opax-mcp", "version": "3.6.0"}})) # BUMPED
if method == "tools/list":
return JSONResponse(_jsonrpc_ok(req_id, {"tools": [
{"name": name, "description": desc, "inputSchema": schema if schema else {"type": "object", "properties": {}}}
for name, (_, desc, schema) in TOOLS.items()
]}))
if method == "tools/call":
tool_name, tool_args = params.get("name") or params.get("tool"), params.get("arguments", params.get("params", {}))
entry = TOOLS.get(tool_name)
if not entry: return JSONResponse(_jsonrpc_err(req_id, -32601, f"Unknown tool: {tool_name}"))
handler, _, _ = entry
try:
if tool_name == "call_emma" and is_call_emma_local_mode_enabled():
caller = get_call_emma_caller_context()
store = get_emma_store()
result = await handler(tool_args, caller, store)
else:
result = await handler(tool_args)
# For call_emma, the result is already the final response object.
# For other tools, we wrap it.
if tool_name == "call_emma":
return JSONResponse(_jsonrpc_ok(req_id, result.model_dump(mode='json')))
else:
return JSONResponse(_jsonrpc_ok(req_id, {"content": [{"type": "text", "text": json.dumps(result, ensure_ascii=False)}]}))
except Exception as e:
if tool_name == "call_emma":
error_type = type(e).__name__
logger.error(
"[MCP HANDLER ERROR] tool=%s error_type=%s",
tool_name,
error_type,
)
return JSONResponse(
_jsonrpc_err(
req_id,
-32000,
f"Server error in tool 'call_emma': {error_type}",
)
)
logger.error(f"[MCP HANDLER ERROR] tool={tool_name} {type(e).__name__}: {e}", exc_info=True)
return JSONResponse(_jsonrpc_err(req_id, -32000, str(e)))
if method.startswith("notifications/"): return JSONResponse(status_code=202, content={})
return JSONResponse(_jsonrpc_err(req_id, -32601, f"Method not found: {method}"), status_code=404)
# ---------------------------------------------------------------------------
# Health (keepalive)
# ---------------------------------------------------------------------------
@app.get("/health")
async def health():
return {"status": "ok", "service": "opax-mcp", "version": "3.6.0"}