feat(ml): add ML layer — Dask DAG, Parameter Server state, telemetry, AGENTS.md [ML-1]
This commit is contained in:
parent
a163ca8c3e
commit
d5ce0227f4
91
AGENTS.md
Normal file
91
AGENTS.md
Normal file
|
|
@ -0,0 +1,91 @@
|
||||||
|
# AGENTS.md — ML-roller i Vauco OS
|
||||||
|
|
||||||
|
**Sist oppdatert:** 2026-05-26
|
||||||
|
**Formål:** Definerer hver agents rolle i ML-laget (Parameter Server-mønster + Dask DAG).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Oversikt
|
||||||
|
|
||||||
|
```
|
||||||
|
┌─────────────────────────────┐
|
||||||
|
│ OPAX Orchestrator │
|
||||||
|
│ (PS Server-node / DAG hub) │
|
||||||
|
└──────────┬──────────────────┘
|
||||||
|
│
|
||||||
|
┌──────────┬──────────┼──────────┬──────────┐
|
||||||
|
▼ ▼ ▼ ▼ ▼
|
||||||
|
billing- rag-agent memory- tui- multi-
|
||||||
|
agent agent bridge agent
|
||||||
|
(worker) (worker) (worker) (worker+ (worker)
|
||||||
|
input-src)
|
||||||
|
```
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Agenttabell
|
||||||
|
|
||||||
|
| Agent | Katalog | PS-rolle | Dask-rolle | Online learning | Heavy mode |
|
||||||
|
|---|---|---|---|---|---|
|
||||||
|
| OPAX Orchestrator | `agents/core-logic/` | **Server-node** (aggregerer state, dispatcher DAG) | DAG scheduler | Nei | Ja (koordinerer heavy sub-agenter) |
|
||||||
|
| billing-agent | `agents/core-logic/` | Worker (pusher kostdata til state store) | `dask.delayed`-task | Nei — strukturert data | Nei |
|
||||||
|
| rag-agent | `agents/rag/` | Worker (pusher retrieval-resultater) | `dask.delayed`-task | **Ja** — `partial_fit` på nye docs (ML-2) | Ja |
|
||||||
|
| memory-agent | `agents/memory/` | Worker (pusher/puller memory-oppslag) | `dask.delayed`-task | Nei — faktahenting | Nei |
|
||||||
|
| eval-agent | `agents/eval/` | Worker (pusher evalueringsmetrikker) | `dask.delayed`-task | **Ja** — kvalitetsscore-feedback (ML-2) | Nei |
|
||||||
|
| multi-agent | `agents/multi_agent/` | Worker + sub-orchestrator | DAG-node med barn | Nei | **Ja — primær heavy-mode bruker** |
|
||||||
|
| tui-bridge | `vauco-gemini-tui-bridge` | Worker + **primær input-kilde** | Streaming input til DAG | **Ja** — real-time feedback-loop (ML-2) | Ja |
|
||||||
|
| vauco-bootstrap | `vauco-bootstrap` | **Config-aggregator** (analog til PS shard) | Trigger for nye DAGer | **Ja** — ASHA-tuning (ML-2) | Nei |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## ML-fase per agent
|
||||||
|
|
||||||
|
| Agent | ML-1 (nå) | ML-2 (neste) | ML-3 (fremtidig) |
|
||||||
|
|---|---|---|---|
|
||||||
|
| OPAX Orchestrator | Kjør DAG via `execute_dag()` | Bytt heavy-mode til ML-predictor | PBT-eksperiment |
|
||||||
|
| billing-agent | `log_agent_call()` + `state.push()` | — | Feature importance-input |
|
||||||
|
| rag-agent | `log_agent_call()` + `state.push()` | `partial_fit` på nye docs | Top-K feature importance |
|
||||||
|
| memory-agent | `log_agent_call()` | — | — |
|
||||||
|
| eval-agent | `log_agent_call()` + `state.push()` | Kvalitetsscore-feedback | Feature importance-target |
|
||||||
|
| multi-agent | `log_agent_call()` | — | PBT worker |
|
||||||
|
| tui-bridge | `log_agent_call()` | `feedback_loop.py` (ML-2.1) | — |
|
||||||
|
| vauco-bootstrap | — | `hypertuner.py` ASHA (ML-2.3) | PBT vinner-konfig → Terraform |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## State store nøkkelkonvensjon
|
||||||
|
|
||||||
|
Format: `"<agent_id>:<key>"`
|
||||||
|
|
||||||
|
| Agent | Nøkkel | Eksempelverdi |
|
||||||
|
|---|---|---|
|
||||||
|
| `billing-agent` | `last_cost_usd` | `12.5` |
|
||||||
|
| `billing-agent` | `last_report_ts` | `"2026-05-26T01:00:00Z"` |
|
||||||
|
| `rag-agent` | `last_retrieval_score` | `0.87` |
|
||||||
|
| `rag-agent` | `chunks_retrieved` | `5` |
|
||||||
|
| `eval-agent` | `quality_score` | `0.91` |
|
||||||
|
| `multi-agent` | `active_sub_agents` | `["billing-agent", "rag-agent"]` |
|
||||||
|
| `tui-bridge` | `last_user_feedback` | `"positive"` |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Telemetri-felter (referanse)
|
||||||
|
|
||||||
|
Alle agenter logger via `ml.telemetry.log_agent_call()`. Feltene er:
|
||||||
|
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"ts": "2026-05-26T01:00:00.000Z",
|
||||||
|
"agent": "billing-agent",
|
||||||
|
"model": "gemini-2.0-flash",
|
||||||
|
"mode": "standard",
|
||||||
|
"duration_s": 0.312,
|
||||||
|
"success": true,
|
||||||
|
"error": null,
|
||||||
|
"input_tokens": 284,
|
||||||
|
"output_tokens": 512
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Filer lagres i: `$TELEMETRY_DIR/YYYY-MM-DD.jsonl` (default: `/tmp/vauco_telemetry/`)
|
||||||
|
For GCS-arkivering: sett `TELEMETRY_DIR=gs://...` eller bruk `10-cost-guard.sh`-cron.
|
||||||
123
docs/ml-integration.md
Normal file
123
docs/ml-integration.md
Normal file
|
|
@ -0,0 +1,123 @@
|
||||||
|
# ML-integrasjon i OSVauco
|
||||||
|
|
||||||
|
**Fase:** ML-1 implementert | ML-2 og ML-3 planlagt
|
||||||
|
**Sist oppdatert:** 2026-05-26
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Hvorfor ML i OPAX?
|
||||||
|
|
||||||
|
OPAX er et multi-agent system. ML-laget gjør agentene sterkere på tre måter:
|
||||||
|
|
||||||
|
1. **Parallell kjøring** (ML-1): Dask DAG lar agenter kjøre samtidig i stedet for sekvensielt.
|
||||||
|
2. **Adaptiv intelligens** (ML-2): Agentene lærer fra faktisk bruk — heavy-mode aktiveres basert på historikk, ikke hardkodede regler.
|
||||||
|
3. **Strategisk innsikt** (ML-3): XGBoost feature importance avslører hvilke agenter og signaler som faktisk driver vellykkede utfall.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Fase ML-1: Grunnlag (implementert)
|
||||||
|
|
||||||
|
### Arkitektur
|
||||||
|
|
||||||
|
```
|
||||||
|
main.py → build_agent_dag([agent_a, agent_b], payload)
|
||||||
|
│
|
||||||
|
Dask task graph
|
||||||
|
│
|
||||||
|
┌──────────┴──────────┐
|
||||||
|
▼ ▼
|
||||||
|
run_agent_task(a) run_agent_task(b) ← parallelt
|
||||||
|
│ │
|
||||||
|
└──────────┬──────────┘
|
||||||
|
▼
|
||||||
|
execute_dag(tasks, scheduler="threads")
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
[result_a, result_b] → AgentStateStore.push()
|
||||||
|
│
|
||||||
|
telemetry.log_agent_call()
|
||||||
|
```
|
||||||
|
|
||||||
|
### Koble ML-laget inn i main.py
|
||||||
|
|
||||||
|
```python
|
||||||
|
# Eksempel: erstatt sekvensielle kall med parallell DAG
|
||||||
|
from ml import build_agent_dag, execute_dag, get_store, log_agent_call
|
||||||
|
import time
|
||||||
|
|
||||||
|
@app.post("/run")
|
||||||
|
def run_agent(req: RunRequest):
|
||||||
|
store = get_store()
|
||||||
|
payload = {"message": req.message, "user_id": req.user_id,
|
||||||
|
"session_id": req.session_id, "mode": req.mode}
|
||||||
|
|
||||||
|
# Velg agenter basert på mode
|
||||||
|
if req.mode == "heavy":
|
||||||
|
agent_fns = [billing_agent_heavy, rag_agent, multi_agent]
|
||||||
|
else:
|
||||||
|
agent_fns = [billing_agent, rag_agent]
|
||||||
|
|
||||||
|
start = time.monotonic()
|
||||||
|
tasks = build_agent_dag(agent_fns, payload)
|
||||||
|
results = execute_dag(tasks, scheduler="threads")
|
||||||
|
total_dur = time.monotonic() - start
|
||||||
|
|
||||||
|
# Push resultater til state store
|
||||||
|
for r in results:
|
||||||
|
if r["success"]:
|
||||||
|
store.push(r["agent"], "last_result", r["result"])
|
||||||
|
store.push(r["agent"], "last_duration_s", r["duration_s"])
|
||||||
|
|
||||||
|
# Logg DAG-kjøring
|
||||||
|
from ml.telemetry import log_dag_execution
|
||||||
|
log_dag_execution(dag_id=req.session_id, agent_results=results,
|
||||||
|
total_duration_s=total_dur)
|
||||||
|
|
||||||
|
# Returner kombinert svar
|
||||||
|
return {"response": [r["result"] for r in results if r["success"]]}
|
||||||
|
```
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Fase ML-2: Adaptiv intelligens (planlagt)
|
||||||
|
|
||||||
|
| Komponent | Fil | Repo |
|
||||||
|
|---|---|---|
|
||||||
|
| Inkrementell feedback-loop | `ml/feedback_loop.py` | `vauco-gemini-tui-bridge` |
|
||||||
|
| ASHA hyperparametertuning | `ml/hypertuner.py` | `vauco-bootstrap` |
|
||||||
|
| Heavy-mode predictor | `ml/heavy_predictor.py` | `OSVauco` |
|
||||||
|
|
||||||
|
**Krav:** `dask-ml>=2024.4.0`, `ray[tune]>=2.10.0` (vauco-bootstrap only)
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Fase ML-3: Strategisk innsikt (planlagt)
|
||||||
|
|
||||||
|
| Komponent | Fil | Repo |
|
||||||
|
|---|---|---|
|
||||||
|
| XGBoost feature importance | `ml/agent_importance.py` | `OSVauco` |
|
||||||
|
| Population-Based Training | `experiments/pbt.py` | `deep-dream` |
|
||||||
|
|
||||||
|
**Krav:** `xgboost>=2.0.0` (legg til requirements.txt i Fase ML-3)
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Telemetri-ruting
|
||||||
|
|
||||||
|
| Miljø | `TELEMETRY_DIR` | Arkivering |
|
||||||
|
|---|---|---|
|
||||||
|
| Lokal dev | `/tmp/vauco_telemetry/` | Manuell |
|
||||||
|
| Cloud Run | `/tmp/vauco_telemetry/` | `10-cost-guard.sh` cron → GCS |
|
||||||
|
| Prod skala | `gs://propane-will-491900-m5-telemetry/` | Auto via GCS sink |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## Cross-repo referanser
|
||||||
|
|
||||||
|
| Repo | ML-relatert fil | Fase |
|
||||||
|
|---|---|---|
|
||||||
|
| `OSVauco` | `ml/`, `AGENTS.md`, `docs/ml-integration.md` | ML-1 |
|
||||||
|
| `vauco-gemini-tui-bridge` | `ml/feedback_loop.py` (planlagt) | ML-2 |
|
||||||
|
| `vauco-bootstrap` | `ml/hypertuner.py` (planlagt) | ML-2 |
|
||||||
|
| `vauco-os` | Konsumerer beste konfig fra state store | ML-2 |
|
||||||
|
| `deep-dream` | `experiments/pbt.py` (planlagt) | ML-3 |
|
||||||
76
ml/README.md
Normal file
76
ml/README.md
Normal file
|
|
@ -0,0 +1,76 @@
|
||||||
|
# OSVauco ML Layer
|
||||||
|
|
||||||
|
ML-laget gir OPAX-agentene parallell kjøring, delt state-kontroll og strukturert
|
||||||
|
telemetri. Implementert i tre faser.
|
||||||
|
|
||||||
|
## Arkitektur
|
||||||
|
|
||||||
|
```
|
||||||
|
OPAX /run-endepunkt
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
build_agent_dag() ← Dask DAG (task_graph.py)
|
||||||
|
│
|
||||||
|
┌────┴────┬──────────────┐
|
||||||
|
▼ ▼ ▼
|
||||||
|
agent_A agent_B ... agent_N ← kjøres parallelt
|
||||||
|
│ │ │
|
||||||
|
└────┬────┴──────────────┘
|
||||||
|
▼
|
||||||
|
AgentStateStore ← Parameter Server state (state_store.py)
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
telemetry JSONL ← rådatagrunnlag for ML-2 og ML-3
|
||||||
|
```
|
||||||
|
|
||||||
|
## Filer
|
||||||
|
|
||||||
|
| Fil | Ansvar |
|
||||||
|
|---|---|
|
||||||
|
| `task_graph.py` | Dask DAG — parallell agent-scheduling |
|
||||||
|
| `state_store.py` | Thread-safe PS-inspirert state store |
|
||||||
|
| `telemetry.py` | Strukturert JSONL-logging per agent-kall |
|
||||||
|
| `__init__.py` | Pakke-eksporter |
|
||||||
|
|
||||||
|
## Scheduler-valg
|
||||||
|
|
||||||
|
| Miljø | `scheduler`-parameter |
|
||||||
|
|---|---|
|
||||||
|
| Dev / unit-test | `'synchronous'` |
|
||||||
|
| Lokal multi-tråd | `'threads'` |
|
||||||
|
| Cloud Run (lett) | `'threads'` |
|
||||||
|
| Full skala | `'distributed'` (krever Dask cluster) |
|
||||||
|
|
||||||
|
## Fase-oversikt
|
||||||
|
|
||||||
|
| Fase | Innhold | Status |
|
||||||
|
|---|---|---|
|
||||||
|
| ML-1 | Dask DAG + PS state + telemetri | ✅ Implementert |
|
||||||
|
| ML-2 | Inkrementell læring (dask-ml) + ASHA-tuning (Ray Tune) | 🔜 Neste |
|
||||||
|
| ML-3 | XGBoost feature importance + PBT | 🔜 Fremtidig |
|
||||||
|
|
||||||
|
## Bruk
|
||||||
|
|
||||||
|
```python
|
||||||
|
from ml import build_agent_dag, execute_dag, get_store, log_agent_call
|
||||||
|
|
||||||
|
# Bygg og kjør DAG
|
||||||
|
tasks = build_agent_dag([billing_agent, rag_agent], payload)
|
||||||
|
results = execute_dag(tasks, scheduler="threads")
|
||||||
|
|
||||||
|
# State store
|
||||||
|
store = get_store()
|
||||||
|
store.push("billing-agent", "last_cost_usd", 12.5)
|
||||||
|
cost = store.pull("billing-agent", "last_cost_usd")
|
||||||
|
|
||||||
|
# Telemetri
|
||||||
|
log_agent_call(
|
||||||
|
agent_id="billing-agent",
|
||||||
|
input_payload=payload,
|
||||||
|
output=result,
|
||||||
|
model_used="gemini-2.0-flash",
|
||||||
|
mode="standard",
|
||||||
|
duration_s=0.4,
|
||||||
|
success=True,
|
||||||
|
)
|
||||||
|
```
|
||||||
14
ml/__init__.py
Normal file
14
ml/__init__.py
Normal file
|
|
@ -0,0 +1,14 @@
|
||||||
|
# OSVauco ML layer
|
||||||
|
# Fase ML-1: Dask DAG scheduling, Parameter Server-inspirert state, telemetri
|
||||||
|
from .task_graph import build_agent_dag, execute_dag, run_agent_task
|
||||||
|
from .state_store import AgentStateStore, get_store
|
||||||
|
from .telemetry import log_agent_call
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"build_agent_dag",
|
||||||
|
"execute_dag",
|
||||||
|
"run_agent_task",
|
||||||
|
"AgentStateStore",
|
||||||
|
"get_store",
|
||||||
|
"log_agent_call",
|
||||||
|
]
|
||||||
124
ml/state_store.py
Normal file
124
ml/state_store.py
Normal file
|
|
@ -0,0 +1,124 @@
|
||||||
|
#!/usr/bin/env python3
|
||||||
|
"""
|
||||||
|
ml/state_store.py — Parameter Server-inspirert state store for OPAX-agenter.
|
||||||
|
|
||||||
|
Konsept (fra ML-dokumentet):
|
||||||
|
Parameter Server (PS) skiller treningsarbeid fra modelltilstand.
|
||||||
|
Server-noder vedlikeholder global state (nøkkel-verdi-store).
|
||||||
|
Worker-noder utfører pull (les parametere) og push (send gradienter).
|
||||||
|
|
||||||
|
OPAX-mapping:
|
||||||
|
AgentStateStore = server-node (global OPAX-tilstand)
|
||||||
|
Agenter (workers) = push resultater, pull delt konteksttilstand
|
||||||
|
aggregate() = analog til gradient-aggregering på PS-server
|
||||||
|
snapshot() = analog til periodisk checkpointing
|
||||||
|
"""
|
||||||
|
|
||||||
|
import threading
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
|
|
||||||
|
class AgentStateStore:
|
||||||
|
"""
|
||||||
|
Thread-safe nøkkel-verdi-store delt av alle OPAX-agenter.
|
||||||
|
|
||||||
|
Nøkkelformat: "<agent_id>:<key>" (f.eks. "billing-agent:last_cost_usd")
|
||||||
|
Historikk: Appendonly-log for telemetri og replay.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self._store: Dict[str, Any] = {}
|
||||||
|
self._lock = threading.RLock()
|
||||||
|
self._history: List[Dict[str, Any]] = []
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------ #
|
||||||
|
# Worker-operasjoner #
|
||||||
|
# ------------------------------------------------------------------ #
|
||||||
|
|
||||||
|
def push(self, agent_id: str, key: str, value: Any) -> None:
|
||||||
|
"""
|
||||||
|
Agent pusher en oppdatering til global state.
|
||||||
|
Analog til PS worker-push av gradienter til server.
|
||||||
|
"""
|
||||||
|
with self._lock:
|
||||||
|
full_key = f"{agent_id}:{key}"
|
||||||
|
self._store[full_key] = value
|
||||||
|
self._history.append({
|
||||||
|
"op": "push",
|
||||||
|
"agent": agent_id,
|
||||||
|
"key": key,
|
||||||
|
"ts": datetime.now(timezone.utc).isoformat(),
|
||||||
|
})
|
||||||
|
|
||||||
|
def pull(self, agent_id: str, key: str) -> Optional[Any]:
|
||||||
|
"""
|
||||||
|
Agent puller gjeldende verdi fra global state.
|
||||||
|
Analog til PS worker-pull av parametere fra server.
|
||||||
|
"""
|
||||||
|
with self._lock:
|
||||||
|
return self._store.get(f"{agent_id}:{key}")
|
||||||
|
|
||||||
|
def pull_any(self, key: str) -> Optional[Any]:
|
||||||
|
"""
|
||||||
|
Pull første treff på key uavhengig av agent_id.
|
||||||
|
Nyttig for delt konfigurasjon (f.eks. "*:active_model").
|
||||||
|
"""
|
||||||
|
with self._lock:
|
||||||
|
for k, v in self._store.items():
|
||||||
|
if k.endswith(f":{key}"):
|
||||||
|
return v
|
||||||
|
return None
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------ #
|
||||||
|
# Server-operasjoner (aggregering og inspeksjon) #
|
||||||
|
# ------------------------------------------------------------------ #
|
||||||
|
|
||||||
|
def aggregate(self, key_pattern: str) -> Dict[str, Any]:
|
||||||
|
"""
|
||||||
|
Aggregerer alle verdier som matcher nøkkelmønsteret på tvers av agenter.
|
||||||
|
Analog til PS-serverens gradient-aggregering over alle workers.
|
||||||
|
|
||||||
|
Eksempel:
|
||||||
|
store.aggregate("last_duration_s")
|
||||||
|
→ {"billing-agent:last_duration_s": 0.3, "rag-agent:last_duration_s": 1.2}
|
||||||
|
"""
|
||||||
|
with self._lock:
|
||||||
|
return {
|
||||||
|
k: v
|
||||||
|
for k, v in self._store.items()
|
||||||
|
if k.endswith(f":{key_pattern}")
|
||||||
|
}
|
||||||
|
|
||||||
|
def list_agents(self) -> List[str]:
|
||||||
|
"""Returnerer unike agent-IDer som har pushet noe."""
|
||||||
|
with self._lock:
|
||||||
|
return list({k.split(":")[0] for k in self._store})
|
||||||
|
|
||||||
|
def snapshot(self) -> Dict[str, Any]:
|
||||||
|
"""
|
||||||
|
Full state-snapshot — brukes til checkpointing og debugging.
|
||||||
|
Analog til PS periodisk modell-checkpoint til distribuert filsystem.
|
||||||
|
"""
|
||||||
|
with self._lock:
|
||||||
|
return dict(self._store)
|
||||||
|
|
||||||
|
def history(self, limit: int = 100) -> List[Dict[str, Any]]:
|
||||||
|
"""Siste N operasjoner fra historikk-loggen."""
|
||||||
|
with self._lock:
|
||||||
|
return list(self._history[-limit:])
|
||||||
|
|
||||||
|
def clear(self) -> None:
|
||||||
|
"""Nullstill state (brukes i tester og teardown)."""
|
||||||
|
with self._lock:
|
||||||
|
self._store.clear()
|
||||||
|
self._history.clear()
|
||||||
|
|
||||||
|
|
||||||
|
# Singleton — deles av alle agenter i én prosess
|
||||||
|
_global_store = AgentStateStore()
|
||||||
|
|
||||||
|
|
||||||
|
def get_store() -> AgentStateStore:
|
||||||
|
"""Returnerer den globale singleton-instansen av AgentStateStore."""
|
||||||
|
return _global_store
|
||||||
90
ml/task_graph.py
Normal file
90
ml/task_graph.py
Normal file
|
|
@ -0,0 +1,90 @@
|
||||||
|
#!/usr/bin/env python3
|
||||||
|
"""
|
||||||
|
ml/task_graph.py — Dask DAG-basert agent-koordinering for OPAX.
|
||||||
|
|
||||||
|
Konsept (fra ML-dokumentet):
|
||||||
|
Dask bygger en task graph (DAG) av operasjoner og avhengigheter.
|
||||||
|
Lazy evaluation: beregninger utføres kun når .compute() kalles.
|
||||||
|
Agent-funksjoner wrappes som @dask.delayed og kjøres parallelt.
|
||||||
|
|
||||||
|
PS-analogi:
|
||||||
|
OPAX orchestrator = DAG scheduler (server-node)
|
||||||
|
Individuelle agenter = Dask-delayed tasks (worker-noder)
|
||||||
|
"""
|
||||||
|
|
||||||
|
import time
|
||||||
|
from typing import Any, Callable, Dict, List, Optional
|
||||||
|
|
||||||
|
import dask
|
||||||
|
|
||||||
|
|
||||||
|
@dask.delayed
|
||||||
|
def run_agent_task(
|
||||||
|
agent_fn: Callable,
|
||||||
|
payload: Dict[str, Any],
|
||||||
|
agent_id: Optional[str] = None,
|
||||||
|
) -> Dict[str, Any]:
|
||||||
|
"""
|
||||||
|
Wrapper: kjører én agent som en Dask-delayed task.
|
||||||
|
Returnerer resultat med metadata for telemetri og state store.
|
||||||
|
"""
|
||||||
|
_id = agent_id or getattr(agent_fn, "__name__", "unknown")
|
||||||
|
start = time.monotonic()
|
||||||
|
try:
|
||||||
|
result = agent_fn(payload)
|
||||||
|
success = True
|
||||||
|
error = None
|
||||||
|
except Exception as exc: # noqa: BLE001
|
||||||
|
result = None
|
||||||
|
success = False
|
||||||
|
error = str(exc)
|
||||||
|
duration = round(time.monotonic() - start, 3)
|
||||||
|
return {
|
||||||
|
"agent": _id,
|
||||||
|
"result": result,
|
||||||
|
"success": success,
|
||||||
|
"error": error,
|
||||||
|
"duration_s": duration,
|
||||||
|
"timestamp": time.time(),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def build_agent_dag(
|
||||||
|
agent_fns: List[Callable],
|
||||||
|
payload: Dict[str, Any],
|
||||||
|
agent_ids: Optional[List[str]] = None,
|
||||||
|
) -> List:
|
||||||
|
"""
|
||||||
|
Bygger en liste av Dask-delayed tasks (én per agent).
|
||||||
|
Agenter uten avhengigheter kjøres parallelt av scheduleren.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
agent_fns: Liste av callable agentfunksjoner.
|
||||||
|
payload: Felles input-dict til alle agenter.
|
||||||
|
agent_ids: Valgfrie ID-er som matcher agent_fns (for logging).
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Liste av Dask-delayed objects — klar for execute_dag().
|
||||||
|
"""
|
||||||
|
ids = agent_ids or [getattr(fn, "__name__", f"agent_{i}") for i, fn in enumerate(agent_fns)]
|
||||||
|
return [run_agent_task(fn, payload, aid) for fn, aid in zip(agent_fns, ids)]
|
||||||
|
|
||||||
|
|
||||||
|
def execute_dag(
|
||||||
|
tasks: List,
|
||||||
|
scheduler: str = "synchronous",
|
||||||
|
) -> List[Dict[str, Any]]:
|
||||||
|
"""
|
||||||
|
Kjører den ferdigbygde DAGen.
|
||||||
|
|
||||||
|
Scheduler-valg:
|
||||||
|
'synchronous' — enkelt-tråd, ingen overhead (dev / unit-test)
|
||||||
|
'threads' — lokal multi-tråd, GIL-vennlig for IO-bound agenter
|
||||||
|
'processes' — multi-prosess, CPU-bound arbeid
|
||||||
|
'distributed' — Dask cluster (full skala, krever dask[distributed])
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Liste av resultat-dicts fra run_agent_task.
|
||||||
|
"""
|
||||||
|
results = dask.compute(*tasks, scheduler=scheduler)
|
||||||
|
return list(results)
|
||||||
98
ml/telemetry.py
Normal file
98
ml/telemetry.py
Normal file
|
|
@ -0,0 +1,98 @@
|
||||||
|
#!/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")
|
||||||
|
|
@ -28,3 +28,15 @@ fastapi>=0.111.0
|
||||||
# Utilities
|
# Utilities
|
||||||
python-dotenv>=1.0.0
|
python-dotenv>=1.0.0
|
||||||
pydantic>=2.7.0
|
pydantic>=2.7.0
|
||||||
|
|
||||||
|
# ── ML Layer (Fase ML-1) ─────────────────────────────────────────────────────
|
||||||
|
# Dask: parallell task scheduling + lazy DAG execution
|
||||||
|
dask[distributed]>=2024.1.0
|
||||||
|
# dask-ml: scikit-learn-kompatibel ML på Dask DataFrames/Arrays (Fase ML-2)
|
||||||
|
dask-ml>=2024.4.0
|
||||||
|
# scikit-learn: base estimatorer med partial_fit for inkrementell læring
|
||||||
|
scikit-learn>=1.4.0
|
||||||
|
# pandas: tabular telemetri-analyse
|
||||||
|
pandas>=2.0.0
|
||||||
|
# XGBoost: distribuert feature importance (Fase ML-3)
|
||||||
|
# xgboost>=2.0.0 ← aktiver i Fase ML-3
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user