Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Streaming

Streaming exists for two different reasons, and the SDK offers a method for each:

  • You want to show tokens as they arrive — use stream() and handle events.
  • You want a complete response but delivered over the streaming transport — use stream_accumulated(), which streams internally and hands you a finished Response. This is useful for long generations that would otherwise sit near a request timeout.

Accumulated streaming

use rai_sdk::{ClientBuilder, Model};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let client = ClientBuilder::new()
        .from_env()
        .model(Model::gpt4o_mini())
        .build()?;

    let response = client
        .request()
        .prompt("Write a short launch announcement.")
        .stream_accumulated()
        .await?;

    println!("{}", response.text());
    Ok(())
}

The result is the same shape as generate() — only the transport differs.

Raw stream events

For incremental output, iterate the stream. This needs futures::StreamExt in scope.

use futures::StreamExt;
use rai_sdk::{ClientBuilder, Model, provider::ProviderStreamEvent};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let client = ClientBuilder::new()
        .from_env()
        .model(Model::gpt4o_mini())
        .build()?;

    let mut stream = client
        .request()
        .prompt("Count from one to five.")
        .stream()
        .await?;

    while let Some(event) = stream.next().await {
        match event? {
            ProviderStreamEvent::Text(text) => print!("{text}"),
            ProviderStreamEvent::Done { .. } => println!(),
            _ => {}
        }
    }

    Ok(())
}

Each item is a Result, because a stream can fail partway through. Do not use while let Some(Ok(event)): that silently swallows mid-stream errors and looks like a clean early finish.

Event kinds

ProviderStreamEvent normalizes each provider’s SSE format:

EventMeaning
Text(String)An incremental text delta. Concatenate them in order.
ToolCallStart { id, name }The model began requesting a tool call.
ToolCallChunk { id, arguments }A fragment of that call’s JSON arguments. Accumulate by id.
Done { finish_reason, usage }The stream ended; usage is reported here when the provider supplies it.

Tool-call arguments arrive as fragments that are not individually valid JSON. Buffer all chunks for an id and only parse once Done arrives.

Match non-exhaustively (_ => {}) so new event kinds do not break your code.

Streaming and tools

The streaming methods reject requests that carry tools, rather than silently ignoring them, with Error::InvalidRequest. Executing a tool loop requires sending follow-up requests, which is incompatible with handing you a single continuous stream.

The check uses the request’s effective tool set, so one client can do both. Opt out per request to stream:

#![allow(unused)]
fn main() {
use rai_sdk::{ClientBuilder, Model};

async fn run(tool: rai_sdk::Tool) -> Result<(), Box<dyn std::error::Error>> {
let client = ClientBuilder::new()
    .from_env()
    .model(Model::gpt4o_mini())
    .tool(tool)
    .build()?;

// Tools run here.
let answer = client.request().prompt("Use a tool if needed.").generate().await?;

// And this request streams, because it opts out of tools.
let stream = client
    .request()
    .no_tools()
    .prompt("Just write prose.")
    .stream()
    .await?;
let _ = (answer, stream);
Ok(())
}
}

The rule applies in both directions: adding a tool with .tool(..) on a request makes it non-streamable even if the client has no tools, instead of quietly dropping it.

If you need incremental output and tool execution in one exchange, drive the loop yourself with generate_once(), executing calls between turns.

Cancellation

Dropping a stream aborts the upstream provider request. Every streaming method is driven entirely by its consumer: the provider’s HTTP response body is polled from inside the returned stream, never from a detached background task. Dropping the stream drops that body, closes the connection, and the provider stops generating.

That holds for the whole family — stream(), generate_stream_events(), stream_wire_events(), and stream_accumulated() — and it holds when the surrounding task is cancelled rather than the stream explicitly dropped, which is what a tokio::time::timeout or a web-framework client disconnect looks like. Nothing keeps running in the background.

Two consequences to plan for:

  • A cancelled generation produces no terminal event, so no usage is reported. Providers still bill for what they generated before the abort, so metering cannot rely on the final usage event alone.
  • Conversely, there is nothing to clean up. You do not need a cancellation token or an abort handle; letting the stream go out of scope is the whole mechanism.

Proxying a stream to your own clients

If your server holds the provider credentials and streams results on to a desktop or browser client, the events have to cross a wire. The wire module covers that case:

  client ──POST──▶ your server ──rai-sdk──▶ provider
         ◀──SSE─── WireStreamEvent ◀────────┘

stream_wire_events() yields WireStreamEvents, which serialize to a tagged JSON object — one SSE data: payload each:

#![allow(unused)]
fn main() {
use futures::StreamExt;
use rai_sdk::{ClientBuilder, Model};

async fn run() -> Result<(), Box<dyn std::error::Error>> {
let client = ClientBuilder::new().from_env().model(Model::gpt4o_mini()).build()?;

let mut events = client
    .request()
    .prompt("Explain SSE in one sentence.")
    .stream_wire_events()
    .await?;

while let Some(event) = events.next().await {
    // event: text_delta
    // data: {"type":"text_delta","text":"Server-sent"}
    println!("event: {}\ndata: {}\n", event.tag(), serde_json::to_string(&event)?);
}
Ok(())
}
}

On the receiving side, StreamAccumulator is the client-side counterpart of stream_accumulated(): feed it the parsed events and it hands back one Response, tool calls included.

use rai_sdk::wire::{StreamAccumulator, WireStreamEvent};

fn main() -> Result<(), Box<dyn std::error::Error>> {
let payloads = [
    r#"{"type":"message_start","protocol_version":1,"model":"gpt-4o-mini","provider":"openai"}"#,
    r#"{"type":"text_delta","text":"Hello, "}"#,
    r#"{"type":"text_delta","text":"world."}"#,
    r#"{"type":"usage","usage":{"prompt_tokens":9,"completion_tokens":3,"total_tokens":12}}"#,
    r#"{"type":"message_stop","finish_reason":"stop"}"#,
];

let mut accumulator = StreamAccumulator::new();
for payload in payloads {
    accumulator.push(serde_json::from_str::<WireStreamEvent>(payload)?)?;
}

let response = accumulator.finish()?;
assert_eq!(response.text(), "Hello, world.");
assert_eq!(response.usage.unwrap().total_tokens, Some(12));
Ok(())
}

examples/sse_proxy.rs runs the whole loop — an axum handler, SSE re-emission, and client-side reassembly — in one process.

Wire events

Unlike the other streaming methods, stream_wire_events() items are not Results. Once the stream is open, every outcome is an event:

"type"Meaning
message_startFirst event of every stream; names the protocol version, model, and provider.
text_deltaAppend this text to the output so far.
tool_call_start / tool_call_delta / tool_call_endA tool call, first incrementally and then assembled.
tool_resultThe output of executing a tool call. Only a proxy that runs tools itself emits this.
usageToken counts, emitted once just before the terminal event.
message_stopTerminal event of a successful stream.
turn_completeAn assembled ConversationTurn, for history.
errorTerminal event of a failed stream.

A mid-stream provider failure arrives as error rather than as a truncated response, which is the point: a client that receives no terminal event at all knows its connection died instead. StreamAccumulator::finish() enforces the distinction — it returns the carried error for the first case and a stream-kind error naming the truncation for the second.

Versioning

The "type" strings and each event’s field names are a compatibility surface: a server and a client can be built from different rai-sdk versions. Renaming or removing one is a breaking change and will be called out in the changelog. Adding a variant is not, so match with a catch-all arm — WireStreamEvent and WireErrorKind are both #[non_exhaustive], and an unrecognized error kind deserializes into WireErrorKind::Other instead of failing.

WIRE_PROTOCOL_VERSION names the current revision of the framing and rides on every message_start. It is bumped only when a client must react to a framing change, never for additive variants.

Timeouts

Streaming does not exempt a request from the configured timeout. A long generation can still exceed AI_TIMEOUT_SECONDS; raise it for workloads that legitimately run long. See Configuration.