mirror of
https://github.com/actions-runner-controller/actions-runner-controller.git
synced 2026-09-30 14:01:32 +02:00
Hand a finished runner's pod back as soon as its job is over (#4673)
This commit is contained in:
@@ -59,6 +59,19 @@ type EphemeralRunnerReconciler struct {
|
||||
// SetupWithManager creates one and registers it with the manager.
|
||||
UnregistrationQueue *RunnerUnregistrationQueue
|
||||
|
||||
// TerminatedPodGracePeriodSeconds is the grace period used when deleting a
|
||||
// runner pod whose containers have all exited. It is zero by default, so
|
||||
// the pod leaves the API as soon as the delete is issued instead of sitting
|
||||
// in Terminating while the kubelet cleans up locally.
|
||||
//
|
||||
// Raise it to keep those pods around for longer, for example to give a log
|
||||
// collector time to read them. A negative value asks for no override at
|
||||
// all, leaving the deletion to the pod's own terminationGracePeriodSeconds.
|
||||
//
|
||||
// It is only ever applied to a pod with nothing left running in it. Pods
|
||||
// that are still alive are always deleted gracefully.
|
||||
TerminatedPodGracePeriodSeconds int64
|
||||
|
||||
ResourceBuilder
|
||||
}
|
||||
|
||||
@@ -527,7 +540,7 @@ func (r *EphemeralRunnerReconciler) cleanupResources(ctx context.Context, epheme
|
||||
case err == nil:
|
||||
if pod.DeletionTimestamp.IsZero() {
|
||||
log.Info("Deleting the runner pod")
|
||||
if err := r.Delete(ctx, pod); err != nil && !kerrors.IsNotFound(err) {
|
||||
if err := r.Delete(ctx, pod, r.deletePodOptions(pod)...); err != nil && !kerrors.IsNotFound(err) {
|
||||
return fmt.Errorf("failed to delete pod: %w", err)
|
||||
}
|
||||
log.Info("Deleted the runner pod")
|
||||
@@ -604,7 +617,7 @@ func (r *EphemeralRunnerReconciler) cleanupRunnerLinkedPods(ctx context.Context,
|
||||
}
|
||||
|
||||
log.Info("Deleting container hooks runner-linked pod", "name", linkedPod.Name)
|
||||
if err := r.Delete(ctx, linkedPod); err != nil && !kerrors.IsNotFound(err) {
|
||||
if err := r.Delete(ctx, linkedPod, r.deletePodOptions(linkedPod)...); err != nil && !kerrors.IsNotFound(err) {
|
||||
errs = append(errs, fmt.Errorf("failed to delete runner linked pod %q: %w", linkedPod.Name, err))
|
||||
}
|
||||
}
|
||||
@@ -759,7 +772,7 @@ func (r *EphemeralRunnerReconciler) markAsSucceeded(ctx context.Context, ephemer
|
||||
func (r *EphemeralRunnerReconciler) deletePodAsFailed(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, pod *corev1.Pod, log logr.Logger) error {
|
||||
if pod.DeletionTimestamp.IsZero() {
|
||||
log.Info("Deleting the ephemeral runner pod", "podId", pod.UID)
|
||||
if err := r.Delete(ctx, pod); err != nil && !kerrors.IsNotFound(err) {
|
||||
if err := r.Delete(ctx, pod, r.deletePodOptions(pod)...); err != nil && !kerrors.IsNotFound(err) {
|
||||
return fmt.Errorf("failed to delete pod with status failed: %w", err)
|
||||
}
|
||||
}
|
||||
@@ -1135,6 +1148,97 @@ func (r *EphemeralRunnerReconciler) SetupWithManager(mgr ctrl.Manager, opts ...O
|
||||
).Complete(r)
|
||||
}
|
||||
|
||||
// podTerminated reports whether every container in the pod has stopped.
|
||||
//
|
||||
// No container the kubelet has reported on may still be running. That check is
|
||||
// made whatever the pod phase says, because the phase is not always the
|
||||
// kubelet's account of the containers: a pod is moved to Failed by the control
|
||||
// plane when its node is lost or shut down, while the last status the kubelet
|
||||
// managed to send still shows a container running on the other side of the
|
||||
// partition. Believing the phase there would drop the pod out of the API while
|
||||
// something is still alive under it.
|
||||
//
|
||||
// The phase is what says whether the containers that have not reported are
|
||||
// still to come. A pod that has reached Succeeded or Failed is not going to
|
||||
// start anything else, so a container missing from the status is one that never
|
||||
// ran, which is how a pod whose init container failed is still terminated. Short
|
||||
// of a terminal phase every container has to have reported, or the runner that
|
||||
// is about to be reported as started would be missed.
|
||||
//
|
||||
// Native sidecars run as init containers that outlive the regular ones, so they
|
||||
// are checked too. A pod still running one of those, or a legacy sidecar
|
||||
// alongside the runner, is not terminated no matter what the runner container
|
||||
// did.
|
||||
func podTerminated(pod *corev1.Pod) bool {
|
||||
for i := range pod.Status.ContainerStatuses {
|
||||
if pod.Status.ContainerStatuses[i].State.Terminated == nil {
|
||||
return false
|
||||
}
|
||||
}
|
||||
for i := range pod.Status.InitContainerStatuses {
|
||||
if pod.Status.InitContainerStatuses[i].State.Terminated == nil {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
switch pod.Status.Phase {
|
||||
case corev1.PodSucceeded, corev1.PodFailed:
|
||||
return true
|
||||
}
|
||||
|
||||
return len(pod.Status.ContainerStatuses) == len(pod.Spec.Containers)
|
||||
}
|
||||
|
||||
// deletePodOptions asks for an immediate deletion of a pod that has nothing
|
||||
// left running in it.
|
||||
//
|
||||
// A graceful deletion exists to give containers their terminationGracePeriod to
|
||||
// shut down, and the API object survives until the kubelet reports that they
|
||||
// have. For a pod whose containers have all terminated there is nothing to
|
||||
// shut down and nothing to protect: the grace period is spent waiting on the
|
||||
// kubelet to finish unmounting volumes and tearing down the sandbox, which it
|
||||
// does whether or not the object is still there.
|
||||
//
|
||||
// That wait is what fills a cluster with Terminating runner pods during a burst
|
||||
// of jobs. They hold their name, their scheduling slot, and their share of any
|
||||
// ResourceQuota, so the runners waiting to replace them cannot start. Dropping
|
||||
// the object as the delete is issued hands those back immediately.
|
||||
//
|
||||
// How long to wait is TerminatedPodGracePeriodSeconds, zero by default. A
|
||||
// negative value leaves the deletion alone, which restores whatever the pod
|
||||
// asks for in its own spec.
|
||||
//
|
||||
// A pod that is still running is deleted normally. Skipping the grace period
|
||||
// there would drop the object while its containers were still alive, leaving
|
||||
// the kubelet to kill them with nothing in the API to account for the resources
|
||||
// they hold in the meantime.
|
||||
//
|
||||
// The deletion is pinned to the pod the decision was made about. Pods are read
|
||||
// through the informer cache and every generation of a runner's pod carries the
|
||||
// same name, so a delete by name is a delete of whatever holds that name when
|
||||
// the API server reads the request, not of the pod whose containers were
|
||||
// observed to have stopped. A single controller cannot get that wrong, since it
|
||||
// is the only thing creating that name and it only creates after a read says the
|
||||
// name is free, but that argument is worth exactly as much as the single writer
|
||||
// it assumes: during a leader election handover the outgoing leader's reconcile
|
||||
// is still in flight while the new leader is already replacing pods. Naming the
|
||||
// UID turns the delete into a conflict when it lands on a pod the controller
|
||||
// never looked at. Without the grace period there is nothing to catch it
|
||||
// afterwards: the pod would be gone the moment the request was accepted,
|
||||
// killing a job instead of handing its container a SIGTERM to deregister with.
|
||||
func (r *EphemeralRunnerReconciler) deletePodOptions(pod *corev1.Pod) []client.DeleteOption {
|
||||
if !podTerminated(pod) || r.TerminatedPodGracePeriodSeconds < 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
opts := []client.DeleteOption{client.GracePeriodSeconds(r.TerminatedPodGracePeriodSeconds)}
|
||||
if pod.UID != "" {
|
||||
uid := pod.UID
|
||||
opts = append(opts, client.Preconditions{UID: &uid})
|
||||
}
|
||||
return opts
|
||||
}
|
||||
|
||||
func runnerContainerStatus(pod *corev1.Pod) *corev1.ContainerStatus {
|
||||
for i := range pod.Status.ContainerStatuses {
|
||||
cs := &pod.Status.ContainerStatuses[i]
|
||||
|
||||
@@ -0,0 +1,295 @@
|
||||
package actionsgithubcom
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
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/runtime"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
ctrlfake "sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
|
||||
)
|
||||
|
||||
func terminatedRunnerPod(exitCode int32) *corev1.Pod {
|
||||
return &corev1.Pod{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "test-runner", Namespace: "default", UID: "d6f1e2b4-0f3a-4a1e-9a1f-2c9f0a3b7c55"},
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName}},
|
||||
},
|
||||
Status: corev1.PodStatus{
|
||||
Phase: corev1.PodRunning,
|
||||
ContainerStatuses: []corev1.ContainerStatus{
|
||||
{
|
||||
Name: v1alpha1.EphemeralRunnerContainerName,
|
||||
State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: exitCode}},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeletePodOptionsSkipsGracePeriodOnlyWhenNothingIsRunning pins which pods
|
||||
// are deleted without a grace period.
|
||||
//
|
||||
// A pod that is done holds its name, its place on the node and its share of any
|
||||
// ResourceQuota until the API server stops waiting on the kubelet, which is
|
||||
// what leaves a burst of finished jobs sitting in Terminating while the runners
|
||||
// replacing them have nowhere to start. Skipping the wait is only safe once
|
||||
// every container has stopped, so the cases below cover both the pod phase and
|
||||
// the container states the phase lags behind, including the sidecars that
|
||||
// outlive the runner container and the lost node whose phase says the pod
|
||||
// failed while the kubelet's last word was that the runner is still up.
|
||||
func TestDeletePodOptionsSkipsGracePeriodOnlyWhenNothingIsRunning(t *testing.T) {
|
||||
sidecarRunning := terminatedRunnerPod(0)
|
||||
sidecarRunning.Spec.Containers = append(sidecarRunning.Spec.Containers, corev1.Container{Name: "dind"})
|
||||
sidecarRunning.Status.ContainerStatuses = append(sidecarRunning.Status.ContainerStatuses, corev1.ContainerStatus{
|
||||
Name: "dind",
|
||||
State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}},
|
||||
})
|
||||
|
||||
nativeSidecarRunning := terminatedRunnerPod(0)
|
||||
nativeSidecarRunning.Status.InitContainerStatuses = []corev1.ContainerStatus{
|
||||
{Name: "dind", State: corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}},
|
||||
}
|
||||
|
||||
statusNotReportedYet := terminatedRunnerPod(0)
|
||||
statusNotReportedYet.Spec.Containers = append(statusNotReportedYet.Spec.Containers, corev1.Container{Name: "dind"})
|
||||
|
||||
runnerStillRunning := terminatedRunnerPod(0)
|
||||
runnerStillRunning.Status.ContainerStatuses[0].State = corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}
|
||||
|
||||
succeeded := terminatedRunnerPod(0)
|
||||
succeeded.Status.Phase = corev1.PodSucceeded
|
||||
|
||||
failed := terminatedRunnerPod(1)
|
||||
failed.Status.Phase = corev1.PodFailed
|
||||
|
||||
nodeLost := terminatedRunnerPod(0)
|
||||
nodeLost.Status.Phase = corev1.PodFailed
|
||||
nodeLost.Status.ContainerStatuses[0].State = corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}
|
||||
|
||||
failedBeforeStarting := terminatedRunnerPod(0)
|
||||
failedBeforeStarting.Status.Phase = corev1.PodFailed
|
||||
failedBeforeStarting.Status.ContainerStatuses = nil
|
||||
failedBeforeStarting.Status.InitContainerStatuses = []corev1.ContainerStatus{
|
||||
{Name: "init", State: corev1.ContainerState{Terminated: &corev1.ContainerStateTerminated{ExitCode: 1}}},
|
||||
}
|
||||
|
||||
tt := map[string]struct {
|
||||
pod *corev1.Pod
|
||||
immediate bool
|
||||
}{
|
||||
"pod succeeded": {pod: succeeded, immediate: true},
|
||||
"pod failed": {pod: failed, immediate: true},
|
||||
"runner exited but the phase has not moved": {pod: terminatedRunnerPod(0), immediate: true},
|
||||
"pod failed before its containers started": {pod: failedBeforeStarting, immediate: true},
|
||||
"runner is still running": {pod: runnerStillRunning},
|
||||
"sidecar is still running": {pod: sidecarRunning},
|
||||
"native sidecar is still running": {pod: nativeSidecarRunning},
|
||||
"a container has not reported a status yet": {pod: statusNotReportedYet},
|
||||
"the phase failed the pod but a container is still reported as running": {pod: nodeLost},
|
||||
}
|
||||
|
||||
for name, tc := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
r := &EphemeralRunnerReconciler{}
|
||||
opts := r.deletePodOptions(tc.pod)
|
||||
|
||||
var deleteOptions client.DeleteOptions
|
||||
for _, opt := range opts {
|
||||
opt.ApplyToDelete(&deleteOptions)
|
||||
}
|
||||
|
||||
if !tc.immediate {
|
||||
assert.Empty(t, opts)
|
||||
assert.Nil(t, deleteOptions.GracePeriodSeconds)
|
||||
return
|
||||
}
|
||||
|
||||
require.NotNil(t, deleteOptions.GracePeriodSeconds)
|
||||
assert.Equal(t, int64(0), *deleteOptions.GracePeriodSeconds)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeletePodOptionsHonorsTheConfiguredGracePeriod pins the escape hatch for
|
||||
// clusters that want finished pods to stick around.
|
||||
//
|
||||
// Removing a pod the moment its job is over is what hands the cluster back to
|
||||
// the runners waiting for it, but it also takes the pod away from anything that
|
||||
// reads it afterwards, such as a log collector scraping from the API. The grace
|
||||
// period is configurable for those, and a negative value hands the decision
|
||||
// back to the pod, which is how the deletion behaved before it was skipped.
|
||||
func TestDeletePodOptionsHonorsTheConfiguredGracePeriod(t *testing.T) {
|
||||
pod := terminatedRunnerPod(0)
|
||||
|
||||
t.Run("a configured grace period is applied to a finished pod", func(t *testing.T) {
|
||||
r := &EphemeralRunnerReconciler{TerminatedPodGracePeriodSeconds: 30}
|
||||
|
||||
var deleteOptions client.DeleteOptions
|
||||
for _, opt := range r.deletePodOptions(pod) {
|
||||
opt.ApplyToDelete(&deleteOptions)
|
||||
}
|
||||
|
||||
require.NotNil(t, deleteOptions.GracePeriodSeconds)
|
||||
assert.Equal(t, int64(30), *deleteOptions.GracePeriodSeconds)
|
||||
})
|
||||
|
||||
t.Run("a negative grace period leaves the deletion alone", func(t *testing.T) {
|
||||
r := &EphemeralRunnerReconciler{TerminatedPodGracePeriodSeconds: -1}
|
||||
|
||||
assert.Empty(t, r.deletePodOptions(pod))
|
||||
})
|
||||
}
|
||||
|
||||
// TestDeletePodOptionsDeletesOnlyThePodItLookedAt pins that the force delete
|
||||
// names the pod it was handed, not the name that pod happens to hold.
|
||||
//
|
||||
// Every generation of a runner's pod is named after the EphemeralRunner, and
|
||||
// pods are read through the informer cache, so a delete by name can outlive the
|
||||
// object it was decided on: an outgoing leader whose reconcile is still running
|
||||
// after the lease moved would otherwise delete the replacement pod the new
|
||||
// leader has already started a job in. The UID makes that a conflict instead.
|
||||
// This has to be asserted on the options, because the fake client honours only
|
||||
// the ResourceVersion precondition and would delete the pod either way.
|
||||
func TestDeletePodOptionsDeletesOnlyThePodItLookedAt(t *testing.T) {
|
||||
r := &EphemeralRunnerReconciler{}
|
||||
|
||||
t.Run("a finished pod is deleted by identity", func(t *testing.T) {
|
||||
pod := terminatedRunnerPod(0)
|
||||
|
||||
var deleteOptions client.DeleteOptions
|
||||
for _, opt := range r.deletePodOptions(pod) {
|
||||
opt.ApplyToDelete(&deleteOptions)
|
||||
}
|
||||
|
||||
require.NotNil(t, deleteOptions.Preconditions)
|
||||
require.NotNil(t, deleteOptions.Preconditions.UID)
|
||||
assert.Equal(t, pod.UID, *deleteOptions.Preconditions.UID)
|
||||
})
|
||||
|
||||
t.Run("a running pod is deleted without any options", func(t *testing.T) {
|
||||
pod := terminatedRunnerPod(0)
|
||||
pod.Status.ContainerStatuses[0].State = corev1.ContainerState{Running: &corev1.ContainerStateRunning{}}
|
||||
|
||||
assert.Empty(t, r.deletePodOptions(pod))
|
||||
})
|
||||
|
||||
t.Run("a pod with no UID carries no precondition", func(t *testing.T) {
|
||||
pod := terminatedRunnerPod(0)
|
||||
pod.UID = ""
|
||||
|
||||
var deleteOptions client.DeleteOptions
|
||||
for _, opt := range r.deletePodOptions(pod) {
|
||||
opt.ApplyToDelete(&deleteOptions)
|
||||
}
|
||||
|
||||
require.NotNil(t, deleteOptions.GracePeriodSeconds)
|
||||
assert.Nil(t, deleteOptions.Preconditions, "an empty UID is a precondition nothing can satisfy")
|
||||
})
|
||||
}
|
||||
|
||||
// TestReconcileReleasesTheRunnerPodWithoutAGracePeriod pins how the pod of a
|
||||
// finished runner goes away.
|
||||
//
|
||||
// The reconcile that observes the clean exit marks the runner Succeeded and
|
||||
// deletes it; the finalizers hold the object until the deletion reconcile,
|
||||
// which is what deletes the pod. What matters is the state that run of
|
||||
// reconciles leaves behind: the runner deleted, and its pod gone rather than
|
||||
// lingering in Terminating for the full terminationGracePeriodSeconds with
|
||||
// nothing left inside it to shut down.
|
||||
//
|
||||
// The pod is deliberately not deleted on the success path itself. Doing that
|
||||
// emits a pod deletion on the pod informer microseconds after the status patch
|
||||
// goes to the runner informer, and the two streams have no ordering between
|
||||
// them, so the reconcile that deletion wakes can read a runner that is not yet
|
||||
// Succeeded and build a replacement pod from a JIT config that has already been
|
||||
// used.
|
||||
func TestReconcileReleasesTheRunnerPodWithoutAGracePeriod(t *testing.T) {
|
||||
scheme := runtime.NewScheme()
|
||||
require.NoError(t, corev1.AddToScheme(scheme))
|
||||
require.NoError(t, v1alpha1.AddToScheme(scheme))
|
||||
|
||||
runner := &v1alpha1.EphemeralRunner{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: "test-runner",
|
||||
Namespace: "default",
|
||||
Finalizers: []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName},
|
||||
},
|
||||
Spec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/org/repo",
|
||||
PodTemplateSpec: corev1.PodTemplateSpec{
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName}},
|
||||
},
|
||||
},
|
||||
},
|
||||
Status: v1alpha1.EphemeralRunnerStatus{
|
||||
Phase: v1alpha1.EphemeralRunnerPhaseRunning,
|
||||
RunnerID: 42,
|
||||
RunnerName: "test-runner",
|
||||
},
|
||||
}
|
||||
|
||||
secret := &corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "test-runner", Namespace: "default"},
|
||||
Data: map[string][]byte{"runnerId": []byte("42"), "runnerName": []byte("test-runner")},
|
||||
}
|
||||
|
||||
var podDeleteOptions []client.DeleteOptions
|
||||
c := ctrlfake.NewClientBuilder().
|
||||
WithScheme(scheme).
|
||||
WithObjects(runner, secret, terminatedRunnerPod(0)).
|
||||
WithStatusSubresource(&v1alpha1.EphemeralRunner{}).
|
||||
WithInterceptorFuncs(interceptor.Funcs{
|
||||
Delete: func(ctx context.Context, c client.WithWatch, obj client.Object, opts ...client.DeleteOption) error {
|
||||
if _, ok := obj.(*corev1.Pod); ok {
|
||||
var applied client.DeleteOptions
|
||||
for _, opt := range opts {
|
||||
opt.ApplyToDelete(&applied)
|
||||
}
|
||||
podDeleteOptions = append(podDeleteOptions, applied)
|
||||
}
|
||||
return c.Delete(ctx, obj, opts...)
|
||||
},
|
||||
}).
|
||||
Build()
|
||||
|
||||
reconciler := &EphemeralRunnerReconciler{
|
||||
Client: c,
|
||||
Scheme: scheme,
|
||||
ResourceBuilder: ResourceBuilder{ResourceCache: newTestResourceCache()},
|
||||
}
|
||||
|
||||
key := types.NamespacedName{Namespace: "default", Name: "test-runner"}
|
||||
|
||||
// The reconcile that sees the exit records it and deletes the runner. The
|
||||
// pod is still the deletion reconcile's to clean up.
|
||||
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: key})
|
||||
require.NoError(t, err)
|
||||
|
||||
var got v1alpha1.EphemeralRunner
|
||||
require.NoError(t, c.Get(t.Context(), key, &got), "the runner is deleted, but its finalizers keep it until the deletion reconcile runs")
|
||||
assert.False(t, got.DeletionTimestamp.IsZero(), "the runner must be deleted by the reconcile that observes the exit")
|
||||
assert.Equal(t, v1alpha1.EphemeralRunnerPhaseSucceeded, got.Status.Phase)
|
||||
assert.Empty(t, podDeleteOptions, "the success path must not delete the pod itself")
|
||||
|
||||
// The deletion reconcile is what hands the pod back.
|
||||
_, err = reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: key})
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Len(t, podDeleteOptions, 1, "the pod of a finished runner must be deleted once it is being finalized")
|
||||
require.NotNil(t, podDeleteOptions[0].GracePeriodSeconds)
|
||||
assert.Equal(t, int64(0), *podDeleteOptions[0].GracePeriodSeconds, "a pod with nothing left running in it must not wait out a grace period")
|
||||
|
||||
err = c.Get(t.Context(), key, new(corev1.Pod))
|
||||
assert.True(t, kerrors.IsNotFound(err), "the pod must be gone, got %v", err)
|
||||
}
|
||||
@@ -25,6 +25,7 @@ import (
|
||||
"slices"
|
||||
"sort"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
@@ -32,6 +33,7 @@ import (
|
||||
"github.com/actions/scaleset"
|
||||
"github.com/go-logr/logr"
|
||||
"go.uber.org/multierr"
|
||||
"golang.org/x/sync/errgroup"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
kerrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
@@ -47,6 +49,23 @@ import (
|
||||
const (
|
||||
// EphemeralRunnerSetFinalizerName is the finalizer name used in EphemeralRunnerSet resource to protect the cleanup process of the child ephemeral runners and proxy secret.
|
||||
EphemeralRunnerSetFinalizerName = "ephemeralrunnerset.actions.github.com/finalizer"
|
||||
|
||||
// runnerBatchConcurrency is how many runners of a single EphemeralRunnerSet
|
||||
// are created or deleted at the same time.
|
||||
//
|
||||
// A set reconciles as one object, and controller-runtime serialises
|
||||
// reconciles per object, so every runner a burst of jobs asks for is created
|
||||
// by one reconcile and every finished runner it leaves behind is deleted by
|
||||
// one reconcile. Done one at a time, the round trip to the API server is
|
||||
// paid once per runner in sequence, and the last runner of a scale up waits
|
||||
// for all the ones before it. The requests are independent, so they are
|
||||
// issued in a batch instead.
|
||||
//
|
||||
// The bound exists because these are writes, and an unbounded fan-out would
|
||||
// hand the whole burst to the client rate limiter at once, where it would
|
||||
// queue in front of the reconciles of every other controller rather than in
|
||||
// front of itself.
|
||||
runnerBatchConcurrency = 8
|
||||
)
|
||||
|
||||
// EphemeralRunnerSetReconciler reconciles a EphemeralRunnerSet object
|
||||
@@ -650,16 +669,40 @@ func (r *EphemeralRunnerSetReconciler) updateStatus(ctx context.Context, ephemer
|
||||
|
||||
// deleteTerminatedEphemeralRunners deletes runners that have reached a terminal
|
||||
// state and are no longer useful, so that the scaling logic can replace them.
|
||||
//
|
||||
// The deletions are issued in a batch. Each one is an independent request, and
|
||||
// the pods they release are what the jobs waiting behind them need, so there is
|
||||
// nothing to be gained by making the hundredth runner of a burst wait for the
|
||||
// ninety-nine round trips before it.
|
||||
func (r *EphemeralRunnerSetReconciler) deleteTerminatedEphemeralRunners(ctx context.Context, ephemeralRunners []*v1alpha1.EphemeralRunner, log logr.Logger) error {
|
||||
return r.deleteEphemeralRunnersInBatches(ctx, ephemeralRunners, log)
|
||||
}
|
||||
|
||||
// deleteEphemeralRunnersInBatches deletes the given runners with a bounded
|
||||
// number of requests in flight, and reports every failure rather than the first.
|
||||
func (r *EphemeralRunnerSetReconciler) deleteEphemeralRunnersInBatches(ctx context.Context, ephemeralRunners []*v1alpha1.EphemeralRunner, log logr.Logger) error {
|
||||
if len(ephemeralRunners) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
g, ctx := errgroup.WithContext(ctx)
|
||||
g.SetLimit(runnerBatchConcurrency)
|
||||
|
||||
var mu sync.Mutex
|
||||
var errs []error
|
||||
for i := range ephemeralRunners {
|
||||
log.Info("Deleting terminated ephemeral runner", "name", ephemeralRunners[i].Name, "phase", ephemeralRunners[i].Status.Phase)
|
||||
if err := r.Delete(ctx, ephemeralRunners[i]); err != nil {
|
||||
if !kerrors.IsNotFound(err) {
|
||||
ephemeralRunner := ephemeralRunners[i]
|
||||
g.Go(func() error {
|
||||
log.Info("Deleting terminated ephemeral runner", "name", ephemeralRunner.Name, "phase", ephemeralRunner.Status.Phase)
|
||||
if err := r.Delete(ctx, ephemeralRunner); err != nil && !kerrors.IsNotFound(err) {
|
||||
mu.Lock()
|
||||
errs = append(errs, err)
|
||||
mu.Unlock()
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
_ = g.Wait()
|
||||
|
||||
return multierr.Combine(errs...)
|
||||
}
|
||||
@@ -713,18 +756,9 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte
|
||||
)
|
||||
|
||||
log.Info("Cleanup terminated ephemeral runners")
|
||||
var errs []error
|
||||
for _, ephemeralRunner := range ephemeralRunnerState.terminated() {
|
||||
log.Info("Deleting ephemeral runner", "name", ephemeralRunner.Name)
|
||||
if err := r.Delete(ctx, ephemeralRunner); err != nil && !kerrors.IsNotFound(err) {
|
||||
errs = append(errs, err)
|
||||
}
|
||||
}
|
||||
|
||||
if len(errs) > 0 {
|
||||
mergedErrs := multierr.Combine(errs...)
|
||||
log.Error(mergedErrs, "Failed to delete ephemeral runners")
|
||||
return false, mergedErrs
|
||||
if err := r.deleteEphemeralRunnersInBatches(ctx, ephemeralRunnerState.terminated(), log); err != nil {
|
||||
log.Error(err, "Failed to delete ephemeral runners")
|
||||
return false, err
|
||||
}
|
||||
|
||||
// avoid fetching the client if we have nothing left to do
|
||||
@@ -737,8 +771,8 @@ func (r *EphemeralRunnerSetReconciler) cleanUpEphemeralRunners(ctx context.Conte
|
||||
return false, err
|
||||
}
|
||||
|
||||
var errs []error
|
||||
log.Info("Cleanup pending or running ephemeral runners")
|
||||
errs = errs[0:0]
|
||||
for _, ephemeralRunner := range ephemeralRunnerState.pending {
|
||||
log.Info("Removing the ephemeral runner from the service", "name", ephemeralRunner.Name)
|
||||
_, err := r.deleteEphemeralRunnerWithActionsClient(ctx, ephemeralRunner, actionsClient, log)
|
||||
@@ -878,30 +912,53 @@ func (r *EphemeralRunnerSetReconciler) reconcileEphemeralRunnerSetProxySecret(ct
|
||||
}
|
||||
|
||||
// createEphemeralRunners provisions `count` number of v1alpha1.EphemeralRunner resources in the cluster.
|
||||
//
|
||||
// The creations are issued in a batch for the same reason the deletions are:
|
||||
// they are independent of one another, they all belong to one reconcile of one
|
||||
// object, and a job waiting for the last runner of a scale up should not also
|
||||
// be waiting for the round trips of every runner created before it.
|
||||
func (r *EphemeralRunnerSetReconciler) createEphemeralRunners(ctx context.Context, runnerSet *v1alpha1.EphemeralRunnerSet, count int, log logr.Logger) error {
|
||||
// Track multiple errors at once and return the bundle.
|
||||
errs := make([]error, 0)
|
||||
for i := range count {
|
||||
ephemeralRunner, err := r.newEphemeralRunner(runnerSet)
|
||||
if err != nil {
|
||||
log.Error(err, "failed to build ephemeral runner")
|
||||
errs = append(errs, err)
|
||||
continue
|
||||
}
|
||||
if runnerSet.Spec.EphemeralRunnerSpec.Proxy != nil {
|
||||
ephemeralRunner.Spec.ProxySecretRef = proxyEphemeralRunnerSetSecretName(runnerSet)
|
||||
}
|
||||
|
||||
log.Info("Creating new ephemeral runner", "progress", i+1, "total", count)
|
||||
if err := r.Create(ctx, ephemeralRunner); err != nil {
|
||||
log.Error(err, "failed to make ephemeral runner")
|
||||
errs = append(errs, err)
|
||||
continue
|
||||
}
|
||||
|
||||
log.Info("Created new ephemeral runner", "runner", ephemeralRunner.Name)
|
||||
if count <= 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
g, ctx := errgroup.WithContext(ctx)
|
||||
g.SetLimit(runnerBatchConcurrency)
|
||||
|
||||
// Track multiple errors at once and return the bundle.
|
||||
var mu sync.Mutex
|
||||
var errs []error
|
||||
addErr := func(err error) {
|
||||
mu.Lock()
|
||||
errs = append(errs, err)
|
||||
mu.Unlock()
|
||||
}
|
||||
|
||||
for i := range count {
|
||||
g.Go(func() error {
|
||||
ephemeralRunner, err := r.newEphemeralRunner(runnerSet)
|
||||
if err != nil {
|
||||
log.Error(err, "failed to build ephemeral runner")
|
||||
addErr(err)
|
||||
return nil
|
||||
}
|
||||
if runnerSet.Spec.EphemeralRunnerSpec.Proxy != nil {
|
||||
ephemeralRunner.Spec.ProxySecretRef = proxyEphemeralRunnerSetSecretName(runnerSet)
|
||||
}
|
||||
|
||||
log.Info("Creating new ephemeral runner", "progress", i+1, "total", count)
|
||||
if err := r.Create(ctx, ephemeralRunner); err != nil {
|
||||
log.Error(err, "failed to make ephemeral runner")
|
||||
addErr(err)
|
||||
return nil
|
||||
}
|
||||
|
||||
log.Info("Created new ephemeral runner", "runner", ephemeralRunner.Name)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
_ = g.Wait()
|
||||
|
||||
return multierr.Combine(errs...)
|
||||
}
|
||||
|
||||
|
||||
@@ -1872,6 +1872,30 @@ var _ = Describe("EphemeralRunner phase metrics", func() {
|
||||
expectEphemeralRunnerPhase(ctx, ephemeralRunner, v1alpha1.EphemeralRunnerPhaseSucceeded)
|
||||
expectEphemeralRunnerPhaseMetric(ephemeralRunner, v1alpha1.EphemeralRunnerPhaseRunning, 0)
|
||||
expectEphemeralRunnerPhaseMetric(ephemeralRunner, v1alpha1.EphemeralRunnerPhaseSucceeded, 1)
|
||||
|
||||
// The pod is handed back by the reconcile that finalizes the deleted
|
||||
// runner, not by the one that observes the exit. The runner itself is
|
||||
// still there to be read: it is deleted, but its finalizers hold it until
|
||||
// that reconcile, which is also what lets the EphemeralRunnerSet see that
|
||||
// the job finished.
|
||||
//
|
||||
// Deleting the pod on the success path instead would emit a pod deletion
|
||||
// on the pod informer just after the status patch goes to the runner
|
||||
// informer, and nothing orders those two streams, so the reconcile that
|
||||
// deletion wakes could read a runner that is not yet Succeeded and build
|
||||
// a replacement pod from a JIT config that has already been used.
|
||||
runnerFinishing := new(v1alpha1.EphemeralRunner)
|
||||
err = k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, runnerFinishing)
|
||||
Expect(err).NotTo(HaveOccurred(), "the finalizers must keep the finished runner readable")
|
||||
Expect(runnerFinishing.DeletionTimestamp.IsZero()).To(BeFalse(), "the reconcile that observes the exit must delete the runner")
|
||||
|
||||
_, err = controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred(), "failed to reconcile the deleted runner")
|
||||
|
||||
Eventually(func() bool {
|
||||
err := k8sClient.Get(ctx, client.ObjectKey{Name: ephemeralRunner.Name, Namespace: ephemeralRunner.Namespace}, pod)
|
||||
return kerrors.IsNotFound(err)
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "expected the pod of the finished runner to be gone")
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -96,15 +96,24 @@ func (b *ResourceBuilder) setSchemeIfUnset(scheme *runtime.Scheme) {
|
||||
}
|
||||
}
|
||||
|
||||
// setControllerReference marks object as owned by owner.
|
||||
//
|
||||
// The scheme is not memoised when the builder was built without one. Runners
|
||||
// are built concurrently now, and a lazily assigned field is a write shared
|
||||
// with every goroutine reading it: they would race on the pointer, and one of
|
||||
// them could read a scheme the other had allocated but not yet registered the
|
||||
// types on, failing the ownership call with an unknown kind. Building a local
|
||||
// one costs an allocation on a path no caller with a scheme ever takes.
|
||||
func (b *ResourceBuilder) setControllerReference(owner client.Object, object client.Object) error {
|
||||
if b.Scheme == nil {
|
||||
b.Scheme = runtime.NewScheme()
|
||||
if err := v1alpha1.AddToScheme(b.Scheme); err != nil {
|
||||
scheme := b.Scheme
|
||||
if scheme == nil {
|
||||
scheme = runtime.NewScheme()
|
||||
if err := v1alpha1.AddToScheme(scheme); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return ctrl.SetControllerReference(owner, object, b.Scheme)
|
||||
return ctrl.SetControllerReference(owner, object, scheme)
|
||||
}
|
||||
|
||||
func (b *ResourceBuilder) newAutoscalingListener(autoscalingRunnerSet *v1alpha1.AutoscalingRunnerSet, ephemeralRunnerSet *v1alpha1.EphemeralRunnerSet, namespace, image string, imagePullSecrets []corev1.LocalObjectReference) (*v1alpha1.AutoscalingListener, error) {
|
||||
@@ -849,7 +858,13 @@ func (b *ResourceBuilder) newEphemeralRunner(ephemeralRunnerSet *v1alpha1.Epheme
|
||||
ephemeralRunnerActionsFinalizerName,
|
||||
},
|
||||
},
|
||||
Spec: ephemeralRunnerSet.Spec.EphemeralRunnerSpec,
|
||||
// Copied rather than shared. A plain assignment is a shallow copy, which
|
||||
// leaves every runner built from this set pointing at the same container,
|
||||
// volume and map values. Creating a runner writes the API server's
|
||||
// response back into the object it was given, and the decoder reuses the
|
||||
// maps and slices it finds there, so runners built in parallel would be
|
||||
// writing into each other. Concurrently written maps end the process.
|
||||
Spec: *ephemeralRunnerSet.Spec.EphemeralRunnerSpec.DeepCopy(),
|
||||
}
|
||||
if err := b.setControllerReference(ephemeralRunnerSet, ephemeralRunner); err != nil {
|
||||
return nil, fmt.Errorf("failed to set controller reference for ephemeral runner: %w", err)
|
||||
|
||||
@@ -3,6 +3,7 @@ package actionsgithubcom
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
@@ -597,3 +598,143 @@ func TestNewEphemeralRunnerStampsActionableRevision(t *testing.T) {
|
||||
assert.Equal(t, "7", runner.Annotations[AnnotationKeyActionableRevision])
|
||||
})
|
||||
}
|
||||
|
||||
// TestNewEphemeralRunnerDoesNotShareItsSpec pins that every runner built from a
|
||||
// set owns its spec outright.
|
||||
//
|
||||
// Creating a runner hands the object to the API server and decodes the reply
|
||||
// back into it, and the decoder writes into the maps and slice elements it
|
||||
// already finds rather than allocating new ones. Runners are built and created
|
||||
// in parallel, so a spec shared between two of them is memory two goroutines
|
||||
// write at the same time, which takes the process down rather than failing a
|
||||
// request. The set is read by all of them at once, so its own copy has to come
|
||||
// through untouched too.
|
||||
func TestNewEphemeralRunnerDoesNotShareItsSpec(t *testing.T) {
|
||||
b := &ResourceBuilder{}
|
||||
set := &v1alpha1.EphemeralRunnerSet{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: "test-set",
|
||||
Namespace: "test-ns",
|
||||
Labels: map[string]string{"set-label": "original"},
|
||||
Annotations: map[string]string{"set-annotation": "original"},
|
||||
},
|
||||
Spec: v1alpha1.EphemeralRunnerSetSpec{
|
||||
Replicas: 2,
|
||||
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/org/repo",
|
||||
Proxy: &v1alpha1.ProxyConfig{
|
||||
HTTP: &v1alpha1.ProxyServerConfig{Url: "http://original"},
|
||||
NoProxy: []string{"original"},
|
||||
},
|
||||
EphemeralRunnerConfigSecretMetadata: &v1alpha1.ResourceMeta{
|
||||
Labels: map[string]string{"secret-label": "original"},
|
||||
},
|
||||
PodTemplateSpec: corev1.PodTemplateSpec{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Labels: map[string]string{"pod-label": "original"},
|
||||
Annotations: map[string]string{"pod-annotation": "original"},
|
||||
},
|
||||
Spec: corev1.PodSpec{
|
||||
NodeSelector: map[string]string{"node": "original"},
|
||||
Volumes: []corev1.Volume{{Name: "original"}},
|
||||
Containers: []corev1.Container{{
|
||||
Name: v1alpha1.EphemeralRunnerContainerName,
|
||||
Image: "original",
|
||||
Env: []corev1.EnvVar{{Name: "KEY", Value: "original"}},
|
||||
}},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
unchanged := set.DeepCopy()
|
||||
|
||||
first, err := b.newEphemeralRunner(set)
|
||||
require.NoError(t, err)
|
||||
second, err := b.newEphemeralRunner(set)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Write over everything the decoder would reach on its way through the
|
||||
// reply, as it would for the runner that happened to be created first.
|
||||
first.Spec.Spec.Containers[0].Image = "decoded"
|
||||
first.Spec.Spec.Containers[0].Env[0].Value = "decoded"
|
||||
first.Spec.Spec.Volumes[0].Name = "decoded"
|
||||
first.Spec.Spec.NodeSelector["node"] = "decoded"
|
||||
first.Spec.Labels["pod-label"] = "decoded"
|
||||
first.Spec.Annotations["pod-annotation"] = "decoded"
|
||||
first.Spec.Proxy.HTTP.Url = "decoded"
|
||||
first.Spec.Proxy.NoProxy[0] = "decoded"
|
||||
first.Spec.EphemeralRunnerConfigSecretMetadata.Labels["secret-label"] = "decoded"
|
||||
first.Labels["set-label"] = "decoded"
|
||||
first.Annotations["set-annotation"] = "decoded"
|
||||
|
||||
assert.Equal(t, "original", second.Spec.Spec.Containers[0].Image)
|
||||
assert.Equal(t, "original", second.Spec.Spec.Containers[0].Env[0].Value)
|
||||
assert.Equal(t, "original", second.Spec.Spec.Volumes[0].Name)
|
||||
assert.Equal(t, "original", second.Spec.Spec.NodeSelector["node"])
|
||||
assert.Equal(t, "original", second.Spec.Labels["pod-label"])
|
||||
assert.Equal(t, "original", second.Spec.Annotations["pod-annotation"])
|
||||
assert.Equal(t, "http://original", second.Spec.Proxy.HTTP.Url)
|
||||
assert.Equal(t, "original", second.Spec.Proxy.NoProxy[0])
|
||||
assert.Equal(t, "original", second.Spec.EphemeralRunnerConfigSecretMetadata.Labels["secret-label"])
|
||||
assert.Equal(t, "original", second.Labels["set-label"])
|
||||
assert.Equal(t, "original", second.Annotations["set-annotation"])
|
||||
|
||||
assert.Equal(t, unchanged.Spec, set.Spec, "the set a runner was built from was written into")
|
||||
assert.Equal(t, unchanged.Labels, set.Labels)
|
||||
assert.Equal(t, unchanged.Annotations, set.Annotations)
|
||||
}
|
||||
|
||||
// TestNewEphemeralRunnerIsSafeToBuildConcurrentlyWithoutAScheme pins that a
|
||||
// builder that was never given a scheme can still build runners in parallel.
|
||||
//
|
||||
// Runners are built concurrently, and the ownership reference needs a scheme to
|
||||
// resolve the owner's kind. A builder without one falls back to a scheme it
|
||||
// makes itself, and doing that by assigning to the builder would be a write
|
||||
// every other goroutine is reading at the same time: they would race on the
|
||||
// field, and one could pick up a scheme another had allocated but not yet
|
||||
// registered the types on, which fails the build with an unknown kind rather
|
||||
// than racing quietly. Every runner here has to come back owned, whichever
|
||||
// goroutine got there first. The race itself is only reported under -race.
|
||||
func TestNewEphemeralRunnerIsSafeToBuildConcurrentlyWithoutAScheme(t *testing.T) {
|
||||
b := &ResourceBuilder{}
|
||||
require.Nil(t, b.Scheme, "the fallback only runs for a builder without a scheme")
|
||||
|
||||
set := &v1alpha1.EphemeralRunnerSet{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: "test-set", Namespace: "test-ns"},
|
||||
Spec: v1alpha1.EphemeralRunnerSetSpec{
|
||||
EphemeralRunnerSpec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/org/repo",
|
||||
PodTemplateSpec: corev1.PodTemplateSpec{
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName}},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
const runners = 32
|
||||
var wg sync.WaitGroup
|
||||
built := make([]*v1alpha1.EphemeralRunner, runners)
|
||||
errs := make([]error, runners)
|
||||
|
||||
start := make(chan struct{})
|
||||
for i := range runners {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
built[i], errs[i] = b.newEphemeralRunner(set)
|
||||
}()
|
||||
}
|
||||
close(start)
|
||||
wg.Wait()
|
||||
|
||||
for i := range runners {
|
||||
require.NoError(t, errs[i])
|
||||
require.Len(t, built[i].OwnerReferences, 1, "the runner has to come back owned by the set")
|
||||
assert.Equal(t, set.Name, built[i].OwnerReferences[0].Name)
|
||||
assert.Equal(t, "EphemeralRunnerSet", built[i].OwnerReferences[0].Kind)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user