Implementing Bidirectional Streaming RPCs
1. Receiving and Sending Concurrently
Example: Bidi handler with goroutines
func (s *server) Chat(stream chatv1.ChatService_ChatServer) error {
errCh := make(chan error, 2)
go func() { // sender
for msg := range s.broker.Subscribe() {
if err := stream.Send(msg); err != nil { errCh <- err; return }
}
errCh <- nil
}()
go func() { // receiver
for {
m, err := stream.Recv()
if err == io.EOF { errCh <- nil; return }
if err != nil { errCh <- err; return }
s.broker.Publish(m)
}
}()
return <-errCh
}
| Rule | Detail |
|---|---|
| Send / Recv | Safe to call concurrently from different goroutines |
| Multiple senders | Not allowed — serialize via a channel |
2. Creating Receive Loop
| Pattern | Detail |
|---|---|
| For loop | Until io.EOF or error |
| Forward to channel | Decouple processing |
3. Creating Send Loop
| Pattern | Detail |
|---|---|
| Range over channel | Send each event |
| Check ctx.Done() | Exit on cancel |
4. Coordinating Goroutines
| Tool | Use |
|---|---|
errgroup.Group | First error cancels siblings |
sync.WaitGroup | Wait for both to finish |
Buffered errCh | Non-blocking error report |
5. Handling Stream Completion
| End From | Indicator |
|---|---|
| Client done | Recv returns io.EOF |
| Server done | Handler returns; client Recv gets io.EOF |
| Either errored | Counterpart gets status error |
6. Detecting Stream Errors
| Symptom | Cause |
|---|---|
Send returns io.EOF | Server closed stream — read final status |
| Recv returns non-EOF | RPC failed |
| ctx canceled | Local/remote cancellation |
7. Creating Client Bidirectional Stream
Example: Client bidi loop
stream, _ := client.Chat(ctx)
go func() {
for input := range outbox {
if err := stream.Send(input); err != nil { return }
}
_ = stream.CloseSend()
}()
for {
m, err := stream.Recv()
if err == io.EOF { break }
if err != nil { return err }
handle(m)
}
8. Implementing Request-Response Pattern
| Approach | Detail |
|---|---|
| Correlation ID | Include request_id in each message |
| Pending map | Match responses to waiting callers |
| Timeout per request | Per-request context with deadline |
9. Implementing Full-Duplex Pattern
| Use Case | Pattern |
|---|---|
| Pub/sub | Both sides push independent streams |
| Chat | Echo + broadcast model |
| Market data | Subscribe requests + price ticks |
10. Closing Bidirectional Stream
| Side | Method |
|---|---|
| Client send-side | stream.CloseSend() |
| Server | Return from handler |
| Abort either side | Cancel ctx → status Canceled |