Expand description
§Module: Streaming
§Responsibility
Emits agent reasoning events (thought, action, observation, result, error)
as Server-Sent Events (SSE) over a tokio::sync::broadcast channel so that
multiple HTTP clients can subscribe to a single agent run in real time.
§Design
AgentEvent— the typed event enum; each variant serialises to JSON and is wrapped in the SSEdata: …\n\nenvelope.StreamBroadcaster— thin wrapper around abroadcast::Sender<String>. CallStreamBroadcaster::sendto publish an event; callStreamBroadcaster::subscribeto obtain abroadcast::Receiver<String>from which a subscriber reads raw SSE lines.AgentEventStream— convenience helper that owns aStreamBroadcasterand exposes typedemit_*methods.
§Guarantees
- Non-panicking: all operations return
Result - No
unwrap/expect/panicin production paths - Thread-safe:
StreamBroadcasterandAgentEventStreamareClone,Send, andSync
§NOT Responsible For
- HTTP transport (callers wire the
Receiveroutput into their own HTTP layer — axum, actix-web, hyper, etc.) - Persistence of streamed events
Structs§
- Agent
Event Stream - High-level typed event stream for a single agent run.
- Stream
Broadcaster - A broadcast publisher for raw SSE strings.
Enums§
- Agent
Event - A single event emitted by a running agent.