125 lines
4.3 KiB
Python
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
|