Streaming
Streaming lets you process a response while it is being generated instead of waiting for all of it. It is what makes chat UIs feel responsive, and it lets you show tool calls as the model writes them. Rig streams at two levels: a whole agent run (every model turn plus tool calls and results), and a single model call with no agent loop.
Streaming an agent
Section titled “Streaming an agent”Replace .await on a prompt with .stream(). The stream yields
MultiTurnStreamItem values;
match on them to print text as it arrives and to pick up the final response:
use futures::StreamExt;use rig::agent::MultiTurnStreamItem;use rig::prelude::*;use rig::providers::openai::{self, OpenAI};use rig::streaming::{Item, StreamEvent};use std::io::Write;
#[tokio::main]async fn main() -> anyhow::Result<()> { let agent = AgentBuilder::new(OpenAI::from_env()?.completion(openai::GPT_5_5)) .preamble("You are a storyteller.") .build();
let mut stream = agent.prompt("Tell me a short story about a robot.").stream();
while let Some(item) = stream.next().await { match item? { MultiTurnStreamItem::StreamAssistantItem(Item::Event(StreamEvent::Text { text, .. })) => { print!("{text}"); std::io::stdout().flush()?; } MultiTurnStreamItem::FinalResponse(response) => { println!("\n\n{} tokens used", response.usage().total_tokens.unwrap_or(0)); } _ => {} } }
Ok(())}Once, in a quiet workshop, a small robot named Bolt woke to the hum of morninglight. It had one task left unfinished: to water the single flower on the bench.
312 tokens used.stream() accepts everything .await does: .history(&messages), .max_turns(n),
.add_hook(..), .tool_context(..). It runs the same loop, fires the same hooks, and executes tools
the same way; the only difference is that you see each step as it happens. The stream is lazy:
nothing is sent until you poll it.
What the stream yields
Section titled “What the stream yields”MultiTurnStreamItem has these variants. It is #[non_exhaustive], so keep a _ => {} arm.
| Variant | When it arrives |
|---|---|
StreamAssistantItem(Item<StreamEvent>) | Live output from the model: text, reasoning, and tool-call fragments (see below). |
ToolCall { tool_call } | A tool call the model made, reported once its turn is complete and the call is about to run. |
ToolExecutionCommitted { tool_call } | A tool call Rig executed, with any hook rewrite applied. Arrives with the results, after the whole batch of calls finished. |
ToolResult { tool_result } | The result sent back to the model. tool_result.call is the id of the call it answers. |
CompletionCall(CompletionCall) | One model call finished. Carries that call’s usage, finish_reason and response id. |
ModelTurnRetried { turn } | A hook rejected a finished turn and asked for a retry. Text streamed for turn should be discarded. |
FinalResponse(PromptResponse) | The run is done. The same PromptResponse that .await returns: output(), usage(), messages, completion_calls(). |
The stream’s item type is Result<MultiTurnStreamItem, PromptError>. An Err is the last item: the
run failed with the same PromptError that .await would have
returned, and nothing follows it.
Model output: StreamEvent
Section titled “Model output: StreamEvent”StreamAssistantItem wraps an Item: either
Item::Event(StreamEvent) or Item::Unknown(..), a provider payload Rig doesn’t model (safe to
ignore). A response is made of numbered parts (a text block, a reasoning block, one tool call),
and each part streams as a start, some fragments, and an end:
StreamEvent | Meaning |
|---|---|
Start { part, kind, name } | A part opened. kind is a PartKind (Text, Reasoning, ToolCall, …); for a tool call, name is the tool. |
Text { part, text } | A text fragment. |
Reasoning { part, text } | A reasoning fragment, for models that stream their thinking. |
Arguments { part, json } | A fragment of a tool call’s argument JSON. |
End { part, content } | The part is complete. content is the finished AssistantContent, for a tool call the full ToolCall with its id. |
Tool calls stream their arguments piece by piece, so a UI can show a long argument (a file being
written, a long query) while the model is still producing it. The fragments of one call, joined, are
its argument JSON. rig::streaming::parse_partial_arguments(&so_far) reads an incomplete prefix
into a best-effort JSON object:
use rig::streaming::{Item, PartKind, StreamEvent, parse_partial_arguments};use std::collections::HashMap;
let mut stream = agent.prompt("Write a haiku to notes.md").max_turns(3).stream();let mut args_so_far: HashMap<usize, String> = HashMap::new();
while let Some(item) = stream.next().await { match item? { MultiTurnStreamItem::StreamAssistantItem(Item::Event(event)) => match event { StreamEvent::Start { part, kind: PartKind::ToolCall, name } => { println!("calling {name:?}"); args_so_far.insert(part.index(), String::new()); } StreamEvent::Arguments { part, json } => { if let Some(args) = args_so_far.get_mut(&part.index()) { args.push_str(&json); println!(" args so far: {:?}", parse_partial_arguments(args)); } } _ => {} }, MultiTurnStreamItem::ToolResult { tool_result } => { println!("result for {}", tool_result.call); } _ => {} }}A tool only runs after its call has ended, never on partial arguments. The
agent_tool_call_streaming
example prints every event with timings for several providers.
Usage per model call
Section titled “Usage per model call”FinalResponse’s usage() is the total for the whole run. A run with tools makes several model
calls, and each one re-sends the growing conversation, so the total says little about how large the
context got. Each CompletionCall item carries the usage of one call:
let mut stream = agent.prompt("Hello!").stream();while let Some(item) = stream.next().await { match item? { MultiTurnStreamItem::CompletionCall(call) => { println!("call {}: {:?} input tokens", call.call_index, call.usage.input_tokens); } MultiTurnStreamItem::FinalResponse(response) => { println!("total: {:?} tokens", response.usage().total_tokens); } _ => {} }}Usage counters are Option<u64>: None means the provider didn’t report that number. The same
per-call list is available after the run as response.completion_calls().
Printing to stdout
Section titled “Printing to stdout”For quick programs, rig::agent::stream_to_stdout prints text and reasoning as they arrive and
returns the final response:
use rig::agent::stream_to_stdout;
let mut stream = agent.prompt("Hello!").stream();let response = stream_to_stdout(&mut stream).await?;println!("\n{} model calls", response.requests());Events without a stream: run_channel
Section titled “Events without a stream: run_channel”When the code that drives the run isn’t the code that renders it (a game loop, a UI frame, an ECS
system), run_channel() splits the run into a future that resolves to the final response and a
RunEvents feed of the same
MultiTurnStreamItems. Spawn the future anywhere and drain the feed with the non-blocking
try_next():
let (run, mut events) = agent.prompt("Tell me a joke").run_channel();let handle = tokio::spawn(run);
// Later, from a synchronous tick:while let Some(item) = events.try_next() { if let MultiTurnStreamItem::FinalResponse(response) = item { println!("{}", response.output()); }}
let response = handle.await??;The feed is bounded: when the consumer falls behind, the run waits instead of dropping events. Dropping the feed doesn’t cancel the run.
Streaming a model call directly
Section titled “Streaming a model call directly”To stream one completion with no agent loop (no tools executed, no hooks), call stream on a model
with a CompletionRequest. It yields Item<StreamEvent> values, the
same events an agent forwards, and finish() returns the assembled response:
use rig::completion::CompletionRequest;use rig::streaming::{Item, StreamEvent};
let model = OpenAI::from_env()?.completion(openai::GPT_5_5);let request = CompletionRequest::new(Message::user("Name three rivers.")) .preamble("Answer briefly.") .max_tokens(200);
let mut stream = model.stream(request)?;while let Some(item) = stream.next().await { if let Item::Event(StreamEvent::Text { text, .. }) = item? { print!("{text}"); }}
let response = stream.finish().await?;println!("\nfinish reason: {:?}, usage: {:?}", response.finish_reason(), response.usage);Errors here are ProviderErrors. A stream that ends before the
provider finished the reply is an error (ProviderError::Truncated), not a short answer. To pause a
stream, stop polling it; to cancel it, drop it.
Practical notes
Section titled “Practical notes”- Handle errors in the loop. Starting a stream can’t fail; every failure arrives as an
Erritem, and it is always the last one. - Discard retried turns. If you use hooks that retry turns, reset the text you showed for that
turn when
ModelTurnRetriedarrives. - Read tool calls from
EndorToolCall, not from joined fragments, when you need the final, validated call.
