ai.stream.TurnStream

The raw text-delta stream of ONE model turn. Pull-based: nothing is read from the wire until `next()` or `final_turn()` is called.

Reference version

Signature

class ai.stream.TurnStream

The raw text-delta stream of ONE model turn. Pull-based: nothing is read from the wire until next() or final_turn() is called.

next() yields DELTA chunks (each item is only the newly arrived text, never the accumulated prefix) and returns Done when the turn is over. final_turn() drains the rest of the stream and returns the completed turn: the accumulated text as one Text block, the folded stop_reason, and usage where the wire provided token counts.

Source:<builtin>/ai/ns_stream/stream.bamlbytes 6904–28729

Fields

_sse

baml.http.SseStream | null

_chunks

string[]

_cursor

int

_pcursor

int

_text

string

_input_tokens

int | null

_output_tokens

int | null

_cached_input_tokens

int | null

_cache_write_input_tokens

int | null

_uncached_input_tokens

int | null

_reasoning_tokens

int | null

_done

bool

_require_terminal

string | null

_saw_terminal

bool

_call_started

baml.time.Instant | null

_ttft_from

baml.time.Instant | null

_first_token_timeout

baml.time.Duration | null

_stream_token_timeout

baml.time.Duration | null

_content_deltas

int

_token_waited

baml.time.Duration

_capture_sse

bool

Static methods

function

from_chunks

(chunks: string[]) -> ai.stream.TurnStream throws never

A turn stream fed from literal delta chunks — the scripted counterpart of a live stream, for tests. There is no wire to be cut, so a scripted stream is never strict: the turn completes as StopReason.Complete after the last chunk.

function

from_event_source

(next_event: ai.stream.EventSource) -> ai.stream.TurnStream throws never

A stream backed by one normalized-event callback. Reliability clients use this to relay a chosen member's deltas and final metadata while retaining TurnStream's ordinary accumulation and termination behavior. The source itself owns any retry/fallback policy; once it returns a TextDelta, that delta is immediately visible to the caller.

function

from_sse

(
sse: baml.http.SseStream,
decode: (string) -> ai.stream.Event[] throws ai.errors.InvokeError,
require_terminal: string | null = …,
timeout_options: ai.LlmTimeoutOptions | null = …,
dispatched: baml.time.Instant | null = …
) -> ai.stream.TurnStream throws baml.errors.Timeout

A turn stream over a live SSE connection. decode receives the raw batch string from sse.next() (a JSON array of SseEvent objects; see decode_sse_batch) and yields the events it contains.

require_terminal opts into the strict termination contract: pass the provider id ("anthropic", "openai", …) when the protocol has an explicit end-of-turn event that the decoder yields as TurnDone. The stream then throws ai.errors.NetworkFailure if the connection ends before that event arrives, rather than reporting the truncated text as a completed turn. Leave it null (the default) only for a protocol whose sole end-of-turn signal IS the socket closing.

timeout_options supplies the two token deadlines, which this stream enforces: next() throws baml.errors.Timeout with timeout_type "first_token_timeout" when no non-empty content delta arrives in time, and "stream_token_timeout" when the gap between two of them is too long. Content is a non-empty TextDelta or a ContentSeen. dispatched is when the request attempt was sent (before connecting), which is where the first-token clock starts; it defaults to now. If the first-token deadline has already passed, the stream is closed and this throws that Timeout. timeout_options.timeout and timeout_options.connect_timeout are the transport's and are not read here.

Instance methods

function

_capture_batch

(self, batch: string) -> void throws never

(internal) Append a raw SSE batch to the wire record when capture is on. Telemetry must never fail the stream, so a malformed batch is skipped.

function

_check_terminal

(self) -> void throws ai.errors.NetworkFailure

(internal) The strict-termination check: a no-op unless this stream required an explicit terminal event and the stream ended without one, in which case the truncated turn is reported as the transport failure it is instead of being passed off as complete. Called on every end-of-stream path.

function

_finish

(self) -> void throws never

(internal) Marks the turn over and releases the SSE connection behind it, if this stream has one.

function

_pull

(
self,
sse: baml.http.SseStream
) -> string | null throws baml.errors.Io | baml.errors.Timeout

(internal) One wire pull, bounded by whichever token deadline is running: first_token_timeout until the first content delta, then stream_token_timeout. The clock only runs here, while waiting on the provider. A batch with no content in it does not reset it, so the next pull gets what is left of the same budget.

The transport's own Timeout (the request's end-to-end limit) is carried out of the timed body as a value so it is never mistaken for the token deadline, then rethrown unchanged.

function

_saw_content

(self) -> void throws never

(internal) A non-empty content delta was delivered: the first-token deadline is met, and the stream-token clock starts over.

function

_stamp_first_delta

(self) -> void throws never

(internal) Stamp time-to-first-token on this stream's wire record, once.

function

final_turn

(
self
) -> ai.ModelTurn throws baml.errors.Io | baml.errors.Timeout | baml.errors.UnknownError | reflect.errors.CompilationError | ai.errors.Failure

Drain the stream and return the completed turn. Safe to call at any point (including after Done); calling twice returns an equal turn.

On a strict stream (from_sse's require_terminal) that ended without the provider's terminal event this throws ai.errors.NetworkFailure instead of returning a turn — the stop_reason default below must never turn a truncated stream into a StopReason.Complete answer.

function

next

(
self
) -> string[] | ai.stream.Done throws baml.errors.Io | baml.errors.Timeout | baml.errors.UnknownError | reflect.errors.CompilationError | ai.errors.Failure

All READY text deltas, or Done when the turn is over. Metadata events are folded silently; calling next() after Done keeps returning Done.

Drain-ready-then-deliver: every event already decoded (the wire layer returns whole batches — everything that queued while the caller was busy) is returned as ONE batch, and the wire is only pulled when no deliverable text is in hand. Delta boundaries are preserved — every element is one provider text delta, so consumers that want each SSE delta individually iterate the batch. A consumer that keeps up sees one-element batches; a consumer that falls behind (e.g. a partial parse slower than the token interval) gets the whole backlog in one return instead of replaying it delta by delta — the pull-model equivalent of the legacy engine's latest-snapshot watch channel.