Skip to main content

Stopping Streams

Cancelling ongoing streams is often needed in production applications. For example, users might want to stop a stream when they realize that the response is not what they want, or you may need to implement timeout policies.

The Go AI SDK uses Go's standard context.Context for cancellation, providing a native and idiomatic way to stop streams.

Context-Based Cancellation​

All AI SDK functions accept a context.Context parameter that can be used to cancel operations. When the context is cancelled, the stream will stop immediately and clean up resources.

Basic Example​

package main

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

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

func main() {
// Create a cancellable context
ctx, cancel := context.WithCancel(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 very long story...",
})
if err != nil {
log.Fatal(err)
}

// Cancel the stream after 5 seconds
go func() {
time.Sleep(5 * time.Second)
fmt.Println("\n[Cancelling stream...]")
cancel()
}()

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

// Check why the stream ended
if err := stream.Err(); err != nil {
if ctx.Err() == context.Canceled {
fmt.Println("\nStream was cancelled successfully")
} else {
log.Printf("Stream error: %v", err)
}
}
}

Timeout-Based Cancellation​

Use context.WithTimeout to automatically cancel a stream after a specified duration:

package main

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

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

func main() {
// Create a context with 10-second timeout
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()

provider := anthropic.New(anthropic.Config{
APIKey: os.Getenv("ANTHROPIC_API_KEY"),
})
model, _ := provider.LanguageModel("claude-sonnet-4-5")

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

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

// Check if timeout occurred
if err := stream.Err(); err != nil {
if ctx.Err() == context.DeadlineExceeded {
fmt.Println("\n\nStream stopped due to timeout")
} else {
log.Printf("Stream error: %v", err)
}
}
}

HTTP Request Cancellation​

In web servers, you can forward the request context to automatically cancel streams when clients disconnect:

package main

import (
"context"
"fmt"
"log"
"net/http"
"os"

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

func streamHandler(w http.ResponseWriter, r *http.Request) {
// Use request context - automatically cancelled when client disconnects
ctx := r.Context()

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 comprehensive guide...",
})
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}

// Set up SSE headers
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")

flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "Streaming not supported", http.StatusInternalServerError)
return
}

// Stream to client
for chunk := range stream.Chunks() {
select {
case <-ctx.Done():
// Client disconnected
log.Println("Client disconnected, cleaning up")
return
default:
fmt.Fprintf(w, "data: %s\n\n", chunk.Text)
flusher.Flush()
}
}

// Log cancellation reason
if err := stream.Err(); err != nil {
if ctx.Err() == context.Canceled {
log.Println("Stream cancelled by client disconnect")
} else {
log.Printf("Stream error: %v", err)
}
}
}

func main() {
http.HandleFunc("/stream", streamHandler)

fmt.Println("Server listening on :8080")
fmt.Println("Try: curl http://localhost:8080/stream")
fmt.Println("Press Ctrl+C in curl to test cancellation")

log.Fatal(http.ListenAndServe(":8080", nil))
}

Handling Cancellation with GenerateText​

Non-streaming operations also respect context cancellation:

package main

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

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

func main() {
// Create a context with 5-second timeout
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

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

result, err := ai.GenerateText(ctx, ai.GenerateTextOptions{
Model: model,
Prompt: "Write a comprehensive analysis of...",
})

if err != nil {
if ctx.Err() == context.DeadlineExceeded {
log.Println("Operation timed out after 5 seconds")
} else {
log.Printf("Error: %v", err)
}
return
}

fmt.Println(result.Text)
}

Graceful Shutdown​

For long-running services, implement graceful shutdown with context cancellation:

package main

import (
"context"
"fmt"
"log"
"os"
"os/signal"
"syscall"

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

func processStream(ctx context.Context) error {
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: "Generate a report...",
})
if err != nil {
return err
}

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

// Check context periodically
select {
case <-ctx.Done():
fmt.Println("\n[Shutting down gracefully...]")
return ctx.Err()
default:
}
}

return stream.Err()
}

func main() {
// Create root context
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

// Handle OS signals for graceful shutdown
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM)

go func() {
sig := <-sigChan
fmt.Printf("\nReceived signal: %v\n", sig)
cancel() // Cancel all operations
}()

// Process stream with cancellation support
if err := processStream(ctx); err != nil {
if ctx.Err() == context.Canceled {
log.Println("Operation cancelled successfully")
} else {
log.Printf("Error: %v", err)
}
}
}

Best Practices​

1. Always Defer Cancel​

// Good - Always defer cancel to prevent context leaks
ctx, cancel := context.WithCancel(parent)
defer cancel()

result, err := ai.GenerateText(ctx, options)

2. Check Context in Long Operations​

for chunk := range stream.Chunks() {
// Check context regularly
select {
case <-ctx.Done():
return ctx.Err()
default:
process(chunk)
}
}

3. Distinguish Cancellation from Errors​

if err := stream.Err(); err != nil {
switch ctx.Err() {
case context.Canceled:
log.Println("User cancelled the operation")
case context.DeadlineExceeded:
log.Println("Operation timed out")
default:
log.Printf("Stream error: %v", err)
}
}

4. Use Request Context in HTTP Handlers​

func handler(w http.ResponseWriter, r *http.Request) {
// Use request context - automatic cancellation on disconnect
ctx := r.Context()

stream, err := ai.StreamText(ctx, options)
// Stream will automatically stop if client disconnects
}

5. Set Appropriate Timeouts​

// Short timeout for user-facing requests
ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second)
defer cancel()

// Longer timeout for batch processing
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()

Context Propagation​

Contexts propagate cancellation through your application:

func processRequest(ctx context.Context, prompt string) error {
// Cancellation from parent context propagates here
return generateAndProcess(ctx, prompt)
}

func generateAndProcess(ctx context.Context, prompt string) error {
// Cancellation propagates to AI SDK calls
stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: prompt,
})
// ... Stream will stop if parent context is cancelled
}

Comparison with JavaScript​

JavaScript (AbortSignal)​

const controller = new AbortController();

const result = streamText({
model,
prompt: "...",
abortSignal: controller.signal,
});

// Cancel the stream
controller.abort();

Go (Context)​

ctx, cancel := context.WithCancel(context.Background())

stream, err := ai.StreamText(ctx, ai.StreamTextOptions{
Model: model,
Prompt: "...",
})

// Cancel the stream
cancel()

Go's context pattern is more powerful because:

  • Automatic propagation through call chains
  • Built-in deadline and timeout support
  • Native integration with standard library (HTTP, database, etc.)
  • Type-safe cancellation reason checking

Next Steps​