StreamText
Performs streaming text generation where tokens are returned incrementally as they are generated.
Signature
func StreamText(ctx context.Context, opts StreamTextOptions) (*StreamTextResult, error)
Parameters
StreamTextOptions
| Field | Type | Required | Description |
|---|---|---|---|
| Model | provider.LanguageModel | Yes | Language model to use for generation |
| Prompt | string | No | Simple string prompt (alternative to Messages) |
| Messages | []types.Message | No | List of conversation messages |
| System | string | No | System instructions |
| Temperature | *float64 | No | Sampling temperature (0.0 to 2.0) |
| MaxTokens | *int | No | Maximum tokens to generate |
| TopP | *float64 | No | Nucleus sampling parameter |
| TopK | *int | No | Top-K sampling parameter |
| FrequencyPenalty | *float64 | No | Frequency penalty (-2.0 to 2.0) |
| PresencePenalty | *float64 | No | Presence penalty (-2.0 to 2.0) |
| StopSequences | []string | No | Sequences that stop generation |
| Seed | *int | No | Random seed for reproducibility |
| Tools | []types.Tool | No | Tools available for the model to call |
| StopWhen | []ai.StopCondition | No | Conditions that stop the tool-calling loop (see IsStepCount). Default: IsStepCount(1), a single step (tool calls in it still execute). Set e.g. IsStepCount(5) to let the model answer after tool results. |
| ToolChoice | types.ToolChoice | No | How the model should choose tools |
| ResponseFormat | *provider.ResponseFormat | No | Response format specification |
| Timeout | *TimeoutConfig | No | Timeout configuration |
| ExperimentalRetention | *types.RetentionSettings | No | Data retention settings |
| ProviderOptions | map[string]interface | No | Provider-specific options |
| OnChunk | func(provider.StreamChunk) | No | Called for each chunk |
| OnFinish | func(*StreamTextResult) | No | Called when stream completes |
Return Value
StreamTextResult
The result provides methods to consume the stream:
| Method | Returns | Description |
|---|---|---|
| Stream() | provider.TextStream | Underlying text stream |
| FullStream() | provider.TextStream | Every chunk type, including lifecycle chunks (see below) |
| Text() | string | Accumulated text so far |
| FinishReason() | types.FinishReason | Finish reason (after stream ends) |
| Usage() | types.Usage | Usage info (after stream ends) |
| ContextManagement() | interface | Context management info |
| Err() | error | Error that occurred during streaming |
| Close() | error | Close the stream |
| ReadAll() | (string, error) | Read all chunks and return complete text |
| Chunks() | <-chan provider.StreamChunk | Channel of chunks |
| Steps() | []types.StepResult | Completed steps, including per-step Performance statistics |
| FinalStep() | types.StepResult | Last completed step with StepTimeMs, ResponseTimeMs, ToolExecutionMs, EffectiveOutputTokensPerSecond, EffectiveTotalTokensPerSecond, and streaming throughput fields |
| Content() | []types.ContentPart | Ordered generated content from all completed steps |
Examples
Basic Streaming
package main
import (
"context"
"fmt"
"io"
"log"
"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)
func main() {
p := openai.New(openai.Config{
APIKey: "your-api-key",
})
model, err := p.LanguageModel("gpt-4")
if err != nil {
log.Fatal(err)
}
// Start streaming
result, err := ai.StreamText(context.Background(), ai.StreamTextOptions{
Model: model,
Prompt: "Write a story about a robot",
})
if err != nil {
log.Fatal(err)
}
defer result.Close()
// Read chunks manually
for {
chunk, err := result.Stream().Next()
if err == io.EOF {
break
}
if err != nil {
log.Fatal(err)
}
if chunk.Type == provider.ChunkTypeText {
fmt.Print(chunk.Text)
}
}
fmt.Printf("\n\nTotal tokens: %d\n", result.Usage().GetTotalTokens())
}
Using Callbacks
result, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Explain machine learning",
OnChunk: func(chunk provider.StreamChunk) {
if chunk.Type == provider.ChunkTypeText {
fmt.Print(chunk.Text)
}
},
OnFinish: func(result *ai.StreamTextResult) {
fmt.Printf("\n\nDone! Tokens: %d\n", result.Usage().GetTotalTokens())
},
})
if err != nil {
log.Fatal(err)
}
Using ReadAll
result, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a poem about nature",
})
if err != nil {
log.Fatal(err)
}
defer result.Close()
// Read all at once
text, err := result.ReadAll()
if err != nil {
log.Fatal(err)
}
fmt.Println(text)
fmt.Printf("Tokens: %d\n", result.Usage().GetTotalTokens())
Using Channel-based Consumption
result, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Describe quantum physics",
})
if err != nil {
log.Fatal(err)
}
defer result.Close()
// Consume chunks from channel
for chunk := range result.Chunks() {
if chunk.Type == provider.ChunkTypeText {
fmt.Print(chunk.Text)
}
}
fmt.Printf("\n\nFinish reason: %s\n", result.FinishReason())
With Per-Chunk Timeout
timeout := &ai.TimeoutConfig{
PerChunk: 5 * time.Second,
}
result, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a long essay",
Timeout: timeout,
})
if err != nil {
log.Fatal(err)
}
defer result.Close()
text, err := result.ReadAll()
if err != nil {
if strings.Contains(err.Error(), "chunk timeout") {
log.Println("Stream timed out waiting for chunk")
} else {
log.Fatal(err)
}
}
Full Stream Lifecycle
FullStream() (and OnChunk) receive every chunk the call produces,
including lifecycle chunks that bracket the call and each step:
provider.ChunkTypeStartopens the stream — emitted exactly once, before the first step's provider chunks.provider.ChunkTypeStartStep/provider.ChunkTypeFinishStepbracket each step.ChunkTypeFinishStepcarries that step's request, response, usage, finish reason, and provider metadata.provider.ChunkTypeFinishcloses the call — emitted exactly once, after the last step, carrying total usage across all steps.
This differs from earlier SDK versions, where ChunkTypeFinish was forwarded
once per step. Code that used to detect step boundaries by watching for
ChunkTypeFinish must switch to ChunkTypeFinishStep; only the very last
ChunkTypeFinish in a multi-step run signals the whole call is done.
for chunk := range result.Chunks() {
switch chunk.Type {
case provider.ChunkTypeStart:
fmt.Println("stream started")
case provider.ChunkTypeStartStep:
fmt.Println("step started")
case provider.ChunkTypeText:
fmt.Print(chunk.Text)
case provider.ChunkTypeFinishStep:
fmt.Printf("\nstep finished: %s\n", chunk.FinishReason)
case provider.ChunkTypeFinish:
fmt.Printf("call finished, total usage: %d tokens\n", chunk.Usage.GetTotalTokens())
}
}
Error Handling
result, err := ai.StreamText(ctx, opts)
if err != nil {
log.Fatal("Failed to start stream:", err)
}
defer result.Close()
for {
chunk, err := result.Stream().Next()
if err == io.EOF {
break
}
if err != nil {
if strings.Contains(err.Error(), "chunk timeout") {
log.Println("Chunk timeout, continuing...")
continue
}
log.Fatal("Stream error:", err)
}
// Process chunk
fmt.Print(chunk.Text)
}
// Check for errors during streaming
if result.Err() != nil {
log.Println("Error during streaming:", result.Err())
}
See Also
- GenerateText - Non-streaming text generation
- StreamObject - Streaming structured output
- Streaming Guide