Skip to content

Commit 234e3c7

Browse files
authored
* feat: usage limit policy plugin Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * formatting Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * remove leftover runner_test.go Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * missing boilerplate header Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * replace Noop policy with StaticUsageLimitPolicy Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * hoist check to the top of dispatchCycle() Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * drop UsageLimitConfig wrapper Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * formatting Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * linting Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> * apply suggestions Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com> --------- Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com>
1 parent 13ad8f1 commit 234e3c7

14 files changed

Lines changed: 334 additions & 21 deletions

File tree

apix/config/v1alpha1/endpointpickerconfig_types.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -314,6 +314,12 @@ type FlowControlConfig struct {
314314
// priority levels. Traffic matching these priorities will be handled according to these rules.
315315
// If a priority band is not specified, it uses specific defaults.
316316
PriorityBands []PriorityBandConfig `json:"priorityBands,omitempty"`
317+
318+
// +optional
319+
// UsageLimitPolicyPluginRef specifies the UsageLimitPolicy plugin to use for adaptive capacity management.
320+
// Must reference a named plugin instance defined in the top-level Plugins section.
321+
// If omitted, a default static policy (threshold=1.0, no gating) is used.
322+
UsageLimitPolicyPluginRef string `json:"usageLimitPolicyPluginRef,omitempty"`
317323
}
318324

319325
func (fcc *FlowControlConfig) String() string {
@@ -340,6 +346,10 @@ func (fcc *FlowControlConfig) String() string {
340346
parts = append(parts, fmt.Sprintf("PriorityBands: %v", fcc.PriorityBands))
341347
}
342348

349+
if fcc.UsageLimitPolicyPluginRef != "" {
350+
parts = append(parts, "UsageLimitPolicyRef: "+fcc.UsageLimitPolicyPluginRef)
351+
}
352+
343353
return "{" + strings.Join(parts, ", ") + "}"
344354
}
345355

cmd/epp/runner/runner.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@ import (
6565
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/ordering"
6666
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/saturationdetector/concurrency"
6767
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/saturationdetector/utilization"
68+
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/usagelimits"
6869
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/requestcontrol/requestattributereporter"
6970
testresponsereceived "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/requestcontrol/test/responsereceived"
7071
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/requesthandling/parsers/openai"
@@ -342,7 +343,7 @@ func (r *Runner) setup(ctx context.Context, cfg *rest.Config, opts *runserver.Op
342343
opts.PoolName,
343344
eppConfig.FlowControlConfig.Controller,
344345
registry, eppConfig.SaturationDetector,
345-
locator,
346+
locator, eppConfig.FlowControlConfig.UsageLimitPolicy,
346347
)
347348
if err != nil {
348349
return nil, nil, fmt.Errorf("failed to initialize Flow Controller: %w", err)
@@ -458,6 +459,7 @@ func (r *Runner) registerInTreePlugins() {
458459
fwkplugin.Register(ordering.FCFSOrderingPolicyType, ordering.FCFSOrderingPolicyFactory)
459460
fwkplugin.Register(ordering.EDFOrderingPolicyType, ordering.EDFOrderingPolicyFactory)
460461
fwkplugin.Register(ordering.SLODeadlineOrderingPolicyType, ordering.SLODeadlineOrderingPolicyFactory)
462+
fwkplugin.Register(usagelimits.StaticUsageLimitPolicyType, usagelimits.StaticPolicyFactory)
461463
// Latency predictor plugins
462464
fwkplugin.Register(predictedlatency.PredictedLatencyPluginType, predictedlatency.PredictedLatencyFactory)
463465
// register filter for test purpose only (used in conformance tests)

pkg/epp/config/loader/configloader_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ import (
4242
framework "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/scheduling"
4343
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/fairness"
4444
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/ordering"
45+
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/usagelimits"
4546
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/requesthandling/parsers/openai"
4647
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/scheduling/picker"
4748
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/scheduling/profile"
@@ -712,6 +713,7 @@ func registerTestPlugins(t *testing.T) {
712713
fwkplugin.Register(picker.MaxScorePickerType, picker.MaxScorePickerFactory)
713714
fwkplugin.Register(profile.SingleProfileHandlerType, profile.SingleProfileHandlerFactory)
714715
fwkplugin.Register(openai.OpenAIParserType, openai.OpenAIParserPluginFactory)
716+
fwkplugin.Register(usagelimits.StaticUsageLimitPolicyType, usagelimits.StaticPolicyFactory)
715717
}
716718

717719
func TestValidateSaturationDetector(t *testing.T) {

pkg/epp/config/loader/defaults.go

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -205,7 +205,14 @@ func ensureFlowControlLayer(
205205
}
206206
}
207207
if _, ok := allPlugins[registry.DefaultFairnessPolicyRef]; !ok {
208-
return registerDefaultPlugin(cfg, handle, registry.DefaultFairnessPolicyRef)
208+
if err := registerDefaultPlugin(cfg, handle, registry.DefaultFairnessPolicyRef); err != nil {
209+
return err
210+
}
211+
}
212+
if _, ok := allPlugins[registry.DefaultUsageLimitPolicyRef]; !ok {
213+
if err := registerDefaultPlugin(cfg, handle, registry.DefaultUsageLimitPolicyRef); err != nil {
214+
return err
215+
}
209216
}
210217
return nil
211218
}

pkg/epp/flowcontrol/benchmark/benchmark.go

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

6868
"github.com/go-logr/logr"
6969
"sigs.k8s.io/controller-runtime/pkg/log"
70+
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/usagelimits"
7071

7172
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/contracts"
7273
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/contracts/mocks"
@@ -275,7 +276,7 @@ func setupBenchmarkHarness(
275276
}
276277
}
277278

278-
fc, err := controller.NewFlowController(ctx, "benchmark", cfg, reg, detector, &mocks.MockPodLocator{})
279+
fc, err := controller.NewFlowController(ctx, "benchmark", cfg, reg, detector, &mocks.MockPodLocator{}, usagelimits.DefaultPolicy())
279280
if err != nil {
280281
b.Fatalf("Failed to init FlowController: %v", err)
281282
}

pkg/epp/flowcontrol/config.go

Lines changed: 30 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222
configapi "sigs.k8s.io/gateway-api-inference-extension/apix/config/v1alpha1"
2323
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/controller"
2424
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/registry"
25+
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/flowcontrol"
2526
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/plugin"
2627
)
2728

@@ -31,8 +32,9 @@ const FeatureGate = "flowControl"
3132
// It embeds the configurations for the controller and the registry, providing a single point of entry for validation
3233
// and initialization.
3334
type Config struct {
34-
Controller *controller.Config
35-
Registry *registry.Config
35+
Controller *controller.Config
36+
Registry *registry.Config
37+
UsageLimitPolicy flowcontrol.UsageLimitPolicy
3638
}
3739

3840
// NewConfigFromAPI creates a new Config by translating the top-level API configuration.
@@ -45,8 +47,30 @@ func NewConfigFromAPI(apiConfig *configapi.FlowControlConfig, handle plugin.Hand
4547
if err != nil {
4648
return nil, fmt.Errorf("failed to create controller config: %w", err)
4749
}
48-
return &Config{
49-
Controller: ctrlCfg,
50-
Registry: registryConfig,
51-
}, nil
50+
usageLimitPolicy, err := ensureUsageLimitPolicy(apiConfig, handle)
51+
if err != nil {
52+
return nil, err
53+
}
54+
cfg := &Config{
55+
Controller: ctrlCfg,
56+
Registry: registryConfig,
57+
UsageLimitPolicy: usageLimitPolicy,
58+
}
59+
return cfg, nil
60+
}
61+
62+
func ensureUsageLimitPolicy(apiConfig *configapi.FlowControlConfig, handle plugin.Handle) (flowcontrol.UsageLimitPolicy, error) {
63+
ref := registry.DefaultUsageLimitPolicyRef
64+
if apiConfig != nil && apiConfig.UsageLimitPolicyPluginRef != "" {
65+
ref = apiConfig.UsageLimitPolicyPluginRef
66+
}
67+
p := handle.Plugin(ref)
68+
if p == nil {
69+
return nil, fmt.Errorf("usage limit policy plugin '%s' not found", ref)
70+
}
71+
usageLimitPolicy, ok := p.(flowcontrol.UsageLimitPolicy)
72+
if !ok {
73+
return nil, fmt.Errorf("plugin '%s' does not implement UsageLimitPolicy", ref)
74+
}
75+
return usageLimitPolicy, nil
5276
}

pkg/epp/flowcontrol/config_test.go

Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,10 +26,12 @@ import (
2626
"k8s.io/utils/ptr"
2727

2828
configapi "sigs.k8s.io/gateway-api-inference-extension/apix/config/v1alpha1"
29+
flowcontrolif "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/flowcontrol"
2930
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/flowcontrol/mocks"
3031
fwkplugin "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/plugin"
3132
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/fairness"
3233
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/ordering"
34+
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/flowcontrol/usagelimits"
3335
"sigs.k8s.io/gateway-api-inference-extension/test/utils"
3436
)
3537

@@ -49,6 +51,23 @@ func TestNewConfigFromAPI(t *testing.T) {
4951
Type: ordering.FCFSOrderingPolicyType,
5052
},
5153
})
54+
handle.AddPlugin(usagelimits.StaticUsageLimitPolicyType, usagelimits.DefaultPolicy())
55+
56+
// A func-based custom policy that always returns 0.8 — demonstrates that users can define
57+
// their policy in the standard plugins section and reference it via UsageLimit.PluginRef.
58+
const funcPolicyName = "func-policy"
59+
handle.AddPlugin(funcPolicyName, usagelimits.NewPolicyFunc(funcPolicyName, func(_ context.Context, _ float64, priorities []int) []float64 {
60+
result := make([]float64, len(priorities))
61+
for i := range result {
62+
result[i] = 0.8
63+
}
64+
return result
65+
}))
66+
67+
// A hand-rolled struct implementing UsageLimitPolicy — demonstrates that any struct
68+
// satisfying the interface can be registered and resolved, without relying on usagelimits helpers.
69+
const structPolicyName = "struct-policy"
70+
handle.AddPlugin(structPolicyName, &constantPointEightPolicy{})
5271

5372
testCases := []struct {
5473
name string
@@ -79,6 +98,70 @@ func TestNewConfigFromAPI(t *testing.T) {
7998
"MaxBytes should be correctly translated from resource.Quantity in API to uint64 in internal config")
8099
},
81100
},
101+
{
102+
name: "Success - Default UsageLimitPolicy when UsageLimit is nil",
103+
apiConfig: nil,
104+
assertion: func(t *testing.T, cfg *Config) {
105+
require.NotNil(t, cfg.UsageLimitPolicy, "UsageLimitPolicy should be resolved even when not explicitly configured")
106+
ceilings := cfg.UsageLimitPolicy.ComputeLimit(context.Background(), 0.5, []int{0})
107+
assert.Equal(t, []float64{1.0}, ceilings, "Default noop policy should return 1.0 (no gating)")
108+
},
109+
},
110+
{
111+
name: "Success - UsageLimitPolicyPluginRef is resolved",
112+
apiConfig: &configapi.FlowControlConfig{
113+
UsageLimitPolicyPluginRef: usagelimits.StaticUsageLimitPolicyType,
114+
},
115+
assertion: func(t *testing.T, cfg *Config) {
116+
require.NotNil(t, cfg.UsageLimitPolicy, "UsageLimitPolicy should be resolved from the handle")
117+
ceilings := cfg.UsageLimitPolicy.ComputeLimit(context.Background(), 0.5, []int{0})
118+
assert.Equal(t, []float64{1.0}, ceilings, "Noop policy should return 1.0 (no gating)")
119+
},
120+
},
121+
{
122+
name: "Success - Func-based UsageLimitPolicy resolved via PluginRef",
123+
apiConfig: &configapi.FlowControlConfig{
124+
UsageLimitPolicyPluginRef: funcPolicyName,
125+
},
126+
assertion: func(t *testing.T, cfg *Config) {
127+
require.NotNil(t, cfg.UsageLimitPolicy)
128+
ctx := context.Background()
129+
for _, tc := range []struct {
130+
name string
131+
priority int
132+
saturation float64
133+
}{
134+
{"zero saturation", 0, 0.0},
135+
{"half saturation", 1, 0.5},
136+
{"full saturation", 5, 1.0},
137+
} {
138+
assert.Equal(t, []float64{0.8}, cfg.UsageLimitPolicy.ComputeLimit(ctx, tc.saturation, []int{tc.priority}),
139+
"func-based policy should return 0.8 at %s", tc.name)
140+
}
141+
},
142+
},
143+
{
144+
name: "Success - Struct-based UsageLimitPolicy resolved via PluginRef",
145+
apiConfig: &configapi.FlowControlConfig{
146+
UsageLimitPolicyPluginRef: structPolicyName,
147+
},
148+
assertion: func(t *testing.T, cfg *Config) {
149+
require.NotNil(t, cfg.UsageLimitPolicy)
150+
ctx := context.Background()
151+
for _, tc := range []struct {
152+
name string
153+
priority int
154+
saturation float64
155+
}{
156+
{"zero saturation", 0, 0.0},
157+
{"half saturation", 1, 0.5},
158+
{"full saturation", 5, 1.0},
159+
} {
160+
assert.Equal(t, []float64{0.8}, cfg.UsageLimitPolicy.ComputeLimit(ctx, tc.saturation, []int{tc.priority}),
161+
"struct-based policy should return 0.8 at %s", tc.name)
162+
}
163+
},
164+
},
82165
}
83166

84167
for _, tc := range testCases {
@@ -96,3 +179,26 @@ func TestNewConfigFromAPI(t *testing.T) {
96179
})
97180
}
98181
}
182+
183+
// constantPointEightPolicy is a hand-rolled UsageLimitPolicy implementation that always returns 0.8.
184+
// It exists to show that any struct satisfying the interface can be registered and resolved,
185+
// without relying on the usagelimits.NewPolicyFunc helper.
186+
type constantPointEightPolicy struct{}
187+
188+
func (p *constantPointEightPolicy) TypedName() fwkplugin.TypedName {
189+
return fwkplugin.TypedName{
190+
Type: "constant-point-eight-policy-type",
191+
Name: "constant-point-eight-policy",
192+
}
193+
}
194+
195+
func (p *constantPointEightPolicy) ComputeLimit(_ context.Context, _ float64, priorities []int) []float64 {
196+
result := make([]float64, len(priorities))
197+
for i := range result {
198+
result[i] = 0.8
199+
}
200+
return result
201+
}
202+
203+
// compile-time check that constantPointEightPolicy satisfies the interface.
204+
var _ flowcontrolif.UsageLimitPolicy = (*constantPointEightPolicy)(nil)

pkg/epp/flowcontrol/controller/controller.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ type shardProcessorFactory func(
6464
shard contracts.RegistryShard,
6565
saturationDetector flowcontrol.SaturationDetector,
6666
podLocator contracts.PodLocator,
67+
usageLimitPolicy flowcontrol.UsageLimitPolicy,
6768
clock clock.WithTicker,
6869
cleanupSweepInterval time.Duration,
6970
enqueueChannelBufferSize int,
@@ -100,6 +101,7 @@ type FlowController struct {
100101
registry registryClient
101102
saturationDetector flowcontrol.SaturationDetector
102103
podLocator contracts.PodLocator
104+
usageLimitPolicy flowcontrol.UsageLimitPolicy
103105
clock clock.WithTicker
104106
logger logr.Logger
105107
shardProcessorFactory shardProcessorFactory
@@ -133,13 +135,15 @@ func NewFlowController(
133135
registry contracts.FlowRegistry,
134136
sd flowcontrol.SaturationDetector,
135137
podLocator contracts.PodLocator,
138+
usageLimitPolicy flowcontrol.UsageLimitPolicy,
136139
opts ...flowControllerOption,
137140
) (*FlowController, error) {
138141
fc := &FlowController{
139142
config: config,
140143
registry: registry,
141144
saturationDetector: sd,
142145
podLocator: podLocator,
146+
usageLimitPolicy: usageLimitPolicy,
143147
clock: clock.RealClock{},
144148
logger: log.FromContext(ctx).WithName("flow-controller"),
145149
parentCtx: ctx,
@@ -150,6 +154,7 @@ func NewFlowController(
150154
shard contracts.RegistryShard,
151155
saturationDetector flowcontrol.SaturationDetector,
152156
podLocator contracts.PodLocator,
157+
usageLimitPolicy flowcontrol.UsageLimitPolicy,
153158
clock clock.WithTicker,
154159
cleanupSweepInterval time.Duration,
155160
enqueueChannelBufferSize int,
@@ -161,6 +166,7 @@ func NewFlowController(
161166
shard,
162167
saturationDetector,
163168
podLocator,
169+
usageLimitPolicy,
164170
clock,
165171
cleanupSweepInterval,
166172
enqueueChannelBufferSize,
@@ -487,6 +493,7 @@ func (fc *FlowController) getOrStartWorker(shard contracts.RegistryShard) *manag
487493
shard,
488494
fc.saturationDetector,
489495
fc.podLocator,
496+
fc.usageLimitPolicy,
490497
fc.clock,
491498
fc.config.ExpiryCleanupInterval,
492499
fc.config.EnqueueChannelBufferSize,

0 commit comments

Comments
 (0)