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)