diff --git a/apis/actions.github.com/v1alpha1/listenerconfig_types.go b/apis/actions.github.com/v1alpha1/listenerconfig_types.go index ec64e614..8117e092 100644 --- a/apis/actions.github.com/v1alpha1/listenerconfig_types.go +++ b/apis/actions.github.com/v1alpha1/listenerconfig_types.go @@ -14,19 +14,58 @@ func (c *ListenerConfig) GetScaler() *ScalerConfig { return c.Scaler } -// ScalerConfig configures the Kubernetes client used by the ghalistener scaler. +// ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. +// +// The scaler talks to the API server over two independent clients, because the +// two kinds of traffic have very different shapes and only one of them is on the +// critical path: +// +// - The scale client publishes the desired runner count. That is one patch per +// message, and it is the only patch that creates runners, so it must never +// queue behind anything. +// - The job client patches job started events. That is two calls per event and +// can be a whole batch at once, so it is the traffic that actually consumes +// the rate limit. +// +// A single shared client lets a batch of job event patches drain the token +// bucket ahead of the scale patch, which delays the only call new jobs are +// waiting on. Splitting them keeps the scale patch clear of that backlog. type ScalerConfig struct { + // QPS is the query per second limit of the client that patches job started + // events. This is the bulk of the scaler's API traffic, at up to two calls + // per job started event. // +optional // +kubebuilder:validation:Minimum:=1 QPS *int `json:"qps,omitempty"` + // Burst is the burst limit of the client that patches job started events. // +optional // +kubebuilder:validation:Minimum:=1 Burst *int `json:"burst,omitempty"` - // Workers is the number of job started and job completed events the scaler - // handles concurrently within a single scale set message. The worker that - // scales the EphemeralRunnerSet runs on top of these. + // ScaleQPS is the query per second limit of the client that publishes the + // desired runner count. The scaler issues at most one such patch per scale + // set message, so this only has to be large enough that the patch never + // waits on a token; it is deliberately a small budget separate from QPS + // rather than a share of it. + // +optional + // +kubebuilder:validation:Minimum:=1 + ScaleQPS *int `json:"scaleQPS,omitempty"` + + // ScaleBurst is the burst limit of the client that publishes the desired + // runner count. + // +optional + // +kubebuilder:validation:Minimum:=1 + ScaleBurst *int `json:"scaleBurst,omitempty"` + + // Workers is the number of job started events the scaler patches + // concurrently. The events are drained from a background queue rather than + // being tied to the message they arrived on, so this bounds how many job + // patches are in flight at any moment, not how many a single message may + // carry. + // + // Raising it past what QPS sustains does nothing, since the rate limiter + // rather than the worker count is what bounds throughput. // +optional // +kubebuilder:validation:Minimum:=1 Workers *int `json:"workers,omitempty"` diff --git a/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go b/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go index 63684022..9c88dc7f 100644 --- a/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go +++ b/apis/actions.github.com/v1alpha1/zz_generated.deepcopy.go @@ -807,6 +807,16 @@ func (in *ScalerConfig) DeepCopyInto(out *ScalerConfig) { *out = new(int) **out = **in } + if in.ScaleQPS != nil { + in, out := &in.ScaleQPS, &out.ScaleQPS + *out = new(int) + **out = **in + } + if in.ScaleBurst != nil { + in, out := &in.ScaleBurst, &out.ScaleBurst + *out = new(int) + **out = **in + } if in.Workers != nil { in, out := &in.Workers, &out.Workers *out = new(int) diff --git a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml index 6551d057..57c87ce1 100644 --- a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml +++ b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalinglisteners.yaml @@ -131,20 +131,61 @@ spec: pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used - by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches + job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object diff --git a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml index 3f8922cf..290ebddc 100644 --- a/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller-experimental/crds/actions.github.com_autoscalingrunnersets.yaml @@ -149,19 +149,60 @@ spec: description: ListenerConfig holds configuration for the ghalistener pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object diff --git a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml index 6551d057..57c87ce1 100644 --- a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml +++ b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalinglisteners.yaml @@ -131,20 +131,61 @@ spec: pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used - by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches + job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object diff --git a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml index 3f8922cf..290ebddc 100644 --- a/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml +++ b/charts/gha-runner-scale-set-controller/crds/actions.github.com_autoscalingrunnersets.yaml @@ -149,19 +149,60 @@ spec: description: ListenerConfig holds configuration for the ghalistener pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object diff --git a/charts/gha-runner-scale-set/tests/template_test.go b/charts/gha-runner-scale-set/tests/template_test.go index fa21e8d5..f9535f05 100644 --- a/charts/gha-runner-scale-set/tests/template_test.go +++ b/charts/gha-runner-scale-set/tests/template_test.go @@ -175,6 +175,16 @@ func TestTemplateListenerScalerValidation(t *testing.T) { setValues: map[string]string{"listenerConfig.scaler.workers": "0"}, wantErrorText: "at '/listenerConfig/scaler/workers': minimum: got 0, want 1", }, + { + name: "zero scaleQPS", + setValues: map[string]string{"listenerConfig.scaler.scaleQPS": "0"}, + wantErrorText: "at '/listenerConfig/scaler/scaleQPS': minimum: got 0, want 1", + }, + { + name: "negative scaleBurst", + setValues: map[string]string{"listenerConfig.scaler.scaleBurst": "-1"}, + wantErrorText: "at '/listenerConfig/scaler/scaleBurst': minimum: got -1, want 1", + }, { name: "fractional qps", setValues: map[string]string{"listenerConfig.scaler.qps": "1.5"}, @@ -248,38 +258,59 @@ func TestTemplateListenerScalerConfig(t *testing.T) { "controllerServiceAccount.namespace": "arc-system", } tests := []struct { - name string - setValues map[string]string - wantQPS *int - wantBurst *int - wantWorkers *int + name string + setValues map[string]string + wantQPS *int + wantBurst *int + wantScaleQPS *int + wantScaleBurst *int + wantWorkers *int }{ { - name: "defaults from values.yaml", - wantQPS: ptr.To(50), - wantBurst: ptr.To(100), - wantWorkers: ptr.To(10), + name: "defaults from values.yaml", + wantQPS: ptr.To(50), + wantBurst: ptr.To(100), + wantScaleQPS: ptr.To(10), + wantScaleBurst: ptr.To(20), + wantWorkers: ptr.To(10), }, { - name: "both overridden", - setValues: map[string]string{"listenerConfig.scaler.qps": "100", "listenerConfig.scaler.burst": "200"}, - wantQPS: ptr.To(100), - wantBurst: ptr.To(200), - wantWorkers: ptr.To(10), + name: "both overridden", + setValues: map[string]string{"listenerConfig.scaler.qps": "100", "listenerConfig.scaler.burst": "200"}, + wantQPS: ptr.To(100), + wantBurst: ptr.To(200), + wantScaleQPS: ptr.To(10), + wantScaleBurst: ptr.To(20), + wantWorkers: ptr.To(10), }, { - name: "qps overridden keeps default burst", - setValues: map[string]string{"listenerConfig.scaler.qps": "75"}, - wantQPS: ptr.To(75), - wantBurst: ptr.To(100), - wantWorkers: ptr.To(10), + name: "qps overridden keeps default burst", + setValues: map[string]string{"listenerConfig.scaler.qps": "75"}, + wantQPS: ptr.To(75), + wantBurst: ptr.To(100), + wantScaleQPS: ptr.To(10), + wantScaleBurst: ptr.To(20), + wantWorkers: ptr.To(10), }, { - name: "workers overridden", - setValues: map[string]string{"listenerConfig.scaler.workers": "32"}, - wantQPS: ptr.To(50), - wantBurst: ptr.To(100), - wantWorkers: ptr.To(32), + name: "workers overridden", + setValues: map[string]string{"listenerConfig.scaler.workers": "32"}, + wantQPS: ptr.To(50), + wantBurst: ptr.To(100), + wantScaleQPS: ptr.To(10), + wantScaleBurst: ptr.To(20), + wantWorkers: ptr.To(32), + }, + { + // The scale client budget is deliberately independent of the job + // client budget: raising qps must not quietly reallocate it. + name: "scale client overridden independently", + setValues: map[string]string{"listenerConfig.scaler.scaleQPS": "5", "listenerConfig.scaler.scaleBurst": "9"}, + wantQPS: ptr.To(50), + wantBurst: ptr.To(100), + wantScaleQPS: ptr.To(5), + wantScaleBurst: ptr.To(9), + wantWorkers: ptr.To(10), }, { name: "listenerConfig disabled", @@ -307,7 +338,8 @@ func TestTemplateListenerScalerConfig(t *testing.T) { var ars v1alpha1.AutoscalingRunnerSet helm.UnmarshalK8SYaml(t, output, &ars) - if tt.wantQPS == nil && tt.wantBurst == nil && tt.wantWorkers == nil { + if tt.wantQPS == nil && tt.wantBurst == nil && tt.wantWorkers == nil && + tt.wantScaleQPS == nil && tt.wantScaleBurst == nil { assert.Nil(t, ars.Spec.ListenerConfig.GetScaler()) return } @@ -316,6 +348,8 @@ func TestTemplateListenerScalerConfig(t *testing.T) { require.NotNil(t, scaler) assert.Equal(t, tt.wantQPS, scaler.QPS) assert.Equal(t, tt.wantBurst, scaler.Burst) + assert.Equal(t, tt.wantScaleQPS, scaler.ScaleQPS) + assert.Equal(t, tt.wantScaleBurst, scaler.ScaleBurst) assert.Equal(t, tt.wantWorkers, scaler.Workers) }) } diff --git a/charts/gha-runner-scale-set/values.schema.json b/charts/gha-runner-scale-set/values.schema.json index 72b26aa5..3319a530 100644 --- a/charts/gha-runner-scale-set/values.schema.json +++ b/charts/gha-runner-scale-set/values.schema.json @@ -10,22 +10,32 @@ "additionalProperties": false, "properties": { "scaler": { - "description": "Configuration for the Kubernetes client used by the ghalistener scaler.", + "description": "Configuration for the Kubernetes clients used by the ghalistener scaler.", "type": ["object", "null"], "additionalProperties": false, "properties": { "qps": { - "description": "Queries per second the listener may issue to the Kubernetes API server.", + "description": "Queries per second the listener may issue to the Kubernetes API server when patching job started events.", "type": "integer", "minimum": 1 }, "burst": { - "description": "Burst of queries the listener may issue to the Kubernetes API server.", + "description": "Burst of queries the listener may issue to the Kubernetes API server when patching job started events.", + "type": "integer", + "minimum": 1 + }, + "scaleQPS": { + "description": "Queries per second the listener may issue to the Kubernetes API server when publishing the desired runner count. This is a separate budget from qps so a batch of job event patches cannot delay it.", + "type": "integer", + "minimum": 1 + }, + "scaleBurst": { + "description": "Burst of queries the listener may issue to the Kubernetes API server when publishing the desired runner count.", "type": "integer", "minimum": 1 }, "workers": { - "description": "Number of job started and job completed events the listener handles concurrently within a single scale set message.", + "description": "Number of job started events the listener patches concurrently from its background queue.", "type": "integer", "minimum": 1 } diff --git a/charts/gha-runner-scale-set/values.yaml b/charts/gha-runner-scale-set/values.yaml index 269983f5..88977634 100644 --- a/charts/gha-runner-scale-set/values.yaml +++ b/charts/gha-runner-scale-set/values.yaml @@ -149,11 +149,25 @@ githubConfigSecret: ## listenerConfig holds configuration for the ghalistener pod. listenerConfig: scaler: + ## The listener talks to the Kubernetes API server over two independent + ## clients. Job started patches are the bulk of the traffic, at up to two + ## calls per event, while publishing the desired runner count is a single + ## patch per message and is the only one that creates runners. Giving them + ## separate budgets keeps a large batch of job patches from draining the + ## rate limiter ahead of the patch new jobs are actually waiting on. + ## + ## qps and burst govern the client patching job started events. qps: 50 burst: 100 - ## Number of job started and job completed events handled concurrently - ## within a single scale set message. The worker patching the - ## EphemeralRunnerSet with the desired runner count runs on top of these. + ## scaleQPS and scaleBurst govern the client publishing the desired runner + ## count. One patch per message needs very little, so this only has to be + ## large enough that the patch never waits on a token. + scaleQPS: 10 + scaleBurst: 20 + ## Number of job started events patched concurrently. They are drained from + ## a background queue rather than tied to the message they arrived on, so + ## this bounds how many patches are in flight at once. Raising it beyond + ## what qps sustains has no effect. workers: 10 ## listenerMetrics are configurable metrics applied to the listener. diff --git a/cmd/ghalistener/main.go b/cmd/ghalistener/main.go index 52a31c70..d8737dcc 100644 --- a/cmd/ghalistener/main.go +++ b/cmd/ghalistener/main.go @@ -7,6 +7,7 @@ import ( "os" "os/signal" "syscall" + "time" "github.com/actions/actions-runner-controller/cmd/ghalistener/config" "github.com/actions/actions-runner-controller/cmd/ghalistener/metrics" @@ -122,6 +123,18 @@ func run(ctx context.Context, config *config.Config) error { return fmt.Errorf("failed to create new kubernetes worker: %w", err) } + // The scaler patches job started events in the background, after the listener + // has already acked the message they arrived on. Draining them on the way out + // keeps a clean shutdown from stranding job information the service believes + // was recorded. + defer func() { + drainCtx, cancelDrain := context.WithTimeout(context.Background(), 30*time.Second) + defer cancelDrain() + if err := scaler.Close(drainCtx); err != nil { + logger.Error("Failed to drain queued job started events", "error", err) + } + }() + g, ctx := errgroup.WithContext(ctx) metricsCtx, cancelMetrics := context.WithCancelCause(ctx) diff --git a/cmd/ghalistener/scaler/queue.go b/cmd/ghalistener/scaler/queue.go new file mode 100644 index 00000000..a5445d3a --- /dev/null +++ b/cmd/ghalistener/scaler/queue.go @@ -0,0 +1,98 @@ +package scaler + +import ( + "sync" + + "github.com/actions/scaleset" +) + +// jobQueue is an unbounded FIFO of job started events waiting to be patched. +// +// It is deliberately unbounded. The alternative, a buffered channel, blocks the +// producer once it fills, and the producer here is Scale: blocking it would put +// the job patches back on the critical path of the next scale decision, which is +// the entire reason they were moved off it. +// +// Unbounded is safe because the queue is not what decides how much work arrives. +// The scaler patches two calls per job started event, so the drain rate it needs +// is twice the rate at which jobs actually start, and that is a property of the +// workload rather than of how fast the listener polls. A faster listener polls +// more often and carries proportionally fewer events per message. The queue +// therefore sits near empty whenever the job client's QPS can sustain the real +// job start rate, and grows only when it genuinely cannot, where the backlog is +// a symptom to observe rather than something a bound would fix. +type jobQueue struct { + mu sync.Mutex + cond *sync.Cond + items []*scaleset.JobStarted + closed bool +} + +func newJobQueue() *jobQueue { + q := &jobQueue{} + q.cond = sync.NewCond(&q.mu) + return q +} + +// push appends events to the queue and never blocks. Events pushed after the +// queue is closed are dropped, which only happens during shutdown, once the +// workers are already on their way out and could not patch them anyway. +func (q *jobQueue) push(items ...*scaleset.JobStarted) { + if len(items) == 0 { + return + } + + q.mu.Lock() + defer q.mu.Unlock() + + if q.closed { + return + } + + q.items = append(q.items, items...) + q.cond.Broadcast() +} + +// pop returns the next event, blocking until one is available. It reports false +// once the queue is both closed and fully drained, which is the signal for a +// worker to exit. Closing does not discard what is already queued, so a shutdown +// still patches everything it accepted. +func (q *jobQueue) pop() (*scaleset.JobStarted, bool) { + q.mu.Lock() + defer q.mu.Unlock() + + for len(q.items) == 0 && !q.closed { + q.cond.Wait() + } + + if len(q.items) == 0 { + return nil, false + } + + item := q.items[0] + // Clear the slot before reslicing so the popped event is not kept alive by + // the backing array until it is overwritten. + q.items[0] = nil + q.items = q.items[1:] + + return item, true +} + +// close stops the queue from accepting new events and wakes every worker so the +// ones with nothing left to do can exit. +func (q *jobQueue) close() { + q.mu.Lock() + defer q.mu.Unlock() + + q.closed = true + q.cond.Broadcast() +} + +// depth is the number of events still waiting to be patched. It is the signal +// that the job client's QPS is not keeping up with the rate jobs are starting. +func (q *jobQueue) depth() int { + q.mu.Lock() + defer q.mu.Unlock() + + return len(q.items) +} diff --git a/cmd/ghalistener/scaler/scale_test.go b/cmd/ghalistener/scaler/scale_test.go index 5e1d6c4e..fca53c8e 100644 --- a/cmd/ghalistener/scaler/scale_test.go +++ b/cmd/ghalistener/scaler/scale_test.go @@ -26,9 +26,16 @@ type fakeAcquirer struct { mu sync.Mutex acquired [][]int64 err error + // onAcquire runs before the call is recorded, so a test can observe what had + // already happened by the time the scaler reached the acquire step. + onAcquire func() } func (f *fakeAcquirer) AcquireJobs(ctx context.Context, requestIDs []int64) ([]int64, error) { + if f.onAcquire != nil { + f.onAcquire() + } + f.mu.Lock() defer f.mu.Unlock() f.acquired = append(f.acquired, requestIDs) @@ -152,23 +159,50 @@ func newScaleScaler(t *testing.T, client JobAcquirer, config Config, workers int } })) t.Cleanup(httpServer.Close) - // Registered last so it runs first: a failed assertion must not leave a - // request parked inside the handler while Close waits for it. - t.Cleanup(server.release) clientset, err := kubernetes.NewForConfig(&rest.Config{Host: httpServer.URL, QPS: -1}) require.NoError(t, err) - return &Scaler{ - clientset: clientset, - client: client, - config: config, - targetRunners: -1, - patchSeq: -1, - logger: discardLogger, - metrics: metrics.Discard, - workers: workers, - }, server + w := &Scaler{ + // The split matters in production, where the two clients carry separate + // rate limits. The tests only care about ordering, so one unthrottled + // client backs both. + scaleClientset: clientset, + jobClientset: clientset, + client: client, + config: config, + targetRunners: -1, + patchSeq: -1, + logger: discardLogger, + metrics: metrics.Discard, + workers: workers, + jobs: newJobQueue(), + } + w.startJobWorkers() + + // Registered after the server so it runs first: the workers have to be gone + // before the server they are calling goes away. + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + _ = w.Close(ctx) + }) + // Registered last so it runs first of all: a failed assertion must not leave + // a request parked inside the handler while the workers are being drained. + t.Cleanup(server.release) + + return w, server +} + +// drain waits for every queued job started event to be patched. Scale +// deliberately does not wait for them, so any assertion about job patches has +// to ask for the drain explicitly. +func drain(t *testing.T, w *Scaler) { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + require.NoError(t, w.Close(ctx)) } // runnerNameFromPath extracts the ephemeral runner name from a request path of @@ -271,6 +305,7 @@ func TestScale_HandlesJobStartedConcurrently(t *testing.T) { } require.NoError(t, w.Scale(t.Context(), msg)) + drain(t, w) server.mu.Lock() peak := server.peak @@ -298,6 +333,7 @@ func TestScale_WorkersBoundConcurrency(t *testing.T) { } require.NoError(t, w.Scale(t.Context(), msg)) + drain(t, w) server.mu.Lock() peak := server.peak @@ -347,18 +383,119 @@ func TestScale_PublishesDesiredCountFirst(t *testing.T) { JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)}, })) + // Safe to assert before the drain: the events are not queued until + // the scale patch has already returned, so no job patch can precede + // it however the workers are scheduled. server.mu.Lock() - defer server.mu.Unlock() - - assert.Equal(t, tt.assigned, w.targetRunners) require.Len(t, server.runnerSetPatches, 1) assert.Equal(t, 0, server.runnerPatchesBeforeRunnerSet, "the desired count is published before any job event patch") + server.mu.Unlock() + + drain(t, w) + + server.mu.Lock() + defer server.mu.Unlock() + assert.Equal(t, tt.assigned, w.targetRunners) assert.Equal(t, 2, server.runnerPatches) }) } } +// TestScale_DoesNotWaitForJobStartedPatches is the property the background queue +// exists for. The desired count is what new jobs wait on; the job started +// patches are bookkeeping, and at two API calls each a full batch of them used +// to sit between one scale decision and the next. +// +// The server holds every ephemeral runner request for the whole test, so a +// scaler that still patched them inline could not return at all. +func TestScale_DoesNotWaitForJobStartedPatches(t *testing.T) { + const jobs = 20 + + w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers) + + // Never released by the test itself; cleanup releases it before draining. + server.block() + + msg := &scaleset.RunnerScaleSetMessage{ + MessageID: 1, + Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: jobs}, + } + for i := range jobs { + msg.JobStartedMessages = append(msg.JobStartedMessages, jobStarted(i)) + } + + done := make(chan error, 1) + go func() { done <- w.Scale(t.Context(), msg) }() + + select { + case err := <-done: + require.NoError(t, err) + case <-time.After(30 * time.Second): + t.Fatal("Scale blocked on job started patches that cannot complete") + } + + server.mu.Lock() + defer server.mu.Unlock() + assert.Len(t, server.runnerSetPatches, 1, + "the desired count is published even while every job patch is stuck") + assert.Equal(t, 0, server.runnerPatches) +} + +// TestScale_AcquiresAfterPublishingDesiredCount pins the other half of the +// reordering. Acquiring is a round trip to the Actions service, and the scale +// decision does not depend on its result: it is derived from the statistics the +// service already put in the message. Publishing first therefore costs nothing +// and keeps that round trip off the path new runners wait on. +func TestScale_AcquiresAfterPublishingDesiredCount(t *testing.T) { + acquirer := &fakeAcquirer{} + w, server := newScaleScaler(t, acquirer, defaultConfig(), defaultWorkers) + + var runnerSetPatchesAtAcquire int + acquirer.onAcquire = func() { + server.mu.Lock() + defer server.mu.Unlock() + runnerSetPatchesAtAcquire = len(server.runnerSetPatches) + } + + require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{ + MessageID: 1, + JobAvailableMessages: []*scaleset.JobAvailable{{JobMessageBase: scaleset.JobMessageBase{RunnerRequestID: 1}}}, + Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 1}, + })) + + assert.Equal(t, 1, runnerSetPatchesAtAcquire, + "the desired count is published before the acquire round trip") +} + +// TestScale_CloseDrainsQueuedJobStartedEvents covers the cost of not waiting. +// The listener acks a message as soon as Scale returns, so by the time these +// patches run the service already believes they were recorded and will never +// redeliver them. Shutting down has to finish them rather than drop them. +func TestScale_CloseDrainsQueuedJobStartedEvents(t *testing.T) { + const jobs = 16 + + // Two workers against sixteen events, so the queue is guaranteed to still + // hold work when Close is called. + w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), 2) + + msg := &scaleset.RunnerScaleSetMessage{ + MessageID: 1, + Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: jobs}, + } + for i := range jobs { + msg.JobStartedMessages = append(msg.JobStartedMessages, jobStarted(i)) + } + require.NoError(t, w.Scale(t.Context(), msg)) + + drain(t, w) + + server.mu.Lock() + defer server.mu.Unlock() + assert.Equal(t, jobs, server.runnerPatches, + "a clean shutdown patches everything it accepted from an acked message") +} + // TestScale_ScaleDownNeverPublishesPatchIDZero is the invariant the ordering // above relies on. Patch ID 0 is the only one the runner set controller acts on // to delete idle runners, so a falling target must never carry it -- otherwise @@ -429,7 +566,52 @@ func TestScale_AcquireFailureIsNotAcked(t *testing.T) { server.mu.Lock() defer server.mu.Unlock() - assert.Empty(t, server.runnerSetPatches, "nothing is published when the jobs were never acquired") + // The desired count is published ahead of the acquire, so it survives the + // failure. That is safe, and deliberate: it is derived from the statistics + // the service put in the message rather than from anything the acquire + // returns, and republishing it when the message is redelivered is a no-op. + assert.Len(t, server.runnerSetPatches, 1, + "the desired count is published before the acquire, so it stands even when the acquire fails") +} + +func TestEffectiveScaleRateLimiterConfig(t *testing.T) { + qps, burst := 3, 7 + + tests := []struct { + name string + config *v1alpha1.ScalerConfig + wantQPS int + wantBurst int + }{ + {name: "nil config", config: nil, wantQPS: defaultScaleQPS, wantBurst: defaultScaleBurst}, + {name: "unset", config: &v1alpha1.ScalerConfig{}, wantQPS: defaultScaleQPS, wantBurst: defaultScaleBurst}, + { + name: "configured", + config: &v1alpha1.ScalerConfig{ScaleQPS: &qps, ScaleBurst: &burst}, + wantQPS: qps, + wantBurst: burst, + }, + { + name: "zero falls back", + config: &v1alpha1.ScalerConfig{ScaleQPS: new(int), ScaleBurst: new(int)}, + wantQPS: defaultScaleQPS, + wantBurst: defaultScaleBurst, + }, + { + name: "independent of the job client budget", + config: &v1alpha1.ScalerConfig{QPS: &qps, Burst: &burst}, + wantQPS: defaultScaleQPS, + wantBurst: defaultScaleBurst, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + gotQPS, gotBurst := effectiveScaleRateLimiterConfig(tt.config, discardLogger) + assert.Equal(t, tt.wantQPS, gotQPS) + assert.Equal(t, tt.wantBurst, gotBurst) + }) + } } func TestEffectiveWorkerCount(t *testing.T) { diff --git a/cmd/ghalistener/scaler/scaler.go b/cmd/ghalistener/scaler/scaler.go index a3ffae4d..268a2fe7 100644 --- a/cmd/ghalistener/scaler/scaler.go +++ b/cmd/ghalistener/scaler/scaler.go @@ -7,13 +7,13 @@ import ( "fmt" "log/slog" "math" + "sync" "github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1" "github.com/actions/actions-runner-controller/cmd/ghalistener/metrics" "github.com/actions/scaleset" "github.com/actions/scaleset/listener" jsonpatch "github.com/evanphx/json-patch" - "golang.org/x/sync/errgroup" kerrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/kubernetes" @@ -53,10 +53,17 @@ type Config struct { const ( defaultQPS = 50 defaultBurst = 100 - // defaultWorkers bounds how many job events are patched at once. Each event - // costs at most a GET and a PATCH, so the default stays well inside the - // default QPS budget while still collapsing a batch of events into a few - // round trips worth of latency. + // defaultScaleQPS and defaultScaleBurst budget the client that publishes the + // desired runner count. That is one patch per message, so a small budget is + // enough for it to never wait on a token. It is separate from the job client + // rather than carved out of it: sharing one bucket is what let a batch of job + // patches delay the scale patch in the first place. + defaultScaleQPS = 10 + defaultScaleBurst = 20 + // defaultWorkers bounds how many job started events are patched at once. + // Each event costs at most a GET and a PATCH, and the rate limiter rather + // than this number is what bounds sustained throughput, so this only has to + // be large enough to keep the job client's tokens spoken for. defaultWorkers = 10 ) @@ -70,7 +77,12 @@ type JobAcquirer interface { // 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 { - clientset *kubernetes.Clientset + // scaleClientset publishes the desired runner count and nothing else, so its + // rate limiter is never drained by job event traffic. + scaleClientset *kubernetes.Clientset + // jobClientset patches job started events. It is used only by the background + // workers, never by Scale. + jobClientset *kubernetes.Clientset client JobAcquirer config Config metrics metrics.Recorder @@ -84,7 +96,19 @@ type Scaler struct { // message at all, so the scaler keeps them to stay able to converge on an // otherwise idle scale set. lastStatistics *scaleset.RunnerScaleSetStatistic - logger *slog.Logger + + // jobs holds job started events accepted from a message but not yet patched. + jobs *jobQueue + // jobWorkers tracks the pool draining jobs. + jobWorkers sync.WaitGroup + // jobCtx scopes the background patches. It is rooted at context.Background() + // rather than at any message's context: the patches outlive the message that + // produced them, so cancelling that message must not abandon them. + jobCtx context.Context + jobCancel context.CancelFunc + closeOnce sync.Once + + logger *slog.Logger } var _ listener.Scaler = (*Scaler)(nil) @@ -99,6 +123,7 @@ func New(client JobAcquirer, config Config, options ...Option) (*Scaler, error) config: config, targetRunners: -1, patchSeq: -1, + jobs: newJobQueue(), } for _, option := range options { option(w) @@ -112,21 +137,101 @@ func New(client JobAcquirer, config Config, options ...Option) (*Scaler, error) return nil, err } - qps, burst := effectiveRateLimiterConfig(config.ScalerConfig, w.logger) - conf.QPS = float32(qps) - conf.Burst = burst w.workers = effectiveWorkerCount(config.ScalerConfig, w.logger) - clientset, err := kubernetes.NewForConfig(conf) + jobQPS, jobBurst := effectiveRateLimiterConfig(config.ScalerConfig, w.logger) + jobConf := rest.CopyConfig(conf) + jobConf.QPS = float32(jobQPS) + jobConf.Burst = jobBurst + jobClientset, err := kubernetes.NewForConfig(jobConf) if err != nil { return nil, err } - w.clientset = clientset + scaleQPS, scaleBurst := effectiveScaleRateLimiterConfig(config.ScalerConfig, w.logger) + scaleConf := rest.CopyConfig(conf) + scaleConf.QPS = float32(scaleQPS) + scaleConf.Burst = scaleBurst + scaleClientset, err := kubernetes.NewForConfig(scaleConf) + if err != nil { + return nil, err + } + + w.jobClientset = jobClientset + w.scaleClientset = scaleClientset + + w.startJobWorkers() return w, nil } +// startJobWorkers brings up the pool that drains the job started queue. +func (w *Scaler) startJobWorkers() { + w.jobCtx, w.jobCancel = context.WithCancel(context.Background()) + + for range w.workers { + w.jobWorkers.Add(1) + go func() { + defer w.jobWorkers.Done() + + for { + jobInfo, ok := w.jobs.pop() + if !ok { + return + } + + // Recorded here rather than at enqueue time so the metric and the + // patch describe the same moment. + w.metrics.RecordJobStarted(jobInfo) + + if err := w.HandleJobStarted(w.jobCtx, jobInfo); err != nil { + // The message this event arrived on has already been acked, so + // there is nobody left to return the error to. Losing the patch + // costs a stale Status.JobID: the runner set may try to delete + // the runner as idle, and the Actions service rejects that while + // the job is still running, so the job itself is not at risk. + w.logger.Error("Failed to patch job started event", + "runnerName", jobInfo.RunnerName, + "requestId", jobInfo.RunnerRequestID, + "error", err.Error(), + ) + } + } + }() + } +} + +// Close drains the job started queue and stops the workers. Everything already +// accepted from an acked message is patched first, so a clean shutdown does not +// strand job information the service believes was recorded. +// +// If ctx expires before the queue drains, the remaining patches are abandoned +// and any in-flight request is cancelled. +func (w *Scaler) Close(ctx context.Context) error { + w.closeOnce.Do(func() { + w.jobs.close() + }) + + drained := make(chan struct{}) + go func() { + w.jobWorkers.Wait() + close(drained) + }() + + select { + case <-drained: + w.jobCancel() + return nil + case <-ctx.Done(): + if depth := w.jobs.depth(); depth > 0 { + w.logger.Error("Abandoning queued job started events", "count", depth) + } + w.jobCancel() + <-drained + return ctx.Err() + } +} + 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) @@ -154,6 +259,35 @@ func effectiveRateLimiterConfig(config *v1alpha1.ScalerConfig, logger *slog.Logg return qps, burst } +// effectiveScaleRateLimiterConfig resolves the budget for the client that +// publishes the desired runner count. +func effectiveScaleRateLimiterConfig(config *v1alpha1.ScalerConfig, logger *slog.Logger) (int, int) { + if config == nil { + logger.Debug("Listener scaler configuration is missing; using defaults", "scaleQPS", defaultScaleQPS, "scaleBurst", defaultScaleBurst) + return defaultScaleQPS, defaultScaleBurst + } + + qps := defaultScaleQPS + if config.ScaleQPS == nil { + logger.Debug("Listener scaler scaleQPS is missing; using default", "default", defaultScaleQPS) + } else if *config.ScaleQPS < 1 { + logger.Warn("Listener scaler scaleQPS must be greater than 0; using default", "configured", *config.ScaleQPS, "default", defaultScaleQPS) + } else { + qps = *config.ScaleQPS + } + + burst := defaultScaleBurst + if config.ScaleBurst == nil { + logger.Debug("Listener scaler scaleBurst is missing; using default", "default", defaultScaleBurst) + } else if *config.ScaleBurst < 1 { + logger.Warn("Listener scaler scaleBurst must be greater than 0; using default", "configured", *config.ScaleBurst, "default", defaultScaleBurst) + } else { + burst = *config.ScaleBurst + } + + return qps, burst +} + func effectiveWorkerCount(config *v1alpha1.ScalerConfig, logger *slog.Logger) int { if config == nil || config.Workers == nil { logger.Debug("Listener scaler workers is missing; using default", "default", defaultWorkers) @@ -191,11 +325,26 @@ func (w *Scaler) applyDefaults() error { // redelivers it otherwise, so every step below is idempotent and safe to repeat // after a partially applied message. // -// The desired runner count is published first, since that is the patch new jobs -// wait on. The job started and job completed events are then patched across a -// bounded worker pool: they touch distinct EphemeralRunners and carry no -// ordering between them, so they run concurrently instead of one API call at a -// time. +// The order is chosen so that nothing the runner set is waiting on sits behind +// work it is not waiting on: +// +// 1. The desired runner count is published. It is the only patch that creates +// runners, so it goes out on its own client before anything else can consume +// a rate limit token. +// 2. The job started events are handed to the background pool. They are +// bookkeeping rather than something new jobs wait on, and at two API calls +// per event they are what a large batch would otherwise spend the whole +// message on. +// 3. The available jobs are acquired. This is a single call to the Actions +// service, and it stays here, ahead of the ack, because a job that is never +// acquired is never assigned; unlike the patches above, losing it is not +// something a later message repairs. +// +// Acquiring last rather than first costs the round trip of the scale patch in +// acquisition delay and saves the entire job patch batch in scale latency. The +// scale decision itself is unaffected either way: it is derived from +// msg.Statistics, a snapshot taken by the service when the message was built, +// and jobs acquired now are reported as assigned in a later message. func (w *Scaler) Scale(ctx context.Context, msg *scaleset.RunnerScaleSetMessage) error { if msg == nil { // The long poll timed out without any activity. There is nothing to @@ -212,62 +361,37 @@ func (w *Scaler) Scale(ctx context.Context, msg *scaleset.RunnerScaleSetMessage) w.metrics.RecordStatistics(msg.Statistics) } - // Acquire first so the jobs are assigned as early as possible. Acquiring a - // job that is already acquired is a no-op, so a redelivered message repeats - // this safely. - if err := w.acquireAvailableJobs(ctx, msg.JobAvailableMessages); err != nil { - return err - } - if len(msg.JobStartedMessages) > 0 || len(msg.JobCompletedMessages) > 0 { w.dirty = true } - // Publish the desired count before anything else. It is the only patch that - // creates runners, so it is what new jobs actually wait on, while the job - // event patches below are bookkeeping. Sending it first also keeps it clear - // of the client rate limiter, which a large batch of event patches would - // otherwise drain ahead of it. - // - // Nothing in the batch has to land first for this to be safe. The runner set - // controller only deletes idle runners under Spec.PatchID == 0, and - // setDesiredWorkerState emits that only when the target is unchanged (or on - // the very first patch, before any target exists), never when the target - // drops. A scale down therefore cannot reach the deletion path in the same - // patch that the job started events are racing. if msg.Statistics != nil { if err := w.patchDesiredRunnerCount(ctx, w.setDesiredWorkerState(msg.Statistics.TotalAssignedJobs)); err != nil { return err } } - // The job events touch distinct runners and carry no ordering between them, - // so they are patched concurrently rather than one round trip at a time. - // - // They must still all land before this returns. The listener acks the message - // the moment Scale succeeds, and nothing other than these patches ever writes - // Status.JobID, so a runner whose patch was dropped after the ack would look - // idle forever. A later message that settles back to MinRunners publishes - // patch ID 0, and that is the one patch the controller does act on to delete - // idle runners. - g, gctx := errgroup.WithContext(ctx) - g.SetLimit(w.workers) - - for _, jobStarted := range msg.JobStartedMessages { - g.Go(func() error { - w.metrics.RecordJobStarted(jobStarted) - return w.HandleJobStarted(gctx, jobStarted) - }) - } - + // A completed job has nothing to patch: the runner is torn down by the + // ephemeral runner controller once its pod exits, and the completion is + // already reflected in the desired count published above. It is handled + // inline because it costs no API call. for _, jobCompleted := range msg.JobCompletedMessages { - g.Go(func() error { - w.metrics.RecordJobCompleted(jobCompleted) - return w.HandleJobCompleted(gctx, jobCompleted) - }) + w.metrics.RecordJobCompleted(jobCompleted) + if err := w.HandleJobCompleted(ctx, jobCompleted); err != nil { + return err + } } - return g.Wait() + // Queued rather than awaited. The listener acks as soon as this returns, so + // these patches outlive the message, which is why the workers run on their + // own context and the queue is drained by Close rather than here. + w.jobs.push(msg.JobStartedMessages...) + + if depth := w.jobs.depth(); depth > 0 { + w.logger.Info("Job started events queued for patching", "depth", depth) + } + + return w.acquireAvailableJobs(ctx, msg.JobAvailableMessages) } // acquireAvailableJobs assigns every available job to this scale set. A job that @@ -301,8 +425,8 @@ func (w *Scaler) acquireAvailableJobs(ctx context.Context, jobsAvailable []*scal // It also transitions the phase to Running if the runner is not in a terminal state. // It returns an error if there is any issue with updating the job information. // -// It is called from a worker goroutine, once per job started event in a message, -// and only ever touches the runner named by its own event. +// It is called from a background worker rather than from Scale, once per job +// started event, and only ever touches the runner named by its own event. func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error { w.logger.Info("Updating job info for the runner", "runnerName", jobInfo.RunnerName, @@ -325,7 +449,7 @@ func (w *Scaler) HandleJobStarted(ctx context.Context, jobInfo *scaleset.JobStar func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStarted) error { // Fetch current EphemeralRunner to check phase and deletion status currentRunner := &v1alpha1.EphemeralRunner{} - err := w.clientset.RESTClient(). + err := w.jobClientset.RESTClient(). Get(). Prefix("apis", v1alpha1.GroupVersion.Group, v1alpha1.GroupVersion.Version). Namespace(w.config.EphemeralRunnerSetNamespace). @@ -389,7 +513,7 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart w.logger.Info("Updating ephemeral runner with merge patch", "json", string(mergePatch)) patchedStatus := &v1alpha1.EphemeralRunner{} - err = w.clientset.RESTClient(). + err = w.jobClientset.RESTClient(). Patch(types.MergePatchType). Prefix("apis", v1alpha1.GroupVersion.Group, v1alpha1.GroupVersion.Version). Namespace(w.config.EphemeralRunnerSetNamespace). @@ -421,7 +545,7 @@ func (w *Scaler) patchJobStarted(ctx context.Context, jobInfo *scaleset.JobStart // is nothing to patch here; the completion only has to be reflected in the // desired count, which Scale already derives from the message. // -// It is called from a worker goroutine, once per job completed event in a message. +// It is called inline from Scale, once per job completed event. func (w *Scaler) HandleJobCompleted(ctx context.Context, msg *scaleset.JobCompleted) error { w.logger.Info("Job completed", "runnerName", msg.RunnerName, @@ -475,7 +599,7 @@ func (w *Scaler) patchDesiredRunnerCount(ctx context.Context, patchID int) error w.logger.Info("Preparing EphemeralRunnerSet update", "json", string(mergePatch)) patchedEphemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{} - err = w.clientset.RESTClient(). + err = w.scaleClientset.RESTClient(). Patch(types.MergePatchType). Prefix("apis", v1alpha1.GroupVersion.Group, v1alpha1.GroupVersion.Version). Namespace(w.config.EphemeralRunnerSetNamespace). diff --git a/cmd/ghalistener/scaler/scaler_apiserver_test.go b/cmd/ghalistener/scaler/scaler_apiserver_test.go index 904a156a..30855285 100644 --- a/cmd/ghalistener/scaler/scaler_apiserver_test.go +++ b/cmd/ghalistener/scaler/scaler_apiserver_test.go @@ -119,13 +119,14 @@ func TestHandleJobStartedAgainstAPIServer(t *testing.T) { require.NoError(t, err) return &Scaler{ - clientset: clientset, - config: Config{EphemeralRunnerSetNamespace: namespace.Name}, - targetRunners: -1, - patchSeq: -1, - logger: discardLogger, - metrics: metrics.Discard, - workers: defaultWorkers, + scaleClientset: clientset, + jobClientset: clientset, + config: Config{EphemeralRunnerSetNamespace: namespace.Name}, + targetRunners: -1, + patchSeq: -1, + logger: discardLogger, + metrics: metrics.Discard, + workers: defaultWorkers, } } diff --git a/cmd/ghalistener/scaler/scaler_test.go b/cmd/ghalistener/scaler/scaler_test.go index 6477d407..e736f277 100644 --- a/cmd/ghalistener/scaler/scaler_test.go +++ b/cmd/ghalistener/scaler/scaler_test.go @@ -333,7 +333,8 @@ func newTestScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, onPatch ...fu require.NoError(t, err) return &Scaler{ - clientset: clientset, + scaleClientset: clientset, + jobClientset: clientset, config: Config{ EphemeralRunnerSetNamespace: runner.Namespace, }, @@ -771,7 +772,8 @@ func newRecordingScaler(t *testing.T, runner *v1alpha1.EphemeralRunner, notFound require.NoError(t, err) return &Scaler{ - clientset: clientset, + scaleClientset: clientset, + jobClientset: clientset, config: Config{ EphemeralRunnerSetNamespace: runner.Namespace, }, diff --git a/config/crd/bases/actions.github.com_autoscalinglisteners.yaml b/config/crd/bases/actions.github.com_autoscalinglisteners.yaml index 6551d057..57c87ce1 100644 --- a/config/crd/bases/actions.github.com_autoscalinglisteners.yaml +++ b/config/crd/bases/actions.github.com_autoscalinglisteners.yaml @@ -131,20 +131,61 @@ spec: pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used - by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches + job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object diff --git a/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml b/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml index 3f8922cf..290ebddc 100644 --- a/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml +++ b/config/crd/bases/actions.github.com_autoscalingrunnersets.yaml @@ -149,19 +149,60 @@ spec: description: ListenerConfig holds configuration for the ghalistener pod. properties: scaler: - description: ScalerConfig configures the Kubernetes client used by the ghalistener scaler. + description: |- + ScalerConfig configures the Kubernetes clients used by the ghalistener scaler. + + The scaler talks to the API server over two independent clients, because the + two kinds of traffic have very different shapes and only one of them is on the + critical path: + + - The scale client publishes the desired runner count. That is one patch per + message, and it is the only patch that creates runners, so it must never + queue behind anything. + - The job client patches job started events. That is two calls per event and + can be a whole batch at once, so it is the traffic that actually consumes + the rate limit. + + A single shared client lets a batch of job event patches drain the token + bucket ahead of the scale patch, which delays the only call new jobs are + waiting on. Splitting them keeps the scale patch clear of that backlog. properties: burst: + description: Burst is the burst limit of the client that patches job started events. minimum: 1 type: integer qps: + description: |- + QPS is the query per second limit of the client that patches job started + events. This is the bulk of the scaler's API traffic, at up to two calls + per job started event. + minimum: 1 + type: integer + scaleBurst: + description: |- + ScaleBurst is the burst limit of the client that publishes the desired + runner count. + minimum: 1 + type: integer + scaleQPS: + description: |- + ScaleQPS is the query per second limit of the client that publishes the + desired runner count. The scaler issues at most one such patch per scale + set message, so this only has to be large enough that the patch never + waits on a token; it is deliberately a small budget separate from QPS + rather than a share of it. minimum: 1 type: integer workers: description: |- - Workers is the number of job started and job completed events the scaler - handles concurrently within a single scale set message. The worker that - scales the EphemeralRunnerSet runs on top of these. + Workers is the number of job started events the scaler patches + concurrently. The events are drained from a background queue rather than + being tied to the message they arrived on, so this bounds how many job + patches are in flight at any moment, not how many a single message may + carry. + + Raising it past what QPS sustains does nothing, since the rate limiter + rather than the worker count is what bounds throughput. minimum: 1 type: integer type: object