Imp.Streaming (Imp v0.5.0)

Copy Markdown View Source

Enumerable-friendly streaming helpers.

Summary

Functions

Collects stream chunks into a string.

Parses provider chunks into incremental typed field updates.

Streams one program call as an Enumerable of chunks.

Functions

collect(program, inputs, opts \\ [])

@spec collect(term(), term(), keyword()) :: String.t() | {:error, term()}

Collects stream chunks into a string.

Successful streams return the collected string. If the stream emits an error chunk, collection stops and returns {:error, reason} instead of partial output.

incremental_fields(chunks, signature)

Parses provider chunks into incremental typed field updates.

The parser follows the ChatAdapter delimiter format: [[ ## field ## ]] followed by field text. A field is emitted when the next field delimiter arrives or when the stream ends.

stream(program, inputs, opts \\ [])

Streams one program call as an Enumerable of chunks.

With provider_stream: true, Imp executes the real program and streams from its named predictors as they are reached. :stream_listeners can select intermediate output fields; without listeners, normalized provider events from every named predictor are yielded. The final typed prediction is yielded by default and can be disabled with include_final_prediction: false.

Composed modules expose their predictor names through the ordinary optimizer predictor callback in Imp.Module. A program with no named predictors returns a terminal {:provider_stream_unsupported, module} error. Without provider_stream: true, the program runs once and the result is chunked locally (grapheme by grapheme, or through the :chunker function).

Either way a failure is the last element of the stream, as %Imp.Streaming.Messages.StreamResponse{chunk: {:error, reason}, done: true}.