Anand Sharma
Clean Deployment v2
76022ae
Raw
History Blame Contribute Delete
3.55 kB
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)