Set listener qps and burst (#4558)

This commit is contained in:
Nikola Jokic
2026-09-08 15:49:58 +02:00
committed by GitHub
parent bc677d306c
commit 7c68e1d318
27 changed files with 1161 additions and 226 deletions
+13 -12
View File
@@ -32,18 +32,19 @@ type Config struct {
// It is initially set to nil if VaultType is set.
// Otherwise, it is populated with the GitHub App credentials from the GitHub secret.
*appconfig.AppConfig
EphemeralRunnerSetNamespace string `json:"ephemeral_runner_set_namespace"`
EphemeralRunnerSetName string `json:"ephemeral_runner_set_name"`
MaxRunners int `json:"max_runners"`
MinRunners int `json:"min_runners"`
RunnerScaleSetID int `json:"runner_scale_set_id"`
RunnerScaleSetName string `json:"runner_scale_set_name"`
ServerRootCA string `json:"server_root_ca"`
LogLevel string `json:"log_level"`
LogFormat string `json:"log_format"`
MetricsAddr string `json:"metrics_addr"`
MetricsEndpoint string `json:"metrics_endpoint"`
Metrics *v1alpha1.MetricsConfig `json:"metrics"`
EphemeralRunnerSetNamespace string `json:"ephemeral_runner_set_namespace"`
EphemeralRunnerSetName string `json:"ephemeral_runner_set_name"`
MaxRunners int `json:"max_runners"`
MinRunners int `json:"min_runners"`
RunnerScaleSetID int `json:"runner_scale_set_id"`
RunnerScaleSetName string `json:"runner_scale_set_name"`
ServerRootCA string `json:"server_root_ca"`
LogLevel string `json:"log_level"`
LogFormat string `json:"log_format"`
MetricsAddr string `json:"metrics_addr"`
MetricsEndpoint string `json:"metrics_endpoint"`
Metrics *v1alpha1.MetricsConfig `json:"metrics"`
ListenerConfig *v1alpha1.ListenerConfig `json:"listener_config"`
}
func Read(ctx context.Context, configPath string) (*Config, error) {
+1
View File
@@ -120,6 +120,7 @@ func run(ctx context.Context, config *config.Config) error {
EphemeralRunnerSetName: config.EphemeralRunnerSetName,
MaxRunners: config.MaxRunners,
MinRunners: config.MinRunners,
ScalerConfig: config.ListenerConfig.GetScaler(),
},
scaler.WithLogger(logger.With("component", "worker")),
)
+46 -10
View File
@@ -30,8 +30,14 @@ type Config struct {
EphemeralRunnerSetName string
MaxRunners int
MinRunners int
ScalerConfig *v1alpha1.ScalerConfig
}
const (
defaultQPS = 50
defaultBurst = 100
)
// The Scaler's role is to process the messages it receives from the listener.
// It then initiates Kubernetes API requests to carry out the necessary actions.
type Scaler struct {
@@ -52,12 +58,22 @@ func New(config Config, options ...Option) (*Scaler, error) {
targetRunners: -1,
patchSeq: -1,
}
for _, option := range options {
option(w)
}
if err := w.applyDefaults(); err != nil {
return nil, err
}
conf, err := rest.InClusterConfig()
if err != nil {
return nil, err
}
qps, burst := effectiveRateLimiterConfig(config.ScalerConfig, w.logger)
conf.QPS = float32(qps)
conf.Burst = burst
clientset, err := kubernetes.NewForConfig(conf)
if err != nil {
return nil, err
@@ -65,17 +81,36 @@ func New(config Config, options ...Option) (*Scaler, error) {
w.clientset = clientset
for _, option := range options {
option(w)
}
if err := w.applyDefaults(); err != nil {
return nil, err
}
return w, nil
}
func effectiveRateLimiterConfig(config *v1alpha1.ScalerConfig, logger *slog.Logger) (int, int) {
if config == nil {
logger.Debug("Listener scaler configuration is missing; using defaults", "qps", defaultQPS, "burst", defaultBurst)
return defaultQPS, defaultBurst
}
qps := defaultQPS
if config.QPS == nil {
logger.Debug("Listener scaler qps is missing; using default", "default", defaultQPS)
} else if *config.QPS < 1 {
logger.Warn("Listener scaler qps must be greater than 0; using default", "configured", *config.QPS, "default", defaultQPS)
} else {
qps = *config.QPS
}
burst := defaultBurst
if config.Burst == nil {
logger.Debug("Listener scaler burst is missing; using default", "default", defaultBurst)
} else if *config.Burst < 1 {
logger.Warn("Listener scaler burst must be greater than 0; using default", "configured", *config.Burst, "default", defaultBurst)
} else {
burst = *config.Burst
}
return qps, burst
}
func (w *Scaler) applyDefaults() error {
if w.logger == nil {
w.logger = slog.New(slog.DiscardHandler)
@@ -213,10 +248,11 @@ func (w *Scaler) HandleDesiredRunnerCount(ctx context.Context, count int) (int,
Do(ctx).
Into(patchedEphemeralRunnerSet)
if err != nil {
return 0, fmt.Errorf("could not patch ephemeral runner set , patch JSON: %s, error: %w", string(mergePatch), err)
return 0, fmt.Errorf("could not patch ephemeral runner set, patch JSON: %s, error: %w", string(mergePatch), err)
}
w.logger.Info("Ephemeral runner set scaled.",
w.logger.Info(
"Ephemeral runner set scaled.",
"namespace", w.config.EphemeralRunnerSetNamespace,
"name", w.config.EphemeralRunnerSetName,
"replicas", patchedEphemeralRunnerSet.Spec.Replicas,
+111
View File
@@ -1,15 +1,126 @@
package scaler
import (
"bytes"
"log/slog"
"math"
"strconv"
"testing"
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
"github.com/stretchr/testify/assert"
)
var discardLogger = slog.New(slog.DiscardHandler)
func TestEffectiveRateLimiterConfig(t *testing.T) {
qps := 75
burst := 150
zero := 0
negative := -1
tests := []struct {
name string
config *v1alpha1.ScalerConfig
wantQPS int
wantBurst int
wantLog string
wantLevel string
}{
{
name: "uses configured values",
config: &v1alpha1.ScalerConfig{
QPS: &qps,
Burst: &burst,
},
wantQPS: qps,
wantBurst: burst,
},
{
name: "defaults missing config",
wantQPS: defaultQPS,
wantBurst: defaultBurst,
wantLog: "Listener scaler configuration is missing; using defaults",
wantLevel: "DEBUG",
},
{
name: "defaults missing qps",
config: &v1alpha1.ScalerConfig{Burst: &burst},
wantQPS: defaultQPS,
wantBurst: burst,
wantLog: "Listener scaler qps is missing; using default",
wantLevel: "DEBUG",
},
{
name: "defaults missing burst",
config: &v1alpha1.ScalerConfig{QPS: &qps},
wantQPS: qps,
wantBurst: defaultBurst,
wantLog: "Listener scaler burst is missing; using default",
wantLevel: "DEBUG",
},
{
name: "defaults zero qps",
config: &v1alpha1.ScalerConfig{QPS: &zero, Burst: &burst},
wantQPS: defaultQPS,
wantBurst: burst,
wantLog: "Listener scaler qps must be greater than 0; using default",
wantLevel: "WARN",
},
{
name: "defaults negative burst",
config: &v1alpha1.ScalerConfig{QPS: &qps, Burst: &negative},
wantQPS: qps,
wantBurst: defaultBurst,
wantLog: "Listener scaler burst must be greater than 0; using default",
wantLevel: "WARN",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
var logs bytes.Buffer
logger := slog.New(slog.NewTextHandler(&logs, &slog.HandlerOptions{Level: slog.LevelDebug}))
qps, burst := effectiveRateLimiterConfig(tt.config, logger)
assert.Equal(t, tt.wantQPS, qps)
assert.Equal(t, tt.wantBurst, burst)
if tt.wantLog == "" {
assert.Empty(t, logs.String())
return
}
assert.Contains(t, logs.String(), "msg="+strconv.Quote(tt.wantLog))
// Missing values are a normal configuration, so they must not be
// logged as warnings; only out-of-range values are.
assert.Contains(t, logs.String(), "level="+tt.wantLevel)
})
}
}
// TestEffectiveRateLimiterConfig_QuietAtInfoLevel asserts that a listener which
// does not configure the scaler produces no output at the default log level.
func TestEffectiveRateLimiterConfig_QuietAtInfoLevel(t *testing.T) {
for _, tt := range []struct {
name string
config *v1alpha1.ScalerConfig
}{
{name: "nil config"},
{name: "empty config", config: &v1alpha1.ScalerConfig{}},
} {
t.Run(tt.name, func(t *testing.T) {
var logs bytes.Buffer
logger := slog.New(slog.NewTextHandler(&logs, &slog.HandlerOptions{Level: slog.LevelInfo}))
qps, burst := effectiveRateLimiterConfig(tt.config, logger)
assert.Equal(t, defaultQPS, qps)
assert.Equal(t, defaultBurst, burst)
assert.Empty(t, logs.String())
})
}
}
func TestSetDesiredWorkerState_MinMaxDefaults(t *testing.T) {
newEmptyWorker := func() *Scaler {
return &Scaler{