Skip to content

Commit 1d91727

Browse files
committed
apply suggestions
Signed-off-by: Edoardo Vacchi <evacchi@users.noreply.github.com>
1 parent b0a8dbf commit 1d91727

10 files changed

Lines changed: 174 additions & 201 deletions

File tree

apix/config/v1alpha1/endpointpickerconfig_types.go

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -302,15 +302,15 @@ type FlowControlConfig struct {
302302
PriorityBands []PriorityBandConfig `json:"priorityBands,omitempty"`
303303

304304
// +optional
305-
// UsageLimitPolicyType is the plugin type name for the UsageLimitPolicy.
306-
// The named plugin must be registered in the plugin registry before the runner starts.
307-
// If omitted or the plugin is not found, a default no-op policy (limit=1.0) is used.
308-
UsageLimitPolicyType string `json:"usageLimitPolicyType,omitempty"`
305+
// UsageLimit specifies the UsageLimitPolicy plugin to use for adaptive capacity management.
306+
// The PluginRef must reference a named plugin instance defined in the top-level Plugins section.
307+
// If omitted, a default no-op policy (limit=1.0) is used.
308+
UsageLimit *UsageLimitConfig `json:"usageLimit,omitempty"`
309309
}
310310

311311
func (fcc *FlowControlConfig) String() string {
312-
return fmt.Sprintf("{MaxBytes: %v, DefaultPriorityBand: %v, PriorityBands: %v, UsageLimitPolicyType: %v}",
313-
fcc.MaxBytes, fcc.DefaultPriorityBand, fcc.PriorityBands, fcc.UsageLimitPolicyType)
312+
return fmt.Sprintf("{MaxBytes: %v, DefaultPriorityBand: %v, PriorityBands: %v, UsageLimit: %v}",
313+
fcc.MaxBytes, fcc.DefaultPriorityBand, fcc.PriorityBands, fcc.UsageLimit)
314314
}
315315

316316
// PriorityBandConfig configures a single priority band.
@@ -340,3 +340,17 @@ func (pbc PriorityBandConfig) String() string {
340340
return fmt.Sprintf("{Priority: %d, MaxBytes: %v, FairnessPolicyRef: %s, OrderingPolicyRef: %s}",
341341
pbc.Priority, pbc.MaxBytes, pbc.FairnessPolicyRef, pbc.OrderingPolicyRef)
342342
}
343+
344+
// UsageLimitConfig contains the configuration for a UsageLimitPolicy plugin.
345+
type UsageLimitConfig struct {
346+
// +required
347+
// +kubebuilder:validation:Required
348+
// PluginRef specifies a particular Plugin instance to be used as the UsageLimitPolicy.
349+
// The reference is to the name of an entry of the Plugins defined in the configuration's
350+
// Plugins section.
351+
PluginRef string `json:"pluginRef"`
352+
}
353+
354+
func (ulc UsageLimitConfig) String() string {
355+
return fmt.Sprintf("{PluginRef: %s}", ulc.PluginRef)
356+
}

cmd/epp/runner/runner.go

Lines changed: 2 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,6 @@ import (
5959
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/framework/plugins/usagelimits"
6060
fcregistry "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/registry"
6161
fwkdl "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/datalayer"
62-
flowcontrolplugins "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/flowcontrol"
6362
fwkplugin "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/plugin"
6463
fwkrh "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/requesthandling"
6564
extractormetrics "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/plugins/datalayer/extractor/metrics"
@@ -329,8 +328,6 @@ func (r *Runner) setup(ctx context.Context, cfg *rest.Config, opts *runserver.Op
329328

330329
saturationDetector := utilizationdetector.NewDetector(eppConfig.SaturationDetectorConfig, setupLog)
331330

332-
usageLimitPolicy := r.setupFlowControlPlugins(eppConfig.FlowControlConfig.UsageLimitPolicyType, setupLog)
333-
334331
// --- Admission Control Initialization ---
335332
var admissionController requestcontrol.AdmissionController
336333
var locator contracts.PodLocator
@@ -347,7 +344,7 @@ func (r *Runner) setup(ctx context.Context, cfg *rest.Config, opts *runserver.Op
347344
opts.PoolName,
348345
eppConfig.FlowControlConfig.Controller,
349346
registry, saturationDetector,
350-
locator, usageLimitPolicy,
347+
locator, eppConfig.FlowControlConfig.UsageLimitPolicy,
351348
)
352349
if err != nil {
353350
return nil, nil, fmt.Errorf("failed to initialize Flow Controller: %w", err)
@@ -463,6 +460,7 @@ func (r *Runner) registerInTreePlugins() {
463460
fwkplugin.Register(fairness.RoundRobinFairnessPolicyType, fairness.RoundRobinFairnessPolicyFactory)
464461
fwkplugin.Register(ordering.FCFSOrderingPolicyType, ordering.FCFSOrderingPolicyFactory)
465462
fwkplugin.Register(ordering.EDFOrderingPolicyType, ordering.EDFOrderingPolicyFactory)
463+
fwkplugin.Register(usagelimits.NoopUsageLimitPolicyType, usagelimits.NoopPolicyFactory)
466464
// Latency predictor plugins
467465
fwkplugin.Register(predictedlatency.PredictedLatencyPluginType, predictedlatency.PredictedLatencyFactory)
468466
// register filter for test purpose only (used in conformance tests)
@@ -639,26 +637,6 @@ func (r *Runner) setupDataLayer(enableNewMetrics bool, cfg *datalayer.Config,
639637
return nil
640638
}
641639

642-
func (r *Runner) setupFlowControlPlugins(policyType string, logger logr.Logger) flowcontrolplugins.UsageLimitPolicy {
643-
if factory, ok := fwkplugin.Registry[policyType]; ok {
644-
p, err := factory(policyType, nil, nil)
645-
if err != nil {
646-
logger.Info("Failed to instantiate usage limit policy plugin, using default", "type", policyType, "error", err)
647-
return usagelimits.DefaultPolicy()
648-
}
649-
policy, ok := p.(flowcontrolplugins.UsageLimitPolicy)
650-
if !ok {
651-
logger.Info("Registered plugin does not implement UsageLimitPolicy, using default", "type", policyType)
652-
return usagelimits.DefaultPolicy()
653-
}
654-
return policy
655-
}
656-
if policyType != "" {
657-
logger.Info("No plugin registered for usage limit policy type, using default", "type", policyType)
658-
}
659-
return usagelimits.DefaultPolicy()
660-
}
661-
662640
func (r *Runner) setupMetricsCollection(enableNewMetrics bool, opts *runserver.Options, pmc backendmetrics.PodMetricsClient) datalayer.EndpointFactory {
663641
if enableNewMetrics {
664642
return datalayer.NewEndpointFactory(nil, opts.RefreshMetricsInterval)

cmd/epp/runner/runner_test.go

Lines changed: 0 additions & 151 deletions
Original file line numberDiff line numberDiff line change
@@ -18,109 +18,13 @@ package runner
1818

1919
import (
2020
"context"
21-
"encoding/json"
2221
"testing"
2322

24-
"github.com/go-logr/logr"
2523
"github.com/stretchr/testify/assert"
26-
"github.com/stretchr/testify/require"
2724

28-
fc "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol"
2925
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/flowcontrol/framework/plugins/usagelimits"
30-
"sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/flowcontrol"
31-
fwkplugin "sigs.k8s.io/gateway-api-inference-extension/pkg/epp/framework/interface/plugin"
3226
)
3327

34-
// TestSetupFlowControlPlugins verifies that setupFlowControlPlugins resolves a registered
35-
// UsageLimitPolicy from the plugin registry and that the resolved policy returns the expected limit.
36-
func TestSetupFlowControlPlugins(t *testing.T) {
37-
// Register the factory before any parallel sub-test starts so the write
38-
// to the global registry map does not race with reads in sub-tests.
39-
const pluginType = "func-plugin"
40-
fwkplugin.Register(pluginType, func(name string, _ json.RawMessage, _ fwkplugin.Handle) (fwkplugin.Plugin, error) {
41-
return usagelimits.NewPolicyFunc(name, func(_ context.Context, _ float64, priorities []int) []float64 {
42-
result := make([]float64, len(priorities))
43-
for i := range result {
44-
result[i] = 0.8
45-
}
46-
return result
47-
}), nil
48-
})
49-
t.Cleanup(func() {
50-
delete(fwkplugin.Registry, pluginType)
51-
})
52-
53-
r := &Runner{}
54-
55-
t.Run("resolves registered plugin and returns 0.8", func(t *testing.T) {
56-
t.Parallel()
57-
58-
policy := r.setupFlowControlPlugins(pluginType, logr.Discard())
59-
require.NotNil(t, policy)
60-
61-
ctx := context.Background()
62-
for _, tc := range []struct {
63-
name string
64-
priority int
65-
saturation float64
66-
}{
67-
{"zero saturation", 0, 0.0},
68-
{"half saturation", 1, 0.5},
69-
{"full saturation", 5, 1.0},
70-
} {
71-
t.Run(tc.name, func(t *testing.T) {
72-
t.Parallel()
73-
assert.Equal(t, []float64{0.8}, policy.ComputeLimit(ctx, tc.saturation, []int{tc.priority}))
74-
})
75-
}
76-
})
77-
78-
t.Run("falls back to default policy when type is empty", func(t *testing.T) {
79-
t.Parallel()
80-
81-
policy := r.setupFlowControlPlugins("", logr.Discard())
82-
require.NotNil(t, policy)
83-
84-
ctx := context.Background()
85-
assert.Equal(t, []float64{1.0}, policy.ComputeLimit(ctx, 0.5, []int{0}))
86-
})
87-
88-
t.Run("falls back to default policy when type is unknown", func(t *testing.T) {
89-
t.Parallel()
90-
91-
policy := r.setupFlowControlPlugins("no-such-policy-type", logr.Discard())
92-
require.NotNil(t, policy)
93-
94-
ctx := context.Background()
95-
assert.Equal(t, []float64{1.0}, policy.ComputeLimit(ctx, 0.5, []int{0}))
96-
})
97-
}
98-
99-
const constantPointEightStructPolicyType = "test-constant-point-eight-struct-usage-limit-policy"
100-
101-
// constantPointEightPolicy is a hand-rolled UsageLimitPolicy implementation that always returns 0.8.
102-
// It exists to show that any struct satisfying the interface can be registered and resolved,
103-
// without relying on the usagelimits.NewPolicyFunc helper.
104-
type constantPointEightPolicy struct{}
105-
106-
func (p *constantPointEightPolicy) TypedName() fwkplugin.TypedName {
107-
return fwkplugin.TypedName{
108-
Type: constantPointEightStructPolicyType + "-type",
109-
Name: constantPointEightStructPolicyType,
110-
}
111-
}
112-
113-
func (p *constantPointEightPolicy) ComputeLimit(_ context.Context, _ float64, priorities []int) []float64 {
114-
result := make([]float64, len(priorities))
115-
for i := range result {
116-
result[i] = 0.8
117-
}
118-
return result
119-
}
120-
121-
// compile-time check that constantPointEightPolicy satisfies the interface.
122-
var _ flowcontrol.UsageLimitPolicy = (*constantPointEightPolicy)(nil)
123-
12428
// TestLinearSpacingPolicy demonstrates a stateless UsageLimitPolicy that dynamically spaces
12529
// ceilings based on the active priority domain. The highest-active priority always gets ceiling
12630
// 1.0, and each subsequent tier drops by a fixed step (0.2). When priorities go idle and the
@@ -154,58 +58,3 @@ func TestLinearSpacingPolicy(t *testing.T) {
15458
got = linearSpacing.ComputeLimit(ctx, 0.5, []int{0})
15559
assert.Equal(t, []float64{1.0}, got, "single active priority should get full ceiling")
15660
}
157-
158-
// TestSetupFlowControlPlugins_WithUsageLimitPolicyType verifies that UsageLimitPolicyType on
159-
// flowcontrol.Config (the field a downstream project would set from a config file) is read by
160-
// setupFlowControlPlugins to resolve and return the correct registered plugin.
161-
func TestSetupFlowControlPlugins_WithUsageLimitPolicyType(t *testing.T) {
162-
const pluginType = "struct-plugin-via-config-option"
163-
164-
fwkplugin.Register(pluginType, func(name string, _ json.RawMessage, _ fwkplugin.Handle) (fwkplugin.Plugin, error) {
165-
return &constantPointEightPolicy{}, nil
166-
})
167-
t.Cleanup(func() {
168-
delete(fwkplugin.Registry, pluginType)
169-
})
170-
171-
cfg := &fc.Config{UsageLimitPolicyType: pluginType}
172-
173-
policy := (&Runner{}).setupFlowControlPlugins(cfg.UsageLimitPolicyType, logr.Discard())
174-
require.NotNil(t, policy)
175-
176-
ctx := context.Background()
177-
assert.Equal(t, []float64{0.8}, policy.ComputeLimit(ctx, 0.5, []int{0}))
178-
}
179-
180-
// TestSetupFlowControlPlugins_StructPlugin verifies that a hand-rolled struct implementing
181-
// UsageLimitPolicy (without using usagelimits.NewPolicyFunc) can be registered and resolved.
182-
func TestSetupFlowControlPlugins_StructPlugin(t *testing.T) {
183-
184-
fwkplugin.Register(constantPointEightStructPolicyType, func(name string, _ json.RawMessage, _ fwkplugin.Handle) (fwkplugin.Plugin, error) {
185-
return &constantPointEightPolicy{}, nil
186-
})
187-
t.Cleanup(func() {
188-
delete(fwkplugin.Registry, constantPointEightStructPolicyType)
189-
})
190-
191-
r := &Runner{}
192-
193-
policy := r.setupFlowControlPlugins(constantPointEightStructPolicyType, logr.Discard())
194-
require.NotNil(t, policy)
195-
196-
ctx := context.Background()
197-
for _, tc := range []struct {
198-
name string
199-
priority int
200-
saturation float64
201-
}{
202-
{"zero saturation", 0, 0.0},
203-
{"half saturation", 1, 0.5},
204-
{"full saturation", 5, 1.0},
205-
} {
206-
t.Run(tc.name, func(t *testing.T) {
207-
t.Parallel()
208-
assert.Equal(t, []float64{0.8}, policy.ComputeLimit(ctx, tc.saturation, []int{tc.priority}))
209-
})
210-
}
211-
}

pkg/epp/config/loader/defaults.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -202,6 +202,9 @@ func ensureFlowControlLayer(
202202
if _, ok := allPlugins[registry.DefaultFairnessPolicyRef]; !ok {
203203
return registerDefaultPlugin(cfg, handle, registry.DefaultFairnessPolicyRef)
204204
}
205+
if _, ok := allPlugins[registry.DefaultUsageLimitPolicyRef]; !ok {
206+
return registerDefaultPlugin(cfg, handle, registry.DefaultUsageLimitPolicyRef)
207+
}
205208
return nil
206209
}
207210

pkg/epp/flowcontrol/config.go

Lines changed: 27 additions & 11 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,12 +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
36-
37-
// UsageLimitPolicyType is the plugin type name used to resolve a UsageLimitPolicy from the plugin registry.
38-
// Optional: Defaults to empty string (uses the default no-op policy).
39-
UsageLimitPolicyType string
35+
Controller *controller.Config
36+
Registry *registry.Config
37+
UsageLimitPolicy flowcontrol.UsageLimitPolicy
4038
}
4139

4240
// NewConfigFromAPI creates a new Config by translating the top-level API configuration.
@@ -49,12 +47,30 @@ func NewConfigFromAPI(apiConfig *configapi.FlowControlConfig, handle plugin.Hand
4947
if err != nil {
5048
return nil, fmt.Errorf("failed to create controller config: %w", err)
5149
}
52-
cfg := &Config{
53-
Controller: ctrlCfg,
54-
Registry: registryConfig,
50+
usageLimitPolicy, err := ensureUsageLimitPolicy(apiConfig, handle)
51+
if err != nil {
52+
return nil, err
5553
}
56-
if apiConfig != nil {
57-
cfg.UsageLimitPolicyType = apiConfig.UsageLimitPolicyType
54+
cfg := &Config{
55+
Controller: ctrlCfg,
56+
Registry: registryConfig,
57+
UsageLimitPolicy: usageLimitPolicy,
5858
}
5959
return cfg, nil
6060
}
61+
62+
func ensureUsageLimitPolicy(apiConfig *configapi.FlowControlConfig, handle plugin.Handle) (flowcontrol.UsageLimitPolicy, error) {
63+
ref := registry.DefaultUsageLimitPolicyRef
64+
if apiConfig != nil && apiConfig.UsageLimit != nil {
65+
ref = apiConfig.UsageLimit.PluginRef
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
76+
}

0 commit comments

Comments
 (0)