From 2d796509ff8bab5674e07d1a9938eb94d20359a7 Mon Sep 17 00:00:00 2001 From: N Rohith Reddy Date: Fri, 17 Jul 2026 14:19:48 -0400 Subject: [PATCH 1/3] test(proxy): add failing OpenAI UsageBypass consumer coverage Lock in bypass weekly-limit fallthrough and subscription-only refuse on ProxyOpenAIChatCompletion before wiring the Messages consumer. Signed-off-by: N Rohith Reddy Co-authored-by: Cursor --- internal/proxy/openai_usage_bypass_test.go | 111 +++++++++++++++++++++ 1 file changed, 111 insertions(+) create mode 100644 internal/proxy/openai_usage_bypass_test.go diff --git a/internal/proxy/openai_usage_bypass_test.go b/internal/proxy/openai_usage_bypass_test.go new file mode 100644 index 00000000..ce4ba22d --- /dev/null +++ b/internal/proxy/openai_usage_bypass_test.go @@ -0,0 +1,111 @@ +package proxy_test + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "workweave/router/internal/billing" + "workweave/router/internal/providers" + "workweave/router/internal/proxy" + "workweave/router/internal/proxy/usage" + "workweave/router/internal/router" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func oaiBypassBody() []byte { + return []byte(`{"model":"` + bypassRequestedMdl + `","messages":[{"role":"user","content":"hi"}]}`) +} + +// TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer: gate on + headroom +// must skip the scorer and serve the requested model on the OpenAI wire. +func TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer(t *testing.T) { + svc, fr, p := bypassFixture(t, 0.20) + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader("")) + + require.NoError(t, svc.ProxyOpenAIChatCompletion(bypassCtx(0.80), oaiBypassBody(), rec, req)) + + assert.Equal(t, 0, fr.routeCalls, "scorer must not run while subscription has headroom") + require.Len(t, p.proxyBodies, 1) + assert.Contains(t, string(p.proxyBodies[0]), `"`+bypassRequestedMdl+`"`) + assert.Equal(t, "usage_bypass", rec.Header().Get(proxy.HeaderRouterDecision)) + assert.Equal(t, bypassRequestedMdl, rec.Header().Get(proxy.HeaderRouterModel)) + assert.Contains(t, rec.Body.String(), `"object":"chat.completion"`) +} + +// TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch: a buffered +// Anthropic 429 on the bypass attempt must not reach the client — clear bypass +// and serve via routeFor (Messages parity). +func TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch(t *testing.T) { + bypassResp := &providers.UpstreamErrorResponse{ + Status: http.StatusTooManyRequests, + Headers: http.Header{ + "anthropic-ratelimit-unified-weekly-limit": []string{"100000"}, + "anthropic-ratelimit-unified-weekly-reset": []string{"2025-12-31T00:00:00Z"}, + "anthropic-ratelimit-unified-weekly-remaining": []string{"0"}, + }, + Body: []byte(`{"type":"error","error":{"type":"rate_limit_error","message":"weekly limit exceeded"}}`), + } + routedResp := func(w http.ResponseWriter) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{"id":"msg_1","type":"message","role":"assistant","model":"` + bypassScorerPickMdl + `","content":[{"type":"text","text":"hi"}],"stop_reason":"end_turn","usage":{"input_tokens":1,"output_tokens":1}}`)) + } + inner := &fakeProvider{proxyResponse: routedResp} + wrappedP := &swapErrProvider{first: bypassResp, second: nil, inner: inner} + fr := &fakeRouter{decision: router.Decision{Provider: providers.ProviderAnthropic, Model: bypassScorerPickMdl, Reason: "cluster:v0.2"}} + obs := usage.NewObserver([]byte("salt"), 10*time.Minute, time.Now) + obs.Record(obs.Key([]byte(bypassSubToken)), usage.Snapshot{ + Primary: usage.Window{UsedPercent: 0.20, WindowMinutes: 300}, + }) + svc := proxy.NewService(fr, map[string]providers.Client{providers.ProviderAnthropic: wrappedP}, nil, false, nil, nil, false, providers.ProviderAnthropic, bypassScorerPickMdl, nil). + WithSubscriptionAwareRouting(obs, 0.05, 2.0) + + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader("")) + ctx := bypassCtx(0.80) + ctx = context.WithValue(ctx, proxy.ExternalIDContextKey{}, "org-oai-bypass-reroute") + ctx = context.WithValue(ctx, proxy.InstallationIDContextKey{}, uuid.New().String()) + + require.NoError(t, svc.ProxyOpenAIChatCompletion(ctx, oaiBypassBody(), rec, req)) + + assert.Equal(t, 1, fr.routeCalls, "scorer must run once on the reroute after bypass 429") + assert.NotEqual(t, http.StatusTooManyRequests, rec.Code, "the 429 must not be flushed to the client") + assert.Equal(t, bypassScorerPickMdl, rec.Header().Get(proxy.HeaderRouterModel)) + assert.Equal(t, "cluster:v0.2", rec.Header().Get(proxy.HeaderRouterDecision)) + assert.Contains(t, rec.Body.String(), `"object":"chat.completion"`) +} + +// TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402: subscription-only +// must refuse a retryable bypass failure instead of paid reroute. +func TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402(t *testing.T) { + bypassResp := &providers.UpstreamErrorResponse{ + Status: http.StatusTooManyRequests, + Body: []byte(`{"type":"error","error":{"type":"rate_limit_error","message":"weekly limit exceeded"}}`), + } + p := &fakeProvider{proxyErr: bypassResp} + wrappedP := &swapErrProvider{first: bypassResp, second: nil, inner: p} + fr := &fakeRouter{decision: router.Decision{Provider: providers.ProviderAnthropic, Model: bypassScorerPickMdl, Reason: "cluster:v0.2"}} + obs := usage.NewObserver([]byte("salt"), 10*time.Minute, time.Now) + obs.Record(obs.Key([]byte(bypassSubToken)), usage.Snapshot{ + Primary: usage.Window{UsedPercent: 0.20, WindowMinutes: 300}, + }) + svc := proxy.NewService(fr, map[string]providers.Client{providers.ProviderAnthropic: wrappedP}, nil, false, nil, nil, false, providers.ProviderAnthropic, bypassScorerPickMdl, nil). + WithSubscriptionAwareRouting(obs, 0.05, 2.0) + + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader("")) + err := svc.ProxyOpenAIChatCompletion(billing.WithSubscriptionOnly(bypassCtx(0.80)), oaiBypassBody(), rec, req) + require.Error(t, err) + assert.True(t, errors.Is(err, proxy.ErrCreditsExhaustedSubscriptionUnavailable)) + assert.Equal(t, 0, fr.routeCalls) + assert.Equal(t, 1, wrappedP.calls) +} From 7994e670e398204e8fbc4a1134274252e34877af Mon Sep 17 00:00:00 2001 From: N Rohith Reddy Date: Fri, 17 Jul 2026 14:22:04 -0400 Subject: [PATCH 2/3] fix(proxy): consume UsageBypass on OpenAI chat completions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire bypassAnthropicOpenAI after runTurnLoop so Anthropic usage-bypass turns get skip-billing pass-through and 429→routeFor fallthrough, matching ProxyMessages. Signed-off-by: N Rohith Reddy Co-authored-by: Cursor --- internal/proxy/openai_usage_bypass_test.go | 10 +- internal/proxy/service.go | 35 +++++- internal/proxy/turnloop.go | 4 +- internal/proxy/usage_bypass.go | 121 +++++++++++++++++++++ 4 files changed, 159 insertions(+), 11 deletions(-) diff --git a/internal/proxy/openai_usage_bypass_test.go b/internal/proxy/openai_usage_bypass_test.go index ce4ba22d..b2dde8ca 100644 --- a/internal/proxy/openai_usage_bypass_test.go +++ b/internal/proxy/openai_usage_bypass_test.go @@ -24,8 +24,7 @@ func oaiBypassBody() []byte { return []byte(`{"model":"` + bypassRequestedMdl + `","messages":[{"role":"user","content":"hi"}]}`) } -// TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer: gate on + headroom -// must skip the scorer and serve the requested model on the OpenAI wire. +// TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer: gate on + headroom skips scorer on OpenAI wire. func TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer(t *testing.T) { svc, fr, p := bypassFixture(t, 0.20) rec := httptest.NewRecorder() @@ -41,9 +40,7 @@ func TestProxyOpenAI_UsageBypass_BelowThreshold_SkipsScorer(t *testing.T) { assert.Contains(t, rec.Body.String(), `"object":"chat.completion"`) } -// TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch: a buffered -// Anthropic 429 on the bypass attempt must not reach the client — clear bypass -// and serve via routeFor (Messages parity). +// TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch: bypass 429 must reroute via routeFor. func TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch(t *testing.T) { bypassResp := &providers.UpstreamErrorResponse{ Status: http.StatusTooManyRequests, @@ -84,8 +81,7 @@ func TestProxyOpenAI_UsageBypass_WeeklyLimit_FallsBackToRoutedDispatch(t *testin assert.Contains(t, rec.Body.String(), `"object":"chat.completion"`) } -// TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402: subscription-only -// must refuse a retryable bypass failure instead of paid reroute. +// TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402: subscription-only refuses retryable bypass failure. func TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402(t *testing.T) { bypassResp := &providers.UpstreamErrorResponse{ Status: http.StatusTooManyRequests, diff --git a/internal/proxy/service.go b/internal/proxy/service.go index 993c0039..48022767 100644 --- a/internal/proxy/service.go +++ b/internal/proxy/service.go @@ -4186,11 +4186,42 @@ func (s *Service) ProxyOpenAIChatCompletion(ctx context.Context, body []byte, w } routeStart := time.Now() routeRes, err := s.runTurnLoop(ctx, env, feats, apiKeyID, installationID, subAgentHint, r.Header, routeRequest) - routeMs := time.Since(routeStart).Milliseconds() if err != nil { - log.Error("Routing failed for OpenAI request", "err", err, "route_ms", routeMs, "requested_model", feats.Model, "total_input_tokens", feats.Tokens) + log.Error("Routing failed for OpenAI request", "err", err, "route_ms", time.Since(routeStart).Milliseconds(), "requested_model", feats.Model, "total_input_tokens", feats.Tokens) return err } + + // Anthropic UsageBypass consumer — same contract as ProxyMessages. + if routeRes.UsageBypass && routeRes.Decision.Provider == providers.ProviderAnthropic { + err := s.bypassAnthropicOpenAI(ctx, env, feats, routeRes.modelSwitched(), requestStart, requestID, externalID, r, w) + if !errors.Is(err, errBypassRetryable) { + s.firePolicyShadowForServingDecision(ctx, routeRes.Decision, routeRequest) + return err + } + if billing.SubscriptionOnlyFromContext(ctx) { + log.Info("Subscription-only OpenAI bypass hit retryable error; refusing instead of paid reroute", + "request_id", requestID, "external_id", externalID) + return ErrCreditsExhaustedSubscriptionUnavailable + } + routeRequest.SubsidizedModelCostFactor = s.subsidyFactors(ctx, r.Header) + if s.pinStore != nil { + role := roleForTier(catalog.TierFor(feats.Model)) + pin, _ := s.loadPin(ctx, sessionKey, role) + hmmHistory := s.loadHMMHistory(ctx, sessionKey, role) + routeRes.SessionKey = sessionKey + routeRes.PriorServedModel, routeRes.SessionEverSwitched = switchHistoryFromPins(pin, hmmHistory) + } + routeRes.UsageBypass = false + decision, rerouteErr := s.routeFor(ctx, routeRequest) + if rerouteErr != nil { + log.Error("Reroute after OpenAI usage-bypass failure failed", "err", rerouteErr) + return rerouteErr + } + routeRes.Decision = decision + routeRes.Fresh = decision + } + + routeMs := time.Since(routeStart).Milliseconds() routeRes.SuggestionMode = r.Header.Get("x-weave-suggestion-mode") == "true" decision := routeRes.Decision s.firePolicyShadowForServingDecision(ctx, decision, routeRequest) diff --git a/internal/proxy/turnloop.go b/internal/proxy/turnloop.go index 562b91e0..b9f651fb 100644 --- a/internal/proxy/turnloop.go +++ b/internal/proxy/turnloop.go @@ -57,8 +57,8 @@ type turnLoopResult struct { StickyHit bool HardPinned bool // UsageBypass is true when the caller's own subscription has headroom: - // ProxyMessages must serve the requested model straight through with no - // billing debit, bypassing Decision's normal dispatch. + // ProxyMessages / ProxyOpenAIChatCompletion must serve the requested model + // straight through with no billing debit, bypassing Decision's normal dispatch. UsageBypass bool PinTier string PinAgeSec int64 diff --git a/internal/proxy/usage_bypass.go b/internal/proxy/usage_bypass.go index 283367ef..514905fc 100644 --- a/internal/proxy/usage_bypass.go +++ b/internal/proxy/usage_bypass.go @@ -367,3 +367,124 @@ func (s *Service) bypassToAnthropic( ) return proxyErr } + +// bypassAnthropicOpenAI is OpenAI-wire bypassToAnthropic: Anthropic upstream, OpenAI response, skip billing. +func (s *Service) bypassAnthropicOpenAI( + ctx context.Context, + env *translate.RequestEnvelope, + feats translate.RoutingFeatures, + modelSwitched bool, + requestStart time.Time, + requestID, externalID string, + r *http.Request, + w http.ResponseWriter, +) error { + log := observability.FromContext(ctx) + decision := router.Decision{ + Provider: providers.ProviderAnthropic, + Model: feats.Model, + Reason: "usage_bypass", + } + w.Header().Set(HeaderRouterDecision, decision.Reason) + w.Header().Set(HeaderRouterProvider, decision.Provider) + w.Header().Set(HeaderRouterModel, decision.Model) + + p, provErr := s.provider(providers.ProviderAnthropic) + if provErr != nil { + return provErr + } + + ctx = resolveAndInjectCredentials(ctx, decision.Provider, r.Header) + + outputReserve := contextWindowOutputReserve + if feats.MaxTokens > outputReserve { + outputReserve = feats.MaxTokens + } + opts := translate.EmitOptions{ + TargetModel: decision.Model, + TargetProvider: decision.Provider, + Capabilities: router.Lookup(decision.Model), + IncludeStreamUsage: s.usageRequired(), + EnableExtendedContext: shouldEnableExtendedContext(env.FullTokenEstimate(), outputReserve), + ModelSwitched: modelSwitched, + } + prep, emitErr := env.PrepareAnthropic(r.Header, opts) + if emitErr != nil { + log.Error("Failed to emit Anthropic body on OpenAI usage-bypass path", "err", emitErr) + return fmt.Errorf("emit bypass body: %w", emitErr) + } + + sink := http.ResponseWriter(w) + if billing.SubscriptionOnlyFromContext(ctx) { + mw := translate.NewOpenAIRoutingMarkerWriter(w, decision.Model, subscriptionOnlyWarningMarker) + if err := mw.Prelude(env.Stream()); err != nil { + log.Error("OpenAI usage-bypass routing-marker prelude failed", "err", err) + } + sink = mw + } + + var extractor *otel.UsageExtractor + var usageSink otel.UsageSink + if s.usageRequired() { + extractor = otel.NewUsageExtractor(nil, providers.ProviderAnthropic) + usageSink = extractor + } + translator := translate.NewSSETranslator(sink, decision.Model, usageSink) + + proxyStart := time.Now() + proxyErr := p.Proxy(ctx, decision, prep, translator, r) + if providers.IsRetryable(proxyErr) { + return errBypassRetryable + } + proxyErr = finalizeAfterProxy(proxyErr, translator.Finalize) + + var upstreamErr *providers.UpstreamErrorResponse + if errors.As(proxyErr, &upstreamErr) { + flushBufferedIfPresent(w, proxyErr) + proxyErr = nil + } + + in, out := extractor.Tokens() + cacheCreation, cacheRead := extractor.CacheTokens() + pricing, _ := catalog.PriceFor(decision.Provider, decision.Model) + inputCost := catalog.EffectiveInputCost(in, cacheCreation, cacheRead, pricing.InputUSDPer1M, pricing, decision.Provider) + outputCost := catalog.EffectiveOutputCost(out, pricing.OutputUSDPer1M) + + clientID := ClientIdentityFrom(ctx) + otel.Record(ctx, otel.Span{ + Name: "router.usage_bypass", + Start: requestStart, + End: time.Now(), + Attrs: otel.NewAttrBuilder(18). + String("request_id", requestID). + String("external_id", externalID). + String("router_user_id", auth.UserIDFrom(ctx)). + String("client.app", clientID.ClientApp). + String("client.session_id", clientID.SessionID). + String("requested.model", decision.Model). + String("decision.model", decision.Model). + String("decision.provider", decision.Provider). + String("decision.reason", decision.Reason). + Bool("cost.subscription_served", servedOnSubscription(ctx)). + Int64("usage.input_tokens", int64(in)). + Int64("usage.output_tokens", int64(out)). + Int64("usage.cache_creation_input_tokens", int64(cacheCreation)). + Int64("usage.cache_read_input_tokens", int64(cacheRead)). + Float64("cost.requested_input_usd", inputCost). + Float64("cost.requested_output_usd", outputCost). + Float64("cost.actual_input_usd", inputCost). + Float64("cost.actual_output_usd", outputCost). + Build(), + }) + otel.Flush(ctx) + log.Info("ProxyOpenAIChatCompletion usage-bypass complete", + "request_id", requestID, + "external_id", externalID, + "requested_model", feats.Model, + "decision_model", decision.Model, + "proxy_ms", time.Since(proxyStart).Milliseconds(), + "total_ms", time.Since(requestStart).Milliseconds(), + "proxy_err", proxyErr, + ) + return proxyErr +} From 6590eb73f806f05a0907025682ab42434d91399a Mon Sep 17 00:00:00 2001 From: N Rohith Reddy Date: Fri, 17 Jul 2026 14:27:41 -0400 Subject: [PATCH 3/3] fix(proxy): buffer OpenAI usage-bypass Prelude behind preludeBuffer Subscription-only streaming marker must not commit SSE before upstream; Discard on retryable bypass so 402 refuse leaves no partial stream. Signed-off-by: N Rohith Reddy Co-authored-by: Cursor --- internal/proxy/openai_usage_bypass_test.go | 28 ++++++++++++++++++++++ internal/proxy/usage_bypass.go | 20 +++++++++++++--- 2 files changed, 45 insertions(+), 3 deletions(-) diff --git a/internal/proxy/openai_usage_bypass_test.go b/internal/proxy/openai_usage_bypass_test.go index b2dde8ca..cc64df16 100644 --- a/internal/proxy/openai_usage_bypass_test.go +++ b/internal/proxy/openai_usage_bypass_test.go @@ -105,3 +105,31 @@ func TestSubscriptionOnly_OpenAI_BypassRetryable_Refuses402(t *testing.T) { assert.Equal(t, 0, fr.routeCalls) assert.Equal(t, 1, wrappedP.calls) } + +// TestSubscriptionOnly_OpenAI_BypassRetryable_Stream_NoPartialCommit: streaming Prelude stays buffered until Discard on 402. +func TestSubscriptionOnly_OpenAI_BypassRetryable_Stream_NoPartialCommit(t *testing.T) { + bypassResp := &providers.UpstreamErrorResponse{ + Status: http.StatusTooManyRequests, + Body: []byte(`{"type":"error","error":{"type":"rate_limit_error","message":"weekly limit exceeded"}}`), + } + p := &fakeProvider{proxyErr: bypassResp} + wrappedP := &swapErrProvider{first: bypassResp, second: nil, inner: p} + fr := &fakeRouter{decision: router.Decision{Provider: providers.ProviderAnthropic, Model: bypassScorerPickMdl, Reason: "cluster:v0.2"}} + obs := usage.NewObserver([]byte("salt"), 10*time.Minute, time.Now) + obs.Record(obs.Key([]byte(bypassSubToken)), usage.Snapshot{ + Primary: usage.Window{UsedPercent: 0.20, WindowMinutes: 300}, + }) + svc := proxy.NewService(fr, map[string]providers.Client{providers.ProviderAnthropic: wrappedP}, nil, false, nil, nil, false, providers.ProviderAnthropic, bypassScorerPickMdl, nil). + WithSubscriptionAwareRouting(obs, 0.05, 2.0) + + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/v1/chat/completions", strings.NewReader("")) + body := []byte(`{"model":"` + bypassRequestedMdl + `","stream":true,"messages":[{"role":"user","content":"hi"}]}`) + err := svc.ProxyOpenAIChatCompletion(billing.WithSubscriptionOnly(bypassCtx(0.80)), body, rec, req) + require.Error(t, err) + assert.True(t, errors.Is(err, proxy.ErrCreditsExhaustedSubscriptionUnavailable)) + assert.Equal(t, 0, fr.routeCalls) + assert.Equal(t, 1, wrappedP.calls) + assert.Empty(t, rec.Body.String(), "retryable bypass must Discard buffered Prelude; client must see no SSE bytes") + assert.False(t, rec.Flushed, "HTTP status must not be committed before the 402 mapping") +} diff --git a/internal/proxy/usage_bypass.go b/internal/proxy/usage_bypass.go index 514905fc..03a72aa9 100644 --- a/internal/proxy/usage_bypass.go +++ b/internal/proxy/usage_bypass.go @@ -414,9 +414,12 @@ func (s *Service) bypassAnthropicOpenAI( return fmt.Errorf("emit bypass body: %w", emitErr) } - sink := http.ResponseWriter(w) + // Buffer synthetic preamble (subscription-only marker) so a retryable bypass + // failure can Discard and return 402 without a partial SSE stream on the wire. + preludeBuf := newPreludeBuffer(w) + sink := http.ResponseWriter(preludeBuf) if billing.SubscriptionOnlyFromContext(ctx) { - mw := translate.NewOpenAIRoutingMarkerWriter(w, decision.Model, subscriptionOnlyWarningMarker) + mw := translate.NewOpenAIRoutingMarkerWriter(preludeBuf, decision.Model, subscriptionOnlyWarningMarker) if err := mw.Prelude(env.Stream()); err != nil { log.Error("OpenAI usage-bypass routing-marker prelude failed", "err", err) } @@ -430,17 +433,28 @@ func (s *Service) bypassAnthropicOpenAI( usageSink = extractor } translator := translate.NewSSETranslator(sink, decision.Model, usageSink) + preludeBuf.Seal() proxyStart := time.Now() proxyErr := p.Proxy(ctx, decision, prep, translator, r) if providers.IsRetryable(proxyErr) { + if !preludeBuf.Committed() { + preludeBuf.Discard() + } return errBypassRetryable } proxyErr = finalizeAfterProxy(proxyErr, translator.Finalize) var upstreamErr *providers.UpstreamErrorResponse if errors.As(proxyErr, &upstreamErr) { - flushBufferedIfPresent(w, proxyErr) + if !preludeBuf.Committed() { + preludeBuf.Discard() + flushBufferedIfPresent(w, proxyErr) + } else if env.Stream() { + _ = emitOpenAISSEErrorEvent(sink, proxyErr) + } else { + flushBufferedIfPresent(w, proxyErr) + } proxyErr = nil }