Streaming Patterns
Primary Primitive: iter.Seq2[T, error]
func (m *Model) Stream(ctx context.Context, msgs []schema.Message) iter.Seq2[schema.StreamChunk, error] {
return func(yield func(schema.StreamChunk, error) bool) {
stream, err := m.client.Stream(ctx, msgs)
if err != nil { yield(schema.StreamChunk{}, err); return }
defer stream.Close()
for {
select {
case <-ctx.Done(): yield(schema.StreamChunk{}, ctx.Err()); return
default:
}
chunk, err := stream.Recv()
if err == io.EOF { return }
if err != nil { yield(schema.StreamChunk{}, err); return }
if !yield(convertChunk(chunk), nil) { return } // consumer stopped
}
}
}
Composition
- Pipe:
func Pipe[A, B any](first iter.Seq2[A, error], transform func(A) (B, error)) iter.Seq2[B, error]
- Collect: Stream to slice —
func Collect[T any](stream iter.Seq2[T, error]) ([]T, error)
- Invoke from Stream: Stream, collect, return last.
- Fan-out:
iter.Pull2() to get next/stop, broadcast to N consumers.
- BufferedStream: Channel-backed buffer for backpressure.
Rules
- Public API:
iter.Seq2[T, error] — never <-chan.
- Internal goroutine communication: channels are fine.
- Always check context cancellation in producers.
yield returning false = consumer stopped — respect immediately.
- Use
iter.Pull2 only when pull semantics are genuinely needed.
Converted and distributed by TomeVault — claim your Tome and manage your conversions.
1---2name: streaming-patterns3description: Go 1.23 iter.Seq2 streaming patterns for Beluga AI v2. Use when implementing streaming, transforms, or backpressure. Use when this capability is needed.4---56# Streaming Patterns78## Primary Primitive: iter.Seq2[T, error]910```go11func (m *Model) Stream(ctx context.Context, msgs []schema.Message) iter.Seq2[schema.StreamChunk, error] {12 return func(yield func(schema.StreamChunk, error) bool) {13 stream, err := m.client.Stream(ctx, msgs)14 if err != nil { yield(schema.StreamChunk{}, err); return }15 defer stream.Close()16 for {17 select {18 case <-ctx.Done(): yield(schema.StreamChunk{}, ctx.Err()); return19 default:20 }21 chunk, err := stream.Recv()22 if err == io.EOF { return }23 if err != nil { yield(schema.StreamChunk{}, err); return }24 if !yield(convertChunk(chunk), nil) { return } // consumer stopped25 }26 }27}28```2930## Composition3132- **Pipe**: `func Pipe[A, B any](first iter.Seq2[A, error], transform func(A) (B, error)) iter.Seq2[B, error]`33- **Collect**: Stream to slice — `func Collect[T any](stream iter.Seq2[T, error]) ([]T, error)`34- **Invoke from Stream**: Stream, collect, return last.35- **Fan-out**: `iter.Pull2()` to get next/stop, broadcast to N consumers.36- **BufferedStream**: Channel-backed buffer for backpressure.3738## Rules39401. Public API: `iter.Seq2[T, error]` — never `<-chan`.412. Internal goroutine communication: channels are fine.423. Always check context cancellation in producers.434. `yield` returning false = consumer stopped — respect immediately.445. Use `iter.Pull2` only when pull semantics are genuinely needed.4546---47> Converted and distributed by [TomeVault](https://tomevault.io/claim/lookatitude) — claim your Tome and manage your conversions.48<!-- tomevault:4.0:skill_md:2026-04-11 -->