Skip to main content

Streaming Issues

This guide covers common problems with streaming AI responses and their solutions.

Error: "Goroutine Leak - Stream Not Consumed"​

Symptoms:

  • Application hangs
  • Memory usage grows continuously
  • runtime.NumGoroutine() keeps increasing

Cause: Not consuming all values from stream channels, causing producer goroutines to block.

Solution:

package main

import (
"context"
"fmt"
"log"
"os"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

// BAD: Not consuming entire stream
func badStreamHandling(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a long story",
})
if err != nil {
return err
}

// Only reading first chunk - LEAK!
chunk := <-stream.Chunks()
fmt.Println(chunk.Text)

// Returning without consuming rest of channel
// Producer goroutine will block forever
return nil
}

// GOOD: Always consume entire stream
func goodStreamHandling(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a long story",
})
if err != nil {
return err
}

// Range loop consumes all values until channel closes
for chunk := range stream.Chunks() {
fmt.Print(chunk.Text)
}

// Check for errors after stream completes
return stream.Err()
}

// ALSO GOOD: Early exit with context cancellation
func earlyExitStreamHandling(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a long story",
})
if err != nil {
return err
}

count := 0
for chunk := range stream.Chunks() {
fmt.Print(chunk.Text)
count++

// If you need to exit early, cancel the context
// This will cause the stream to close cleanly
if count >= 10 {
// Stream will be canceled and channel will close
return nil
}
}

return stream.Err()
}

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

if err := goodStreamHandling(ctx, model); err != nil {
log.Fatal(err)
}
}

Error: "Send on Closed Channel"​

Symptoms:

panic: send on closed channel

Cause: Trying to read from a stream channel multiple times or after it's closed.

Solution:

package main

import (
"context"
"fmt"
"log"
"os"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a short story",
})
if err != nil {
log.Fatal(err)
}

// First read - OK
for chunk := range stream.Chunks() {
fmt.Print(chunk.Text)
}

// Channel is now closed

// BAD: Don't try to read again
// for chunk := range stream.Chunks() { // Channel is closed!
// fmt.Print(chunk.Text)
// }

// GOOD: Only check error after stream closes
if err := stream.Err(); err != nil {
log.Printf("Stream error: %v", err)
}

// Don't reuse the stream object
}

Stream Interruption Issues​

Error: "Stream Stops Prematurely"​

Symptoms:

  • Stream closes before all content received
  • Incomplete responses
  • No error reported

Cause: Network issues, provider errors, or context cancellation.

Solution:

package main

import (
"context"
"fmt"
"log"
"os"
"strings"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

func robustStreaming(ctx context.Context, model provider.LanguageModel, prompt string) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: prompt,
})
if err != nil {
return fmt.Errorf("failed to start stream: %w", err)
}

var fullText strings.Builder
var chunkCount int
var finishReason string

// Use Chunks to get detailed chunk information
for chunk := range stream.Chunks() {
switch chunk.Type {
case provider.ChunkTypeText:
fmt.Print(chunk.Text)
fullText.WriteString(chunk.Text)
chunkCount++

case provider.ChunkTypeToolCall:
log.Printf("Tool call: %s", chunk.ToolCall.ToolName)

case provider.ChunkTypeError:
log.Printf("Error in stream: %v", chunk.Err)

case provider.ChunkTypeFinish:
finishReason = string(chunk.FinishReason)
log.Printf("\nStream finished: %s", finishReason)
}
}

// Check for errors
if err := stream.Err(); err != nil {
log.Printf("Stream error (received %d chunks): %v", chunkCount, err)
return err
}

// Validate stream completed properly
if chunkCount == 0 {
return fmt.Errorf("stream produced no chunks")
}

if finishReason == "" {
log.Println("Warning: Stream ended without finish reason")
}

log.Printf("Received %d chunks, %d bytes", chunkCount, fullText.Len())
return nil
}

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

if err := robustStreaming(ctx, model, "Write a story about AI"); err != nil {
log.Fatalf("Streaming failed: %v", err)
}
}

Error: "Missing Final Chunks"​

Symptoms:

  • Last few chunks not received
  • Response appears incomplete

Cause: Not waiting for channel to close naturally.

Solution:

package main

import (
"context"
"fmt"
"log"
"os"
"time"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

func completeStreamRead(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Count from 1 to 100",
})
if err != nil {
return err
}

timeout := time.After(5 * time.Minute)

for {
select {
case chunk, ok := <-stream.Chunks():
if !ok {
// Channel properly closed, all chunks received
log.Println("\nStream completed")
return stream.Err()
}
fmt.Print(chunk.Text)

case <-timeout:
log.Println("\nTimeout waiting for stream")
return fmt.Errorf("stream timeout")

case <-ctx.Done():
log.Println("\nContext canceled")
return ctx.Err()
}
}
}

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

if err := completeStreamRead(ctx, model); err != nil {
log.Fatalf("Error: %v", err)
}
}

Concurrent Streaming Issues​

Error: "Race Condition in Stream Processing"​

Symptoms:

  • Panic: concurrent map read/write
  • Data corruption
  • Unpredictable behavior

Cause: Multiple goroutines accessing stream data without synchronization.

Solution:

package main

import (
"context"
"fmt"
"log"
"os"
"sync"
"time"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

// SafeStreamProcessor handles concurrent stream processing safely
type SafeStreamProcessor struct {
mu sync.Mutex
chunks []string
done bool
}

func (p *SafeStreamProcessor) AddChunk(chunk string) {
p.mu.Lock()
defer p.mu.Unlock()
p.chunks = append(p.chunks, chunk)
}

func (p *SafeStreamProcessor) GetChunks() []string {
p.mu.Lock()
defer p.mu.Unlock()
result := make([]string, len(p.chunks))
copy(result, p.chunks)
return result
}

func (p *SafeStreamProcessor) MarkDone() {
p.mu.Lock()
defer p.mu.Unlock()
p.done = true
}

func (p *SafeStreamProcessor) IsDone() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.done
}

func processStreamSafely(ctx context.Context, model provider.LanguageModel, prompt string) error {
processor := &SafeStreamProcessor{}

stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: prompt,
})
if err != nil {
return err
}

var wg sync.WaitGroup

// Goroutine 1: Read and store chunks
wg.Add(1)
go func() {
defer wg.Done()
for chunk := range stream.Chunks() {
processor.AddChunk(chunk.Text)
fmt.Print(chunk.Text)
}
processor.MarkDone()
}()

// Goroutine 2: Monitor progress
wg.Add(1)
go func() {
defer wg.Done()
for !processor.IsDone() {
chunks := processor.GetChunks()
log.Printf("Progress: %d chunks received", len(chunks))
select {
case <-ctx.Done():
return
case <-time.After(2 * time.Second):
}
}
}()

wg.Wait()

return stream.Err()
}

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

if err := processStreamSafely(ctx, model, "Write a story"); err != nil {
log.Fatal(err)
}
}

Buffer and Memory Issues​

Error: "Out of Memory with Streaming"​

Symptoms:

  • High memory usage
  • OOM errors
  • Application crashes

Cause: Accumulating all chunks in memory instead of processing them incrementally.

Solution:

package main

import (
"bufio"
"context"
"fmt"
"log"
"os"

"github.com/digitallysavvy/go-ai/pkg/ai"
"github.com/digitallysavvy/go-ai/pkg/provider"
"github.com/digitallysavvy/go-ai/pkg/providers/openai"
)

// BAD: Accumulating all chunks in memory
func memoryHungryStreaming(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a very long essay",
})
if err != nil {
return err
}

var allText string // Accumulates all chunks - can use lots of memory!
for chunk := range stream.Chunks() {
allText += chunk.Text // Bad: growing string
}

fmt.Println(allText)
return stream.Err()
}

// GOOD: Process chunks incrementally
func memoryEfficientStreaming(ctx context.Context, model provider.LanguageModel, outputFile string) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a very long essay",
})
if err != nil {
return err
}

// Write directly to file instead of accumulating in memory
file, err := os.Create(outputFile)
if err != nil {
return err
}
defer file.Close()

writer := bufio.NewWriter(file)
defer writer.Flush()

for chunk := range stream.Chunks() {
// Write chunk immediately, don't accumulate
if _, err := writer.WriteString(chunk.Text); err != nil {
return err
}
// Also print to console
fmt.Print(chunk.Text)
}

return stream.Err()
}

// ALSO GOOD: Process chunks with fixed buffer
func bufferedStreaming(ctx context.Context, model provider.LanguageModel) error {
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "Write a long story",
})
if err != nil {
return err
}

const maxBufferSize = 1000 // Process in chunks of 1000 chars
buffer := make([]byte, 0, maxBufferSize)

for chunk := range stream.Chunks() {
buffer = append(buffer, []byte(chunk.Text)...)

// Process buffer when full
if len(buffer) >= maxBufferSize {
// Process the buffer (e.g., write to file, send to client, etc.)
processBuffer(buffer)
buffer = buffer[:0] // Reset buffer
}
}

// Process remaining buffer
if len(buffer) > 0 {
processBuffer(buffer)
}

return stream.Err()
}

func processBuffer(data []byte) {
// Process data (write to file, database, etc.)
fmt.Print(string(data))
}

func main() {
ctx := context.Background()
provider := openai.New(openai.Config{
APIKey: os.Getenv("OPENAI_API_KEY"),
})
model, _ := provider.LanguageModel("gpt-4")

if err := memoryEfficientStreaming(ctx, model, "output.txt"); err != nil {
log.Fatal(err)
}
}

Best Practices​

1. Always Consume Full Stream​

// Use range to ensure all chunks consumed
for chunk := range stream.Chunks() {
process(chunk)
}
// Channel is now properly closed

2. Check Errors After Stream​

for chunk := range stream.Chunks() {
fmt.Print(chunk.Text)
}

if err := stream.Err(); err != nil {
log.Printf("Stream error: %v", err)
}

3. Use Chunks for Detailed Info​

for chunk := range stream.Chunks() {
switch chunk.Type {
case provider.ChunkTypeText:
// Handle text
case provider.ChunkTypeError:
// Handle errors mid-stream
case provider.ChunkTypeFinish:
// Handle completion
}
}

4. Process Incrementally​

// Don't accumulate all chunks
var allText string // BAD

// Process each chunk as it arrives
for chunk := range stream.Chunks() {
processImmediately(chunk) // GOOD
}

5. Use Timeouts for Streams​

ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()

stream, _ := ai.StreamText(ctx, options)

6. Handle Context Cancellation​

for {
select {
case chunk, ok := <-stream.Chunks():
if !ok {
return stream.Err()
}
process(chunk)
case <-ctx.Done():
return ctx.Err()
}
}

See Also​