99 lines
3.0 KiB
Python
99 lines
3.0 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
ml/telemetry.py — Strukturert telemetri-innsamling for OPAX-agenter.
|
|
|
|
Formål:
|
|
Loggfører hvert agent-kall til JSONL-filer.
|
|
Disse filene er rådata for Fase ML-2 (inkrementell læring)
|
|
og Fase ML-3 (XGBoost feature importance).
|
|
|
|
Felter per record:
|
|
ts — ISO 8601 UTC tidsstempel
|
|
agent — agent-ID (f.eks. "billing-agent")
|
|
model — LLM-modell brukt (f.eks. "gemini-2.0-flash")
|
|
mode — "standard" eller "heavy"
|
|
duration_s — latens i sekunder
|
|
success — bool
|
|
error — feilmelding eller None
|
|
input_tokens — proxy: len(str(input))
|
|
output_tokens — proxy: len(str(output))
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import time
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any, Dict, Optional
|
|
|
|
# Skriv til Cloud Run-vennlig sti — overstyr med TELEMETRY_DIR env var
|
|
_DEFAULT_DIR = os.environ.get("TELEMETRY_DIR", "/tmp/vauco_telemetry")
|
|
|
|
|
|
def log_agent_call(
|
|
agent_id: str,
|
|
input_payload: Dict[str, Any],
|
|
output: Any,
|
|
model_used: str,
|
|
mode: str,
|
|
duration_s: float,
|
|
success: bool,
|
|
error: Optional[str] = None,
|
|
telemetry_dir: Optional[str] = None,
|
|
) -> None:
|
|
"""
|
|
Logg ett agent-kall til JSONL-fil (én fil per dato).
|
|
|
|
Filnavnformat: YYYY-MM-DD.jsonl
|
|
Trygt for parallell skriving (én linje per kall, atomisk append).
|
|
"""
|
|
base = Path(telemetry_dir or _DEFAULT_DIR)
|
|
base.mkdir(parents=True, exist_ok=True)
|
|
|
|
record = {
|
|
"ts": datetime.now(timezone.utc).isoformat(),
|
|
"agent": agent_id,
|
|
"model": model_used,
|
|
"mode": mode,
|
|
"duration_s": round(duration_s, 3),
|
|
"success": success,
|
|
"error": error,
|
|
"input_tokens": len(str(input_payload)),
|
|
"output_tokens": len(str(output)),
|
|
}
|
|
|
|
log_file = base / f"{datetime.now(timezone.utc).strftime('%Y-%m-%d')}.jsonl"
|
|
with open(log_file, "a", encoding="utf-8") as fh:
|
|
fh.write(json.dumps(record, ensure_ascii=False) + "\n")
|
|
|
|
|
|
def log_dag_execution(
|
|
dag_id: str,
|
|
agent_results: list,
|
|
total_duration_s: float,
|
|
telemetry_dir: Optional[str] = None,
|
|
) -> None:
|
|
"""
|
|
Logg ett DAG-kjøring (aggregert over alle agenter i DAGen).
|
|
Brukes av execute_dag() for å gi oversikt over parallelle kjøringer.
|
|
"""
|
|
base = Path(telemetry_dir or _DEFAULT_DIR)
|
|
base.mkdir(parents=True, exist_ok=True)
|
|
|
|
record = {
|
|
"ts": datetime.now(timezone.utc).isoformat(),
|
|
"type": "dag_execution",
|
|
"dag_id": dag_id,
|
|
"agent_count": len(agent_results),
|
|
"success_count": sum(1 for r in agent_results if r.get("success")),
|
|
"total_duration_s": round(total_duration_s, 3),
|
|
"agents": [
|
|
{"agent": r.get("agent"), "duration_s": r.get("duration_s"), "success": r.get("success")}
|
|
for r in agent_results
|
|
],
|
|
}
|
|
|
|
log_file = base / f"{datetime.now(timezone.utc).strftime('%Y-%m-%d')}.jsonl"
|
|
with open(log_file, "a", encoding="utf-8") as fh:
|
|
fh.write(json.dumps(record, ensure_ascii=False) + "\n")
|