From 700b7cd5d672baa70d1e274f3f6a9941e6a03153 Mon Sep 17 00:00:00 2001 From: newkirk Date: Tue, 21 Jul 2026 13:21:39 -0400 Subject: [PATCH] fix(phase-3): enable circuit breaker for streaming requests --- internal/providers/circuit_breaker.go | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) diff --git a/internal/providers/circuit_breaker.go b/internal/providers/circuit_breaker.go index bcfb23c3..701a7b18 100644 --- a/internal/providers/circuit_breaker.go +++ b/internal/providers/circuit_breaker.go @@ -50,9 +50,13 @@ func (cbp *CircuitBreakerProvider) ChatCompletion(ctx context.Context, req *mode } func (cbp *CircuitBreakerProvider) ChatCompletionStream(ctx context.Context, req *models.UnifiedRequest) (<-chan *models.ChatCompletionStreamResponse, error) { - // Circuit breaker for streaming is tricky. We'll just call the provider directly. - // Future: Implement a way to track stream failures in the circuit breaker. - return cbp.provider.ChatCompletionStream(ctx, req) + result, err := cbp.cb.Execute(func() (interface{}, error) { + return cbp.provider.ChatCompletionStream(ctx, req) + }) + if err != nil { + return nil, err + } + return result.(<-chan *models.ChatCompletionStreamResponse), nil } func (cbp *CircuitBreakerProvider) ImageGeneration(ctx context.Context, req *models.ImageGenerationRequest) (*models.ImageGenerationResponse, error) { @@ -76,6 +80,11 @@ func (cbp *CircuitBreakerProvider) Responses(ctx context.Context, req *models.Re } func (cbp *CircuitBreakerProvider) ResponsesStream(ctx context.Context, req *models.ResponsesRequest) (<-chan *models.ResponsesStreamChunk, error) { - // Circuit breaker passthrough for streaming (same pattern as ChatCompletionStream) - return cbp.provider.ResponsesStream(ctx, req) + result, err := cbp.cb.Execute(func() (interface{}, error) { + return cbp.provider.ResponsesStream(ctx, req) + }) + if err != nil { + return nil, err + } + return result.(<-chan *models.ResponsesStreamChunk), nil }