mirror of
https://github.com/actions-runner-controller/actions-runner-controller.git
synced 2026-10-06 12:11:16 +02:00
Scale down runners whose pod never reports the runner container (#4696)
This commit is contained in:
@@ -68,18 +68,18 @@ const (
|
||||
runnerBatchConcurrency = 8
|
||||
)
|
||||
|
||||
// unrecordedRunnerIDGracePeriod is how long cleanup waits for a runner to
|
||||
// record its runner ID before deleting it without one.
|
||||
// unrecordedRunnerIDGracePeriod is how long cleanup and scale down wait for a
|
||||
// runner to record its runner ID before deleting it without one.
|
||||
//
|
||||
// The runner controller records the ID only after it has created the pod,
|
||||
// so for a moment a runner can be executing a job while its status still
|
||||
// says 0, and 0 is not a registration the service can be asked about.
|
||||
// Cleanup leaves such a runner until the ID is recorded, which updates the
|
||||
// runner and so reconciles the set again, and then treats it like any other.
|
||||
// A runner that goes on without one usually cannot register or cannot start
|
||||
// its pod, and waiting on it forever would hold up the cleanup behind it. It
|
||||
// is deleted instead, and finalizing it asks the service before a live pod
|
||||
// goes.
|
||||
// The runner controller records the ID only once the pod reports the runner
|
||||
// container, so for a moment a runner can be executing a job while its status
|
||||
// still says 0, and 0 is not a registration the service can be asked about.
|
||||
// Cleanup and scale down leave such a runner until the ID is recorded, which
|
||||
// updates the runner and so reconciles the set again, and then treat it like
|
||||
// any other. A runner that goes on without one usually cannot register or its
|
||||
// pod cannot start, for instance because it cannot be scheduled, and waiting
|
||||
// on it forever would strand its pod and registration. It is deleted instead,
|
||||
// and finalizing it asks the service before a live pod goes.
|
||||
var unrecordedRunnerIDGracePeriod = time.Minute
|
||||
|
||||
// EphemeralRunnerSetReconciler reconciles a EphemeralRunnerSet object
|
||||
@@ -309,6 +309,7 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
|
||||
}
|
||||
|
||||
total := ephemeralRunnersByState.scaleTotal()
|
||||
var requeueAfter time.Duration
|
||||
if ephemeralRunnerSet.Spec.PatchID == 0 || ephemeralRunnerSet.Spec.PatchID != ephemeralRunnersByState.latestPatchID {
|
||||
// Spec.Replicas is the count the listener asked for when it published
|
||||
// Spec.PatchID. Deleting finished runners here changes the live count that
|
||||
@@ -367,21 +368,23 @@ func (r *EphemeralRunnerSetReconciler) Reconcile(ctx context.Context, req ctrl.R
|
||||
case ephemeralRunnerSet.Spec.PatchID == 0 && total > ephemeralRunnerSet.Spec.Replicas:
|
||||
count := total - ephemeralRunnerSet.Spec.Replicas
|
||||
log.Info("Deleting ephemeral runners (scale down)", "count", count)
|
||||
if err := r.deleteIdleEphemeralRunners(
|
||||
wait, err := r.deleteIdleEphemeralRunners(
|
||||
ctx,
|
||||
&ephemeralRunnerSet,
|
||||
ephemeralRunnersByState.pending,
|
||||
ephemeralRunnersByState.running,
|
||||
count,
|
||||
log,
|
||||
); err != nil {
|
||||
)
|
||||
if err != nil {
|
||||
log.Error(err, "failed to delete idle runners")
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
requeueAfter = wait
|
||||
}
|
||||
}
|
||||
|
||||
return ctrl.Result{}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
|
||||
return ctrl.Result{RequeueAfter: requeueAfter}, r.updateStatus(ctx, &ephemeralRunnerSet, ephemeralRunnersByState, log)
|
||||
}
|
||||
|
||||
// patchAppliedActionableRevisionStatus brings status into line with the runner
|
||||
@@ -1066,32 +1069,39 @@ func (r *EphemeralRunnerSetReconciler) createProxySecret(ctx context.Context, ep
|
||||
}
|
||||
|
||||
// deleteIdleEphemeralRunners try to deletes `count` number of v1alpha1.EphemeralRunner resources in the cluster.
|
||||
// It will only delete `v1alpha1.EphemeralRunner` that has registered with Actions service
|
||||
// which has a `v1alpha1.EphemeralRunner.Status.RunnerId` set.
|
||||
// A runner with a `v1alpha1.EphemeralRunner.Status.RunnerId` is deleted once the Actions service removes it.
|
||||
// A runner without one waits for unrecordedRunnerIDGracePeriod and is then deleted with its
|
||||
// registration finalizer kept, so finalizing it asks the service before a live pod goes.
|
||||
// It never records one when its pod cannot start, and no update notifies the set of that,
|
||||
// so requeueAfter is when the first waiting runner is due.
|
||||
// So, it is possible that this function will not delete enough ephemeral runners
|
||||
// if there are not enough ephemeral runners that have registered with Actions service.
|
||||
// if there are not enough ephemeral runners that can be removed at this time.
|
||||
// When this happens, the next reconcile loop will try to delete the remaining ephemeral runners
|
||||
// after we get notified by any of the `v1alpha1.EphemeralRunner.Status` updates.
|
||||
func (r *EphemeralRunnerSetReconciler) deleteIdleEphemeralRunners(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, pendingEphemeralRunners, runningEphemeralRunners []*v1alpha1.EphemeralRunner, count int, log logr.Logger) error {
|
||||
func (r *EphemeralRunnerSetReconciler) deleteIdleEphemeralRunners(ctx context.Context, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, pendingEphemeralRunners, runningEphemeralRunners []*v1alpha1.EphemeralRunner, count int, log logr.Logger) (requeueAfter time.Duration, err error) {
|
||||
if count <= 0 {
|
||||
return nil
|
||||
return 0, nil
|
||||
}
|
||||
runners := newEphemeralRunnerStepper(pendingEphemeralRunners, runningEphemeralRunners)
|
||||
if runners.len() == 0 {
|
||||
log.Info("No pending or running ephemeral runners running at this time for scale down")
|
||||
return nil
|
||||
return 0, nil
|
||||
}
|
||||
actionsClient, err := r.GetActionsService(ctx, ephemeralRunnerSet)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create actions client for ephemeral runner replica set: %w", err)
|
||||
return 0, fmt.Errorf("failed to create actions client for ephemeral runner replica set: %w", err)
|
||||
}
|
||||
var errs []error
|
||||
deletedCount := 0
|
||||
now := time.Now()
|
||||
for runners.next() {
|
||||
ephemeralRunner := runners.object()
|
||||
isDone := ephemeralRunner.IsDone()
|
||||
if !isDone && ephemeralRunner.Status.RunnerID == 0 {
|
||||
log.Info("Skipping ephemeral runner since it is not registered yet", "name", ephemeralRunner.Name)
|
||||
if wait := unrecordedRunnerIDWait(ephemeralRunner, now); !isDone && wait > 0 {
|
||||
log.Info("Skipping ephemeral runner since its runner ID is not recorded yet", "name", ephemeralRunner.Name, "retryAfter", wait)
|
||||
if requeueAfter == 0 || wait < requeueAfter {
|
||||
requeueAfter = wait
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -1116,11 +1126,11 @@ func (r *EphemeralRunnerSetReconciler) deleteIdleEphemeralRunners(ctx context.Co
|
||||
|
||||
deletedCount++
|
||||
if deletedCount == count {
|
||||
break
|
||||
return 0, multierr.Combine(errs...)
|
||||
}
|
||||
}
|
||||
|
||||
return multierr.Combine(errs...)
|
||||
return requeueAfter, multierr.Combine(errs...)
|
||||
}
|
||||
|
||||
func (r *EphemeralRunnerSetReconciler) deleteEphemeralRunnerWithActionsClient(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, actionsClient multiclient.Client, log logr.Logger) (bool, error) {
|
||||
|
||||
@@ -929,6 +929,18 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
|
||||
ephemeralRunnerSetTestInterval,
|
||||
).Should(BeNil(), "1 EphemeralRunner should be in Pending and 1 in Running phase")
|
||||
|
||||
// Let's say ephemeral runner controller patched these ephemeral runners with the registration.
|
||||
|
||||
updatedRunner = runnerList.Items[0].DeepCopy()
|
||||
updatedRunner.Status.RunnerID = 1
|
||||
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[0]))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner")
|
||||
|
||||
updatedRunner = runnerList.Items[1].DeepCopy()
|
||||
updatedRunner.Status.RunnerID = 2
|
||||
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[1]))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner")
|
||||
|
||||
// Scale down to 0 with patch ID 0. This forces the scale down to self correct on empty batch
|
||||
|
||||
ers = new(v1alpha1.EphemeralRunnerSet)
|
||||
@@ -942,31 +954,6 @@ var _ = Describe("Test EphemeralRunnerSet controller", func() {
|
||||
err = k8sClient.Patch(ctx, updated, client.MergeFrom(ers))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunnerSet")
|
||||
|
||||
runnerList = new(v1alpha1.EphemeralRunnerList)
|
||||
Consistently(
|
||||
func() (int, error) {
|
||||
if err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace); err != nil {
|
||||
return -1, err
|
||||
}
|
||||
|
||||
return len(runnerList.Items), nil
|
||||
},
|
||||
ephemeralRunnerSetTestTimeout,
|
||||
ephemeralRunnerSetTestInterval,
|
||||
).Should(BeEquivalentTo(2), "2 EphemeralRunner should be up since they don't have an ID yet")
|
||||
|
||||
// Now, let's say ephemeral runner controller patched these ephemeral runners with the registration.
|
||||
|
||||
updatedRunner = runnerList.Items[0].DeepCopy()
|
||||
updatedRunner.Status.RunnerID = 1
|
||||
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[0]))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner")
|
||||
|
||||
updatedRunner = runnerList.Items[1].DeepCopy()
|
||||
updatedRunner.Status.RunnerID = 2
|
||||
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runnerList.Items[1]))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner")
|
||||
|
||||
// Now, eventually, they should be deleted
|
||||
runnerList = new(v1alpha1.EphemeralRunnerList)
|
||||
Eventually(
|
||||
@@ -2024,40 +2011,21 @@ var _ = Describe("Test EphemeralRunnerSet controller with proxy settings", func(
|
||||
err = k8sClient.Patch(ctx, ephemeralRunnerSet, patch)
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to patch EphemeralRunnerSet")
|
||||
|
||||
// Set pods to PodSucceeded to simulate an actual EphemeralRunner stopping
|
||||
// The runner has not recorded a runner ID, and the suite does not wait for
|
||||
// it, so scale down retires it.
|
||||
Eventually(
|
||||
func(g Gomega) (int, error) {
|
||||
func() (int, error) {
|
||||
runnerList := new(v1alpha1.EphemeralRunnerList)
|
||||
err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace)
|
||||
if err != nil {
|
||||
return -1, err
|
||||
}
|
||||
|
||||
// Set status to simulate a configured EphemeralRunner
|
||||
refetch := false
|
||||
for i, runner := range runnerList.Items {
|
||||
if runner.Status.RunnerID == 0 {
|
||||
updatedRunner := runner.DeepCopy()
|
||||
updatedRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseSucceeded
|
||||
updatedRunner.Status.RunnerID = i + 100
|
||||
err = k8sClient.Status().Patch(ctx, updatedRunner, client.MergeFrom(&runner))
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to update EphemeralRunner")
|
||||
refetch = true
|
||||
}
|
||||
}
|
||||
|
||||
if refetch {
|
||||
err := listEphemeralRunnersAndRemoveFinalizers(ctx, k8sClient, runnerList, ephemeralRunnerSet.Namespace)
|
||||
if err != nil {
|
||||
return -1, err
|
||||
}
|
||||
}
|
||||
|
||||
return len(runnerList.Items), nil
|
||||
},
|
||||
ephemeralRunnerSetTestTimeout,
|
||||
ephemeralRunnerSetTestInterval,
|
||||
).Should(BeEquivalentTo(1), "1 EphemeralRunner should exist")
|
||||
).Should(BeEquivalentTo(0), "EphemeralRunner should be retired by the scale down")
|
||||
|
||||
// Delete the EphemeralRunnerSet
|
||||
err = k8sClient.Delete(ctx, ephemeralRunnerSet)
|
||||
@@ -2377,6 +2345,8 @@ var _ = Describe("Test EphemeralRunnerSet actionable revision cleanup", func() {
|
||||
ephemeralRunnerSet := &v1alpha1.EphemeralRunnerSet{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "test-actionable-revision-initial", Namespace: autoscalingNS.Name},
|
||||
Spec: v1alpha1.EphemeralRunnerSetSpec{
|
||||
// Wants the runner, so that only cleanup can remove it.
|
||||
Replicas: 1,
|
||||
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/owner/repo",
|
||||
GitHubConfigSecret: configSecret.Name,
|
||||
|
||||
@@ -0,0 +1,275 @@
|
||||
package actionsgithubcom
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
scalefake "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake"
|
||||
"github.com/actions/actions-runner-controller/controllers/actions.github.com/secretresolver"
|
||||
"github.com/actions/scaleset"
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
logf "sigs.k8s.io/controller-runtime/pkg/log"
|
||||
)
|
||||
|
||||
// The set runs under a manager, so it is woken by real watch events. The runner
|
||||
// controller is stepped by hand, which is what leaves a pod in the state under
|
||||
// test: envtest runs no scheduler and no kubelet, so the pod is exactly as the
|
||||
// spec left it. The Actions service is modeled, and answers for registration 7.
|
||||
var _ = Describe("Test EphemeralRunnerSet scale down of a runner without a recorded runner ID", func() {
|
||||
const registeredRunnerID = 7
|
||||
|
||||
var (
|
||||
ctx context.Context
|
||||
ns *corev1.Namespace
|
||||
runnerSet *v1alpha1.EphemeralRunnerSet
|
||||
runnerController *EphemeralRunnerReconciler
|
||||
runnerKey types.NamespacedName
|
||||
|
||||
mu sync.Mutex
|
||||
removeReply error
|
||||
removals []int64
|
||||
)
|
||||
|
||||
setRemoveReply := func(err error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
removeReply = err
|
||||
}
|
||||
removedRunners := func() []int64 {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return append([]int64(nil), removals...)
|
||||
}
|
||||
|
||||
reconcileRunner := func() ctrl.Result {
|
||||
result, err := runnerController.Reconcile(ctx, ctrl.Request{NamespacedName: runnerKey})
|
||||
ExpectWithOffset(1, err).NotTo(HaveOccurred())
|
||||
return result
|
||||
}
|
||||
getRunner := func() (*v1alpha1.EphemeralRunner, error) {
|
||||
runner := new(v1alpha1.EphemeralRunner)
|
||||
return runner, k8sClient.Get(ctx, runnerKey, runner)
|
||||
}
|
||||
getPod := func() (*corev1.Pod, error) {
|
||||
pod := new(corev1.Pod)
|
||||
return pod, k8sClient.Get(ctx, runnerKey, pod)
|
||||
}
|
||||
setPodStatus := func(status corev1.PodStatus) {
|
||||
pod, err := getPod()
|
||||
ExpectWithOffset(1, err).NotTo(HaveOccurred())
|
||||
pod.Status = status
|
||||
ExpectWithOffset(1, k8sClient.Status().Update(ctx, pod)).To(Succeed())
|
||||
}
|
||||
runnerContainerState := func(state corev1.ContainerState) corev1.PodStatus {
|
||||
return corev1.PodStatus{
|
||||
Phase: corev1.PodRunning,
|
||||
ContainerStatuses: []corev1.ContainerStatus{{
|
||||
Name: v1alpha1.EphemeralRunnerContainerName,
|
||||
State: state,
|
||||
}},
|
||||
}
|
||||
}
|
||||
scaleToZero := func() {
|
||||
current := new(v1alpha1.EphemeralRunnerSet)
|
||||
ExpectWithOffset(1, k8sClient.Get(ctx, client.ObjectKeyFromObject(runnerSet), current)).To(Succeed())
|
||||
updated := current.DeepCopy()
|
||||
updated.Spec.Replicas = 0
|
||||
updated.Spec.PatchID = 0
|
||||
ExpectWithOffset(1, k8sClient.Patch(ctx, updated, client.MergeFrom(current))).To(Succeed())
|
||||
}
|
||||
// The runner controller is left alone until the set has deleted the runner,
|
||||
// because it would put back the registration finalizer that a scale down of a
|
||||
// registered runner drops before deleting it.
|
||||
expectRunnerRetired := func() {
|
||||
EventuallyWithOffset(1, func() (bool, error) {
|
||||
runner, err := getRunner()
|
||||
if kerrors.IsNotFound(err) {
|
||||
return true, nil
|
||||
}
|
||||
return err == nil && !runner.DeletionTimestamp.IsZero(), err
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "scale down did not retire the runner")
|
||||
}
|
||||
expectRunnerAndPodGone := func() {
|
||||
EventuallyWithOffset(1, func() bool {
|
||||
reconcileRunner()
|
||||
_, err := getRunner()
|
||||
return kerrors.IsNotFound(err)
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "the runner was not finalized")
|
||||
|
||||
// An unscheduled pod goes as soon as it is deleted.
|
||||
pod, err := getPod()
|
||||
if !kerrors.IsNotFound(err) {
|
||||
ExpectWithOffset(1, err).NotTo(HaveOccurred())
|
||||
ExpectWithOffset(1, pod.DeletionTimestamp.IsZero()).To(BeFalse(), "the runner pod was left behind")
|
||||
}
|
||||
ExpectWithOffset(1, kerrors.IsNotFound(k8sClient.Get(ctx, runnerKey, new(corev1.Secret)))).To(BeTrue(), "the jitconfig secret was left behind")
|
||||
}
|
||||
|
||||
BeforeEach(func() {
|
||||
ctx = context.Background()
|
||||
removeReply, removals = nil, nil
|
||||
|
||||
var mgr ctrl.Manager
|
||||
ns, mgr = createNamespace(GinkgoT(), k8sClient)
|
||||
configSecret := createDefaultSecret(GinkgoT(), k8sClient, ns.Name)
|
||||
|
||||
service := scalefake.NewClient(
|
||||
scalefake.WithGenerateJitRunnerConfig(&scaleset.RunnerScaleSetJitRunnerConfig{
|
||||
Runner: &scaleset.RunnerReference{ID: registeredRunnerID, Name: "runner", RunnerScaleSetID: 100},
|
||||
EncodedJITConfig: "jit",
|
||||
}, nil),
|
||||
scalefake.WithRemoveRunnerFunc(func(_ context.Context, id int64) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
removals = append(removals, id)
|
||||
if id != registeredRunnerID {
|
||||
return fmt.Errorf("%w: runner %d", scaleset.NotFoundError, id)
|
||||
}
|
||||
return removeReply
|
||||
}),
|
||||
)
|
||||
multiClient := scalefake.NewMultiClient(scalefake.WithClient(service))
|
||||
|
||||
setController := &EphemeralRunnerSetReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Scheme: mgr.GetScheme(),
|
||||
Log: logf.Log,
|
||||
ResourceBuilder: ResourceBuilder{
|
||||
ResourceCache: newTestResourceCache(),
|
||||
SecretResolver: secretresolver.New(mgr.GetClient(), multiClient),
|
||||
},
|
||||
}
|
||||
Expect(setController.SetupWithManager(mgr)).To(Succeed())
|
||||
|
||||
runnerController = &EphemeralRunnerReconciler{
|
||||
Client: k8sClient,
|
||||
APIReader: k8sClient,
|
||||
Scheme: mgr.GetScheme(),
|
||||
Log: logf.Log,
|
||||
// The workers are never started, so nothing leaves the queue.
|
||||
UnregistrationQueue: NewRunnerUnregistrationQueue(logf.Log, nil, 0),
|
||||
ResourceBuilder: ResourceBuilder{
|
||||
Scheme: mgr.GetScheme(),
|
||||
ResourceCache: newTestResourceCache(),
|
||||
SecretResolver: secretresolver.New(k8sClient, multiClient),
|
||||
},
|
||||
}
|
||||
|
||||
runnerSet = &v1alpha1.EphemeralRunnerSet{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "test-ers", Namespace: ns.Name},
|
||||
Spec: v1alpha1.EphemeralRunnerSetSpec{
|
||||
Replicas: 1,
|
||||
PatchID: 1,
|
||||
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/owner/repo",
|
||||
GitHubConfigSecret: configSecret.Name,
|
||||
RunnerScaleSetID: 100,
|
||||
PodTemplateSpec: corev1.PodTemplateSpec{Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName, Image: runnerImage}},
|
||||
}},
|
||||
},
|
||||
},
|
||||
}
|
||||
Expect(k8sClient.Create(ctx, runnerSet)).To(Succeed())
|
||||
startManagers(GinkgoT(), mgr)
|
||||
|
||||
Eventually(func() (int, error) {
|
||||
var runners v1alpha1.EphemeralRunnerList
|
||||
if err := k8sClient.List(ctx, &runners, client.InNamespace(ns.Name)); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if len(runners.Items) == 1 {
|
||||
runnerKey = client.ObjectKeyFromObject(&runners.Items[0])
|
||||
}
|
||||
return len(runners.Items), nil
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(Equal(1), "the set did not create its runner")
|
||||
|
||||
// Registers the runner and creates the pod.
|
||||
reconcileRunner()
|
||||
_, err := getPod()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
runner, err := getRunner()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(runner.Status.RunnerID).To(BeZero())
|
||||
})
|
||||
|
||||
It("retires a runner whose pod cannot be scheduled", func() {
|
||||
setPodStatus(corev1.PodStatus{
|
||||
Phase: corev1.PodPending,
|
||||
Conditions: []corev1.PodCondition{{
|
||||
Type: corev1.PodScheduled,
|
||||
Status: corev1.ConditionFalse,
|
||||
Reason: corev1.PodReasonUnschedulable,
|
||||
Message: "0/1 nodes are available",
|
||||
}},
|
||||
})
|
||||
reconcileRunner()
|
||||
runner, err := getRunner()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(runner.Status.RunnerID).To(BeZero(), "a pod without a runner container status must not publish the runner ID")
|
||||
|
||||
scaleToZero()
|
||||
|
||||
expectRunnerRetired()
|
||||
expectRunnerAndPodGone()
|
||||
Expect(removedRunners()).To(Equal([]int64{registeredRunnerID}), "the registration was not removed before the pod")
|
||||
})
|
||||
|
||||
It("retires a runner whose pod is waiting on its container", func() {
|
||||
setPodStatus(runnerContainerState(corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{Reason: "ContainerCreating"}}))
|
||||
reconcileRunner()
|
||||
runner, err := getRunner()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(runner.Status.RunnerID).To(Equal(registeredRunnerID))
|
||||
|
||||
scaleToZero()
|
||||
|
||||
expectRunnerRetired()
|
||||
expectRunnerAndPodGone()
|
||||
Expect(removedRunners()).To(Equal([]int64{registeredRunnerID}))
|
||||
})
|
||||
|
||||
It("keeps the pod of a busy runner that has not recorded its runner ID", func() {
|
||||
// The container is running and the service says it is executing a job,
|
||||
// but the runner controller has not caught up: the runner still says 0 and
|
||||
// has no job.
|
||||
setPodStatus(runnerContainerState(corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}))
|
||||
setRemoveReply(fmt.Errorf("%w: %w", scaleset.ConflictError, scaleset.JobStillRunningError))
|
||||
runner, err := getRunner()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(runner.Status.RunnerID).To(BeZero())
|
||||
Expect(runner.HasJob()).To(BeFalse())
|
||||
|
||||
scaleToZero()
|
||||
|
||||
expectRunnerRetired()
|
||||
Expect(removedRunners()).To(BeEmpty(), "the set asked the service about a runner without a recorded ID")
|
||||
|
||||
// Finalizing asks the service about the registration in the jitconfig
|
||||
// secret before the live pod goes, and is told to keep it.
|
||||
for range 3 {
|
||||
Expect(reconcileRunner().RequeueAfter).To(Equal(busyRunnerRequeueInterval))
|
||||
runner, err = getRunner()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(runner.Finalizers).To(ContainElement(ephemeralRunnerActionsFinalizerName))
|
||||
pod, err := getPod()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(pod.DeletionTimestamp.IsZero()).To(BeTrue(), "the pod of a runner executing a job was deleted")
|
||||
}
|
||||
Expect(removedRunners()).To(HaveEach(int64(registeredRunnerID)))
|
||||
Expect(removedRunners()).NotTo(BeEmpty())
|
||||
|
||||
// The job finished and the service let go of the runner.
|
||||
setRemoveReply(nil)
|
||||
expectRunnerAndPodGone()
|
||||
})
|
||||
})
|
||||
@@ -446,3 +446,204 @@ func TestRunnerFinalizerChecksContainerStatesInTerminalPods(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// scaleToZero is the listener publishing an idle set with no runners wanted.
|
||||
func (f *unrecordedRunnerIDFixture) scaleToZero() {
|
||||
set := new(v1alpha1.EphemeralRunnerSet)
|
||||
require.NoError(f.t, f.c.Get(f.t.Context(), client.ObjectKeyFromObject(f.set), set))
|
||||
set.Spec.Replicas = 0
|
||||
set.Spec.PatchID = 0
|
||||
require.NoError(f.t, f.c.Update(f.t.Context(), set))
|
||||
}
|
||||
|
||||
// makePodUnschedulable leaves the pod as the scheduler does when no node fits
|
||||
// it: pending, with no container status for the runner controller to publish
|
||||
// the runner ID with.
|
||||
func (f *unrecordedRunnerIDFixture) makePodUnschedulable() {
|
||||
pod := f.pod()
|
||||
pod.Status.Phase = corev1.PodPending
|
||||
pod.Status.ContainerStatuses = nil
|
||||
pod.Status.Conditions = []corev1.PodCondition{{
|
||||
Type: corev1.PodScheduled,
|
||||
Status: corev1.ConditionFalse,
|
||||
Reason: corev1.PodReasonUnschedulable,
|
||||
}}
|
||||
require.NoError(f.t, f.c.Status().Update(f.t.Context(), pod))
|
||||
|
||||
_, err := f.reconcileRunner()
|
||||
require.NoError(f.t, err)
|
||||
require.Zero(f.t, f.runner().Status.RunnerID, "a pod without a runner container status must not publish the runner ID")
|
||||
}
|
||||
|
||||
func TestSetScaleDownRetiresRunnerWhosePodCannotStart(t *testing.T) {
|
||||
t.Run("unschedulable pod never records the runner ID", func(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 2*time.Minute)
|
||||
f.reply = nil
|
||||
f.makePodUnschedulable()
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
|
||||
require.Empty(t, f.removals, "a runner without a recorded ID must not be asked about")
|
||||
runner := f.runner()
|
||||
require.NotNil(t, runner)
|
||||
require.False(t, runner.DeletionTimestamp.IsZero(), "scale down left a runner that can never start")
|
||||
require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName)
|
||||
|
||||
// Finalizing asks about the registration in the jitconfig secret before
|
||||
// the pod goes.
|
||||
_, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals)
|
||||
require.Nil(t, f.pod())
|
||||
require.Nil(t, f.runner())
|
||||
require.Empty(t, f.queue.queued(), "the registration was already removed")
|
||||
})
|
||||
|
||||
t.Run("pod waiting on its container records the runner ID", func(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 2*time.Minute)
|
||||
f.reply = nil
|
||||
pod := f.pod()
|
||||
pod.Status.Phase = corev1.PodPending
|
||||
pod.Status.ContainerStatuses = []corev1.ContainerStatus{{
|
||||
Name: v1alpha1.EphemeralRunnerContainerName,
|
||||
State: corev1.ContainerState{Waiting: &corev1.ContainerStateWaiting{Reason: "ContainerCreating"}},
|
||||
}}
|
||||
require.NoError(t, f.c.Status().Update(t.Context(), pod))
|
||||
_, err := f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, unrecordedTestRunnerID, f.runner().Status.RunnerID)
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
|
||||
require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals)
|
||||
runner := f.runner()
|
||||
require.NotNil(t, runner)
|
||||
require.False(t, runner.DeletionTimestamp.IsZero())
|
||||
require.NotContains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName, "the registration was already removed")
|
||||
|
||||
_, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals)
|
||||
require.Nil(t, f.pod())
|
||||
require.Nil(t, f.runner())
|
||||
require.Empty(t, f.queue.queued())
|
||||
})
|
||||
}
|
||||
|
||||
func TestSetScaleDownWaitsForRunnerToRecordItsIDBeforeRetiringIt(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 0)
|
||||
f.reply = nil
|
||||
f.makePodUnschedulable()
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Positive(t, result.RequeueAfter, "nothing else wakes the set for a runner that never records its ID")
|
||||
require.LessOrEqual(t, result.RequeueAfter, unrecordedRunnerIDGracePeriod)
|
||||
require.Empty(t, f.removals)
|
||||
require.True(t, f.runner().DeletionTimestamp.IsZero(), "a runner still within its grace period was retired")
|
||||
f.requirePodKept()
|
||||
|
||||
unrecordedRunnerIDGracePeriod = 0
|
||||
|
||||
result, err = f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
require.False(t, f.runner().DeletionTimestamp.IsZero())
|
||||
|
||||
_, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, f.pod())
|
||||
require.Nil(t, f.runner())
|
||||
}
|
||||
|
||||
func TestSetScaleDownKeepsBusyRunnerWithoutRecordedID(t *testing.T) {
|
||||
for _, phase := range []v1alpha1.EphemeralRunnerPhase{"", v1alpha1.EphemeralRunnerPhaseRunning} {
|
||||
t.Run("phase="+string(phase), func(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 2*time.Minute)
|
||||
if phase != "" {
|
||||
runner := f.runner()
|
||||
runner.Status.Phase = phase
|
||||
require.NoError(t, f.c.Status().Update(t.Context(), runner))
|
||||
}
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
require.Empty(t, f.removals, "a runner without a recorded ID must not be asked about")
|
||||
runner := f.runner()
|
||||
require.NotNil(t, runner)
|
||||
require.Zero(t, runner.Status.RunnerID)
|
||||
require.False(t, runner.HasJob())
|
||||
require.Contains(t, runner.Finalizers, ephemeralRunnerActionsFinalizerName)
|
||||
f.requirePodKept()
|
||||
|
||||
// The service still reports the registration busy, so finalizing
|
||||
// keeps the live pod and asks again later.
|
||||
result, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, busyRunnerRequeueInterval, result.RequeueAfter)
|
||||
require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals)
|
||||
require.NotNil(t, f.runner())
|
||||
f.requirePodKept()
|
||||
require.Empty(t, f.queue.queued())
|
||||
|
||||
// The job finished and the service let go of the runner.
|
||||
f.reply = nil
|
||||
_, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, f.pod())
|
||||
require.Nil(t, f.runner())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestSetScaleDownKeepsBusyRunnerUntilItRecordsItsID(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 0)
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Positive(t, result.RequeueAfter)
|
||||
require.Empty(t, f.removals)
|
||||
require.True(t, f.runner().DeletionTimestamp.IsZero())
|
||||
f.requirePodKept()
|
||||
|
||||
_, err = f.reconcileRunner()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, unrecordedTestRunnerID, f.runner().Status.RunnerID)
|
||||
|
||||
// With the ID recorded the set asks the service itself, and a runner it
|
||||
// reports busy is neither deleted nor left in the deleting state.
|
||||
result, err = f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
require.Equal(t, []int64{unrecordedTestRunnerID}, f.removals)
|
||||
require.True(t, f.runner().DeletionTimestamp.IsZero(), "a runner executing a job was retired")
|
||||
f.requirePodKept()
|
||||
}
|
||||
|
||||
func TestSetScaleDownKeepsRunnerWithReportedJobAndNoRecordedID(t *testing.T) {
|
||||
f := newUnrecordedRunnerIDFixture(t, 2*time.Minute)
|
||||
runner := f.runner()
|
||||
runner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
|
||||
runner.Status.JobID = "job-1"
|
||||
require.NoError(t, f.c.Status().Update(t.Context(), runner))
|
||||
f.scaleToZero()
|
||||
|
||||
result, err := f.reconcileSet()
|
||||
require.NoError(t, err)
|
||||
require.Zero(t, result.RequeueAfter)
|
||||
require.Empty(t, f.removals)
|
||||
runner = f.runner()
|
||||
require.NotNil(t, runner)
|
||||
require.True(t, runner.DeletionTimestamp.IsZero(), "a runner with a reported job was retired")
|
||||
f.requirePodKept()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user