Skip to content

Commit 3758ae0

Browse files
authored
feat: Add Vertex AI parser support (#884)
* Add vertex parser Signed-off-by: bobzetian <bobzetian@google.com> add vertex parser Signed-off-by: bobzetian <bobzetian@google.com> * update readme, comment on vertexai parser, add unittest for the grpc util Signed-off-by: bobzetian <bobzetian@google.com> * update compression todo github issue Signed-off-by: bobzetian <bobzetian@google.com> * move :path to common place under parsers Signed-off-by: bobzetian <bobzetian@google.com> * fix lint Signed-off-by: bobzetian <bobzetian@google.com> --------- Signed-off-by: bobzetian <bobzetian@google.com>
1 parent f8976ee commit 3758ae0

11 files changed

Lines changed: 577 additions & 44 deletions

File tree

cmd/epp/runner/runner.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,7 @@ import (
8080
testresponsereceived "github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requestcontrol/test/responsereceived"
8181
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requesthandling/parsers/openai"
8282
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requesthandling/parsers/passthrough"
83+
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requesthandling/parsers/vertexai"
8384
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requesthandling/parsers/vllmgrpc"
8485
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/scheduling/filter/prefixcacheaffinity"
8586
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/scheduling/filter/sloheadroomtier"
@@ -502,6 +503,7 @@ func (r *Runner) registerInTreePlugins() {
502503
fwkplugin.Register(openai.OpenAIParserType, openai.OpenAIParserPluginFactory)
503504
fwkplugin.Register(vllmgrpc.VllmGRPCParserType, vllmgrpc.VllmGRPCParserPluginFactory)
504505
fwkplugin.Register(passthrough.PassthroughParserType, passthrough.PassthroughParserPluginFactory)
506+
fwkplugin.Register(vertexai.VertexAIParserType, vertexai.VertexAIParserPluginFactory)
505507
// register saturation detector plugins
506508
fwkplugin.Register(concurrency.ConcurrencyDetectorType, concurrency.ConcurrencyDetectorFactory)
507509
fwkplugin.Register(utilization.UtilizationDetectorType, utilization.UtilizationDetectorFactory)

go.mod

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go 1.25.7
66
toolchain go1.25.8
77

88
require (
9+
cloud.google.com/go/aiplatform v1.124.0
910
github.com/cespare/xxhash/v2 v2.3.0
1011
github.com/envoyproxy/go-control-plane/envoy v1.37.0
1112
github.com/fsnotify/fsnotify v1.9.0
@@ -37,6 +38,7 @@ require (
3738
go.uber.org/zap v1.27.1
3839
golang.org/x/sync v0.20.0
3940
golang.org/x/time v0.15.0
41+
google.golang.org/genproto/googleapis/api v0.0.0-20260401024825-9d38bb4040a9
4042
google.golang.org/grpc v1.80.0
4143
google.golang.org/protobuf v1.36.11
4244
k8s.io/api v0.35.4
@@ -54,6 +56,7 @@ require (
5456

5557
require (
5658
cel.dev/expr v0.25.1 // indirect
59+
cloud.google.com/go/longrunning v0.9.0 // indirect
5760
github.com/Masterminds/semver/v3 v3.4.0 // indirect
5861
github.com/antlr4-go/antlr/v4 v4.13.0 // indirect
5962
github.com/beorn7/perks v1.0.1 // indirect
@@ -127,13 +130,13 @@ require (
127130
golang.org/x/exp v0.0.0-20260112195511-716be5621a96 // indirect
128131
golang.org/x/mod v0.33.0 // indirect
129132
golang.org/x/net v0.52.0 // indirect
130-
golang.org/x/oauth2 v0.35.0 // indirect
133+
golang.org/x/oauth2 v0.36.0 // indirect
131134
golang.org/x/sys v0.42.0 // indirect
132135
golang.org/x/term v0.41.0 // indirect
133136
golang.org/x/text v0.35.0 // indirect
134137
golang.org/x/tools v0.42.0 // indirect
135138
gomodules.xyz/jsonpatch/v2 v2.4.0 // indirect
136-
google.golang.org/genproto/googleapis/api v0.0.0-20260401024825-9d38bb4040a9 // indirect
139+
google.golang.org/genproto v0.0.0-20260319201613-d00831a3d3e7 // indirect
137140
google.golang.org/genproto/googleapis/rpc v0.0.0-20260401024825-9d38bb4040a9 // indirect
138141
gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect
139142
gopkg.in/inf.v0 v0.9.1 // indirect

go.sum

Lines changed: 18 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,17 @@
11
cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4=
22
cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4=
3-
cloud.google.com/go/auth v0.18.1 h1:IwTEx92GFUo2pJ6Qea0EU3zYvKnTAeRCODxfA/G5UWs=
4-
cloud.google.com/go/auth v0.18.1/go.mod h1:GfTYoS9G3CWpRA3Va9doKN9mjPGRS+v41jmZAhBzbrA=
3+
cloud.google.com/go v0.123.0 h1:2NAUJwPR47q+E35uaJeYoNhuNEM9kM8SjgRgdeOJUSE=
4+
cloud.google.com/go/aiplatform v1.124.0 h1:77uy+1+G11yP98axztHeVaW85NcXVLA7WP6ynUXmOlE=
5+
cloud.google.com/go/aiplatform v1.124.0/go.mod h1:yWTZiCunYDnyxeWWD14tDo6+BMlvAUCC5VxuxhvbrVI=
6+
cloud.google.com/go/auth v0.18.2 h1:+Nbt5Ev0xEqxlNjd6c+yYUeosQ5TtEUaNcN/3FozlaM=
7+
cloud.google.com/go/auth v0.18.2/go.mod h1:xD+oY7gcahcu7G2SG2DsBerfFxgPAJz17zz2joOFF3M=
58
cloud.google.com/go/auth/oauth2adapt v0.2.8 h1:keo8NaayQZ6wimpNSmW5OPc283g65QNIiLpZnkHRbnc=
69
cloud.google.com/go/auth/oauth2adapt v0.2.8/go.mod h1:XQ9y31RkqZCcwJWNSx2Xvric3RrU88hAYYbjDWYDL+c=
10+
cloud.google.com/go/compute v1.54.0 h1:4CKmnpO+40z44bKG5bdcKxQ7ocNpRtOc9SCLLUzze1w=
711
cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdBtwLoEkH9Zs=
812
cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10=
13+
cloud.google.com/go/longrunning v0.9.0 h1:0EzbDEGsAvOZNbqXopgniY0w0a1phvu5IdUFq8grmqY=
14+
cloud.google.com/go/longrunning v0.9.0/go.mod h1:pkTz846W7bF4o2SzdWJ40Hu0Re+UoNT6Q5t+igIcb8E=
915
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0 h1:fou+2+WFTib47nS+nz/ozhEBnvU96bKHy6LjRsY4E28=
1016
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0/go.mod h1:t76Ruy8AHvUAC8GfMWJMa0ElSbuIcO03NLpynfbgsPA=
1117
github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.13.1 h1:Hk5QBxZQC1jb2Fwj6mpzme37xbCDdNTxU7O9eb5+LB4=
@@ -179,10 +185,10 @@ github.com/google/s2a-go v0.1.9 h1:LGD7gtMgezd8a/Xak7mEWL0PjoTQFvpRudN895yqKW0=
179185
github.com/google/s2a-go v0.1.9/go.mod h1:YA0Ei2ZQL3acow2O62kdp9UlnvMmU7kA6Eutn0dXayM=
180186
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
181187
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
182-
github.com/googleapis/enterprise-certificate-proxy v0.3.11 h1:vAe81Msw+8tKUxi2Dqh/NZMz7475yUvmRIkXr4oN2ao=
183-
github.com/googleapis/enterprise-certificate-proxy v0.3.11/go.mod h1:RFV7MUdlb7AgEq2v7FmMCfeSMCllAzWxFgRdusoGks8=
184-
github.com/googleapis/gax-go/v2 v2.16.0 h1:iHbQmKLLZrexmb0OSsNGTeSTS0HO4YvFOG8g5E4Zd0Y=
185-
github.com/googleapis/gax-go/v2 v2.16.0/go.mod h1:o1vfQjjNZn4+dPnRdl/4ZD7S9414Y4xA+a/6Icj6l14=
188+
github.com/googleapis/enterprise-certificate-proxy v0.3.14 h1:yh8ncqsbUY4shRD5dA6RlzjJaT4hi3kII+zYw8wmLb8=
189+
github.com/googleapis/enterprise-certificate-proxy v0.3.14/go.mod h1:vqVt9yG9480NtzREnTlmGSBmFrA+bzb0yl0TxoBQXOg=
190+
github.com/googleapis/gax-go/v2 v2.21.0 h1:h45NjjzEO3faG9Lg/cFrBh2PgegVVgzqKzuZl/wMbiI=
191+
github.com/googleapis/gax-go/v2 v2.21.0/go.mod h1:But/NJU6TnZsrLai/xBAQLLz+Hc7fHZJt/hsCz3Fih4=
186192
github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 h1:JeSE6pjso5THxAzdVpqr6/geYxZytqFMBCOtn/ujyeo=
187193
github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674/go.mod h1:r4w70xmWCQKmi1ONH4KIaBptdivuRPyosB9RmPlGEwA=
188194
github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 h1:cLN4IBkmkYZNnk7EAJ0BHIethd+J6LqxFNw5mSiI2bM=
@@ -353,8 +359,8 @@ golang.org/x/mod v0.33.0 h1:tHFzIWbBifEmbwtGz65eaWyGiGZatSrT9prnU8DbVL8=
353359
golang.org/x/mod v0.33.0/go.mod h1:swjeQEj+6r7fODbD2cqrnje9PnziFuw4bmLbBZFrQ5w=
354360
golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0=
355361
golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw=
356-
golang.org/x/oauth2 v0.35.0 h1:Mv2mzuHuZuY2+bkyWXIHMfhNdJAdwW3FuWeCPYN5GVQ=
357-
golang.org/x/oauth2 v0.35.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA=
362+
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
363+
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
358364
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
359365
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
360366
golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo=
@@ -371,8 +377,10 @@ gomodules.xyz/jsonpatch/v2 v2.4.0 h1:Ci3iUJyx9UeRx7CeFN8ARgGbkESwJK+KB9lLcWxY/Zw
371377
gomodules.xyz/jsonpatch/v2 v2.4.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY=
372378
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
373379
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
374-
google.golang.org/api v0.265.0 h1:FZvfUdI8nfmuNrE34aOWFPmLC+qRBEiNm3JdivTvAAU=
375-
google.golang.org/api v0.265.0/go.mod h1:uAvfEl3SLUj/7n6k+lJutcswVojHPp2Sp08jWCu8hLY=
380+
google.golang.org/api v0.274.0 h1:aYhycS5QQCwxHLwfEHRRLf9yNsfvp1JadKKWBE54RFA=
381+
google.golang.org/api v0.274.0/go.mod h1:JbAt7mF+XVmWu6xNP8/+CTiGH30ofmCmk9nM8d8fHew=
382+
google.golang.org/genproto v0.0.0-20260319201613-d00831a3d3e7 h1:XzmzkmB14QhVhgnawEVsOn6OFsnpyxNPRY9QV01dNB0=
383+
google.golang.org/genproto v0.0.0-20260319201613-d00831a3d3e7/go.mod h1:L43LFes82YgSonw6iTXTxXUX1OlULt4AQtkik4ULL/I=
376384
google.golang.org/genproto/googleapis/api v0.0.0-20260401024825-9d38bb4040a9 h1:VPWxll4HlMw1Vs/qXtN7BvhZqsS9cdAittCNvVENElA=
377385
google.golang.org/genproto/googleapis/api v0.0.0-20260401024825-9d38bb4040a9/go.mod h1:7QBABkRtR8z+TEnmXTqIqwJLlzrZKVfAUm7tY3yGv0M=
378386
google.golang.org/genproto/googleapis/rpc v0.0.0-20260401024825-9d38bb4040a9 h1:m8qni9SQFH0tJc1X0vmnpw/0t+AImlSvp30sEupozUg=

pkg/epp/framework/plugins/requesthandling/parsers/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ This directory contains parser plugins used to parse and understand the payloads
66

77
* **`openai-parser`**: The default parser, supporting the [OpenAI API](https://developers.openai.com/api/reference/overview). This is used when no parser is explicitly specified in the `EndpointPickerConfig`.
88
* **`vllmgrpc-parser`**: A parser designed to handle requests specifically for the [vLLM gRPC API](https://docs.vllm.ai/en/latest/api/vllm/entrypoints/grpc_server/).
9+
* **`vertexai-parser`**: A parser designed to handle requests for the Vertex AI gRPC API, specifically supporting [PredictionService/ChatCompletions](https://github.com/googleapis/googleapis/blob/89c3153888201c9e80bc5ec78d6ffca0debe6b52/google/cloud/aiplatform/v1beta1/prediction_service.proto#L235). For unsupported Vertex AI APIs, it skips parsing and lets the request pass through without interpretation resulting in routing to a random endpoint.
910
* **`passthrough-parser`**: A model-agnostic parser that supports any request format by passing the request body through without interpretation.
1011
* **Drawback**: EPP cannot parse the payload, so payload-related scheduling scorers (e.g., `prefix-cache-scorer`) are not supported.
1112

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
/*
2+
Copyright 2026 The Kubernetes Authors.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package parsers
18+
19+
const (
20+
// MethodPathKey is the header key for the request path.
21+
MethodPathKey = ":path"
22+
)

pkg/epp/framework/plugins/requesthandling/parsers/openai/openai.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import (
2727

2828
fwkplugin "github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/interface/plugin"
2929
fwkrh "github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/interface/requesthandling"
30+
"github.com/llm-d/llm-d-inference-scheduler/pkg/epp/framework/plugins/requesthandling/parsers"
3031
)
3132

3233
const (
@@ -145,7 +146,7 @@ func (p *OpenAIParser) parseStreamResponse(chunk []byte) (*fwkrh.ParsedResponse,
145146
// getRequestPath extracts the request path from headers with fallback priority
146147
func getRequestPath(headers map[string]string) string {
147148
// Try primary path header
148-
if path := headers[":path"]; path != "" {
149+
if path := headers[parsers.MethodPathKey]; path != "" {
149150
return path
150151
}
151152

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
/*
2+
Copyright 2026 The Kubernetes Authors.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package grpcutil
18+
19+
import (
20+
"encoding/binary"
21+
"errors"
22+
"fmt"
23+
)
24+
25+
const (
26+
gRPCPayloadHeaderLen = 5
27+
)
28+
29+
// ParseGrpcPayload extracts the message payload from a gRPC frame.
30+
// A standard gRPC frame consists of a 1-byte compression flag, a 4-byte message length,
31+
// and the actual message payload.
32+
// It returns an error if the payload is compressed.
33+
func ParseGrpcPayload(data []byte) ([]byte, error) {
34+
if len(data) < gRPCPayloadHeaderLen {
35+
return nil, fmt.Errorf("invalid gRPC frame: expected at least %d bytes for header, got %d", gRPCPayloadHeaderLen, len(data))
36+
}
37+
38+
isCompressed := data[0] == 1
39+
if isCompressed {
40+
// TODO(#895): handle compressed payload.
41+
return nil, errors.New("compressed gRPC payload is not supported")
42+
}
43+
msgLen := binary.BigEndian.Uint32(data[1:5])
44+
45+
if uint32(len(data)) < gRPCPayloadHeaderLen+msgLen {
46+
return nil, fmt.Errorf("incomplete gRPC payload: header indicates %d bytes, but only %d bytes are available", msgLen, uint32(len(data))-gRPCPayloadHeaderLen)
47+
}
48+
return data[gRPCPayloadHeaderLen : gRPCPayloadHeaderLen+msgLen], nil
49+
}
Lines changed: 95 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,95 @@
1+
/*
2+
Copyright 2026 The Kubernetes Authors.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package grpcutil
18+
19+
import (
20+
"encoding/binary"
21+
"testing"
22+
23+
"github.com/google/go-cmp/cmp"
24+
)
25+
26+
func TestParseGrpcPayload(t *testing.T) {
27+
tests := []struct {
28+
name string
29+
data []byte
30+
want []byte
31+
wantErr string
32+
}{
33+
{
34+
name: "too short header",
35+
data: []byte{0, 0, 0, 0},
36+
want: nil,
37+
wantErr: "invalid gRPC frame: expected at least 5 bytes for header, got 4",
38+
},
39+
{
40+
name: "compressed payload not supported",
41+
data: []byte{1, 0, 0, 0, 0},
42+
want: nil,
43+
wantErr: "compressed gRPC payload is not supported",
44+
},
45+
{
46+
name: "incomplete payload",
47+
data: func() []byte {
48+
b := make([]byte, 10)
49+
b[0] = 0 // not compressed
50+
binary.BigEndian.PutUint32(b[1:5], 10) // indicates 10 bytes
51+
copy(b[5:], []byte("12345"))
52+
return b
53+
}(),
54+
want: nil,
55+
wantErr: "incomplete gRPC payload: header indicates 10 bytes, but only 5 bytes are available",
56+
},
57+
{
58+
name: "success exact size",
59+
data: func() []byte {
60+
b := make([]byte, 10)
61+
b[0] = 0 // not compressed
62+
binary.BigEndian.PutUint32(b[1:5], 5)
63+
copy(b[5:], []byte("hello"))
64+
return b
65+
}(),
66+
want: []byte("hello"),
67+
wantErr: "",
68+
},
69+
}
70+
71+
for _, tt := range tests {
72+
t.Run(tt.name, func(t *testing.T) {
73+
got, err := ParseGrpcPayload(tt.data)
74+
75+
if tt.wantErr != "" {
76+
if err == nil {
77+
t.Fatalf("got error nil, want %q", tt.wantErr)
78+
}
79+
if gotErr := err.Error(); gotErr != tt.wantErr {
80+
t.Errorf("got error %q, want %q", gotErr, tt.wantErr)
81+
}
82+
return
83+
}
84+
85+
if err != nil {
86+
t.Fatalf("unexpected error: %v", err)
87+
}
88+
89+
want := tt.want
90+
if diff := cmp.Diff(want, got); diff != "" {
91+
t.Errorf("ParseGrpcPayload() mismatch (-want +got):\n%s", diff)
92+
}
93+
})
94+
}
95+
}

0 commit comments

Comments
 (0)