Patch job started events off the scale decision path

The scaler issued every API call a message asked for before returning, and
the listener acks only once it returns. A full batch of job started events
is 50 of them, two calls each, so the next scale decision sat behind 100
calls of bookkeeping.

Those patches are not what new jobs wait on. Only the desired runner count
creates runners; Status.JobID is a hint the runner set consults when
choosing which idle runner to delete, and the Actions service rejects the
deletion of a runner whose job is still running either way.

Publish the desired count first, hand the job started events to a
background pool, and acquire after. Acquiring last costs the round trip of
the scale patch and saves the whole batch, and the scale decision is
unaffected: it comes from msg.Statistics, a snapshot the service took when
it built the message, so jobs acquired now are reported as assigned in a
later one.

Give the two kinds of traffic their own clients. Sharing one token bucket
is what let the job patches delay the scale patch, so splitting them is
what makes the reordering worth anything; backgrounding alone would just
move the same queue.

Measured against a 5ms API server and a 50ms Actions service, at a full
50 event batch:

  scale patch reaches the API server   1.971s -> 6ms
  listener loop                        2.02s  -> 107ms/msg

The loop no longer spends the rate limit inline, so its cost is now the
two service round trips rather than the batch size.

This does not raise throughput. Job patches still cost two calls per
event, so the job client sustains qps/2 job starts per second, and the
queue is bounded by the real job start rate rather than by how fast the
listener polls: a faster loop polls more often and carries proportionally
fewer events. Measured at qps 40, the queue stays empty through 18
starts/sec and degrades gradually past 20 rather than falling over.

Close drains the queue, since the message these patches came from was
acked long before they run.

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
Nikola Jokic
2026-09-18 15:05:54 +02:00
co-authored by Copilot App
parent 8148c24c37
commit 2877779b72
17 changed files with 928 additions and 155 deletions
+13
View File
@@ -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)
+98
View File
@@ -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)
}
+199 -17
View File
@@ -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) {
+190 -66
View File
@@ -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).
@@ -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,
}
}
+4 -2
View File
@@ -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,
},