Skip to main content

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​

FieldTypeRequiredDescription
Modelprovider.LanguageModelYesLanguage model to use for generation
PromptstringNoSimple string prompt (alternative to Messages)
Messages[]types.MessageNoList of conversation messages
SystemstringNoSystem instructions
Temperature*float64NoSampling temperature (0.0 to 2.0)
MaxTokens*intNoMaximum tokens to generate
TopP*float64NoNucleus sampling parameter
TopK*intNoTop-K sampling parameter
FrequencyPenalty*float64NoFrequency penalty (-2.0 to 2.0)
PresencePenalty*float64NoPresence penalty (-2.0 to 2.0)
StopSequences[]stringNoSequences that stop generation
Seed*intNoRandom seed for reproducibility
Tools[]types.ToolNoTools available for the model to call
StopWhen[]ai.StopConditionNoConditions 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.
ToolChoicetypes.ToolChoiceNoHow the model should choose tools
ResponseFormat*provider.ResponseFormatNoResponse format specification
Timeout*TimeoutConfigNoTimeout configuration
ExperimentalRetention*types.RetentionSettingsNoData retention settings
ProviderOptionsmap[string]interfaceNoProvider-specific options
OnChunkfunc(provider.StreamChunk)NoCalled for each chunk
OnFinishfunc(*StreamTextResult)NoCalled when stream completes

Return Value​

StreamTextResult​

The result provides methods to consume the stream:

MethodReturnsDescription
Stream()provider.TextStreamUnderlying text stream
FullStream()provider.TextStreamEvery chunk type, including lifecycle chunks (see below)
Text()stringAccumulated text so far
FinishReason()types.FinishReasonFinish reason (after stream ends)
Usage()types.UsageUsage info (after stream ends)
ContextManagement()interfaceContext management info
Err()errorError that occurred during streaming
Close()errorClose the stream
ReadAll()(string, error)Read all chunks and return complete text
Chunks()<-chan provider.StreamChunkChannel of chunks
Steps()[]types.StepResultCompleted steps, including per-step Performance statistics
FinalStep()types.StepResultLast completed step with StepTimeMs, ResponseTimeMs, ToolExecutionMs, EffectiveOutputTokensPerSecond, EffectiveTotalTokensPerSecond, and streaming throughput fields
Content()[]types.ContentPartOrdered 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.ChunkTypeStart opens the stream — emitted exactly once, before the first step's provider chunks.
  • provider.ChunkTypeStartStep / provider.ChunkTypeFinishStep bracket each step. ChunkTypeFinishStep carries that step's request, response, usage, finish reason, and provider metadata.
  • provider.ChunkTypeFinish closes 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​