sync: 141 file da Baida98/AI@c6afe8eb (2026-08-07 20:01 UTC) [deploy-all]

#18
by Baida07 - opened
Files changed (1) hide show
  1. 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.memory_manager import MemoryManager
151
- _mem_manager = MemoryManager()
152
- _mem_manager_inited = True
153
- except Exception: _mem_manager = None
 
 
 
 
 
 
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
+