OSVauco/ml/state_store.py

125 lines
4.3 KiB
Python

#!/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