mirror of
https://github.com/actions-runner-controller/actions-runner-controller.git
synced 2026-10-02 09:51:30 +02:00
Publish the desired runner count before the job event patches
Scale used to hold the replica patch back whenever the target dropped, so that every job started patch had landed before the runner set controller could act on a lower count. The reasoning was that deleteIdleEphemeralRunners skips a runner only once it carries a job ID, so a runner that had just picked up a job could otherwise be deleted. That guard was unreachable. The controller only deletes idle runners under Spec.PatchID == 0, and setDesiredWorkerState emits patch ID 0 only when the target is unchanged and equal to MinRunners, or on the very first patch, when no previous target exists. Neither can coincide with a falling target, so a scale down never reaches the deletion path. Exhaustively walking message sequences over every MinRunners/MaxRunners pair finds no state where the two occur together. So the replica patch has no reason to wait, and good reason to go first: it is the only patch that creates runners, and therefore the one new jobs wait on, while the job event patches are bookkeeping the controller reads later. 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. The job events still have to land before Scale returns. The listener acks the message the moment it does, and nothing other than these patches ever writes Status.JobID, so a patch dropped after the ack would leave a busy runner looking idle to the scale down that a later patch ID 0 permits. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
This commit is contained in:
co-authored by
Copilot App
parent
7833438dd7
commit
8148c24c37
@@ -303,78 +303,84 @@ func TestScale_WorkersBoundConcurrency(t *testing.T) {
|
||||
peak := server.peak
|
||||
server.mu.Unlock()
|
||||
|
||||
// workers event slots plus the scaling worker, which runs alongside them
|
||||
// because this message scales up.
|
||||
assert.LessOrEqual(t, peak, workers+1)
|
||||
assert.LessOrEqual(t, peak, workers, "the pool bound is a real limit")
|
||||
}
|
||||
|
||||
// TestScale_ScaleDownWaitsForJobStarted pins the ordering the runner set
|
||||
// controller depends on. It skips a runner during scale down only when that
|
||||
// runner already carries a job request ID, so a patch that lowers the replica
|
||||
// count must not be published while job started patches are still outstanding.
|
||||
func TestScale_ScaleDownWaitsForJobStarted(t *testing.T) {
|
||||
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
|
||||
|
||||
// Establish a target of 4 so the next message scales down.
|
||||
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
|
||||
MessageID: 1,
|
||||
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 4},
|
||||
}))
|
||||
require.Equal(t, 4, w.targetRunners)
|
||||
|
||||
server.mu.Lock()
|
||||
server.runnerSetPatches = nil
|
||||
server.mu.Unlock()
|
||||
|
||||
msg := &scaleset.RunnerScaleSetMessage{
|
||||
MessageID: 2,
|
||||
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 1},
|
||||
JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)},
|
||||
// TestScale_PublishesDesiredCountFirst pins the ordering that matters for
|
||||
// scale up latency: the replica patch is the one that creates runners, so it
|
||||
// goes out before any of the job event bookkeeping patches.
|
||||
//
|
||||
// It is asserted for a scale down as well as a scale up. Nothing requires the
|
||||
// job started patches to land first: the runner set controller only deletes
|
||||
// idle runners under Spec.PatchID == 0, which setDesiredWorkerState never emits
|
||||
// together with a falling target.
|
||||
func TestScale_PublishesDesiredCountFirst(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
settleAt int
|
||||
assigned int
|
||||
}{
|
||||
{name: "scale up", settleAt: 1, assigned: 5},
|
||||
{name: "steady", settleAt: 5, assigned: 5},
|
||||
{name: "scale down", settleAt: 5, assigned: 1},
|
||||
}
|
||||
|
||||
require.NoError(t, w.Scale(t.Context(), msg))
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
|
||||
|
||||
server.mu.Lock()
|
||||
defer server.mu.Unlock()
|
||||
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
|
||||
MessageID: 1,
|
||||
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: tt.settleAt},
|
||||
}))
|
||||
require.Equal(t, tt.settleAt, w.targetRunners)
|
||||
|
||||
assert.Equal(t, 1, w.targetRunners)
|
||||
require.Len(t, server.runnerSetPatches, 1)
|
||||
assert.Equal(t, 2, server.runnerPatchesBeforeRunnerSet,
|
||||
"the scale down patch is published only after every job started patch landed")
|
||||
server.mu.Lock()
|
||||
server.runnerSetPatches = nil
|
||||
server.runnerPatches = 0
|
||||
server.runnerPatchesBeforeRunnerSet = 0
|
||||
server.mu.Unlock()
|
||||
|
||||
require.NoError(t, w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
|
||||
MessageID: 2,
|
||||
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: tt.assigned},
|
||||
JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)},
|
||||
}))
|
||||
|
||||
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")
|
||||
assert.Equal(t, 2, server.runnerPatches)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestScale_ScaleUpRunsAlongsideJobStarted is the counterpart: a patch that
|
||||
// cannot delete anything is published without waiting for the event workers.
|
||||
func TestScale_ScaleUpRunsAlongsideJobStarted(t *testing.T) {
|
||||
w, server := newScaleScaler(t, &fakeAcquirer{}, defaultConfig(), defaultWorkers)
|
||||
// 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
|
||||
// a runner whose job started patch has not landed yet would look idle and be
|
||||
// eligible for deletion.
|
||||
func TestScale_ScaleDownNeverPublishesPatchIDZero(t *testing.T) {
|
||||
for _, minRunners := range []int{0, 1, 2} {
|
||||
config := defaultConfig()
|
||||
config.MinRunners = minRunners
|
||||
|
||||
// Block the ephemeral runner requests so the scale patch can only land first
|
||||
// if it genuinely does not wait for them.
|
||||
server.block()
|
||||
defer server.release()
|
||||
w, _ := newScaleScaler(t, &fakeAcquirer{}, config, defaultWorkers)
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
done <- w.Scale(t.Context(), &scaleset.RunnerScaleSetMessage{
|
||||
MessageID: 1,
|
||||
Statistics: &scaleset.RunnerScaleSetStatistic{TotalAssignedJobs: 5},
|
||||
JobStartedMessages: []*scaleset.JobStarted{jobStarted(0), jobStarted(1)},
|
||||
})
|
||||
}()
|
||||
|
||||
require.Eventually(t, func() bool {
|
||||
server.mu.Lock()
|
||||
defer server.mu.Unlock()
|
||||
return len(server.runnerSetPatches) == 1
|
||||
}, 10*time.Second, time.Millisecond, "scale up patch is published without waiting for job started patches")
|
||||
|
||||
server.release()
|
||||
require.NoError(t, <-done)
|
||||
|
||||
server.mu.Lock()
|
||||
defer server.mu.Unlock()
|
||||
assert.Equal(t, 0, server.runnerPatchesBeforeRunnerSet)
|
||||
assert.Equal(t, 2, server.runnerPatches)
|
||||
previous := -1
|
||||
for _, assigned := range []int{0, 3, 3, 1, 0, 0, 4, 2, 0} {
|
||||
patchID := w.setDesiredWorkerState(assigned)
|
||||
if previous >= 0 && w.targetRunners < previous {
|
||||
assert.NotEqual(t, 0, patchID,
|
||||
"minRunners=%d target %d->%d", minRunners, previous, w.targetRunners)
|
||||
}
|
||||
previous = w.targetRunners
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestScale_NilMessage covers the long poll timing out. The listener stopped
|
||||
|
||||
@@ -191,11 +191,11 @@ func (w *Scaler) applyDefaults() error {
|
||||
// redelivers it otherwise, so every step below is idempotent and safe to repeat
|
||||
// after a partially applied message.
|
||||
//
|
||||
// The work is split across workers: one patches the EphemeralRunnerSet with the
|
||||
// desired replica count, the rest patch the EphemeralRunner behind each job
|
||||
// started or job completed event. The events touch distinct resources and carry
|
||||
// no ordering between them, so they run concurrently instead of being replayed
|
||||
// one API call at a time.
|
||||
// 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.
|
||||
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
|
||||
@@ -223,37 +223,35 @@ func (w *Scaler) Scale(ctx context.Context, msg *scaleset.RunnerScaleSetMessage)
|
||||
w.dirty = true
|
||||
}
|
||||
|
||||
// The scale decision is computed up front, on the goroutine that owns the
|
||||
// scaler state, so the scaling worker never races the event workers for it.
|
||||
scaleRequested := msg.Statistics != nil
|
||||
var patchID int
|
||||
var scalesDown bool
|
||||
if scaleRequested {
|
||||
previousTarget := w.targetRunners
|
||||
patchID = w.setDesiredWorkerState(msg.Statistics.TotalAssignedJobs)
|
||||
scalesDown = previousTarget >= 0 && w.targetRunners < previousTarget
|
||||
// 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
|
||||
}
|
||||
}
|
||||
|
||||
// A patch that lowers the replica count can make the runner set controller
|
||||
// delete idle runners, and it only skips a runner that already carries a job
|
||||
// request ID. Publishing it before the job started patches land could
|
||||
// therefore offer up a runner that just picked up a job, so the scaling
|
||||
// worker waits for them in that case. A patch that scales up or holds cannot
|
||||
// delete anything, so it runs alongside the event workers.
|
||||
scaleConcurrently := scaleRequested && !scalesDown
|
||||
|
||||
// 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)
|
||||
limit := w.workers
|
||||
if scaleConcurrently {
|
||||
limit++ // the scaling worker gets a slot of its own
|
||||
}
|
||||
g.SetLimit(limit)
|
||||
|
||||
if scaleConcurrently {
|
||||
g.Go(func() error {
|
||||
return w.patchDesiredRunnerCount(gctx, patchID)
|
||||
})
|
||||
}
|
||||
g.SetLimit(w.workers)
|
||||
|
||||
for _, jobStarted := range msg.JobStartedMessages {
|
||||
g.Go(func() error {
|
||||
@@ -269,15 +267,7 @@ func (w *Scaler) Scale(ctx context.Context, msg *scaleset.RunnerScaleSetMessage)
|
||||
})
|
||||
}
|
||||
|
||||
if err := g.Wait(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if scaleRequested && !scaleConcurrently {
|
||||
return w.patchDesiredRunnerCount(ctx, patchID)
|
||||
}
|
||||
|
||||
return nil
|
||||
return g.Wait()
|
||||
}
|
||||
|
||||
// acquireAvailableJobs assigns every available job to this scale set. A job that
|
||||
|
||||
Reference in New Issue
Block a user