diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index c23ed9a..2ee9143 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -19,7 +19,7 @@ concurrency: cancel-in-progress: true env: - GO_VERSION: "1.26.5" + GO_VERSION: "1.26.6" GOPRIVATE: "github.com/GrayCodeAI/*" GONOSUMDB: "github.com/GrayCodeAI/*" GONOSUMCHECK: "1" diff --git a/CHANGELOG.md b/CHANGELOG.md index 71c289b..336b8e5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,50 @@ Format: [Keep a Changelog](https://keepachangelog.com/en/1.0.0/) · Versioning: ## [Unreleased] +### Changed — Shared MiMo auth-retry helper (2026-08-16) +- **Deduplicated `doRequestWithMimoAuthRetry`** between the OpenAI and + Anthropic adapters into one `doWithMimoAuthRetry` helper (client/adapters, + next to `mimoAuthHeaders`); the two adapters now differ only in the Bearer + headers they apply to the 401 retry. No behavior change. + +### Fixed — Gemini stream request IDs (2026-08-16) +- **Gemini `StreamChat` now propagates the provider request ID.** The client + captured `X-Goog-Request-Id` from the response headers but passed an empty + string to the stream result, so hosts lost the correlation ID on + successful streams (it was only preserved on errors). Both the shared + parser path and the legacy opt-out parser now carry it. + +### Fixed — Non-fatal stream diagnostics no longer fail the stream (2026-08-16) +- **Stream health diagnostics are now warnings, not terminal errors.** + `client/core`'s OpenAI stream processor emits end-of-stream diagnostics + (reasoning-only responses, empty responses) as error-type events followed + by the terminal `done` — but the engine mapped *every* error event to + `provider_unavailable`, stopped forwarding, and set `Err()` even though + content had been delivered. Diagnostic events are now marked non-fatal via + the existing `EyrieStreamEvent.Warning` field (additive); the engine + forwards them as `warning` events and still delivers the final + `done`/usage event with `Err()` unset. Genuinely fatal stream errors keep + the previous behavior. The deprecated client continuation helper and the + tracing middleware treat warning-marked events the same way. + +### Fixed — Concentrate adapter robustness (2026-08-16) +- **Concentrate Responses client now uses the shared pooled HTTP client** + (`core.NewPooledHTTPClient(core.DefaultTimeout)`) instead of a private + `&http.Client{Timeout: 120s}` literal — long streams are no longer cut off + at 2 minutes and connections reuse the process-wide transport pool like + every other adapter. +- **Concentrate requests are retried via `core.DoWithRetry`** (chat and + stream paths) on 429/500/502/503/529 with backoff and `Retry-After` + support; `SetRetry` previously discarded the config with a comment claiming + the HTTP client handled retries (it never does). +- **Concentrate errors are structured `*core.EyrieError`s** built by + `core.ParseProviderError`/`core.FormatAPIError` (8KB bounded read, + provider/op/status/request-ID preserved), so `IsRetriable()`/`IsAuthError()` + and the engine's error classification work; the captured `X-Request-Id` is + also propagated to stream results. +- **`normalizeToolParams` no longer mutates the caller's tool schema map** — + a shallow copy gets `additionalProperties:false` injected for strict mode. + ### Added — Round 3 ecosystem improvements (2026-06-06) - **Reasoning controls** — `reasoning_effort` and Anthropic extended-thinking `thinking_budget_tokens` passthrough on `ChatOptions` (omitted when unset). diff --git a/client/adapters/anthropic.go b/client/adapters/anthropic.go index 3014243..a8c4abc 100644 --- a/client/adapters/anthropic.go +++ b/client/adapters/anthropic.go @@ -629,27 +629,12 @@ func (c *AnthropicClient) StreamChat(ctx context.Context, messages []core.EyrieM } func (c *AnthropicClient) doRequestWithMimoAuthRetry(ctx context.Context, req *http.Request, body []byte) (*http.Response, error) { - resp, err := core.DoWithRetry(ctx, c.httpClient, req, c.retry, c.logger) - if err != nil { - return nil, err - } - if !c.useMimoAuth || resp.StatusCode != http.StatusUnauthorized { - return resp, nil - } - _ = resp.Body.Close() - req2, err := http.NewRequestWithContext(ctx, req.Method, req.URL.String(), bytes.NewReader(body)) - if err != nil { - return nil, err - } - req2.Header.Set("Content-Type", "application/json") - req2.Header.Set("Authorization", "Bearer "+c.apiKey) - req2.Header.Set("Anthropic-Version", c.version) - req2.Header.Set("User-Agent", core.UserAgent()) - if req.Header.Get("Accept") != "" { - req2.Header.Set("Accept", req.Header.Get("Accept")) - } - req2.GetBody = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(body)), nil } - return core.DoWithRetry(ctx, c.httpClient, req2, c.retry, c.logger) + return doWithMimoAuthRetry(ctx, c.httpClient, c.retry, c.logger, c.useMimoAuth, req, body, func(req2 *http.Request) { + req2.Header.Set("Content-Type", "application/json") + req2.Header.Set("Authorization", "Bearer "+c.apiKey) + req2.Header.Set("Anthropic-Version", c.version) + req2.Header.Set("User-Agent", core.UserAgent()) + }) } // Ping checks connectivity to the Anthropic API using a lightweight GET request. diff --git a/client/adapters/concentrate_responses.go b/client/adapters/concentrate_responses.go index 83cb4b1..a806587 100644 --- a/client/adapters/concentrate_responses.go +++ b/client/adapters/concentrate_responses.go @@ -25,6 +25,7 @@ type ConcentrateResponsesClient struct { baseURL string apiKey string httpClient *http.Client + retry core.RetryConfig logger *slog.Logger } @@ -35,7 +36,8 @@ func NewConcentrateResponsesClient(apiKey, baseURL string, opts ...core.ClientOp c := &ConcentrateResponsesClient{ baseURL: baseURL, apiKey: apiKey, - httpClient: &http.Client{Timeout: 120 * time.Second}, + httpClient: core.NewPooledHTTPClient(core.DefaultTimeout), + retry: core.DefaultRetryConfig(), logger: slog.Default(), } for _, opt := range opts { @@ -179,16 +181,19 @@ func (c *ConcentrateResponsesClient) Chat(ctx context.Context, messages []core.E httpReq.Header.Set("Authorization", "Bearer "+c.apiKey) httpReq.Header.Set("Accept", "application/json") httpReq.Header.Set("User-Agent", "eyrie-model-catalog/1.0") + httpReq.GetBody = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(body)), nil } - resp, err := c.httpClient.Do(httpReq) + resp, err := core.DoWithRetry(ctx, c.httpClient, httpReq, c.retry, c.logger) if err != nil { return nil, fmt.Errorf("concentrate: request failed: %w", err) } defer func() { _ = resp.Body.Close() }() + requestID := resp.Header.Get("X-Request-Id") + if resp.StatusCode != http.StatusOK { - body, _ := io.ReadAll(resp.Body) - return nil, fmt.Errorf("concentrate: request failed (%d): %s", resp.StatusCode, string(body)) + detail, readErr := core.ParseProviderError(resp.Body) + return nil, core.FormatAPIError("concentrate", "chat", resp.StatusCode, requestID, detail, readErr) } var apiResp responsesResponse @@ -221,21 +226,24 @@ func (c *ConcentrateResponsesClient) StreamChat(ctx context.Context, messages [] httpReq.Header.Set("Authorization", "Bearer "+c.apiKey) httpReq.Header.Set("Accept", "text/event-stream") httpReq.Header.Set("User-Agent", "eyrie-model-catalog/1.0") + httpReq.GetBody = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(body)), nil } - resp, err := c.httpClient.Do(httpReq) + resp, err := core.DoWithRetry(streamCtx, c.httpClient, httpReq, c.retry, c.logger) if err != nil { cancel() - return nil, fmt.Errorf("concentrate: request failed: %w", err) + return nil, fmt.Errorf("concentrate: stream request failed: %w", err) } + requestID := resp.Header.Get("X-Request-Id") + if resp.StatusCode != http.StatusOK { - body, _ := io.ReadAll(resp.Body) - resp.Body.Close() + detail, readErr := core.ParseProviderError(resp.Body) + _ = resp.Body.Close() cancel() - return nil, fmt.Errorf("concentrate: stream request failed (%d): %s", resp.StatusCode, string(body)) + return nil, core.FormatAPIError("concentrate", "stream", resp.StatusCode, requestID, detail, readErr) } - return c.handleStream(streamCtx, cancel, resp), nil + return c.handleStream(streamCtx, cancel, resp, requestID), nil } // Ping checks the health of the Concentrate API. @@ -354,6 +362,8 @@ func concentrateToolChoice(choice *core.ToolChoiceOption) interface{} { // normalizeToolParams ensures tool parameters conform to Concentrate's strict mode // requirements: additionalProperties must be false at the top level when strict=true. +// The input map is never mutated: a shallow copy is returned so the caller's +// tool definition (which may be reused across requests) stays intact. // See: https://concentrate.ai/docs/api-reference/endpoint/tool-calling func normalizeToolParams(params map[string]interface{}) map[string]interface{} { if params == nil { @@ -362,7 +372,12 @@ func normalizeToolParams(params map[string]interface{}) map[string]interface{} { // Only enforce for object-typed schemas if t, ok := params["type"]; ok && t == "object" { if _, has := params["additionalProperties"]; !has { - params["additionalProperties"] = false + normalized := make(map[string]interface{}, len(params)+1) + for k, v := range params { + normalized[k] = v + } + normalized["additionalProperties"] = false + return normalized } } return params @@ -517,7 +532,7 @@ type streamEvent struct { ContentIndex int `json:"content_index,omitempty"` } -func (c *ConcentrateResponsesClient) handleStream(ctx context.Context, cancel context.CancelFunc, resp *http.Response) *core.StreamResult { +func (c *ConcentrateResponsesClient) handleStream(ctx context.Context, cancel context.CancelFunc, resp *http.Response, requestID string) *core.StreamResult { events := make(chan core.EyrieStreamEvent, core.StreamChannelBuffer) go func() { @@ -639,7 +654,7 @@ func (c *ConcentrateResponsesClient) handleStream(ctx context.Context, cancel co } }() - return llm.NewStreamResult(events, "", cancel) + return llm.NewStreamResult(events, requestID, cancel) } func sendConcentrateStreamEvent(ctx context.Context, events chan<- core.EyrieStreamEvent, event core.EyrieStreamEvent) bool { @@ -722,7 +737,7 @@ func (c *ConcentrateResponsesClient) SetHTTPClient(hc *http.Client) { // SetRetry implements core.Configurable. func (c *ConcentrateResponsesClient) SetRetry(rc core.RetryConfig) { - // Retries handled at HTTP client level + c.retry = rc } // SetLogger implements core.Configurable. @@ -752,7 +767,7 @@ func (c *ConcentrateResponsesClient) HTTPClient() *http.Client { // Retry implements core.Configurable. func (c *ConcentrateResponsesClient) Retry() core.RetryConfig { - return core.RetryConfig{} + return c.retry } // Logger implements core.Configurable. diff --git a/client/adapters/concentrate_responses_test.go b/client/adapters/concentrate_responses_test.go index c98148a..254f2e5 100644 --- a/client/adapters/concentrate_responses_test.go +++ b/client/adapters/concentrate_responses_test.go @@ -3,8 +3,11 @@ package adapters import ( "context" "encoding/json" + "errors" + "fmt" "io" "net/http" + "net/http/httptest" "strings" "testing" "time" @@ -26,6 +29,20 @@ func TestNewConcentrateResponsesClient(t *testing.T) { } } +func TestNewConcentrateResponsesClient_UsesSharedPooledHTTPClient(t *testing.T) { + t.Parallel() + client := NewConcentrateResponsesClient("cn-key", "https://api.concentrate.ai/v1") + if client.httpClient.Timeout != core.DefaultTimeout { + t.Errorf("timeout = %v, want default %v", client.httpClient.Timeout, core.DefaultTimeout) + } + if client.httpClient.Transport != core.NewPooledHTTPClient(0).Transport { + t.Error("client does not use the shared pooled transport") + } + if client.Retry().MaxRetries != core.DefaultRetryConfig().MaxRetries { + t.Error("client does not default to the shared retry config") + } +} + func TestConcentrateResponsesClient_ChatUsesResponsesContract(t *testing.T) { t.Parallel() transport := roundTripFunc(func(req *http.Request) (*http.Response, error) { @@ -311,3 +328,193 @@ func TestConcentrateResponsesClient_PingDoesNotRequireAuth(t *testing.T) { t.Fatal(err) } } + +func TestConcentrateResponsesClient_ChatRetriesOn500ThenSucceeds(t *testing.T) { + t.Parallel() + attempts := 0 + transport := roundTripFunc(func(req *http.Request) (*http.Response, error) { + attempts++ + if attempts == 1 { + return jsonResponse(http.StatusInternalServerError, map[string]any{ + "error": map[string]string{"message": "upstream exploded"}, + }), nil + } + // The retried request must still carry the full body (GetBody path). + var body map[string]interface{} + if err := jsonDecodeRequest(req, &body); err != nil { + t.Fatalf("decode retried body: %v", err) + } + if body["model"] != "gpt-5" { + t.Fatalf("retried body = %#v", body) + } + return jsonResponse(http.StatusOK, map[string]any{ + "id": "resp_retry", + "status": "completed", + "output": []map[string]any{{ + "type": "message", "role": "assistant", + "content": []map[string]any{{"type": "output_text", "text": "ok"}}, + }}, + }), nil + }) + client := NewConcentrateResponsesClient("cn-key", "https://api.concentrate.ai/v1") + client.httpClient = &http.Client{Transport: transport} + client.SetRetry(core.NewRetryConfig(2, time.Millisecond, 2*time.Millisecond, 500)) + + resp, err := client.Chat(context.Background(), []core.EyrieMessage{{Role: "user", Content: "Hi"}}, core.ChatOptions{Model: "gpt-5"}) + if err != nil { + t.Fatalf("Chat: %v", err) + } + if resp.Content != "ok" { + t.Fatalf("content = %q", resp.Content) + } + if attempts != 2 { + t.Fatalf("attempts = %d, want 2 (one 500 + one success)", attempts) + } +} + +func TestConcentrateResponsesClient_ChatErrorIsStructuredEyrieError(t *testing.T) { + t.Parallel() + transport := roundTripFunc(func(*http.Request) (*http.Response, error) { + resp := jsonResponse(http.StatusUnauthorized, map[string]any{ + "error": map[string]string{"code": "invalid_api_key", "message": "bad key"}, + }) + resp.Header.Set("X-Request-Id", "req_abc") + return resp, nil + }) + client := NewConcentrateResponsesClient("cn-key", "https://api.concentrate.ai/v1") + client.httpClient = &http.Client{Transport: transport} + client.SetRetry(core.RetryConfig{}) // no retries: classify the terminal error + + _, err := client.Chat(context.Background(), []core.EyrieMessage{{Role: "user", Content: "Hi"}}, core.ChatOptions{Model: "gpt-5"}) + if err == nil { + t.Fatal("expected error") + } + var eyrieErr *core.EyrieError + if !errors.As(err, &eyrieErr) { + t.Fatalf("error is %T, want *core.EyrieError (%v)", err, err) + } + if eyrieErr.Provider != "concentrate" || eyrieErr.Op != "chat" { + t.Fatalf("provider/op = %s/%s", eyrieErr.Provider, eyrieErr.Op) + } + if eyrieErr.StatusCode != http.StatusUnauthorized { + t.Fatalf("status = %d, want 401", eyrieErr.StatusCode) + } + if !eyrieErr.IsAuthError() { + t.Error("IsAuthError() = false, want true") + } + if eyrieErr.IsRetriable() { + t.Error("IsRetriable() = true for 401, want false") + } + if eyrieErr.RequestID != "req_abc" { + t.Errorf("request id = %q, want req_abc", eyrieErr.RequestID) + } + if !strings.Contains(eyrieErr.Message, "bad key") { + t.Errorf("message = %q, want it to carry the provider detail", eyrieErr.Message) + } +} + +func TestConcentrateResponsesClient_StreamErrorIsStructuredEyrieError(t *testing.T) { + t.Parallel() + transport := roundTripFunc(func(*http.Request) (*http.Response, error) { + resp := jsonResponse(http.StatusTooManyRequests, map[string]any{ + "error": map[string]string{"type": "rate_limit_error", "message": "slow down"}, + }) + resp.Header.Set("X-Request-Id", "req_429") + return resp, nil + }) + client := NewConcentrateResponsesClient("cn-key", "https://api.concentrate.ai/v1") + client.httpClient = &http.Client{Transport: transport} + client.SetRetry(core.RetryConfig{}) // no retries: classify the terminal error + + _, err := client.StreamChat(context.Background(), nil, core.ChatOptions{Model: "gpt-5"}) + if err == nil { + t.Fatal("expected error") + } + var eyrieErr *core.EyrieError + if !errors.As(err, &eyrieErr) { + t.Fatalf("error is %T, want *core.EyrieError (%v)", err, err) + } + if eyrieErr.Op != "stream" { + t.Fatalf("op = %s, want stream", eyrieErr.Op) + } + if eyrieErr.StatusCode != http.StatusTooManyRequests || !eyrieErr.IsRateLimited() || !eyrieErr.IsRetriable() { + t.Fatalf("status = %d (rate-limited=%v, retriable=%v), want 429/true/true", + eyrieErr.StatusCode, eyrieErr.IsRateLimited(), eyrieErr.IsRetriable()) + } + if eyrieErr.RequestID != "req_429" { + t.Errorf("request id = %q, want req_429", eyrieErr.RequestID) + } +} + +// The adapter must not impose a short whole-response timeout: a stream whose +// wall time exceeds the previous hard-coded 120s-class wiring still delivers +// every event plus the terminal done. Gaps are kept small so the test stays +// fast; the default pooled client (core.DefaultTimeout) is exercised as-is. +func TestConcentrateResponsesClient_StreamSurvivesSlowServer(t *testing.T) { + t.Parallel() + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("X-Request-Id", "req_slow") + flusher := w.(http.Flusher) + send := func(payload string) { + fmt.Fprintf(w, "data: %s\n\n", payload) + flusher.Flush() + } + send(`{"type":"response.output_text.delta","delta":"slow"}`) + time.Sleep(150 * time.Millisecond) + send(`{"type":"response.output_text.delta","delta":" but steady"}`) + time.Sleep(150 * time.Millisecond) + send(`{"type":"response.completed","response":{"status":"completed"}}`) + })) + defer server.Close() + + client := NewConcentrateResponsesClient("cn-key", server.URL) + result, err := client.StreamChat(context.Background(), nil, core.ChatOptions{Model: "gpt-5"}) + if err != nil { + t.Fatalf("StreamChat: %v", err) + } + defer result.Close() + + var content strings.Builder + sawDone := false + for evt := range result.Events { + switch evt.Type { + case "content": + content.WriteString(evt.Content) + case "done": + sawDone = true + } + } + if got := content.String(); got != "slow but steady" { + t.Fatalf("content = %q", got) + } + if !sawDone { + t.Fatal("stream never delivered the terminal done event") + } + if result.RequestID != "req_slow" { + t.Fatalf("stream request id = %q, want req_slow", result.RequestID) + } +} + +func TestNormalizeToolParamsDoesNotMutateCallerMap(t *testing.T) { + t.Parallel() + original := map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{"path": map[string]interface{}{"type": "string"}}, + } + normalized := normalizeToolParams(original) + if _, has := original["additionalProperties"]; has { + t.Fatal("caller's parameter map was mutated in place") + } + if got := normalized["additionalProperties"]; got != false { + t.Fatalf("normalized additionalProperties = %#v, want false", got) + } + if _, ok := normalized["properties"].(map[string]interface{}); !ok { + t.Fatalf("normalized properties = %#v, want the original nested schema", normalized["properties"]) + } + + explicit := map[string]interface{}{"type": "object", "additionalProperties": true} + if got := normalizeToolParams(explicit); got["additionalProperties"] != true { + t.Fatalf("explicit additionalProperties = %#v, want preserved true", got["additionalProperties"]) + } +} diff --git a/client/adapters/gemini.go b/client/adapters/gemini.go index bda21d2..6839680 100644 --- a/client/adapters/gemini.go +++ b/client/adapters/gemini.go @@ -161,12 +161,12 @@ func (c *GeminiClient) StreamChat(ctx context.Context, messages []core.EyrieMess if geminiSharedParserEnabled() { sseEvents := core.ParseSSEStream(streamCtx, resp.Body, c.logger) events := processGeminiStream(streamCtx, sseEvents, c.logger) - return llm.NewStreamResult(events, "", cancel), nil + return llm.NewStreamResult(events, requestID, cancel), nil } // Fallback (opt-out via EYRIE_GEMINI_SHARED_PARSER=0): old bespoke parser. events := make(chan core.EyrieStreamEvent, 64) go c.streamLoop(streamCtx, resp.Body, events) - return llm.NewStreamResult(events, "", cancel), nil + return llm.NewStreamResult(events, requestID, cancel), nil } func (c *GeminiClient) Ping(ctx context.Context) error { diff --git a/client/adapters/gemini_test.go b/client/adapters/gemini_test.go index 4faea87..ed7f61b 100644 --- a/client/adapters/gemini_test.go +++ b/client/adapters/gemini_test.go @@ -149,6 +149,37 @@ func TestGeminiClient_StreamChat_EmptyModel(t *testing.T) { } } +// Regression (audit E5): StreamChat captured X-Goog-Request-Id for error +// reporting but passed "" to NewStreamResult, dropping the provider request +// ID on successful streams. +func TestGeminiClient_StreamChat_PropagatesRequestID(t *testing.T) { + t.Parallel() + transport := roundTripFunc(func(req *http.Request) (*http.Response, error) { + body := "data: {\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"hi\"}]},\"finishReason\":\"STOP\"}]}\n\n" + return &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{ + "Content-Type": []string{"text/event-stream"}, + "X-Goog-Request-Id": []string{"goog-req-42"}, + }, + Body: io.NopCloser(strings.NewReader(body)), + }, nil + }) + c := NewGeminiClient("key", "") + c.retry = core.RetryConfig{RetryConfig: types.RetryConfig{MaxRetries: 0}} + c.httpClient = &http.Client{Transport: transport} + result, err := c.StreamChat(context.Background(), []core.EyrieMessage{{Role: "user", Content: "Hi"}}, core.ChatOptions{Model: "gemini-2.0-flash"}) + if err != nil { + t.Fatalf("StreamChat: %v", err) + } + defer result.Close() + for range result.Events { + } + if result.RequestID != "goog-req-42" { + t.Fatalf("stream request id = %q, want goog-req-42", result.RequestID) + } +} + func TestGeminiClient_Ping_Success(t *testing.T) { t.Parallel() transport := roundTripFunc(func(req *http.Request) (*http.Response, error) { diff --git a/client/adapters/mimo.go b/client/adapters/mimo.go index 7012ef9..6493cbb 100644 --- a/client/adapters/mimo.go +++ b/client/adapters/mimo.go @@ -1,7 +1,10 @@ package adapters import ( + "bytes" "context" + "io" + "log/slog" "net/http" "strconv" "strings" @@ -66,5 +69,39 @@ func mimoAuthHeaders(req *http.Request, apiKey string) { xiaomi.SetMimoRequestAuth(req, apiKey) } +// doWithMimoAuthRetry runs the HTTP request through core.DoWithRetry; on 401 +// with MiMo api-key auth it rebuilds the request with the provider's Bearer +// headers (set by setRetryHeaders) and retries once. Shared by the OpenAI and +// Anthropic adapters, which differ only in the headers applied to the retry. +func doWithMimoAuthRetry( + ctx context.Context, + httpClient *http.Client, + retry core.RetryConfig, + logger *slog.Logger, + useMimoAuth bool, + req *http.Request, + body []byte, + setRetryHeaders func(*http.Request), +) (*http.Response, error) { + resp, err := core.DoWithRetry(ctx, httpClient, req, retry, logger) + if err != nil { + return nil, err + } + if !useMimoAuth || resp.StatusCode != http.StatusUnauthorized { + return resp, nil + } + _ = resp.Body.Close() + req2, err := http.NewRequestWithContext(ctx, req.Method, req.URL.String(), bytes.NewReader(body)) + if err != nil { + return nil, err + } + setRetryHeaders(req2) + if req.Header.Get("Accept") != "" { + req2.Header.Set("Accept", req.Header.Get("Accept")) + } + req2.GetBody = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(body)), nil } + return core.DoWithRetry(ctx, httpClient, req2, retry, logger) +} + // ProviderID reports the configured MiMo gateway identity. func (c *MiMoClient) ProviderID() string { return c.providerID } diff --git a/client/adapters/openai.go b/client/adapters/openai.go index a68a3f8..d533450 100644 --- a/client/adapters/openai.go +++ b/client/adapters/openai.go @@ -82,24 +82,7 @@ func (c *OpenAIClient) setBearerHeaders(req *http.Request) { // doRequestWithMimoAuthRetry runs the HTTP request; on 401 with MiMo api-key auth, retries once with Bearer. func (c *OpenAIClient) doRequestWithMimoAuthRetry(ctx context.Context, req *http.Request, body []byte) (*http.Response, error) { - resp, err := core.DoWithRetry(ctx, c.httpClient, req, c.retry, c.logger) - if err != nil { - return nil, err - } - if !c.useMimoAuth || resp.StatusCode != http.StatusUnauthorized { - return resp, nil - } - _ = resp.Body.Close() - req2, err := http.NewRequestWithContext(ctx, req.Method, req.URL.String(), bytes.NewReader(body)) - if err != nil { - return nil, err - } - c.setBearerHeaders(req2) - if req.Header.Get("Accept") != "" { - req2.Header.Set("Accept", req.Header.Get("Accept")) - } - req2.GetBody = func() (io.ReadCloser, error) { return io.NopCloser(bytes.NewReader(body)), nil } - return core.DoWithRetry(ctx, c.httpClient, req2, c.retry, c.logger) + return doWithMimoAuthRetry(ctx, c.httpClient, c.retry, c.logger, c.useMimoAuth, req, body, c.setBearerHeaders) } type ( diff --git a/client/continuation.go b/client/continuation.go index ad679d7..ff97bf0 100644 --- a/client/continuation.go +++ b/client/continuation.go @@ -148,7 +148,12 @@ func StreamChatWithContinuation(ctx context.Context, p Provider, messages []Eyri stopReason = evt.StopReason case "error": emit(cancelCtx, outCh, evt) - return + // Warning-marked error events are non-fatal health + // diagnostics emitted just before the terminal done; + // keep consuming so that done event is observed. + if evt.Warning == "" { + return + } default: emit(cancelCtx, outCh, evt) } diff --git a/client/core/response_health_test.go b/client/core/response_health_test.go index 9e50ec9..a0c4e01 100644 --- a/client/core/response_health_test.go +++ b/client/core/response_health_test.go @@ -62,7 +62,9 @@ func TestHealthFromResponse(t *testing.T) { } // Integration: a stream that emits only reasoning_content then finishes must -// produce a diagnostic error event before the terminal done. +// produce a diagnostic error event before the terminal done. The diagnostic +// must carry the Warning marker (non-fatal) so downstream consumers keep +// streaming and still observe the terminal done event. func TestStreamEmitsErrorOnlyReasoningDiagnostic(t *testing.T) { t.Parallel() events := make(chan SSEEvent, 10) @@ -72,7 +74,9 @@ func TestStreamEmitsErrorOnlyReasoningDiagnostic(t *testing.T) { ch := ProcessOpenAIStream(context.Background(), events, testLogger()) - var sawThinking, sawDiagnostic, sawContent bool + var sawThinking, sawDiagnostic, sawContent, sawDone bool + diagnosticIdx, doneIdx := -1, -1 + var i int for evt := range ch { switch evt.Type { case "thinking": @@ -82,8 +86,16 @@ func TestStreamEmitsErrorOnlyReasoningDiagnostic(t *testing.T) { case "error": if strings.Contains(evt.Error, "reasoning") { sawDiagnostic = true + diagnosticIdx = i + if evt.Warning == "" { + t.Error("diagnostic error event must set the Warning marker") + } } + case "done": + sawDone = true + doneIdx = i } + i++ } if !sawThinking { t.Error("expected a thinking event from reasoning_content") @@ -94,6 +106,12 @@ func TestStreamEmitsErrorOnlyReasoningDiagnostic(t *testing.T) { if !sawDiagnostic { t.Error("expected an error-only-reasoning diagnostic event") } + if !sawDone { + t.Fatal("expected a terminal done event after the diagnostic") + } + if diagnosticIdx < 0 || doneIdx < diagnosticIdx { + t.Fatalf("diagnostic (%d) must precede done (%d)", diagnosticIdx, doneIdx) + } } // A normal stream with content must NOT Emit a health diagnostic. diff --git a/client/core/stream.go b/client/core/stream.go index 2200881..6141e8d 100644 --- a/client/core/stream.go +++ b/client/core/stream.go @@ -453,8 +453,11 @@ func ProcessOpenAIStreamWithOpts(ctx context.Context, sseEvents <-chan SSEEvent, } // finish emits the terminal done event, preceded by a non-fatal health - // diagnostic (as an error-type event the consumer can log) when the - // response produced no usable content/tool-calls. + // diagnostic when the response produced no usable content/tool-calls. + // The diagnostic rides an error-type event (so consumers that only + // read .Error still see it) but is marked non-fatal via the Warning + // field: downstream layers must surface it without terminating the + // stream, because the terminal done event follows immediately. finish := func(stopReason string) { emitTools() health := DetectResponseHealth(ResponseSignals{ @@ -465,7 +468,7 @@ func ProcessOpenAIStreamWithOpts(ctx context.Context, sseEvents <-chan SSEEvent, StreamEnded: true, }) if d := health.Diagnostic(); d != "" { - Emit(ctx, ch, EyrieStreamEvent{Type: "error", Error: d}) + Emit(ctx, ch, EyrieStreamEvent{Type: "error", Error: d, Warning: d}) } Emit(ctx, ch, EyrieStreamEvent{Type: "done", StopReason: stopReason, TTFTms: ttftMs}) } diff --git a/client/tracing.go b/client/tracing.go index 96e177d..6a8573d 100644 --- a/client/tracing.go +++ b/client/tracing.go @@ -97,8 +97,14 @@ func (tp *TracingProvider) StreamChat(ctx context.Context, messages []EyrieMessa for evt := range origEvents { switch evt.Type { case "error": - span.SetStatus(codes.Error, evt.Error) - span.SetAttributes(attribute.Bool("error", true)) + if evt.Warning != "" { + // Non-fatal health diagnostic: record it without + // failing the span (the stream still completes). + span.SetAttributes(attribute.String("warning", evt.Warning)) + } else { + span.SetStatus(codes.Error, evt.Error) + span.SetAttributes(attribute.Bool("error", true)) + } case "usage": // Token usage is delivered on the "usage" event, not "done". if evt.Usage != nil { diff --git a/engine/engine_test.go b/engine/engine_test.go index 525ce75..9cd6692 100644 --- a/engine/engine_test.go +++ b/engine/engine_test.go @@ -177,6 +177,69 @@ func TestNormalizedStreamContract(t *testing.T) { } } +// Regression (audit E4): client/core emits end-of-stream health diagnostics +// (e.g. reasoning-only responses) as error-type events marked non-fatal via +// the Warning field, followed by the terminal done. The engine must forward +// them as warning events and still deliver the done/usage event without +// setting Err(). +func TestStreamDiagnosticErrorEventIsNonFatal(t *testing.T) { + sourceEvents := make(chan client.EyrieStreamEvent, 3) + sourceEvents <- client.EyrieStreamEvent{Type: "content", Content: "answer"} + sourceEvents <- client.EyrieStreamEvent{Type: "error", Error: "model produced reasoning tokens but no answer", Warning: "model produced reasoning tokens but no answer"} + sourceEvents <- client.EyrieStreamEvent{Type: "done", StopReason: "stop", Usage: &client.EyrieUsage{PromptTokens: 1, CompletionTokens: 2, TotalTokens: 3}} + close(sourceEvents) + + ctx, cancel := context.WithCancel(context.Background()) + stream := newStream(ctx, cancel, client.NewStreamResult(sourceEvents, nil), Route{Provider: "mock", Model: "mock/model"}) + defer stream.Close() + + var events []Event + for stream.Next() { + events = append(events, stream.Event()) + } + if err := stream.Err(); err != nil { + t.Fatalf("diagnostic must not set Err(): %v", err) + } + if len(events) != 4 { + t.Fatalf("events = %d, want 4: %+v", len(events), events) + } + if events[1].Type != EventContentDelta || events[1].Content != "answer" { + t.Fatalf("content delta lost: %+v", events[1]) + } + if events[2].Type != EventWarning || events[2].Warning == "" { + t.Fatalf("diagnostic should surface as a warning event: %+v", events[2]) + } + if events[3].Type != EventDone || events[3].Usage == nil || events[3].Usage.TotalTokens != 3 { + t.Fatalf("terminal done/usage lost: %+v", events[3]) + } +} + +// Genuinely fatal error events (no Warning marker) keep the previous +// behavior: the stream terminates and Err() carries the classified error. +func TestStreamFatalErrorEventStillTerminal(t *testing.T) { + sourceEvents := make(chan client.EyrieStreamEvent, 2) + sourceEvents <- client.EyrieStreamEvent{Type: "content", Content: "partial"} + sourceEvents <- client.EyrieStreamEvent{Type: "error", Error: "connection reset"} + close(sourceEvents) + + ctx, cancel := context.WithCancel(context.Background()) + stream := newStream(ctx, cancel, client.NewStreamResult(sourceEvents, nil), Route{Provider: "mock", Model: "mock/model"}) + defer stream.Close() + + var events []Event + for stream.Next() { + events = append(events, stream.Event()) + } + if err := stream.Err(); err == nil { + t.Fatal("fatal error event must set Err()") + } else if !IsCode(err, ErrorProviderUnavailable) { + t.Fatalf("error code = %v, want provider_unavailable", err) + } + if len(events) != 2 { + t.Fatalf("events = %d, want 2 (route_selected + content_delta): %+v", len(events), events) + } +} + func TestSnapshotPublishesCapabilities(t *testing.T) { compiled := &catalog.CompiledCatalog{ ModelsByID: map[string]catalog.Model{ diff --git a/engine/stream.go b/engine/stream.go index fea7b24..76f81ac 100644 --- a/engine/stream.go +++ b/engine/stream.go @@ -148,6 +148,13 @@ func normalizeEvent(event client.EyrieStreamEvent) (Event, error) { case "continuation": out.Type = EventContinuation case "error": + if event.Warning != "" { + // Non-fatal health diagnostic (e.g. a reasoning-only response): + // client/core marks these with Warning so they can be surfaced + // without terminating the stream — the terminal done/usage event + // follows. Forward as a warning event; do not set Err()/stop. + return Event{Type: EventWarning, Warning: event.Warning}, nil + } return Event{}, &Error{Code: ErrorProviderUnavailable, Operation: "stream", Message: event.Error} default: out.Type = event.Type diff --git a/go.mod b/go.mod index 3becbbf..e242d77 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module github.com/GrayCodeAI/eyrie -go 1.26.5 +go 1.26.6 require ( github.com/GrayCodeAI/hawk-core-contracts v0.1.12