Skip to content

Commit b58867f

Browse files
authored
Merge branch 'main' into mcp-claim-headers
2 parents 5b02445 + 7cfd3c3 commit b58867f

21 files changed

Lines changed: 2808 additions & 17 deletions

cmd/aigw/.env.otel.otel-tui

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,6 @@
11
# otel-tui configuration - Terminal UI for OpenTelemetry
22
OTEL_SERVICE_NAME=aigw
3-
# TODO: Next release after Envoy Gateway v1.7.0-rc.1 revert to normal service name
4-
# OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-tui:4317
5-
OTEL_EXPORTER_OTLP_ENDPOINT=http://host.docker.internal:4317
3+
OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-tui:4317
64
OTEL_EXPORTER_OTLP_PROTOCOL=grpc
75
# Reduce trace and metrics export delay for demo purposes
86
OTEL_BSP_SCHEDULE_DELAY=100

cmd/aigw/.env.otel.phoenix

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,7 @@
11
# Phoenix configuration - LLM-specific observability
22
OTEL_SERVICE_NAME=aigw
33
# Phoenix uses port 4317 for OTLP gRPC, port 6006 for web UI
4-
# TODO: Next release after Envoy Gateway v1.7.0-rc.1 revert to normal service name
5-
# OTEL_EXPORTER_OTLP_ENDPOINT=http://phoenix:4317
6-
OTEL_EXPORTER_OTLP_ENDPOINT=http://host.docker.internal:4317
4+
OTEL_EXPORTER_OTLP_ENDPOINT=http://phoenix:4317
75
OTEL_EXPORTER_OTLP_PROTOCOL=grpc
86
# Phoenix only supports traces, not metrics or logs
97
OTEL_METRICS_EXPORTER=none
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
# Copyright Envoy AI Gateway Authors
2+
# SPDX-License-Identifier: Apache-2.0
3+
# The full text of the Apache license is available in the LICENSE file at
4+
# the root of the repo.
5+
6+
# This example routes /v1/messages (Anthropic Messages API) requests to a
7+
# local vLLM backend.
8+
# The AIServiceBackend schema is set to OpenAI so the gateway passes through
9+
# the Anthropic request format with translation.
10+
11+
apiVersion: aigateway.envoyproxy.io/v1alpha1
12+
kind: AIGatewayRoute
13+
metadata:
14+
name: envoy-ai-gateway-basic-anthropic-openai
15+
namespace: default
16+
spec:
17+
parentRefs:
18+
- name: envoy-ai-gateway-basic
19+
kind: Gateway
20+
group: gateway.networking.k8s.io
21+
rules:
22+
- matches:
23+
- headers:
24+
- type: Exact
25+
name: x-ai-eg-model
26+
value: Qwen/Qwen2.5-0.5B-Instruct # Replace with the model name served by your vLLM instance.
27+
backendRefs:
28+
- name: envoy-ai-gateway-basic-anthropic-openai
29+
---
30+
apiVersion: aigateway.envoyproxy.io/v1alpha1
31+
kind: AIServiceBackend
32+
metadata:
33+
name: envoy-ai-gateway-basic-anthropic-openai
34+
namespace: default
35+
spec:
36+
schema:
37+
name: OpenAI # vLLM exposes an OpenAI-compatible API, so Anthropic Messages Requests are translated
38+
backendRef:
39+
name: envoy-ai-gateway-basic-openai
40+
kind: Backend
41+
group: gateway.envoyproxy.io
42+
---
43+
apiVersion: gateway.envoyproxy.io/v1alpha1
44+
kind: Backend
45+
metadata:
46+
name: envoy-ai-gateway-basic-anthropic-openai
47+
namespace: default
48+
spec:
49+
endpoints:
50+
- ip:
51+
address: 0.0.0.0 # Replace with your vLLM service hostname or IP (e.g. localhost's internal IP from kind cluster).
52+
port: 8000 # Replace with the port your vLLM instance listens on (default: 8000).

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ require (
1919
github.com/cohere-ai/cohere-go/v2 v2.16.2
2020
github.com/coreos/go-oidc/v3 v3.17.0
2121
github.com/docker/docker v28.5.2+incompatible
22-
github.com/envoyproxy/gateway v1.7.0-rc.1
22+
github.com/envoyproxy/gateway v1.7.0
2323
github.com/envoyproxy/go-control-plane v0.14.0
2424
github.com/envoyproxy/go-control-plane/envoy v1.36.1-0.20260115164926-066cbd5b3989
2525
github.com/go-logr/logr v1.4.3

go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -156,8 +156,8 @@ github.com/ebitengine/purego v0.9.1 h1:a/k2f2HQU3Pi399RPW1MOaZyhKJL9w/xFpKAg4q1s
156156
github.com/ebitengine/purego v0.9.1/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ=
157157
github.com/emicklei/go-restful/v3 v3.13.0 h1:C4Bl2xDndpU6nJ4bc1jXd+uTmYPVUwkD6bFY/oTyCes=
158158
github.com/emicklei/go-restful/v3 v3.13.0/go.mod h1:6n3XBCmQQb25CM2LCACGz8ukIrRry+4bhvbpWn3mrbc=
159-
github.com/envoyproxy/gateway v1.7.0-rc.1 h1:82R4u70KdM1Bx1lwliDcdTgkPOYOcFUtL1vipoM6/ow=
160-
github.com/envoyproxy/gateway v1.7.0-rc.1/go.mod h1:R/4TXMCrgBmpvROacQ9boJksCAhik8lTY5DFYHRJvFE=
159+
github.com/envoyproxy/gateway v1.7.0 h1:noVz4fADhljSKpQFaspb25+C3kC8bBBisQZle9W0ei4=
160+
github.com/envoyproxy/gateway v1.7.0/go.mod h1:U/UXS8G6Q5SldcKubP8PZazD9OsNT0Mf1XzCVig0KkQ=
161161
github.com/envoyproxy/go-control-plane v0.14.0 h1:hbG2kr4RuFj222B6+7T83thSPqLjwBIfQawTkC++2HA=
162162
github.com/envoyproxy/go-control-plane v0.14.0/go.mod h1:NcS5X47pLl/hfqxU70yPwL9ZMkUlwlKxtAohpi2wBEU=
163163
github.com/envoyproxy/go-control-plane/contrib v1.36.1-0.20260115164926-066cbd5b3989 h1:KTd1TJym7dgV1L1XlxXeJNct7rJI3xTV+iuArq40wm0=

internal/endpointspec/endpointspec.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -343,8 +343,11 @@ func (MessagesEndpointSpec) GetTranslator(schema filterapi.VersionedAPISchema, m
343343
return translator.NewAnthropicToAWSAnthropicTranslator(schema.Version, modelNameOverride), nil
344344
case filterapi.APISchemaAnthropic:
345345
return translator.NewAnthropicToAnthropicTranslator(schema.Version, modelNameOverride), nil
346+
case filterapi.APISchemaOpenAI:
347+
// The Anthropic prefix can be altered using values.yaml if necessary
348+
return translator.NewAnthropicToChatCompletionOpenAITranslator(schema.Version, modelNameOverride), nil
346349
default:
347-
return nil, fmt.Errorf("/v1/messages endpoint only supports backends that return native Anthropic format (Anthropic, GCPAnthropic, AWSAnthropic). Backend %s uses different model format", schema.Name)
350+
return nil, fmt.Errorf("/v1/messages endpoint only supports backends that return native Anthropic format (Anthropic, GCPAnthropic, AWSAnthropic). OpenAI translation is also supported. Backend %s uses different model format", schema.Name)
348351
}
349352
}
350353

internal/endpointspec/endpointspec_test.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -236,13 +236,14 @@ func TestMessagesEndpointSpec_GetTranslator(t *testing.T) {
236236
{Name: filterapi.APISchemaGCPAnthropic},
237237
{Name: filterapi.APISchemaAWSAnthropic},
238238
{Name: filterapi.APISchemaAnthropic},
239+
{Name: filterapi.APISchemaOpenAI}, // This is for OpenAI-schema backends like vLLM that support the /v1/messages endpoint
239240
} {
240241
translator, err := spec.GetTranslator(schema, "override")
241242
require.NoError(t, err)
242243
require.NotNil(t, translator)
243244
}
244245

245-
_, err := spec.GetTranslator(filterapi.VersionedAPISchema{Name: filterapi.APISchemaOpenAI}, "override")
246+
_, err := spec.GetTranslator(filterapi.VersionedAPISchema{Name: filterapi.APISchemaCohere}, "override")
246247
require.ErrorContains(t, err, "only supports")
247248
}
248249

Lines changed: 275 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,275 @@
1+
// Copyright Envoy AI Gateway Authors
2+
// SPDX-License-Identifier: Apache-2.0
3+
// The full text of the Apache license is available in the LICENSE file at
4+
// the root of the repo.
5+
6+
package translator
7+
8+
import (
9+
"cmp"
10+
"fmt"
11+
"io"
12+
"log/slog"
13+
"strconv"
14+
"strings"
15+
16+
"github.com/tidwall/sjson"
17+
18+
"github.com/envoyproxy/ai-gateway/internal/apischema/anthropic"
19+
"github.com/envoyproxy/ai-gateway/internal/apischema/openai"
20+
"github.com/envoyproxy/ai-gateway/internal/internalapi"
21+
"github.com/envoyproxy/ai-gateway/internal/json"
22+
"github.com/envoyproxy/ai-gateway/internal/metrics"
23+
"github.com/envoyproxy/ai-gateway/internal/tracing/tracingapi"
24+
)
25+
26+
// NewAnthropicToChatCompletionOpenAITranslator implements [Factory] for Anthropic to OpenAI ChatCompletion translation.
27+
// This translator converts Anthropic API format to OpenAI ChatCompletion API requests.
28+
func NewAnthropicToChatCompletionOpenAITranslator(version string, modelNameOverride internalapi.ModelNameOverride) AnthropicMessagesTranslator {
29+
// TODO: use "version" in APISchema struct to set the specific prefix if needed like OpenAI does. However, two questions:
30+
// * Is there any "Anthropic compatible" API that uses a different prefix like OpenAI does?
31+
// * Even if there is, we should refactor the APISchema struct to have "prefix" field instead of abusing "version" field.
32+
_ = version
33+
passthroughTranslator := NewAnthropicToAnthropicTranslator(version, modelNameOverride)
34+
return &anthropicToOpenAIV1ChatCompletionTranslator{passthroughTranslator: &passthroughTranslator, modelNameOverride: modelNameOverride}
35+
}
36+
37+
type anthropicToOpenAIV1ChatCompletionTranslator struct {
38+
passthroughTranslator *AnthropicMessagesTranslator
39+
modelNameOverride internalapi.ModelNameOverride
40+
requestModel internalapi.RequestModel
41+
stream bool
42+
streamState *openAIStreamToAnthropicState
43+
// Redaction configuration for debug logging
44+
debugLogEnabled bool
45+
enableRedaction bool
46+
logger *slog.Logger
47+
}
48+
49+
// RequestBody implements [AnthropicMessagesTranslator.RequestBody].
50+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) RequestBody(_ []byte, body *anthropic.MessagesRequest, _ bool) (
51+
newHeaders []internalapi.Header, newBody []byte, err error,
52+
) {
53+
// Set translator config based on Anthropic message request
54+
a.stream = body.Stream
55+
// Store the request model to use as fallback for response model
56+
a.requestModel = cmp.Or(a.modelNameOverride, body.Model)
57+
58+
// Convert Anthropic message request body to OpenAI format.
59+
openAIReq := buildOpenAIChatCompletionRequest(body, a.modelNameOverride)
60+
61+
newBody, err = json.Marshal(openAIReq)
62+
if err != nil {
63+
return nil, nil, fmt.Errorf("failed to marshal OpenAI request: %w", err)
64+
}
65+
66+
// Add stop sequences via sjson because ChatCompletionNewParamsStopUnion (from the external openai-go SDK)
67+
// requires importing the external package. Using sjson avoids that dependency.
68+
if len(body.StopSequences) > 0 {
69+
newBody, err = sjson.SetBytesOptions(newBody, "stop", body.StopSequences, sjsonOptions)
70+
if err != nil {
71+
return nil, nil, fmt.Errorf("failed to set stop sequences: %w", err)
72+
}
73+
}
74+
75+
if body.Stream {
76+
a.streamState = &openAIStreamToAnthropicState{
77+
activeTools: make(map[int64]*streamToolCall),
78+
requestModel: a.requestModel,
79+
}
80+
}
81+
82+
newHeaders = []internalapi.Header{
83+
{pathHeaderName, "/v1/chat/completions"},
84+
{contentLengthHeaderName, strconv.Itoa(len(newBody))},
85+
}
86+
return
87+
}
88+
89+
// ResponseHeaders implements [AnthropicMessagesTranslator.ResponseHeaders].
90+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) ResponseHeaders(_ map[string]string) (
91+
newHeaders []internalapi.Header, err error,
92+
) {
93+
return nil, nil
94+
}
95+
96+
// ResponseBody implements [AnthropicMessagesTranslator.ResponseBody].
97+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) ResponseBody(_ map[string]string, body io.Reader, endOfStream bool, span tracingapi.MessageSpan) (
98+
newHeaders []internalapi.Header, newBody []byte, tokenUsage metrics.TokenUsage, responseModel string, err error,
99+
) {
100+
if a.stream {
101+
return a.responseBodyStreaming(body, endOfStream)
102+
}
103+
return a.responseBodyNonStreaming(body, span)
104+
}
105+
106+
// responseBodyNonStreaming converts an OpenAI ChatCompletionResponse to Anthropic MessagesResponse format.
107+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) responseBodyNonStreaming(body io.Reader, span tracingapi.MessageSpan) (
108+
newHeaders []internalapi.Header, newBody []byte, tokenUsage metrics.TokenUsage, responseModel string, err error,
109+
) {
110+
responseModel = a.requestModel
111+
112+
openAIResp := &openai.ChatCompletionResponse{}
113+
if err = json.NewDecoder(body).Decode(openAIResp); err != nil {
114+
return nil, nil, tokenUsage, responseModel, fmt.Errorf("failed to unmarshal OpenAI response: %w", err)
115+
}
116+
117+
responseModel = cmp.Or(openAIResp.Model, a.requestModel)
118+
119+
tokenUsage = metrics.ExtractTokenUsageFromExplicitCaching(
120+
int64(openAIResp.Usage.PromptTokens),
121+
int64(openAIResp.Usage.CompletionTokens),
122+
nil,
123+
nil,
124+
)
125+
126+
anthropicResp := openAIResponseToAnthropic(openAIResp, responseModel)
127+
128+
// Redact and log response when enabled
129+
if a.debugLogEnabled && a.enableRedaction && a.logger != nil {
130+
redactedResp := a.RedactAnthropicBody(anthropicResp)
131+
if jsonBody, marshalErr := json.Marshal(redactedResp); marshalErr == nil {
132+
a.logger.Debug("response body processing", slog.Any("response", string(jsonBody)))
133+
}
134+
}
135+
136+
if span != nil {
137+
span.RecordResponse(anthropicResp)
138+
}
139+
140+
newBody, err = json.Marshal(anthropicResp)
141+
if err != nil {
142+
return nil, nil, tokenUsage, responseModel, fmt.Errorf("failed to marshal Anthropic response: %w", err)
143+
}
144+
newHeaders = []internalapi.Header{{contentLengthHeaderName, strconv.Itoa(len(newBody))}}
145+
return
146+
}
147+
148+
// responseBodyStreaming handles converting OpenAI SSE chunks to Anthropic SSE events.
149+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) responseBodyStreaming(body io.Reader, endOfStream bool) (
150+
newHeaders []internalapi.Header, newBody []byte, tokenUsage metrics.TokenUsage, responseModel string, err error,
151+
) {
152+
responseModel = a.requestModel
153+
154+
if a.streamState == nil {
155+
return nil, nil, tokenUsage, responseModel, fmt.Errorf("stream state not initialized")
156+
}
157+
158+
// Read body into streamState's buffer
159+
if _, err = a.streamState.buffer.ReadFrom(body); err != nil {
160+
return nil, nil, tokenUsage, responseModel, fmt.Errorf("failed to read stream body: %w", err)
161+
}
162+
163+
// Initialize out as a non-nil empty slice so that if no Anthropic events are emitted
164+
// (e.g., for finish_reason-only chunks or [DONE]), we still return a non-nil newBody.
165+
// A non-nil empty body tells Envoy to replace the chunk with nothing, suppressing the
166+
// raw upstream bytes instead of passing them through unchanged.
167+
out := make([]byte, 0)
168+
if err = a.streamState.processBuffer(&out, endOfStream); err != nil {
169+
return nil, nil, tokenUsage, responseModel, err
170+
}
171+
172+
// Update responseModel if updated in streamState or take requested model
173+
responseModel = cmp.Or(a.streamState.model, a.requestModel)
174+
tokenUsage = a.streamState.tokenUsage
175+
176+
// Always return newBody (even if empty) to suppress the original upstream chunk.
177+
newBody = out
178+
return
179+
}
180+
181+
// ResponseError implements [AnthropicMessagesTranslator] for Anthropic to OpenAI translation.
182+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) ResponseError(respHeaders map[string]string, r io.Reader) (
183+
newHeaders []internalapi.Header,
184+
mutatedBody []byte,
185+
err error,
186+
) {
187+
statusCode := respHeaders[statusHeaderName]
188+
var anthropicError anthropic.ErrorResponse
189+
190+
if strings.Contains(respHeaders[contentTypeHeaderName], jsonContentType) {
191+
// OpenAI backend returned a structured JSON error; translate to Anthropic error format.
192+
var openaiErr openai.Error
193+
if err = json.NewDecoder(r).Decode(&openaiErr); err != nil {
194+
return nil, nil, fmt.Errorf("failed to unmarshal OpenAI error body: %w", err)
195+
}
196+
anthropicError = anthropic.ErrorResponse{
197+
Type: "error",
198+
Error: anthropic.ErrorResponseMessage{
199+
Type: openaiErr.Error.Type,
200+
Message: openaiErr.Error.Message,
201+
},
202+
}
203+
} else {
204+
var buf []byte
205+
buf, err = io.ReadAll(r)
206+
if err != nil {
207+
return nil, nil, fmt.Errorf("failed to read error body: %w", err)
208+
}
209+
var typ string
210+
switch statusCode {
211+
case "400":
212+
typ = "invalid_request_error"
213+
case "401":
214+
typ = "authentication_error"
215+
case "403":
216+
typ = "permission_error"
217+
case "404":
218+
typ = "not_found_error"
219+
case "413":
220+
typ = "request_too_large"
221+
case "429":
222+
typ = "rate_limit_error"
223+
case "500":
224+
typ = "internal_server_error"
225+
case "503":
226+
typ = "service_unavailable_error"
227+
case "529":
228+
typ = "overloaded_error"
229+
default:
230+
typ = "internal_server_error"
231+
}
232+
anthropicError = anthropic.ErrorResponse{
233+
Type: "error", // Always "error" at the top level.
234+
Error: anthropic.ErrorResponseMessage{Type: typ, Message: string(buf)},
235+
}
236+
}
237+
238+
mutatedBody, err = json.Marshal(anthropicError)
239+
if err != nil {
240+
return nil, nil, fmt.Errorf("failed to marshal error body: %w", err)
241+
}
242+
newHeaders = append(newHeaders,
243+
internalapi.Header{contentTypeHeaderName, jsonContentType},
244+
internalapi.Header{contentLengthHeaderName, strconv.Itoa(len(mutatedBody))},
245+
)
246+
return
247+
}
248+
249+
// SetRedactionConfig implements [AnthropicResponseRedactor.SetRedactionConfig].
250+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) SetRedactionConfig(debugLogEnabled, enableRedaction bool, logger *slog.Logger) {
251+
a.debugLogEnabled = debugLogEnabled
252+
a.enableRedaction = enableRedaction
253+
a.logger = logger
254+
}
255+
256+
// RedactAnthropicBody implements [AnthropicResponseRedactor.RedactAnthropicBody].
257+
// Creates a redacted copy of the Anthropic response for safe logging without modifying the original.
258+
func (a *anthropicToOpenAIV1ChatCompletionTranslator) RedactAnthropicBody(resp *anthropic.MessagesResponse) *anthropic.MessagesResponse {
259+
if resp == nil {
260+
return nil
261+
}
262+
263+
// Create a shallow copy of the response
264+
redacted := *resp
265+
266+
// Redact content blocks (contains AI-generated content)
267+
if len(resp.Content) > 0 {
268+
redacted.Content = make([]anthropic.MessagesContentBlock, len(resp.Content))
269+
for i := range resp.Content {
270+
redacted.Content[i] = redactAnthropicContent(&resp.Content[i])
271+
}
272+
}
273+
274+
return &redacted
275+
}

0 commit comments

Comments
 (0)