Spaces:
Sleeping
Sleeping
| import time | |
| import uuid | |
| import asyncio | |
| from typing import Any, Dict, List, Literal, Optional, Set | |
| from pydantic import BaseModel, Field | |
| # ----------------------------------------------------------------------------- | |
| # 1. THE CONTRACT: UNIVERSAL DEBUG EVENT SCHEMA | |
| # ----------------------------------------------------------------------------- | |
| DebugEventType = Literal[ | |
| "run_started", | |
| "input_received", | |
| "prompt_built", | |
| "retrieval_started", | |
| "retrieval_result", | |
| "tool_called", | |
| "tool_result", | |
| "llm_request", | |
| "llm_response", | |
| "policy_check", | |
| "error", | |
| "run_finished" | |
| ] | |
| class DebugEvent(BaseModel): | |
| runId: str | |
| sessionId: Optional[str] = None | |
| userId: Optional[str] = None | |
| timestamp: int = Field(default_factory=lambda: int(time.time() * 1000)) | |
| type: DebugEventType | |
| title: str | |
| payload: Dict[str, Any] = Field(default_factory=dict) | |
| # ----------------------------------------------------------------------------- | |
| # 2. THE BUS: PORTABLE TRANSPORT & STORAGE | |
| # ----------------------------------------------------------------------------- | |
| class DebugBus: | |
| def __init__(self): | |
| self.listeners: Set[asyncio.Queue] = set() | |
| self.session_subscriptions: Dict[str, Set[asyncio.Queue]] = {} | |
| self.runs: Dict[str, List[DebugEvent]] = {} # Memory storage for replay | |
| def subscribe(self, session_id: str, queue: asyncio.Queue): | |
| if session_id not in self.session_subscriptions: | |
| self.session_subscriptions[session_id] = set() | |
| self.session_subscriptions[session_id].add(queue) | |
| self.listeners.add(queue) | |
| def unsubscribe(self, session_id: str, queue: asyncio.Queue): | |
| if session_id in self.session_subscriptions: | |
| self.session_subscriptions[session_id].remove(queue) | |
| if not self.session_subscriptions[session_id]: | |
| del self.session_subscriptions[session_id] | |
| self.listeners.discard(queue) | |
| async def emit(self, event: DebugEvent): | |
| # 1. Store locally for replay | |
| if event.runId not in self.runs: | |
| self.runs[event.runId] = [] | |
| self.runs[event.runId].append(event) | |
| # 2. Push to active subscribers | |
| # Broadcasters listen to specific session_id streams | |
| session_id = event.sessionId | |
| data = event.model_dump() | |
| if session_id and session_id in self.session_subscriptions: | |
| for q in self.session_subscriptions[session_id]: | |
| await q.put(data) | |
| # Global listeners (if any) | |
| # for q in self.listeners: | |
| # await q.put(data) | |
| def get_run(self, run_id: str) -> List[DebugEvent]: | |
| return self.runs.get(run_id, []) | |
| def get_all_runs(self) -> List[Dict[str, Any]]: | |
| return [{"runId": rid, "events": [e.model_dump() for e in evs]} for rid, evs in self.runs.items()] | |
| # Singleton instance | |
| debug_bus = DebugBus() | |
| # ----------------------------------------------------------------------------- | |
| # 3. THE EMITTER: CONVENIENCE WRAPPER | |
| # ----------------------------------------------------------------------------- | |
| async def emit_debug_event( | |
| runId: str, | |
| title: str, | |
| type: DebugEventType, | |
| sessionId: Optional[str] = None, | |
| userId: Optional[str] = None, | |
| payload: Optional[Dict[str, Any]] = None | |
| ): | |
| event = DebugEvent( | |
| runId=runId, | |
| sessionId=sessionId, | |
| userId=userId, | |
| type=type, | |
| title=title, | |
| payload=payload or {} | |
| ) | |
| await debug_bus.emit(event) | |