diff --git a/apps/edge-api/cmd/server/main.go b/apps/edge-api/cmd/server/main.go index 913c54703..5fdf0ace3 100644 --- a/apps/edge-api/cmd/server/main.go +++ b/apps/edge-api/cmd/server/main.go @@ -196,7 +196,20 @@ func main() { // reservation and settlement, upstream retry, tracing and audit are shared // with that surface rather than reimplemented. It previously POSTed straight // to LiteLLM, which let a caller address a raw route id and skip all of it. - anthropicHandler := anthropic.NewHandler(anthropic.Deps{OpenAIChat: openAIChatHandler}) + anthropicHandler := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: openAIChatHandler, + // count_tokens is the one route on this surface that does not delegate, + // so it is the one route that needs its own API-key authority. Without + // it the handler could only see a JWT session user and refused every + // Anthropic SDK caller, which authenticates with an API key (#1261). + // Zero cost arguments: the estimate is computed locally and bills + // nothing, so this resolves and rate-limits the key without reserving + // credit against it. + AuthorizeAPIKey: func(ctx context.Context, authHeader string) (*apierrors.OpenAIError, map[string]string) { + _, headers, authErr := authorizer.Authorize(ctx, authHeader, "", 0, 0, 0) + return authErr, headers + }, + }) mux.Handle("/v1/messages", anthropic.APIKeyNormalizer(anthropicHandler)) mux.Handle("/v1/messages/", anthropic.APIKeyNormalizer(anthropicHandler)) @@ -476,7 +489,7 @@ func main() { } // API routes - mux.Handle("/v1/models", handleModels(catalogClient, authorizer)) + mux.Handle("/v1/models", modelsHandler(catalogClient, authorizer)) mux.Handle("/catalog/models", handleCatalogModels(catalogClient)) // Feature-gate read seam for Open WebUI (issue #293). OWUI has no in-repo @@ -843,6 +856,26 @@ func voiceGateForAPIKeys(gate func(http.Handler) http.Handler) func(http.Handler } } +// modelsHandler is what GET /v1/models is actually registered as: the +// OpenAI-shaped handler below, wrapped so a real Anthropic SDK client works +// against the same route (issue #1259). +// +// APIKeyNormalizer is applied here at the leaf as well as in +// authSelectorMiddleware, and that is not redundant. The selector wrapper only +// exists when JWT auth is wired (jwtMW != nil); on a deployment where Supabase +// JWT config is absent, edge-api logs "JWT auth wiring skipped" and mounts no +// selector at all, so nothing normalizes x-api-key and handleModels reads an +// empty Authorization header for every Anthropic SDK caller. POST /v1/messages +// has always carried the same leaf wrapper for the same reason, which is +// precisely why it kept working on that deployment while this route 401'd. +// +// ModelsCompat then re-shapes the answer, but only for a caller that +// identified itself as Anthropic-shaped; an OpenAI-shaped caller, Open WebUI +// included, still gets the byte-identical OpenAI list it always did. +func modelsHandler(client *catalog.Client, authorizer *authz.Authorizer) http.Handler { + return anthropic.APIKeyNormalizer(anthropic.ModelsCompat(handleModels(client, authorizer))) +} + // handleModels serves the OpenAI-compatible model list. // // Every caller needs a credential this service can resolve: a signed-in diff --git a/apps/edge-api/cmd/server/main_test.go b/apps/edge-api/cmd/server/main_test.go index 3dbc81cd3..e4be7842e 100644 --- a/apps/edge-api/cmd/server/main_test.go +++ b/apps/edge-api/cmd/server/main_test.go @@ -3,6 +3,7 @@ package main import ( "bytes" "context" + "encoding/json" "errors" "fmt" "log" @@ -1292,3 +1293,94 @@ func TestDegradedHealthBodyNamesNoInternalComponent(t *testing.T) { } } } + +// modelsHandlerTestFixtures builds a catalog client and an authorizer that both +// accept the one API key these two tests present. +func modelsHandlerTestFixtures(t *testing.T) (*edgecatalog.Client, *authz.Authorizer) { + t.Helper() + var sawPath string + seeded := `{"models":[{"id":"hive-default","object":"model","created":1716935002,"owned_by":"hive"}],"catalog":[]}` + client := edgecatalog.NewClient(newTenantCatalogSnapshotServer(t, seeded, seeded, &sawPath)) + authorizer := newTestAuthorizer(t, http.StatusOK, `{ + "key_id":"key-1", + "account_id":"acc-1", + "tenant_id":"`+uuid.New().String()+`", + "status":"active", + "allow_all_models":true, + "allowed_aliases":["hive-default"], + "budget_kind":"none", + "budget_consumed_credits":0, + "budget_reserved_credits":0, + "policy_version":1 + }`) + return client, authorizer +} + +// TestModelsHandlerServesAnAnthropicSDKClient is the issue #1259 wiring guard, +// exercising GET /v1/models exactly as it is registered rather than the inner +// OpenAI handler alone. +// +// A real Anthropic SDK client sends the credential on x-api-key and never on +// Authorization. Two things then have to happen at this route and neither did: +// the leaf APIKeyNormalizer has to rewrite the header (authSelectorMiddleware's +// copy of it only exists when JWT auth is wired, so it cannot be the only one), +// and the answer has to come back in the Anthropic list shape. +func TestModelsHandlerServesAnAnthropicSDKClient(t *testing.T) { + client, authorizer := modelsHandlerTestFixtures(t) + + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("x-api-key", "hk_test") + req.Header.Set("anthropic-version", "2023-06-01") + rr := httptest.NewRecorder() + + modelsHandler(client, authorizer).ServeHTTP(rr, req) + + if rr.Code != http.StatusOK { + t.Fatalf("x-api-key caller: want 200 got %d: %s", rr.Code, rr.Body.String()) + } + var got struct { + Data []struct { + Type string `json:"type"` + ID string `json:"id"` + DisplayName string `json:"display_name"` + CreatedAt string `json:"created_at"` + } `json:"data"` + HasMore bool `json:"has_more"` + } + if err := json.Unmarshal(rr.Body.Bytes(), &got); err != nil { + t.Fatalf("decode: %v (body=%s)", err, rr.Body.String()) + } + if len(got.Data) != 1 { + t.Fatalf("data: want 1 entry got %d (%s)", len(got.Data), rr.Body.String()) + } + if got.Data[0].Type != "model" || got.Data[0].ID != "hive-default" || got.Data[0].CreatedAt == "" { + t.Fatalf("entry is not Anthropic-shaped: %+v", got.Data[0]) + } + if strings.Contains(rr.Body.String(), "owned_by") { + t.Errorf("Anthropic list body still carries OpenAI keys: %s", rr.Body.String()) + } +} + +// TestModelsHandlerKeepsTheOpenAIShapeForOpenAIClients is the non-regression +// half of the route wiring: Open WebUI's model picker reads this same route +// with a plain bearer token and must keep getting the OpenAI list. +func TestModelsHandlerKeepsTheOpenAIShapeForOpenAIClients(t *testing.T) { + client, authorizer := modelsHandlerTestFixtures(t) + + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("Authorization", "Bearer hk_test") + rr := httptest.NewRecorder() + + modelsHandler(client, authorizer).ServeHTTP(rr, req) + + if rr.Code != http.StatusOK { + t.Fatalf("bearer caller: want 200 got %d: %s", rr.Code, rr.Body.String()) + } + body := rr.Body.String() + if !strings.Contains(body, `"object":"list"`) || !strings.Contains(body, "hive-default") { + t.Fatalf("OpenAI-shaped caller must keep the OpenAI list: %s", body) + } + if strings.Contains(body, "display_name") || strings.Contains(body, "has_more") { + t.Fatalf("OpenAI-shaped caller was served the Anthropic shape: %s", body) + } +} diff --git a/apps/edge-api/internal/anthropic/handler.go b/apps/edge-api/internal/anthropic/handler.go index efea2f3cc..166059627 100644 --- a/apps/edge-api/internal/anthropic/handler.go +++ b/apps/edge-api/internal/anthropic/handler.go @@ -2,6 +2,7 @@ package anthropic import ( "bytes" + "context" "encoding/json" "io" "log/slog" @@ -12,6 +13,7 @@ import ( "github.com/google/uuid" "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/auth" "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/authz" + apierr "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/errors" ) const maxBodyBytes = 4 << 20 // 4 MiB @@ -34,6 +36,22 @@ type Deps struct { // id, that direct POST let a caller name a route instead of an alias and // skip entitlement and metering in one move. OpenAIChat http.Handler + + // AuthorizeAPIKey resolves a "Bearer hk_..." Authorization header to a Hive + // API-key principal, returning the already-sanitized OpenAI-shaped refusal + // (and any headers that must ride with it) when it cannot. + // + // It exists for POST /v1/messages/count_tokens alone. Every other route on + // this surface delegates to OpenAIChat, which is itself the authority for + // an API-key principal; count_tokens never dispatches anywhere, so without + // this it could only see a session-cookie principal and 401'd every + // programmatic caller -- which is to say, essentially every real Anthropic + // SDK integration, since an API key is how they all authenticate (issue + // #1261). + // + // Nil leaves count_tokens session-only and fail-closed, the pre-existing + // behaviour. + AuthorizeAPIKey func(ctx context.Context, authHeader string) (*apierr.OpenAIError, map[string]string) } // Handler accepts Anthropic Messages requests, translates them to the internal @@ -159,17 +177,7 @@ func (h *Handler) handleMessages(w http.ResponseWriter, r *http.Request) { // handleCountTokens returns a local token count estimate for the request body. func (h *Handler) handleCountTokens(w http.ResponseWriter, r *http.Request) { - user, ok := auth.UserFrom(r.Context()) - if !ok || user == nil { - writeAnthropicError(w, http.StatusUnauthorized, "missing user", "") - return - } - if user.TenantID == uuid.Nil { - writeAnthropicError(w, http.StatusForbidden, "no tenant for user", "") - return - } - if !authz.RoleHas(authz.Role(user.Role), authz.PermChatInvoke) { - writeAnthropicError(w, http.StatusForbidden, "chat not allowed", "") + if !h.authorizeCountTokens(w, r) { return } @@ -206,6 +214,48 @@ func (h *Handler) handleCountTokens(w http.ResponseWriter, r *http.Request) { } } +// authorizeCountTokens accepts either principal type this surface serves: a +// JWT session user (checked for tenant and chat permission, as before) or a +// Hive API key resolved through Deps.AuthorizeAPIKey. It writes the refusal +// itself and reports whether the request may proceed. +// +// The two are checked in that order because the JWT middleware is what +// populates auth.UserFrom; an "hk_" request is routed past it by auth.Selector +// and therefore carries no session user at all, which is exactly why the +// session-only guard this replaces rejected every API-key caller. +func (h *Handler) authorizeCountTokens(w http.ResponseWriter, r *http.Request) bool { + if user, ok := auth.UserFrom(r.Context()); ok && user != nil { + if user.TenantID == uuid.Nil { + writeAnthropicError(w, http.StatusForbidden, "no tenant for user", "") + return false + } + if !authz.RoleHas(authz.Role(user.Role), authz.PermChatInvoke) { + writeAnthropicError(w, http.StatusForbidden, "chat not allowed", "") + return false + } + return true + } + + if h.deps.AuthorizeAPIKey == nil { + writeAnthropicError(w, http.StatusUnauthorized, "missing user", "") + return false + } + authErr, headers := h.deps.AuthorizeAPIKey(r.Context(), r.Header.Get("Authorization")) + if authErr == nil { + return true + } + // Round-trip through the shared OpenAI writer so the status mapping + // (401 vs 403 vs 429 vs 503) stays the single implementation the rest of + // edge-api uses, then reshape the envelope for an Anthropic client. This + // never re-sanitizes: the authorizer's refusals are already customer-safe. + // reshapeInto, not a bare reshape, so the retry metadata WriteAuthFailure + // sets on a 429 or a 503 reaches the client instead of dying in the recorder. + rec := &headerlessRecorder{} + apierr.WriteAuthFailure(rec, authErr, headers) + rec.reshapeInto(w) + return false +} + // normalizeAPIKeyHeader rewrites an Anthropic x-api-key header to a standard // Authorization: Bearer header so downstream auth middleware works uniformly. func normalizeAPIKeyHeader(r *http.Request) *http.Request { diff --git a/apps/edge-api/internal/anthropic/handler_test.go b/apps/edge-api/internal/anthropic/handler_test.go index 5aa1a8d2a..4984752b1 100644 --- a/apps/edge-api/internal/anthropic/handler_test.go +++ b/apps/edge-api/internal/anthropic/handler_test.go @@ -1,6 +1,7 @@ package anthropic_test import ( + "context" "encoding/json" "fmt" "io" @@ -12,6 +13,7 @@ import ( "github.com/google/uuid" "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/anthropic" "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/auth" + apierr "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/errors" ) // fakeChat stands in for the wired POST /v1/chat/completions handler chain that @@ -686,3 +688,209 @@ func TestAPIKeyNormalizer_NoKey_PassesThrough(t *testing.T) { t.Errorf("Authorization: want empty got %q", captured) } } + +// newCountTokensRequest builds an unauthenticated (no session user) +// count_tokens request carrying an API-key credential, which is how every real +// Anthropic SDK integration reaches this route. +func newCountTokensRequest(authHeader string) *http.Request { + req := httptest.NewRequest(http.MethodPost, "/v1/messages/count_tokens", + strings.NewReader(`{"model":"m","messages":[{"role":"user","content":"hello world"}],"max_tokens":5}`)) + req.Header.Set("Content-Type", "application/json") + if authHeader != "" { + req.Header.Set("Authorization", authHeader) + } + return req +} + +// TestHandler_CountTokens_AcceptsAPIKeyPrincipal is the issue #1261 guard. +// count_tokens is the only route on this surface that does not delegate to the +// chat chain, so it used to recognize a JWT session principal and nothing +// else. An "hk_" request is routed past the JWT middleware by auth.Selector and +// therefore carries no session user at all, which made every programmatic +// caller a 401 on this one route while the sibling /v1/messages accepted the +// identical key on the identical connection. +func TestHandler_CountTokens_AcceptsAPIKeyPrincipal(t *testing.T) { + var sawHeader string + calls := 0 + h := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: &fakeChat{}, + AuthorizeAPIKey: func(_ context.Context, authHeader string) (*apierr.OpenAIError, map[string]string) { + calls++ + sawHeader = authHeader + return nil, nil + }, + }) + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, newCountTokensRequest("Bearer hk_live_test")) + + if rec.Code != http.StatusOK { + t.Fatalf("API-key count_tokens: want 200 got %d body=%s", rec.Code, rec.Body.String()) + } + if calls != 1 { + t.Fatalf("AuthorizeAPIKey called %d times, want exactly 1", calls) + } + if sawHeader != "Bearer hk_live_test" { + t.Errorf("authorizer saw header %q, want the request's own Authorization value", sawHeader) + } + var got anthropic.CountTokensResponse + if err := json.NewDecoder(rec.Body).Decode(&got); err != nil { + t.Fatalf("decode: %v", err) + } + if got.InputTokens <= 0 { + t.Errorf("input_tokens: want a positive estimate got %d", got.InputTokens) + } +} + +// TestHandler_CountTokens_RejectedAPIKeyKeepsTheAuthorizersOwnRefusal proves +// the refusal a caller sees is the authorizer's verdict reshaped into the +// Anthropic envelope, not the old blanket "missing user" string. The status +// alone cannot prove that (both are 401), so the message is what this asserts. +func TestHandler_CountTokens_RejectedAPIKeyKeepsTheAuthorizersOwnRefusal(t *testing.T) { + code := "invalid_api_key" + h := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: &fakeChat{}, + AuthorizeAPIKey: func(_ context.Context, _ string) (*apierr.OpenAIError, map[string]string) { + refusal := apierr.NewError("invalid_request_error", "Invalid API key.", &code) + return &refusal, nil + }, + }) + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, newCountTokensRequest("Bearer hk_revoked")) + + if rec.Code != http.StatusUnauthorized { + t.Fatalf("revoked key: want 401 got %d", rec.Code) + } + var body map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatalf("decode: %v", err) + } + if body["type"] != "error" { + t.Errorf("envelope: want top-level type=error got %v", body["type"]) + } + errObj, _ := body["error"].(map[string]any) + if errObj["type"] != "authentication_error" { + t.Errorf("error.type: want authentication_error got %v", errObj["type"]) + } + if errObj["message"] != "Invalid API key." { + t.Errorf("error.message: want the authorizer's own refusal got %v", errObj["message"]) + } + if errObj["code"] != "invalid_api_key" { + t.Errorf("error.code: want invalid_api_key got %v", errObj["code"]) + } +} + +// TestHandler_CountTokens_WithoutAnAPIKeyAuthorityFailsClosed pins the +// deliberate degradation: a Handler wired without AuthorizeAPIKey stays +// session-only and refuses, rather than admitting an unauthenticated caller. +func TestHandler_CountTokens_WithoutAnAPIKeyAuthorityFailsClosed(t *testing.T) { + h := anthropic.NewHandler(anthropic.Deps{OpenAIChat: &fakeChat{}}) + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, newCountTokensRequest("Bearer hk_live_test")) + + if rec.Code != http.StatusUnauthorized { + t.Fatalf("no API-key authority wired: want 401 got %d body=%s", rec.Code, rec.Body.String()) + } +} + +// TestHandler_CountTokens_SessionPrincipalDoesNotConsultTheAPIKeyAuthority +// keeps the two principals independent: a signed-in user must still be served +// from the session checks alone, so a deployment whose authorizer is degraded +// does not start refusing browser traffic. +func TestHandler_CountTokens_SessionPrincipalDoesNotConsultTheAPIKeyAuthority(t *testing.T) { + calls := 0 + h := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: &fakeChat{}, + AuthorizeAPIKey: func(_ context.Context, _ string) (*apierr.OpenAIError, map[string]string) { + calls++ + refusal := apierr.NewError("invalid_request_error", "Invalid API key.", nil) + return &refusal, nil + }, + }) + req := newAuthedRequest(t, `{"model":"m","messages":[{"role":"user","content":"hello world"}],"max_tokens":5}`) + req.URL.Path = "/v1/messages/count_tokens" + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("session count_tokens: want 200 got %d body=%s", rec.Code, rec.Body.String()) + } + if calls != 0 { + t.Errorf("API-key authority consulted %d times for a session principal, want 0", calls) + } +} + +// TestHandler_CountTokens_RefusalCarriesTheAuthorizersRetryHeaders is the +// header half of the same refusal TestHandler_CountTokens_RejectedAPIKeyKeepsTheAuthorizersOwnRefusal +// checks the body half of. WriteAuthFailure is the shared source of truth +// precisely so a retryable 429 is never collapsed into a bare non-retryable +// refusal, and it delivers the retryable part through headers: recording the +// body and discarding those headers restores exactly the collapse it exists to +// prevent, since the Anthropic SDK reads retry-after for its backoff. +func TestHandler_CountTokens_RefusalCarriesTheAuthorizersRetryHeaders(t *testing.T) { + code := "rate_limit_exceeded" + h := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: &fakeChat{}, + AuthorizeAPIKey: func(_ context.Context, _ string) (*apierr.OpenAIError, map[string]string) { + refusal := apierr.NewError("rate_limit_error", "Rate limit reached.", &code) + return &refusal, map[string]string{ + "retry-after": "30", + "x-ratelimit-limit-requests": "100", + "x-ratelimit-remaining-requests": "0", + } + }, + }) + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, newCountTokensRequest("Bearer hk_throttled")) + + if rec.Code != http.StatusTooManyRequests { + t.Fatalf("throttled key: want 429 got %d body=%s", rec.Code, rec.Body.String()) + } + if got := rec.Header().Get("Retry-After"); got != "30" { + t.Errorf("retry-after: want 30 got %q", got) + } + if got := rec.Header().Get("X-Ratelimit-Limit-Requests"); got != "100" { + t.Errorf("x-ratelimit-limit-requests: want 100 got %q", got) + } + if got := rec.Header().Get("X-Ratelimit-Remaining-Requests"); got != "0" { + t.Errorf("x-ratelimit-remaining-requests: want 0 got %q", got) + } + var body map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatalf("decode: %v", err) + } + errObj, _ := body["error"].(map[string]any) + if errObj["type"] != "rate_limit_error" { + t.Errorf("error.type: want rate_limit_error got %v", errObj["type"]) + } +} + +// TestHandler_CountTokens_UpstreamUnavailableCarriesRetryAfter covers the other +// branch that populates the header map: the authorizer refusing because the +// control plane is unreachable. A 503 telling the caller to retry without +// saying when sends every SDK retry layer back to its own short backoff, +// against a dependency that is by construction already unable to answer. +func TestHandler_CountTokens_UpstreamUnavailableCarriesRetryAfter(t *testing.T) { + code := "upstream_unavailable" + h := anthropic.NewHandler(anthropic.Deps{ + OpenAIChat: &fakeChat{}, + AuthorizeAPIKey: func(_ context.Context, _ string) (*apierr.OpenAIError, map[string]string) { + refusal := apierr.NewError("api_error", "Authorization is temporarily unavailable.", &code) + return &refusal, map[string]string{"retry-after": "5"} + }, + }) + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, newCountTokensRequest("Bearer hk_live_test")) + + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("degraded authorizer: want 503 got %d body=%s", rec.Code, rec.Body.String()) + } + if got := rec.Header().Get("Retry-After"); got != "5" { + t.Errorf("retry-after: want 5 got %q", got) + } +} diff --git a/apps/edge-api/internal/anthropic/models.go b/apps/edge-api/internal/anthropic/models.go new file mode 100644 index 000000000..a5d5daca2 --- /dev/null +++ b/apps/edge-api/internal/anthropic/models.go @@ -0,0 +1,145 @@ +package anthropic + +import ( + "encoding/json" + "net/http" + "time" +) + +// ModelsListResponse is the Anthropic GET /v1/models response envelope. +// Verified against the live specification on 2026-08-28 +// (https://docs.claude.com/en/api/models-list): a data array of model objects +// each carrying type/id/display_name/created_at, plus has_more and the +// nullable first_id / last_id pagination cursors. +// +// FirstID and LastID are pointers because the specification types them as +// "string or null" and a real client distinguishes an empty page from a page +// whose cursor happens to be the empty string. +type ModelsListResponse struct { + Data []ModelInfo `json:"data"` + HasMore bool `json:"has_more"` + FirstID *string `json:"first_id"` + LastID *string `json:"last_id"` +} + +// ModelInfo is one entry in an Anthropic model list. +type ModelInfo struct { + Type string `json:"type"` // always "model" + ID string `json:"id"` + DisplayName string `json:"display_name"` + CreatedAt string `json:"created_at"` // RFC 3339 +} + +// openAIModelList is the OpenAI-shaped body this package translates FROM. Only +// the three fields that have an Anthropic counterpart are read; owned_by and +// the rest are deliberately dropped rather than passed through, which also +// keeps this translation strictly leak-reducing with respect to the +// provider-blind invariant. +type openAIModelList struct { + Data []struct { + ID string `json:"id"` + Created int64 `json:"created"` + Name string `json:"name"` + } `json:"data"` +} + +// IsAnthropicClient reports whether a request came from an Anthropic-shaped +// client rather than an OpenAI-shaped one. +// +// Both headers are load-bearing. anthropic-version is sent by every official +// Anthropic SDK on every request and by nothing else, so it identifies the +// dialect even when the caller authenticates with Authorization: Bearer. +// x-api-key is the SDK's default credential header and catches a hand-rolled +// client that copied the documented curl invocation. An OpenAI-shaped caller +// sends neither, so the OpenAI response shape stays the default and Open WebUI +// (which is what actually consumes GET /v1/models on this deployment) is +// unaffected. +func IsAnthropicClient(r *http.Request) bool { + return r.Header.Get("anthropic-version") != "" || r.Header.Get("x-api-key") != "" +} + +// ModelsCompat wraps the OpenAI-shaped GET /v1/models handler so an Anthropic +// SDK client gets an Anthropic-shaped answer: the documented list envelope on +// success, and the Anthropic error envelope ({"type":"error","error":{...}}) +// instead of the OpenAI one on a refusal, which is what +// anthropic.APIStatusError actually reads its .type from (issue #1259). +// +// It wraps rather than modifies the underlying handler so authentication, +// tenant filtering and the OpenAI-shaped contract every other caller depends +// on all stay exactly where they are: this only re-shapes bytes that handler +// already decided to emit. Buffering the whole body is safe here and only +// here, because a model list is small and never streamed. +func ModelsCompat(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + // One URL, two representations, chosen by request headers. Nothing on + // this route sets Cache-Control and the route needs a credential, so no + // correct cache stores it today, but any intermediary keying on URL + // alone (or an edge cache added in front of Caddy later) could hand an + // Anthropic-shaped body to an OpenAI-shaped caller and empty Open + // WebUI's model picker. Declared on both branches, before either + // writes, since a Vary set after WriteHeader is a Vary nobody sends. + // Add rather than Set: an outer middleware may already have declared a + // Vary of its own (Origin, for a CORS layer), and Set would delete it. + w.Header().Add("Vary", "anthropic-version, x-api-key") + + if !IsAnthropicClient(r) { + next.ServeHTTP(w, r) + return + } + + rec := &headerlessRecorder{} + next.ServeHTTP(rec, r) + if rec.status == 0 { + rec.status = http.StatusOK + } + + if rec.status < 200 || rec.status > 299 { + // reshapeInto, not a bare reshape: a refusal from + // authorizeAliasRequest carries retry-after and the x-ratelimit-* + // family on the recorder, and an Anthropic client reads them for + // its backoff exactly as an OpenAI-shaped one does. + rec.reshapeInto(w) + return + } + + var list openAIModelList + if err := json.Unmarshal(rec.body.Bytes(), &list); err != nil { + writeAnthropicError(w, http.StatusBadGateway, "the model list could not be returned", "upstream_error") + return + } + + out := ModelsListResponse{Data: make([]ModelInfo, 0, len(list.Data))} + for _, m := range list.Data { + display := m.Name + if display == "" { + display = m.ID + } + out.Data = append(out.Data, ModelInfo{ + Type: "model", + ID: m.ID, + DisplayName: display, + CreatedAt: time.Unix(m.Created, 0).UTC().Format(time.RFC3339), + }) + } + if len(out.Data) > 0 { + first := out.Data[0].ID + last := out.Data[len(out.Data)-1].ID + out.FirstID = &first + out.LastID = &last + } + // HasMore stays false: this route serves the caller's whole entitled + // catalog in one response and honours no pagination cursor, so + // claiming another page exists would send a client into a loop. + + // The success path re-encodes rather than reshapes, but it owes the + // client the same headers the delegated handler set: this route is + // where a 2xx rate-limit budget would surface, and dropping those is + // the same defect as dropping a refusal's retry metadata. + rec.copyHeadersTo(w) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + if err := json.NewEncoder(w).Encode(out); err != nil { + return + } + }) +} diff --git a/apps/edge-api/internal/anthropic/models_test.go b/apps/edge-api/internal/anthropic/models_test.go new file mode 100644 index 000000000..7678ae75a --- /dev/null +++ b/apps/edge-api/internal/anthropic/models_test.go @@ -0,0 +1,337 @@ +package anthropic_test + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/sakibsadmanshajib/hive/apps/edge-api/internal/anthropic" +) + +// openAIModelListBody is the exact body edge-api's own GET /v1/models writes. +const openAIModelListBody = `{"object":"list","data":[` + + `{"id":"hive-default","object":"model","created":1716935002,"owned_by":"hive","name":"Hive Default"},` + + `{"id":"hive-fast","object":"model","created":1716935003,"owned_by":"hive"}` + + `]}` + +func openAIModelsHandler(status int, body string) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _, _ = w.Write([]byte(body)) + }) +} + +// TestModelsCompat_AnthropicClientGetsTheAnthropicListShape is the issue #1259 +// body-shape guard. Verified against https://docs.claude.com/en/api/models-list +// on 2026-08-28: each entry is {"type":"model","id":...,"display_name":..., +// "created_at":""} and the envelope carries has_more plus the +// nullable first_id / last_id cursors, none of which the OpenAI shape has. +func TestModelsCompat_AnthropicClientGetsTheAnthropicListShape(t *testing.T) { + h := anthropic.ModelsCompat(openAIModelsHandler(http.StatusOK, openAIModelListBody)) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("x-api-key", "hk_live_test") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("status: want 200 got %d body=%s", rec.Code, rec.Body.String()) + } + var got anthropic.ModelsListResponse + if err := json.Unmarshal(rec.Body.Bytes(), &got); err != nil { + t.Fatalf("decode: %v (body=%s)", err, rec.Body.String()) + } + if len(got.Data) != 2 { + t.Fatalf("data: want 2 entries got %d (%s)", len(got.Data), rec.Body.String()) + } + if got.Data[0].Type != "model" { + t.Errorf("data[0].type: want model got %q", got.Data[0].Type) + } + if got.Data[0].ID != "hive-default" { + t.Errorf("data[0].id: want hive-default got %q", got.Data[0].ID) + } + if got.Data[0].DisplayName != "Hive Default" { + t.Errorf("data[0].display_name: want Hive Default got %q", got.Data[0].DisplayName) + } + // An entry with no name falls back to its id rather than an empty string: + // display_name is a required field a client renders directly. + if got.Data[1].DisplayName != "hive-fast" { + t.Errorf("data[1].display_name: want the id as fallback got %q", got.Data[1].DisplayName) + } + if _, err := time.Parse(time.RFC3339, got.Data[0].CreatedAt); err != nil { + t.Errorf("data[0].created_at %q is not RFC 3339: %v", got.Data[0].CreatedAt, err) + } + if got.HasMore { + t.Error("has_more: this route serves the whole entitled list, so it must be false") + } + if got.FirstID == nil || *got.FirstID != "hive-default" { + t.Errorf("first_id: want hive-default got %v", got.FirstID) + } + if got.LastID == nil || *got.LastID != "hive-fast" { + t.Errorf("last_id: want hive-fast got %v", got.LastID) + } + + // The OpenAI-only keys must be gone, not merely supplemented: a real SDK + // parses this into a typed page and an "object":"list" envelope has no + // data array it can read. + raw := rec.Body.String() + for _, absent := range []string{`"object"`, `"owned_by"`} { + if strings.Contains(raw, absent) { + t.Errorf("Anthropic list body still carries the OpenAI key %s: %s", absent, raw) + } + } +} + +// TestModelsCompat_EmptyListStillCarriesNullCursors pins the shape of a page +// with nothing on it: data is an array (never null, same defect class as issue +// #1260) and both cursors are explicitly null rather than absent. +func TestModelsCompat_EmptyListStillCarriesNullCursors(t *testing.T) { + h := anthropic.ModelsCompat(openAIModelsHandler(http.StatusOK, `{"object":"list","data":[]}`)) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("anthropic-version", "2023-06-01") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + raw := rec.Body.String() + if !strings.Contains(raw, `"data":[]`) { + t.Errorf("empty list must serialize data as []: %s", raw) + } + if !strings.Contains(raw, `"first_id":null`) || !strings.Contains(raw, `"last_id":null`) { + t.Errorf("empty list must carry explicit null cursors: %s", raw) + } +} + +// TestModelsCompat_ErrorUsesTheAnthropicEnvelope is the issue #1259 error-shape +// guard. anthropic.APIStatusError reads its .type from body["error"]["type"], +// and the OpenAI envelope has no top-level "type":"error" and carries a value +// there that is not a member of Anthropic's error enum, so a client inspecting +// either field got nothing usable. +func TestModelsCompat_ErrorUsesTheAnthropicEnvelope(t *testing.T) { + openAIRefusal := `{"error":{"message":"You didn't provide an API key.","type":"invalid_request_error","param":null,"code":"invalid_api_key"}}` + h := anthropic.ModelsCompat(openAIModelsHandler(http.StatusUnauthorized, openAIRefusal)) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("x-api-key", "hk_bad") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusUnauthorized { + t.Fatalf("status: want 401 got %d", rec.Code) + } + var body map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatalf("decode: %v", err) + } + if body["type"] != "error" { + t.Errorf("envelope: want top-level type=error got %v (%s)", body["type"], rec.Body.String()) + } + errObj, _ := body["error"].(map[string]any) + if errObj["type"] != "authentication_error" { + t.Errorf("error.type: want authentication_error got %v", errObj["type"]) + } + if errObj["message"] != "You didn't provide an API key." { + t.Errorf("error.message: want the underlying refusal preserved got %v", errObj["message"]) + } +} + +// TestModelsCompat_OpenAIClientIsUntouched is the non-regression half. Open +// WebUI's model picker is built from this route and parses the OpenAI shape; +// re-shaping unconditionally would empty it. +func TestModelsCompat_OpenAIClientIsUntouched(t *testing.T) { + h := anthropic.ModelsCompat(openAIModelsHandler(http.StatusOK, openAIModelListBody)) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("Authorization", "Bearer hk_live_test") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + if rec.Body.String() != openAIModelListBody { + t.Fatalf("an OpenAI-shaped caller must get the byte-identical OpenAI body\nwant: %s\ngot: %s", + openAIModelListBody, rec.Body.String()) + } +} + +func TestIsAnthropicClient(t *testing.T) { + cases := []struct { + name string + headers map[string]string + want bool + }{ + {"x-api-key, the SDK default credential", map[string]string{"x-api-key": "hk_1"}, true}, + {"anthropic-version, sent by every official SDK", map[string]string{"anthropic-version": "2023-06-01"}, true}, + {"bearer only, an OpenAI-shaped caller", map[string]string{"Authorization": "Bearer hk_1"}, false}, + {"no credential at all", nil, false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + for k, v := range tc.headers { + req.Header.Set(k, v) + } + if got := anthropic.IsAnthropicClient(req); got != tc.want { + t.Errorf("IsAnthropicClient = %v, want %v", got, tc.want) + } + }) + } +} + +// TestModelsCompat_RefusalCarriesTheRetryHeaders is the regression guard for +// the header half of a refusal. handleModels answers an unauthorized or +// throttled API-key caller through apierrors.WriteAuthFailure, which delivers +// the status and message in the body and the retry metadata in the headers. +// Buffering the body for reshaping and dropping the headers would leave an +// Anthropic SDK backing off on its own default schedule against a gateway that +// just told it exactly how long to wait. +func TestModelsCompat_RefusalCarriesTheRetryHeaders(t *testing.T) { + upstream := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.Header().Set("Retry-After", "30") + w.Header().Set("X-Ratelimit-Limit-Requests", "100") + w.Header().Set("X-Ratelimit-Remaining-Requests", "0") + // A stale length for the OpenAI-shaped body, which must not follow the + // shorter Anthropic envelope onto the wire. + w.Header().Set("Content-Length", "4096") + w.WriteHeader(http.StatusTooManyRequests) + _, _ = w.Write([]byte(`{"error":{"message":"Rate limit reached.","type":"rate_limit_error","code":"rate_limit_exceeded"}}`)) + }) + + h := anthropic.ModelsCompat(upstream) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("anthropic-version", "2023-06-01") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusTooManyRequests { + t.Fatalf("status: want 429 got %d body=%s", rec.Code, rec.Body.String()) + } + if got := rec.Header().Get("Retry-After"); got != "30" { + t.Errorf("retry-after: want 30 got %q", got) + } + if got := rec.Header().Get("X-Ratelimit-Limit-Requests"); got != "100" { + t.Errorf("x-ratelimit-limit-requests: want 100 got %q", got) + } + if got := rec.Header().Get("X-Ratelimit-Remaining-Requests"); got != "0" { + t.Errorf("x-ratelimit-remaining-requests: want 0 got %q", got) + } + if got := rec.Header().Get("Content-Length"); got == "4096" { + t.Errorf("content-length: the delegated body length must not describe the reshaped one, got %q", got) + } + if got := rec.Header().Get("Content-Type"); got != "application/json" { + t.Errorf("content-type: want application/json got %q", got) + } + var body map[string]any + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatalf("decode: %v (body=%s)", err, rec.Body.String()) + } + if body["type"] != "error" { + t.Errorf("envelope: want top-level type=error got %v", body["type"]) + } +} + +// TestModelsCompat_DeclaresVary pins the cache-correctness half of serving two +// representations from one URL. Nothing on this route sets Cache-Control and +// the route needs a credential, so no correct cache stores it today; the +// declaration is what keeps a future edge cache, or an intermediary keying on +// URL alone, from handing an Anthropic-shaped body to Open WebUI and emptying +// its model picker. +func TestModelsCompat_DeclaresVary(t *testing.T) { + h := anthropic.ModelsCompat(openAIModelsHandler(http.StatusOK, openAIModelListBody)) + + for _, tc := range []struct { + name string + headers map[string]string + }{ + {name: "anthropic client", headers: map[string]string{"x-api-key": "hk_live_test"}}, + {name: "openai client", headers: map[string]string{"Authorization": "Bearer hk_live_test"}}, + } { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + for k, v := range tc.headers { + req.Header.Set(k, v) + } + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + got := rec.Header().Get("Vary") + if !strings.Contains(got, "anthropic-version") || !strings.Contains(got, "x-api-key") { + t.Errorf("vary: want both request headers that select the representation, got %q", got) + } + }) + } +} + +// TestModelsCompat_VaryDoesNotClobberAnExistingDeclaration keeps the wrapper +// additive. Nothing in edge-api declares a Vary on this route today, so Set +// would be harmless right now and wrong the moment a CORS layer or any other +// outer middleware declares Vary: Origin ahead of it. +func TestModelsCompat_VaryDoesNotClobberAnExistingDeclaration(t *testing.T) { + outer := func(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Vary", "Origin") + next.ServeHTTP(w, r) + }) + } + + h := outer(anthropic.ModelsCompat(openAIModelsHandler(http.StatusOK, openAIModelListBody))) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("x-api-key", "hk_live_test") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + got := strings.Join(rec.Header().Values("Vary"), ", ") + for _, want := range []string{"Origin", "anthropic-version", "x-api-key"} { + if !strings.Contains(got, want) { + t.Errorf("vary: want %q preserved, got %q", want, got) + } + } +} + +// TestModelsCompat_SuccessCarriesTheDelegatedHeaders covers the 2xx half of the +// same carry-over the refusal path needs. The success path re-encodes the body +// rather than reshaping it, which made it easy to leave out, and a header a +// delegated handler sets alongside a 200 is as much part of that response as +// the body is. +func TestModelsCompat_SuccessCarriesTheDelegatedHeaders(t *testing.T) { + upstream := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.Header().Set("X-Ratelimit-Remaining-Requests", "42") + w.Header().Set("Content-Length", "4096") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(openAIModelListBody)) + }) + + h := anthropic.ModelsCompat(upstream) + req := httptest.NewRequest(http.MethodGet, "/v1/models", nil) + req.Header.Set("x-api-key", "hk_live_test") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusOK { + t.Fatalf("status: want 200 got %d body=%s", rec.Code, rec.Body.String()) + } + if got := rec.Header().Get("X-Ratelimit-Remaining-Requests"); got != "42" { + t.Errorf("x-ratelimit-remaining-requests: want 42 got %q", got) + } + if got := rec.Header().Get("Content-Length"); got == "4096" { + t.Errorf("content-length: the delegated body length must not describe the re-encoded one, got %q", got) + } + if got := rec.Header().Get("Content-Type"); got != "application/json" { + t.Errorf("content-type: want application/json got %q", got) + } + var body anthropic.ModelsListResponse + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatalf("decode: %v (body=%s)", err, rec.Body.String()) + } + if len(body.Data) != 2 { + t.Errorf("data: want 2 entries got %d", len(body.Data)) + } +} diff --git a/apps/edge-api/internal/anthropic/stream.go b/apps/edge-api/internal/anthropic/stream.go index d89cea8a7..d9b81f989 100644 --- a/apps/edge-api/internal/anthropic/stream.go +++ b/apps/edge-api/internal/anthropic/stream.go @@ -145,17 +145,7 @@ func (t *SSETranslator) FeedLine(line []byte) bool { return false } - if !t.started { - t.started = true - // Always mint a fresh gateway-owned id, one per stream, reused on - // every subsequent chunk via t.messageID. Never derive from - // chunk.ID: that carries the upstream's own id verbatim (OpenRouter - // "gen-*", Groq "chatcmpl-*"), and the old "msg_"-prefix behavior - // still shipped that raw upstream id to the client -- see - // FromOAIResponse for the non-streaming twin of this fix. - t.messageID = "msg_" + uuid.New().String() - t.emitMessageStart() - } + t.ensureStarted() if chunk.Usage != nil { cacheRead, cacheWrite := 0, 0 @@ -204,6 +194,14 @@ func (t *SSETranslator) Finish() { } t.done = true + // An upstream stream that yielded no parseable chunk at all never reached + // emitMessageStart, and the terminal sequence alone is not a stream any + // Anthropic client can read: the SDK accumulator raises a RuntimeError for + // unexpected event order when a message_delta arrives while its snapshot is + // still nil. Opening the message here costs two events and turns that crash + // into an empty turn. + t.ensureStarted() + if t.hasOpenBlock { t.emitContentBlockStop(t.openBlockIndex) t.hasOpenBlock = false @@ -212,6 +210,23 @@ func (t *SSETranslator) Finish() { t.emitMessageStop() } +// ensureStarted emits message_start exactly once per stream, whether the first +// upstream chunk or the terminal sequence gets there first. +func (t *SSETranslator) ensureStarted() { + if t.started { + return + } + t.started = true + // Always mint a fresh gateway-owned id, one per stream, reused on every + // subsequent chunk via t.messageID. Never derive from chunk.ID: that + // carries the upstream own id verbatim (OpenRouter "gen-*", Groq + // "chatcmpl-*"), and the old "msg_"-prefix behavior still shipped that raw + // upstream id to the client -- see FromOAIResponse for the non-streaming + // twin of this fix. + t.messageID = "msg_" + uuid.New().String() + t.emitMessageStart() +} + // WriteErr returns the first write error the translator hit, if any. func (t *SSETranslator) WriteErr() error { return t.writeErr } @@ -256,7 +271,7 @@ func (t *SSETranslator) openTextBlock() { Index: indexPtr(0), ContentBlock: &StreamContentBlock{ Type: "text", - Text: "", + Text: emptyTextPtr(), }, }) } @@ -310,9 +325,10 @@ func (t *SSETranslator) handleToolCallDelta(tc OAIToolCallDelta) { Type: "content_block_start", Index: indexPtr(nextIndex), ContentBlock: &StreamContentBlock{ - Type: "tool_use", - ID: tc.ID, - Name: tc.Function.Name, + Type: "tool_use", + ID: tc.ID, + Name: tc.Function.Name, + Input: json.RawMessage(`{}`), }, }) } @@ -342,6 +358,11 @@ func (t *SSETranslator) emitContentBlockStop(index int) { // content_block_* call site. func indexPtr(i int) *int { return &i } +// emptyTextPtr exists for the same reason indexPtr does: StreamContentBlock.Text +// is a pointer so the empty string still serializes (see its doc comment), and a +// composite literal cannot take the address of a string constant inline. +func emptyTextPtr() *string { s := ""; return &s } + func (t *SSETranslator) emitMessageDelta() { if t.stopReason == "" { t.stopReason = "end_turn" diff --git a/apps/edge-api/internal/anthropic/stream_test.go b/apps/edge-api/internal/anthropic/stream_test.go index 621848df9..8de67ea58 100644 --- a/apps/edge-api/internal/anthropic/stream_test.go +++ b/apps/edge-api/internal/anthropic/stream_test.go @@ -486,3 +486,170 @@ func TestSSETranslator_FirstBlockIndexIsPresentAndZero(t *testing.T) { t.Fatal("no content_block_* events were found to check") } } + +// TestSSETranslator_TextBlockStartCarriesExplicitEmptyText pins the exact wire +// bytes issue #1274 was about, not merely the presence of the event. +// +// The real Anthropic API emits +// {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}} +// (verified against https://docs.claude.com/en/docs/build-with-claude/streaming +// on 2026-08-28). The Anthropic SDK's stream accumulator seeds its snapshot +// from that "text" key and then runs `content.text += event.delta.text` on the +// first content_block_delta, so dropping the key for the empty string made +// every streamed text response a TypeError inside the SDK's own documented +// helper. The assertion is therefore on KEY PRESENCE, which is the thing that +// was actually wrong; asserting the value alone would pass just as happily +// against a missing key decoded as a nil interface. +func TestSSETranslator_TextBlockStartCarriesExplicitEmptyText(t *testing.T) { + stream := buildOAIStream( + `{"id":"chatcmpl-1","model":"route-upstream","choices":[{"index":0,"delta":{"content":"Hello"},"finish_reason":null}]}`, + `{"id":"chatcmpl-1","model":"route-upstream","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}`, + ) + rec := httptest.NewRecorder() + tr := anthropic.NewSSETranslator(rec, "claude-3-haiku") + if err := tr.Translate(strings.NewReader(stream)); err != nil { + t.Fatalf("translate error: %v", err) + } + + var seen bool + for _, ev := range parseSSEEvents(t, rec.Body.String()) { + if ev["type"] != "content_block_start" { + continue + } + block, ok := ev["content_block"].(map[string]interface{}) + if !ok { + t.Fatalf("content_block_start has no content_block object: %v", ev) + } + if block["type"] != "text" { + continue + } + seen = true + text, present := block["text"] + if !present { + t.Fatalf("text content_block_start omits the \"text\" key entirely, which crashes "+ + "the Anthropic SDK accumulator on the first delta: %v", block) + } + if text != "" { + t.Errorf("content_block.text: want empty string got %v", text) + } + } + if !seen { + t.Fatal("no text content_block_start was emitted at all") + } +} + +// TestSSETranslator_ToolUseBlockStartCarriesExplicitEmptyInput is the same +// defect one block type over: Anthropic's documented tool_use block start +// carries "input":{}, and a nil json.RawMessage under `omitempty` drops it. +// The Python SDK's accumulator happens to survive this one (it assigns +// content.input from its own partial-JSON buffer rather than reading the +// block-start value), but a client that types input as a required object does +// not, and the wire shape is wrong either way. +func TestSSETranslator_ToolUseBlockStartCarriesExplicitEmptyInput(t *testing.T) { + stream := buildOAIStream( + `{"id":"chatcmpl-1","model":"route-upstream","choices":[{"index":0,"delta":{"tool_calls":[{"index":0,"id":"call_1","type":"function","function":{"name":"get_weather","arguments":""}}]},"finish_reason":null}]}`, + `{"id":"chatcmpl-1","model":"route-upstream","choices":[{"index":0,"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{\"city\":\"Dhaka\"}"}}]},"finish_reason":null}]}`, + `{"id":"chatcmpl-1","model":"route-upstream","choices":[{"index":0,"delta":{},"finish_reason":"tool_calls"}]}`, + ) + rec := httptest.NewRecorder() + tr := anthropic.NewSSETranslator(rec, "claude-3-haiku") + if err := tr.Translate(strings.NewReader(stream)); err != nil { + t.Fatalf("translate error: %v", err) + } + + var seen bool + for _, ev := range parseSSEEvents(t, rec.Body.String()) { + if ev["type"] != "content_block_start" { + continue + } + block, ok := ev["content_block"].(map[string]interface{}) + if !ok || block["type"] != "tool_use" { + continue + } + seen = true + input, present := block["input"] + if !present { + t.Fatalf("tool_use content_block_start omits the \"input\" key: %v", block) + } + obj, ok := input.(map[string]interface{}) + if !ok || len(obj) != 0 { + t.Errorf("content_block.input: want empty object got %v", input) + } + if _, hasText := block["text"]; hasText { + t.Errorf("tool_use content_block_start must not carry a text key: %v", block) + } + } + if !seen { + t.Fatal("no tool_use content_block_start was emitted at all") + } +} + +// TestSSETranslator_EmptyStreamStillOpensTheMessage pins event ORDER on a +// stream that yielded no parseable chunk, which the sibling +// TestSSETranslator_EmptyStream does not: it asserts message_delta and +// message_stop are present but not that anything precedes them. Finish used to +// emit the terminal pair unconditionally, so a caller that got no chunk at all +// received message_delta as the first event of the stream, and the SDK +// accumulator raises an unexpected-event-order RuntimeError when a +// message_delta arrives while its snapshot is still nil. +func TestSSETranslator_EmptyStreamStillOpensTheMessage(t *testing.T) { + rec := httptest.NewRecorder() + tr := anthropic.NewSSETranslator(rec, "m") + if err := tr.Translate(strings.NewReader("data: [DONE]\n\n")); err != nil { + t.Fatalf("translate error: %v", err) + } + + events := parseSSEEvents(t, rec.Body.String()) + if len(events) == 0 { + t.Fatal("no events emitted for an empty stream") + } + if events[0]["type"] != "message_start" { + t.Fatalf("first event: want message_start got %v (body=%s)", events[0]["type"], rec.Body.String()) + } + + var startAt, deltaAt = -1, -1 + for i, ev := range events { + switch ev["type"] { + case "message_start": + if startAt == -1 { + startAt = i + } else { + t.Errorf("message_start emitted more than once, at %d and %d", startAt, i) + } + case "message_delta": + if deltaAt == -1 { + deltaAt = i + } + } + } + if deltaAt == -1 { + t.Fatal("missing message_delta on empty stream") + } + if startAt > deltaAt { + t.Errorf("message_start at %d must precede message_delta at %d", startAt, deltaAt) + } +} + +// TestSSETranslator_FinishAfterAStartedStreamDoesNotRepeatMessageStart is the +// other half of that guard: opening the message from Finish must not add a +// second message_start to a stream that already carried one. +func TestSSETranslator_FinishAfterAStartedStreamDoesNotRepeatMessageStart(t *testing.T) { + stream := buildOAIStream( + `{"id":"chatcmpl-x","model":"m","choices":[{"index":0,"delta":{"content":"hi"},"finish_reason":"stop"}]}`, + ) + rec := httptest.NewRecorder() + tr := anthropic.NewSSETranslator(rec, "m") + if err := tr.Translate(strings.NewReader(stream)); err != nil { + t.Fatalf("translate error: %v", err) + } + + starts := 0 + for _, ev := range parseSSEEvents(t, rec.Body.String()) { + if ev["type"] == "message_start" { + starts++ + } + } + if starts != 1 { + t.Errorf("message_start count: want exactly 1 got %d (body=%s)", starts, rec.Body.String()) + } +} diff --git a/apps/edge-api/internal/anthropic/translate_response.go b/apps/edge-api/internal/anthropic/translate_response.go index f484c8426..65b5c4c85 100644 --- a/apps/edge-api/internal/anthropic/translate_response.go +++ b/apps/edge-api/internal/anthropic/translate_response.go @@ -34,7 +34,15 @@ func FromOAIResponse(resp OAIResponse, clientAlias string) MessagesResponse { Type: "message", Role: "assistant", Model: model, - Usage: anthropicUsage(resp.Usage, model, resp.Model), + // Content is seeded with an empty, non-nil slice, and appended to + // below, so that a turn which produced neither text nor a tool call + // still serializes as "content":[]. A nil Go slice marshals to JSON + // null, and Anthropic's contract is that content is ALWAYS an array: + // every typed client iterates it unconditionally, so null is a + // TypeError rather than an empty turn (issue #1260). Both exits below + // are covered by seeding it here rather than at one of them. + Content: []ResponseBlock{}, + Usage: anthropicUsage(resp.Usage, model, resp.Model), } if len(resp.Choices) == 0 { @@ -45,10 +53,8 @@ func FromOAIResponse(resp OAIResponse, clientAlias string) MessagesResponse { choice := resp.Choices[0] out.StopReason = mapFinishReason(choice.FinishReason) - var blocks []ResponseBlock - if choice.Message.Content != "" { - blocks = append(blocks, ResponseBlock{ + out.Content = append(out.Content, ResponseBlock{ Type: "text", Text: choice.Message.Content, }) @@ -59,7 +65,7 @@ func FromOAIResponse(resp OAIResponse, clientAlias string) MessagesResponse { if err != nil { input = json.RawMessage(fmt.Sprintf(`{"_raw":%q}`, tc.Function.Arguments)) } - blocks = append(blocks, ResponseBlock{ + out.Content = append(out.Content, ResponseBlock{ Type: "tool_use", ID: tc.ID, Name: tc.Function.Name, @@ -67,7 +73,6 @@ func FromOAIResponse(resp OAIResponse, clientAlias string) MessagesResponse { }) } - out.Content = blocks return out } diff --git a/apps/edge-api/internal/anthropic/translate_response_test.go b/apps/edge-api/internal/anthropic/translate_response_test.go index c1a64cce0..8c1231731 100644 --- a/apps/edge-api/internal/anthropic/translate_response_test.go +++ b/apps/edge-api/internal/anthropic/translate_response_test.go @@ -229,3 +229,60 @@ func TestFromOAIResponse_ToolArgumentsInvalidJSON(t *testing.T) { t.Error("input should not be empty on fallback") } } + +// TestFromOAIResponse_EmptyCompletionSerializesContentAsArray is the issue +// #1260 guard. A turn that produced neither text nor a tool call left Content +// as a nil Go slice, and a nil slice with no `omitempty` marshals to JSON +// null, not []. Anthropic's contract is that content is ALWAYS an array, so +// every typed client iterates it unconditionally and gets a TypeError on null +// instead of an empty turn. +// +// The assertion runs on the MARSHALED BYTES, not on the struct. Checking +// len(got.Content) == 0 in Go passes identically for a nil slice and an empty +// one, so it cannot see this defect at all; the difference exists only once +// encoding/json runs. +func TestFromOAIResponse_EmptyCompletionSerializesContentAsArray(t *testing.T) { + cases := map[string]anthropic.OAIResponse{ + "no choices at all": { + ID: "chatcmpl-empty", + Model: "m", + }, + "choice with empty content and no tool calls": { + ID: "chatcmpl-truncated", + Model: "m", + Choices: []anthropic.OAIChoice{ + {FinishReason: "length", Message: anthropic.OAIMsg{Role: "assistant", Content: ""}}, + }, + }, + } + + for name, oai := range cases { + t.Run(name, func(t *testing.T) { + raw, err := json.Marshal(anthropic.FromOAIResponse(oai, "claude-3-haiku")) + if err != nil { + t.Fatalf("marshal: %v", err) + } + if !strings.Contains(string(raw), `"content":[]`) { + t.Fatalf("empty completion must serialize content as [], got %s", raw) + } + + // Decode back the way a real client does, and prove the value is + // an array rather than null. + var decoded map[string]json.RawMessage + if err := json.Unmarshal(raw, &decoded); err != nil { + t.Fatalf("unmarshal: %v", err) + } + content, present := decoded["content"] + if !present { + t.Fatal("content key absent from the response body") + } + var blocks []anthropic.ResponseBlock + if err := json.Unmarshal(content, &blocks); err != nil { + t.Fatalf("content is not an array: %v", err) + } + if blocks == nil { + t.Fatalf("content decoded to null rather than an empty array: %s", content) + } + }) + } +} diff --git a/apps/edge-api/internal/anthropic/translate_writer.go b/apps/edge-api/internal/anthropic/translate_writer.go index d28576e1b..6c419fca9 100644 --- a/apps/edge-api/internal/anthropic/translate_writer.go +++ b/apps/edge-api/internal/anthropic/translate_writer.go @@ -33,6 +33,42 @@ func (r *headerlessRecorder) WriteHeader(status int) { r.status = status } func (r *headerlessRecorder) Write(p []byte) (int, error) { return r.body.Write(p) } +// reshapeInto re-emits the recorded response to the real client writer in the +// Anthropic error envelope, carrying over the headers the recorded writer set. +// +// The carry-over is the point of the method existing at all. apierr.WriteAuthFailure +// delivers half of a refusal through the body and the other half through headers: +// a 429 arrives with retry-after and the x-ratelimit-* family, and both the +// degraded-limiter and upstream_unavailable branches arrive with retry-after. +// Recording the body and dropping the headers turns a retryable refusal into one +// the Anthropic SDK backs off from on its own default schedule, which is exactly +// the collapse that shared writer exists to prevent. Pairing the copy with the +// reshape in one call keeps a call site from taking one half without the other. +// +// Content-Type and Content-Length are deliberately not carried over: the +// reshaped body is a different envelope of a different length, and a stale +// Content-Length would truncate it on the wire. Everything else forwarded here +// is retry metadata that carries no provider identity. +func (r *headerlessRecorder) reshapeInto(w http.ResponseWriter) { + r.copyHeadersTo(w) + reshapeToAnthropicError(w, r.status, r.body.Bytes()) +} + +// copyHeadersTo carries the recorded headers onto the real client writer. It is +// separate from reshapeInto because a recorded 2xx is re-encoded rather than +// reshaped and needs the same carry-over: whatever a delegated handler sets +// alongside a success is as much part of that response as the body is. +func (r *headerlessRecorder) copyHeadersTo(w http.ResponseWriter) { + for key, values := range r.header { + if key == "Content-Type" || key == "Content-Length" { + continue + } + for _, value := range values { + w.Header().Add(key, value) + } + } +} + // maxTranslatedBodyBytes bounds how much of a buffered response the translator // holds in memory, matching the ceiling the OpenAI sync path already applies // when reading an upstream body. diff --git a/apps/edge-api/internal/anthropic/types.go b/apps/edge-api/internal/anthropic/types.go index a932dbd8c..078271f2f 100644 --- a/apps/edge-api/internal/anthropic/types.go +++ b/apps/edge-api/internal/anthropic/types.go @@ -320,11 +320,32 @@ type StreamMessage struct { } // StreamContentBlock is the block header in content_block_start. +// +// Text and Input are pointers for the same reason StreamEvent.Index is: the +// real Anthropic protocol requires the empty value to be PRESENT on the wire, +// and `omitempty` on a plain Go string or a nil json.RawMessage drops the key +// for the empty value exactly as readily as for "unset". Verified against the +// live streaming specification on 2026-08-28 +// (https://docs.claude.com/en/docs/build-with-claude/streaming), which emits +// +// {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}} +// {"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"toolu_...","name":"get_weather","input":{}}} +// +// on every text and tool_use block respectively. The text case was not +// cosmetic: the Anthropic SDK's own stream accumulator seeds its snapshot from +// this field and then does `content.text += delta.text` on the first +// content_block_delta, so a missing key made that a None += str TypeError on +// every streamed text response (issue #1274). +// +// A text block never carries id/name/input and a tool_use block never carries +// text, so each stays absent by leaving the corresponding field at its zero +// value; only the block's own required-but-empty field is made explicit. type StreamContentBlock struct { - Type string `json:"type"` // "text" | "tool_use" - ID string `json:"id,omitempty"` - Name string `json:"name,omitempty"` - Text string `json:"text,omitempty"` + Type string `json:"type"` // "text" | "tool_use" + ID string `json:"id,omitempty"` + Name string `json:"name,omitempty"` + Text *string `json:"text,omitempty"` + Input json.RawMessage `json:"input,omitempty"` } // StreamDelta is the delta payload in content_block_delta or message_delta.