Spaces:
Running
Running
sync: 141 file da Baida98/AI@c6afe8eb (2026-08-07 20:01 UTC) [deploy-all]
#18
by Baida07 - opened
- api/state.py +121 -4
api/state.py
CHANGED
|
@@ -137,6 +137,23 @@ _heartbeat_state: dict = {
|
|
| 137 |
"runs": 0,
|
| 138 |
}
|
| 139 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 140 |
# ββ Singleton Getters βββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 141 |
def get_supabase() -> Optional[Any]:
|
| 142 |
return _sb
|
|
@@ -147,10 +164,16 @@ def _get_mem_manager() -> Any:
|
|
| 147 |
global _mem_manager, _mem_manager_inited
|
| 148 |
if _mem_manager_inited: return _mem_manager
|
| 149 |
try:
|
| 150 |
-
from memory.
|
| 151 |
-
_mem_manager = MemoryManager()
|
| 152 |
-
|
| 153 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 154 |
return _mem_manager
|
| 155 |
|
| 156 |
_executor: Any = None
|
|
@@ -175,3 +198,97 @@ def _get_ai_client() -> Any:
|
|
| 175 |
|
| 176 |
async def _get_mem_manager_async() -> Any:
|
| 177 |
return _get_mem_manager()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 137 |
"runs": 0,
|
| 138 |
}
|
| 139 |
|
| 140 |
+
# ββ Telemetry & Timing ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 141 |
+
# Shared by the agent loop and the provider diagnostics endpoint. Keep this
|
| 142 |
+
# bounded so long-running workers cannot grow without limit.
|
| 143 |
+
_TIMING_STORE: dict[str, list[float]] = {}
|
| 144 |
+
_REPAIR_STATS: dict[str, int] = {}
|
| 145 |
+
|
| 146 |
+
def record_timing(key: str, duration_ms: float) -> None:
|
| 147 |
+
"""Record a bounded latency sample for agent/provider diagnostics."""
|
| 148 |
+
samples = _TIMING_STORE.setdefault(key, [])
|
| 149 |
+
samples.append(duration_ms)
|
| 150 |
+
if len(samples) > 100:
|
| 151 |
+
samples.pop(0)
|
| 152 |
+
|
| 153 |
+
def increment_stat(key: str, delta: int = 1) -> None:
|
| 154 |
+
"""Increment an aggregated agent quality/recovery counter."""
|
| 155 |
+
_REPAIR_STATS[key] = _REPAIR_STATS.get(key, 0) + delta
|
| 156 |
+
|
| 157 |
# ββ Singleton Getters βββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 158 |
def get_supabase() -> Optional[Any]:
|
| 159 |
return _sb
|
|
|
|
| 164 |
global _mem_manager, _mem_manager_inited
|
| 165 |
if _mem_manager_inited: return _mem_manager
|
| 166 |
try:
|
| 167 |
+
from memory.manager import MemoryManager
|
| 168 |
+
_mem_manager = MemoryManager(sb_client=_get_sb())
|
| 169 |
+
try:
|
| 170 |
+
_asyncio_mod.create_task(_mem_manager.init())
|
| 171 |
+
_mem_manager_inited = True
|
| 172 |
+
except RuntimeError:
|
| 173 |
+
# No running event loop during import; the async getter initializes it.
|
| 174 |
+
pass
|
| 175 |
+
except Exception:
|
| 176 |
+
_mem_manager = None
|
| 177 |
return _mem_manager
|
| 178 |
|
| 179 |
_executor: Any = None
|
|
|
|
| 198 |
|
| 199 |
async def _get_mem_manager_async() -> Any:
|
| 200 |
return _get_mem_manager()
|
| 201 |
+
|
| 202 |
+
_planner: Any = None
|
| 203 |
+
def _get_planner() -> Any:
|
| 204 |
+
global _planner
|
| 205 |
+
if _planner is not None: return _planner
|
| 206 |
+
try:
|
| 207 |
+
from agents.planner import Planner
|
| 208 |
+
_planner = Planner(llm_client=_get_ai_client())
|
| 209 |
+
except Exception:
|
| 210 |
+
_planner = None
|
| 211 |
+
return _planner
|
| 212 |
+
|
| 213 |
+
# ββ Prune helpers βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 214 |
+
def _prune_checkpoints() -> None:
|
| 215 |
+
now = int(time.time() * 1000)
|
| 216 |
+
expired = [k for k, v in list(_task_checkpoints.items())
|
| 217 |
+
if now - v.get('savedAt', 0) > _CHECKPOINT_TTL_MS]
|
| 218 |
+
for k in expired:
|
| 219 |
+
_task_checkpoints.pop(k, None)
|
| 220 |
+
if len(_task_checkpoints) > _CHECKPOINT_MAX:
|
| 221 |
+
oldest = sorted(_task_checkpoints.items(), key=lambda x: x[1].get('savedAt', 0))
|
| 222 |
+
for k, _ in oldest[:len(_task_checkpoints) - _CHECKPOINT_MAX]:
|
| 223 |
+
_task_checkpoints.pop(k, None)
|
| 224 |
+
|
| 225 |
+
def _prune_agent_tasks() -> None:
|
| 226 |
+
now = int(time.time() * 1000)
|
| 227 |
+
expired = [k for k, v in list(_agent_tasks.items())
|
| 228 |
+
if v.get('status') in ('SUCCESS', 'ERROR', 'CANCELLED')
|
| 229 |
+
and now - v.get('created_at', 0) > _AGENT_TASK_TTL_MS]
|
| 230 |
+
for k in expired:
|
| 231 |
+
_agent_tasks.pop(k, None)
|
| 232 |
+
if len(_agent_tasks) > _AGENT_TASK_MAX:
|
| 233 |
+
oldest = sorted(_agent_tasks.items(), key=lambda x: x[1].get('created_at', 0))
|
| 234 |
+
for k, _ in oldest[:len(_agent_tasks) - _AGENT_TASK_MAX]:
|
| 235 |
+
_agent_tasks.pop(k, None)
|
| 236 |
+
|
| 237 |
+
def _prune_loop_registry() -> None:
|
| 238 |
+
now = time.time()
|
| 239 |
+
stale = [k for k, v in list(_loop_registry.items())
|
| 240 |
+
if v.get('done') and now - v.get('finished_at', 0.0) > _LOOP_REGISTRY_TTL_S]
|
| 241 |
+
for k in stale:
|
| 242 |
+
_loop_registry.pop(k, None)
|
| 243 |
+
|
| 244 |
+
# ββ Shared Pydantic models ββββββββββββββββββββββββββββββββββββββββββββββββββββ
|
| 245 |
+
class ReasonLoopIn(BaseModel):
|
| 246 |
+
goal: str
|
| 247 |
+
context: list[dict] = []
|
| 248 |
+
max_steps: int = 8
|
| 249 |
+
project_context: str = ""
|
| 250 |
+
learning_hints: list[str] = []
|
| 251 |
+
session_id: Optional[str] = None
|
| 252 |
+
negative_constraints: Optional[str] = ""
|
| 253 |
+
|
| 254 |
+
@field_validator('goal', mode='before')
|
| 255 |
+
@classmethod
|
| 256 |
+
def validate_goal(cls, v: object) -> str:
|
| 257 |
+
if not isinstance(v, str) or not v.strip():
|
| 258 |
+
raise ValueError('goal must be a non-empty string')
|
| 259 |
+
return v.strip()
|
| 260 |
+
|
| 261 |
+
@field_validator('context', 'learning_hints', mode='before')
|
| 262 |
+
@classmethod
|
| 263 |
+
def coerce_list(cls, v: object) -> list:
|
| 264 |
+
return v if isinstance(v, list) else []
|
| 265 |
+
|
| 266 |
+
@field_validator('project_context', mode='before')
|
| 267 |
+
@classmethod
|
| 268 |
+
def coerce_str(cls, v: object) -> str:
|
| 269 |
+
return str(v).strip()[:2000] if v else ""
|
| 270 |
+
|
| 271 |
+
class AgentTaskIn(BaseModel):
|
| 272 |
+
goal: str
|
| 273 |
+
context: list[dict] = []
|
| 274 |
+
max_steps: int = 8
|
| 275 |
+
taskId: Optional[str] = None
|
| 276 |
+
project_context: str = ""
|
| 277 |
+
learning_hints: list[str] = []
|
| 278 |
+
session_id: Optional[str] = None
|
| 279 |
+
resume_from_step: Optional[int] = None
|
| 280 |
+
persona: Optional[str] = None
|
| 281 |
+
negative_constraints: Optional[str] = ""
|
| 282 |
+
|
| 283 |
+
@field_validator('goal', mode='before')
|
| 284 |
+
@classmethod
|
| 285 |
+
def validate_goal(cls, v: object) -> str:
|
| 286 |
+
if not isinstance(v, str) or not v.strip():
|
| 287 |
+
raise ValueError('goal must be a non-empty string')
|
| 288 |
+
return v.strip()
|
| 289 |
+
|
| 290 |
+
@field_validator('context', 'learning_hints', mode='before')
|
| 291 |
+
@classmethod
|
| 292 |
+
def coerce_list(cls, v: object) -> list:
|
| 293 |
+
return v if isinstance(v, list) else []
|
| 294 |
+
|