EventTranspiler -- converts LangGraph astream_events into granular StreamEvents.
Captures the exact millisecond each tool starts, its arguments, its result,
and each LLM chunk, producing typed StreamEvent objects for the frontend.
EventTranspiler
Converts raw LangGraph astream_events(v2) into typed StreamEvents.
Usage::
transpiler = EventTranspiler()
async for event in transpiler.transpile(runtime.astream_events("v2")):
# event is a TimestampedEvent with ms precision
Source code in src/symfonic/core/streaming/transpiler.py
| def __init__(self) -> None:
self._start_time: float = 0.0
self._active_tool_calls: dict[str, str] = {} # call_id -> tool_name
|
transpile
async
transpile(
raw_events: AsyncIterator[dict[str, Any]],
) -> AsyncIterator[TimestampedEvent]
Transform raw LangGraph events into typed TimestampedEvents.
Source code in src/symfonic/core/streaming/transpiler.py
| async def transpile(
self,
raw_events: AsyncIterator[dict[str, Any]],
) -> AsyncIterator[TimestampedEvent]:
"""Transform raw LangGraph events into typed TimestampedEvents."""
self._start_time = time.monotonic()
async for raw in raw_events:
event_type = raw.get("event", "")
data = raw.get("data", {})
for stream_event in self._map_events(event_type, data, raw):
elapsed = (time.monotonic() - self._start_time) * 1000
yield TimestampedEvent(event=stream_event, elapsed_ms=elapsed)
|
TimestampedEvent
dataclass
TimestampedEvent(event: StreamEvent, elapsed_ms: float)
StreamEvent wrapped with elapsed milliseconds from stream start.
to_dict
to_dict() -> dict[str, Any]
Serialize for SSE transmission.
Source code in src/symfonic/core/streaming/transpiler.py
| def to_dict(self) -> dict[str, Any]:
"""Serialize for SSE transmission."""
return {
"elapsed_ms": round(self.elapsed_ms, 1),
"event_type": type(self.event).__name__,
"data": _event_to_dict(self.event),
}
|