mirror of
https://github.com/actions-runner-controller/actions-runner-controller.git
synced 2026-10-04 16:13:36 +02:00
Deregister runners from the Actions service in the background (#4664)
This commit is contained in:
@@ -53,6 +53,12 @@ type EphemeralRunnerReconciler struct {
|
||||
Log logr.Logger
|
||||
Scheme *runtime.Scheme
|
||||
PublishMetrics bool
|
||||
|
||||
// UnregistrationQueue takes the removal of runner registrations from the
|
||||
// Actions service off the reconcile path. When it is left unset,
|
||||
// SetupWithManager creates one and registers it with the manager.
|
||||
UnregistrationQueue *RunnerUnregistrationQueue
|
||||
|
||||
ResourceBuilder
|
||||
}
|
||||
|
||||
@@ -105,26 +111,65 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
|
||||
}
|
||||
|
||||
if controllerutil.ContainsFinalizer(&ephemeralRunner, ephemeralRunnerActionsFinalizerName) {
|
||||
log.Info("Trying to clean up runner from the service")
|
||||
ok, err := r.cleanupRunnerFromService(ctx, &ephemeralRunner, log)
|
||||
if err != nil {
|
||||
log.Error(err, "Failed to clean up runner from service")
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
if !ok {
|
||||
log.Info("Runner is not finished yet, retrying in 30s")
|
||||
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
|
||||
// This finalizer exists to release the runner's registration with the
|
||||
// Actions service. There are two ways that happens.
|
||||
//
|
||||
// A runner that exited with code 0 already removed its own
|
||||
// registration on the way out. Runners are ephemeral, so a clean exit
|
||||
// means the agent deregistered itself before it stopped, and there is
|
||||
// nothing left to ask the service to remove. That is the path every
|
||||
// completed job takes, and it costs no API call at all.
|
||||
//
|
||||
// Every other runner may still hold a registration. Removing it is
|
||||
// handed to background workers rather than done here, so that deleting
|
||||
// the pod and the secret below is never held up by an external API.
|
||||
// See RunnerUnregistrationQueue for what that costs.
|
||||
//
|
||||
// Queueing is also what stops holding the pod alive when the service
|
||||
// reports that the runner is still executing a job. That used to keep
|
||||
// the runner pod of a job that is still running from being deleted out
|
||||
// from under it, but only for deletions that reach this branch
|
||||
// directly. The EphemeralRunnerSet does not rely on it: it refuses to
|
||||
// delete a runner that has a job assigned, and removes a runner from
|
||||
// the service before deleting it when it scales down. What is left is
|
||||
// an EphemeralRunner deleted by hand, and there the deletion is taken
|
||||
// at face value: the pod goes now, and the workers keep retrying the
|
||||
// removal until the service accepts it.
|
||||
var runnerID int
|
||||
if runnerSelfDeregistered(&ephemeralRunner) {
|
||||
log.Info("Runner exited successfully and deregistered itself, skipping its removal from the service")
|
||||
} else {
|
||||
// Resolved before the finalizer goes, because recovering an ID the
|
||||
// status never recorded reads the jitconfig secret, which the
|
||||
// cleanup below deletes.
|
||||
id, err := r.registeredRunnerID(ctx, &ephemeralRunner, log)
|
||||
if err != nil {
|
||||
log.Error(err, "Failed to resolve the registration of an ephemeral runner being deleted")
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
runnerID = id
|
||||
}
|
||||
|
||||
log.Info("Runner is cleaned up from the service, removing finalizer")
|
||||
log.Info(
|
||||
"Removing the runner registration finalizer",
|
||||
"unregisterFromService", runnerID != 0,
|
||||
"phase", ephemeralRunner.Status.Phase,
|
||||
)
|
||||
|
||||
if controllerutil.RemoveFinalizer(runner.Mutate(), ephemeralRunnerActionsFinalizerName) {
|
||||
log.Info("Removed finalizer from ephemeral runner")
|
||||
if err := r.Patch(ctx, &ephemeralRunner, runner.MergeFrom()); err != nil {
|
||||
log.Error(err, "Failed to update ephemeral runner after removing finalizer")
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
log.Info("Removed finalizer from ephemeral runner")
|
||||
|
||||
// Queued only once the finalizer is actually gone. A failed patch
|
||||
// above sends the reconcile back through this branch, and queueing
|
||||
// first would ask the service to remove the same runner twice.
|
||||
if runnerID != 0 {
|
||||
r.UnregistrationQueue.Push(&ephemeralRunner, runnerID)
|
||||
}
|
||||
log.Info("Removed the runner registration finalizer from ephemeral runner")
|
||||
}
|
||||
|
||||
log.Info("Finalizing ephemeral runner")
|
||||
@@ -160,6 +205,18 @@ func (r *EphemeralRunnerReconciler) Reconcile(ctx context.Context, req ctrl.Requ
|
||||
|
||||
if ephemeralRunner.IsDone() {
|
||||
log.Info("Cleaning up resources after after ephemeral runner termination", "phase", ephemeralRunner.Status.Phase)
|
||||
|
||||
// markAsFailed and markAsOutdated release the registration as they record
|
||||
// the terminal phase, but the patch that does it can fail after the phase
|
||||
// is already recorded, and a retry lands here rather than back in them.
|
||||
// Repeated here so that error costs a reconcile instead of leaving the
|
||||
// registration held until the set gets around to deleting the runner.
|
||||
// Does nothing once the registration is released.
|
||||
if err := r.queueUnregistration(ctx, &ephemeralRunner, log); err != nil {
|
||||
log.Error(err, "Failed to release the registration of a terminated ephemeral runner")
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
err := r.cleanupResources(ctx, &ephemeralRunner, log)
|
||||
if err != nil {
|
||||
log.Error(err, "Failed to clean up ephemeral runner owned resources")
|
||||
@@ -432,18 +489,10 @@ func (r *EphemeralRunnerReconciler) deleteEphemeralRunnerOrPod(ctx context.Conte
|
||||
return err
|
||||
}
|
||||
|
||||
// The runner is gone, and its pod failed with a job assigned, so the
|
||||
// registration is still held. The delete above runs the finalizer, which
|
||||
// queues its removal.
|
||||
log.Info("Deleted the ephemeral runner that has a job assigned but the pod has failed")
|
||||
log.Info("Trying to remove the runner from the service")
|
||||
actionsClient, err := r.GetActionsService(ctx, ephemeralRunner)
|
||||
if err != nil {
|
||||
log.Error(err, "Failed to get actions client for removing the runner from the service")
|
||||
return nil
|
||||
}
|
||||
if err := actionsClient.RemoveRunner(ctx, int64(ephemeralRunner.Status.RunnerID)); err != nil {
|
||||
log.Error(err, "Failed to remove the runner from the service")
|
||||
return nil
|
||||
}
|
||||
log.Info("Removed the runner from the service")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -455,19 +504,6 @@ func (r *EphemeralRunnerReconciler) deleteEphemeralRunnerOrPod(ctx context.Conte
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *EphemeralRunnerReconciler) cleanupRunnerFromService(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) (ok bool, err error) {
|
||||
if err := r.deleteRunnerFromService(ctx, ephemeralRunner, log); err != nil {
|
||||
if errors.Is(err, scaleset.JobStillRunningError) {
|
||||
log.Info("Runner job is still running, cannot remove the runner from the service yet")
|
||||
return false, nil
|
||||
}
|
||||
|
||||
return false, err
|
||||
}
|
||||
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (r *EphemeralRunnerReconciler) cleanupResources(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) error {
|
||||
log.Info("Cleaning up the runner pod")
|
||||
pod := new(corev1.Pod)
|
||||
@@ -607,12 +643,56 @@ func (r *EphemeralRunnerReconciler) markAsFailed(ctx context.Context, ephemeralR
|
||||
}
|
||||
r.publishEphemeralRunnerPhaseMetric(ephemeralRunner, ephemeralRunner.Status.Phase, log)
|
||||
|
||||
log.Info("Removing the runner from the service")
|
||||
if err := r.deleteRunnerFromService(ctx, ephemeralRunner, log); err != nil {
|
||||
return fmt.Errorf("failed to remove the runner from service: %w", err)
|
||||
// A failed runner is not deleted here; it stays until the EphemeralRunnerSet
|
||||
// cleans it up, which can be a long time, so the registration is released now
|
||||
// rather than waiting for the finalizer.
|
||||
if err := r.queueUnregistration(ctx, ephemeralRunner, log); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
log.Info("EphemeralRunner is marked as Failed and deleted from the service")
|
||||
log.Info("EphemeralRunner is marked as Failed and queued for removal from the service")
|
||||
return nil
|
||||
}
|
||||
|
||||
// queueUnregistration releases the runner's registration with the Actions
|
||||
// service: it hands the removal to the background workers and drops the
|
||||
// finalizer that exists to make it happen.
|
||||
//
|
||||
// A runner that exited with code 0 deregistered itself, so it has nothing to
|
||||
// hand over and only the finalizer goes.
|
||||
//
|
||||
// Dropping the finalizer is also what keeps this to a single removal. Without
|
||||
// it the deletion that eventually follows would queue the same runner again.
|
||||
// It doubles as the guard that makes this safe to call repeatedly: a runner
|
||||
// whose registration is already released is left alone.
|
||||
func (r *EphemeralRunnerReconciler) queueUnregistration(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) error {
|
||||
if !controllerutil.ContainsFinalizer(ephemeralRunner, ephemeralRunnerActionsFinalizerName) {
|
||||
return nil
|
||||
}
|
||||
|
||||
var runnerID int
|
||||
if runnerSelfDeregistered(ephemeralRunner) {
|
||||
log.Info("Runner exited successfully and deregistered itself, skipping its removal from the service")
|
||||
} else {
|
||||
id, err := r.registeredRunnerID(ctx, ephemeralRunner, log)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
runnerID = id
|
||||
}
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
controllerutil.RemoveFinalizer(ephemeralRunner, ephemeralRunnerActionsFinalizerName)
|
||||
if err := r.Patch(ctx, ephemeralRunner, client.MergeFrom(original)); err != nil && !kerrors.IsNotFound(err) {
|
||||
return fmt.Errorf("failed to remove the runner registration finalizer: %w", err)
|
||||
}
|
||||
|
||||
// Queued only once the finalizer is gone. A NotFound patch means another
|
||||
// actor already removed it and the runner finished deletion, while any other
|
||||
// failed patch leaves the removal to the retry rather than queueing it twice.
|
||||
if runnerID != 0 {
|
||||
r.UnregistrationQueue.Push(ephemeralRunner, runnerID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -629,10 +709,14 @@ func (r *EphemeralRunnerReconciler) markAsOutdated(ctx context.Context, ephemera
|
||||
}
|
||||
r.publishEphemeralRunnerPhaseMetric(ephemeralRunner, ephemeralRunner.Status.Phase, log)
|
||||
|
||||
log.Info("Removing the runner from the service")
|
||||
if err := r.deleteRunnerFromService(ctx, ephemeralRunner, log); err != nil {
|
||||
return fmt.Errorf("failed to remove the runner from service: %w", err)
|
||||
// Queued rather than removed here, for the same reason as markAsFailed: an
|
||||
// outdated runner waits on the EphemeralRunnerSet to delete it, and the
|
||||
// phase transition has no reason to wait on the service.
|
||||
if err := r.queueUnregistration(ctx, ephemeralRunner, log); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
log.Info("EphemeralRunner is marked as Outdated and queued for removal from the service")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -942,26 +1026,91 @@ func ephemeralRunnerMetricLabels(ephemeralRunner *v1alpha1.EphemeralRunner) (met
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (r *EphemeralRunnerReconciler) deleteRunnerFromService(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) error {
|
||||
client, err := r.GetActionsService(ctx, ephemeralRunner)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to get actions client for runner: %w", err)
|
||||
// registeredRunnerID returns the ID of the registration the runner holds with
|
||||
// the Actions service, or 0 when it never got one.
|
||||
//
|
||||
// Callers decide whether a removal is needed at all; this only names the
|
||||
// registration to remove. See runnerSelfDeregistered for the runners that do
|
||||
// not need one.
|
||||
//
|
||||
// The common answer comes from the runner's own status and costs nothing. The
|
||||
// exception is a runner whose status never recorded an ID: the registration is
|
||||
// created by GenerateJitRunnerConfig, and the ID it returns reaches the
|
||||
// jitconfig secret before the status patch that publishes it. A runner deleted
|
||||
// in that window holds a registration the status cannot name, so the secret is
|
||||
// read to recover it. That read only happens for a runner that got that far and
|
||||
// no further, never on the path a finishing job takes.
|
||||
//
|
||||
// A secret that cannot be read is an error rather than an answer. If it is
|
||||
// absent, the service is checked by name before concluding the runner was
|
||||
// never registered: GenerateJitRunnerConfig can register it just before
|
||||
// createSecret persists the ID. Any other uncertainty leaves the question
|
||||
// open, because answering 0 would drop the finalizer and lose the last record
|
||||
// of a registration that does exist.
|
||||
func (r *EphemeralRunnerReconciler) registeredRunnerID(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, log logr.Logger) (int, error) {
|
||||
if ephemeralRunner.Status.RunnerID != 0 {
|
||||
return ephemeralRunner.Status.RunnerID, nil
|
||||
}
|
||||
|
||||
log.Info("Removing runner from the service", "runnerId", ephemeralRunner.Status.RunnerID)
|
||||
err = client.RemoveRunner(ctx, int64(ephemeralRunner.Status.RunnerID))
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to remove runner from the service: %w", err)
|
||||
secret := new(corev1.Secret)
|
||||
if err := r.Get(ctx, types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name}, secret); err != nil {
|
||||
if !kerrors.IsNotFound(err) {
|
||||
return 0, fmt.Errorf("failed to read the jitconfig secret of a runner without a recorded ID: %w", err)
|
||||
}
|
||||
|
||||
actionsClient, err := r.GetActionsService(ctx, ephemeralRunner)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to get actions client for a runner without a recorded ID or jitconfig secret: %w", err)
|
||||
}
|
||||
|
||||
existingRunner, err := actionsClient.GetRunnerByName(ctx, ephemeralRunner.Name)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to get runner by name for a runner without a recorded ID or jitconfig secret: %w", err)
|
||||
}
|
||||
if existingRunner == nil {
|
||||
log.Info("No runner registration found for a runner without a recorded ID or jitconfig secret")
|
||||
return 0, nil
|
||||
}
|
||||
if existingRunner.RunnerScaleSetID != ephemeralRunner.Spec.RunnerScaleSetID {
|
||||
return 0, fmt.Errorf(
|
||||
"runner registration %d found by name belongs to runner scale set %d, expected %d",
|
||||
existingRunner.ID,
|
||||
existingRunner.RunnerScaleSetID,
|
||||
ephemeralRunner.Spec.RunnerScaleSetID,
|
||||
)
|
||||
}
|
||||
|
||||
log.Info("Recovered the runner ID from the Actions service", "runnerId", existingRunner.ID)
|
||||
return existingRunner.ID, nil
|
||||
}
|
||||
|
||||
log.Info("Removed runner from the service", "runnerId", ephemeralRunner.Status.RunnerID)
|
||||
return nil
|
||||
runnerID, err := strconv.Atoi(string(secret.Data["runnerId"]))
|
||||
if err != nil {
|
||||
// Not retried, unlike a failed read. Nothing about waiting makes the
|
||||
// value parse, and there is no other record of the registration.
|
||||
log.Error(err, "Jitconfig secret of a runner without a recorded ID is corrupted; leaving the runner for the service to clean up")
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
log.Info("Recovered the runner ID from the jitconfig secret", "runnerId", runnerID)
|
||||
return runnerID, nil
|
||||
}
|
||||
|
||||
// SetupWithManager sets up the controller with the Manager.
|
||||
func (r *EphemeralRunnerReconciler) SetupWithManager(mgr ctrl.Manager, opts ...Option) error {
|
||||
r.ResourceBuilder.setSchemeIfUnset(r.Scheme)
|
||||
|
||||
if r.UnregistrationQueue == nil {
|
||||
r.UnregistrationQueue = NewRunnerUnregistrationQueue(
|
||||
r.Log.WithName("runner-unregistration"),
|
||||
r.SecretResolver,
|
||||
0,
|
||||
)
|
||||
if err := mgr.Add(r.UnregistrationQueue); err != nil {
|
||||
return fmt.Errorf("failed to add the runner unregistration workers to the manager: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return builderWithOptions(
|
||||
ctrl.NewControllerManagedBy(mgr).
|
||||
For(&v1alpha1.EphemeralRunner{}).
|
||||
|
||||
@@ -1654,4 +1654,283 @@ var _ = Describe("EphemeralRunner", func() {
|
||||
).Should(BeTrue(), "failed to contact server")
|
||||
})
|
||||
})
|
||||
|
||||
Describe("Unregistering the runner from the service", func() {
|
||||
var ctx context.Context
|
||||
var mgr ctrl.Manager
|
||||
var autoscalingNS *corev1.Namespace
|
||||
var configSecret *corev1.Secret
|
||||
var controller *EphemeralRunnerReconciler
|
||||
var queue *RunnerUnregistrationQueue
|
||||
|
||||
BeforeEach(func() {
|
||||
ctx = context.Background()
|
||||
autoscalingNS, mgr = createNamespace(GinkgoT(), k8sClient)
|
||||
configSecret = createDefaultSecret(GinkgoT(), k8sClient, autoscalingNS.Name)
|
||||
|
||||
// The workers are deliberately never started. These tests are about
|
||||
// what the reconciler hands over to the queue, and leaving the pool
|
||||
// out keeps the queue readable once the reconcile returns.
|
||||
queue = NewRunnerUnregistrationQueue(logf.Log, nil, 0)
|
||||
|
||||
controller = &EphemeralRunnerReconciler{
|
||||
Client: k8sClient,
|
||||
Scheme: mgr.GetScheme(),
|
||||
Log: logf.Log,
|
||||
UnregistrationQueue: queue,
|
||||
ResourceBuilder: ResourceBuilder{
|
||||
ResourceCache: newTestResourceCache(),
|
||||
SecretResolver: secretresolver.New(k8sClient, scalefake.NewMultiClient(
|
||||
scalefake.WithClient(scalefake.NewClient()),
|
||||
)),
|
||||
},
|
||||
}
|
||||
})
|
||||
|
||||
// finalizeRunner drives a runner that has reached phase through deletion,
|
||||
// and returns what the reconciler left on the unregistration queue.
|
||||
finalizeRunner := func(name string, runnerID int, phase v1alpha1.EphemeralRunnerPhase) []runnerUnregistration {
|
||||
ephemeralRunner := newExampleRunner(name, autoscalingNS.Name, configSecret.Name)
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName}
|
||||
Expect(k8sClient.Create(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
ephemeralRunner.Status.RunnerID = runnerID
|
||||
ephemeralRunner.Status.Phase = phase
|
||||
Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed())
|
||||
|
||||
Expect(k8sClient.Delete(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
request := ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ephemeralRunner)}
|
||||
Eventually(func() bool {
|
||||
_, err := controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
err = k8sClient.Get(ctx, request.NamespacedName, new(v1alpha1.EphemeralRunner))
|
||||
return kerrors.IsNotFound(err)
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "ephemeral runner was not finalized")
|
||||
|
||||
return queue.queued()
|
||||
}
|
||||
It("skips the service for a runner that exited successfully", func() {
|
||||
// A runner that exits with code 0 removed its own registration on the
|
||||
// way out, so the deletion costs no API call at all. This is the path
|
||||
// every completed job takes.
|
||||
Expect(finalizeRunner("succeeded-runner", 1, v1alpha1.EphemeralRunnerPhaseSucceeded)).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("skips the service for a runner that exited successfully before recording its ID", func() {
|
||||
// The skip is decided on the exit alone. A succeeded runner is not
|
||||
// chased through the jitconfig secret looking for a registration to
|
||||
// remove, because it already removed its own.
|
||||
name := "succeeded-unrecorded-runner"
|
||||
Expect(k8sClient.Create(ctx, &corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name},
|
||||
Data: map[string][]byte{"runnerId": []byte("7"), "runnerName": []byte(name)},
|
||||
})).To(Succeed())
|
||||
|
||||
Expect(finalizeRunner(name, 0, v1alpha1.EphemeralRunnerPhaseSucceeded)).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("skips the service for a runner that was never registered", func() {
|
||||
Expect(finalizeRunner("unregistered-runner", 0, v1alpha1.EphemeralRunnerPhaseRunning)).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("queues the ID from the jitconfig secret when the status never recorded one", func() {
|
||||
// The registration is created before the status can publish its ID, so
|
||||
// a runner deleted in that window is registered under an ID only the
|
||||
// secret knows.
|
||||
name := "unrecorded-runner"
|
||||
Expect(k8sClient.Create(ctx, &corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name},
|
||||
Data: map[string][]byte{"runnerId": []byte("7"), "runnerName": []byte(name)},
|
||||
})).To(Succeed())
|
||||
|
||||
queued := finalizeRunner(name, 0, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
|
||||
Expect(queued).To(HaveLen(1))
|
||||
Expect(queued[0].runnerID).To(Equal(7))
|
||||
})
|
||||
|
||||
It("deletes the pod and the secret while the service call is still in flight", func() {
|
||||
blocked := make(chan struct{})
|
||||
removing := make(chan struct{})
|
||||
defer close(blocked)
|
||||
|
||||
workerCtx, stopWorkers := context.WithCancel(ctx)
|
||||
defer stopWorkers()
|
||||
|
||||
controller.UnregistrationQueue = NewRunnerUnregistrationQueue(
|
||||
logf.Log,
|
||||
secretresolver.New(k8sClient, scalefake.NewMultiClient(
|
||||
scalefake.WithClient(scalefake.NewClient(
|
||||
scalefake.WithRemoveRunnerFunc(func(context.Context, int64) error {
|
||||
close(removing)
|
||||
<-blocked
|
||||
return nil
|
||||
}),
|
||||
)),
|
||||
)),
|
||||
0,
|
||||
)
|
||||
go func() {
|
||||
defer GinkgoRecover()
|
||||
Expect(controller.UnregistrationQueue.Start(workerCtx)).To(Succeed())
|
||||
}()
|
||||
|
||||
name := "in-flight-runner"
|
||||
ephemeralRunner := newExampleRunner(name, autoscalingNS.Name, configSecret.Name)
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName}
|
||||
Expect(k8sClient.Create(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
ephemeralRunner.Status.RunnerID = 42
|
||||
ephemeralRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
|
||||
Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed())
|
||||
|
||||
Expect(k8sClient.Create(ctx, &corev1.Pod{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name},
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName, Image: "ghcr.io/actions/actions-runner"}},
|
||||
},
|
||||
})).To(Succeed())
|
||||
Expect(k8sClient.Create(ctx, &corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: autoscalingNS.Name},
|
||||
Data: map[string][]byte{jitTokenKey: []byte("jit")},
|
||||
})).To(Succeed())
|
||||
|
||||
Expect(k8sClient.Delete(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
request := ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ephemeralRunner)}
|
||||
Eventually(func() bool {
|
||||
_, err := controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
err = k8sClient.Get(ctx, request.NamespacedName, new(v1alpha1.EphemeralRunner))
|
||||
return kerrors.IsNotFound(err)
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "ephemeral runner was not finalized")
|
||||
|
||||
// The worker picked the removal up and is sitting in the service call,
|
||||
// which is the case the reconcile above used to be serialised behind.
|
||||
Eventually(removing, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeClosed())
|
||||
|
||||
Expect(kerrors.IsNotFound(k8sClient.Get(ctx, request.NamespacedName, new(corev1.Secret)))).
|
||||
To(BeTrue(), "jitconfig secret outlived the runner")
|
||||
|
||||
// envtest runs no kubelet, so nothing confirms the pod is gone and it
|
||||
// stays behind terminating.
|
||||
pod := new(corev1.Pod)
|
||||
err := k8sClient.Get(ctx, request.NamespacedName, pod)
|
||||
if !kerrors.IsNotFound(err) {
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(pod.DeletionTimestamp.IsZero()).To(BeFalse(), "runner pod was not deleted")
|
||||
}
|
||||
})
|
||||
|
||||
for _, phase := range []v1alpha1.EphemeralRunnerPhase{
|
||||
v1alpha1.EphemeralRunnerPhasePending,
|
||||
v1alpha1.EphemeralRunnerPhaseRunning,
|
||||
v1alpha1.EphemeralRunnerPhaseFailed,
|
||||
v1alpha1.EphemeralRunnerPhaseOutdated,
|
||||
} {
|
||||
It(fmt.Sprintf("queues the removal of a runner in phase %s", phase), func() {
|
||||
queued := finalizeRunner(fmt.Sprintf("%s-runner", strings.ToLower(string(phase))), 42, phase)
|
||||
|
||||
Expect(queued).To(HaveLen(1))
|
||||
Expect(queued[0].runnerID).To(Equal(42))
|
||||
// Queued rather than called, so the pod and the secret above are
|
||||
// deleted without waiting on the service.
|
||||
Expect(queued[0].readyAt.IsZero()).To(BeTrue())
|
||||
})
|
||||
}
|
||||
|
||||
It("skips the service for a runner the EphemeralRunnerSet already deregistered", func() {
|
||||
// The set removes the registration before deleting a runner it is
|
||||
// scaling down, and drops this finalizer to say so. Queueing here
|
||||
// would be a second removal for a runner the service has forgotten.
|
||||
name := "already-deregistered-runner"
|
||||
ephemeralRunner := newExampleRunner(name, autoscalingNS.Name, configSecret.Name)
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName}
|
||||
Expect(k8sClient.Create(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
ephemeralRunner.Status.RunnerID = 42
|
||||
ephemeralRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseRunning
|
||||
Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed())
|
||||
|
||||
Expect(k8sClient.Delete(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
request := ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ephemeralRunner)}
|
||||
Eventually(func() bool {
|
||||
_, err := controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
err = k8sClient.Get(ctx, request.NamespacedName, new(v1alpha1.EphemeralRunner))
|
||||
return kerrors.IsNotFound(err)
|
||||
}, ephemeralRunnerTimeout, ephemeralRunnerInterval).Should(BeTrue(), "ephemeral runner was not finalized")
|
||||
|
||||
Expect(queue.queued()).To(BeEmpty())
|
||||
})
|
||||
|
||||
It("releases the registration of a terminated runner that still holds the finalizer", func() {
|
||||
// markAsFailed and markAsOutdated release the registration as they
|
||||
// record the phase, but that patch can fail once the phase is already
|
||||
// recorded, and the retry does not land back in them: the runner is
|
||||
// terminal by then, so the reconcile takes the IsDone path instead.
|
||||
// This is that path picking the release back up.
|
||||
name := "terminated-runner-still-registered"
|
||||
ephemeralRunner := newExampleRunner(name, autoscalingNS.Name, configSecret.Name)
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName}
|
||||
Expect(k8sClient.Create(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
ephemeralRunner.Status.RunnerID = 42
|
||||
ephemeralRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseFailed
|
||||
Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed())
|
||||
|
||||
// Not deleted. The set has not got to it yet, which is the whole
|
||||
// reason the registration should not wait for that.
|
||||
request := ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ephemeralRunner)}
|
||||
_, err := controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
queued := queue.queued()
|
||||
Expect(queued).To(HaveLen(1))
|
||||
Expect(queued[0].runnerID).To(Equal(42))
|
||||
|
||||
updated := new(v1alpha1.EphemeralRunner)
|
||||
Expect(k8sClient.Get(ctx, request.NamespacedName, updated)).To(Succeed())
|
||||
Expect(updated.Finalizers).NotTo(ContainElement(ephemeralRunnerActionsFinalizerName))
|
||||
|
||||
// Reconciling a terminal runner again must not ask the service to
|
||||
// remove the same registration a second time.
|
||||
_, err = controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(queue.queued()).To(HaveLen(1))
|
||||
})
|
||||
|
||||
It("leaves a succeeded runner alone on the terminated path", func() {
|
||||
// Same path, but the runner exited cleanly, so there is nothing to
|
||||
// hand over and only the finalizer goes.
|
||||
name := "terminated-succeeded-runner"
|
||||
ephemeralRunner := newExampleRunner(name, autoscalingNS.Name, configSecret.Name)
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName}
|
||||
Expect(k8sClient.Create(ctx, ephemeralRunner)).To(Succeed())
|
||||
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
ephemeralRunner.Status.RunnerID = 42
|
||||
ephemeralRunner.Status.Phase = v1alpha1.EphemeralRunnerPhaseSucceeded
|
||||
Expect(k8sClient.Status().Patch(ctx, ephemeralRunner, client.MergeFrom(original))).To(Succeed())
|
||||
|
||||
request := ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ephemeralRunner)}
|
||||
_, err := controller.Reconcile(ctx, request)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
Expect(queue.queued()).To(BeEmpty())
|
||||
|
||||
updated := new(v1alpha1.EphemeralRunner)
|
||||
Expect(k8sClient.Get(ctx, request.NamespacedName, updated)).To(Succeed())
|
||||
Expect(updated.Finalizers).NotTo(ContainElement(ephemeralRunnerActionsFinalizerName))
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
@@ -991,12 +991,33 @@ func (r *EphemeralRunnerSetReconciler) deleteIdleEphemeralRunners(ctx context.Co
|
||||
|
||||
func (r *EphemeralRunnerSetReconciler) deleteEphemeralRunnerWithActionsClient(ctx context.Context, ephemeralRunner *v1alpha1.EphemeralRunner, actionsClient multiclient.Client, log logr.Logger) (bool, error) {
|
||||
if err := actionsClient.RemoveRunner(ctx, int64(ephemeralRunner.Status.RunnerID)); err != nil {
|
||||
if errors.Is(err, scaleset.JobStillRunningError) {
|
||||
switch {
|
||||
case errors.Is(err, scaleset.JobStillRunningError):
|
||||
log.Info("Runner is still running a job, skipping deletion", "name", ephemeralRunner.Name, "runnerId", ephemeralRunner.Status.RunnerID)
|
||||
return false, nil
|
||||
}
|
||||
|
||||
return false, err
|
||||
case errors.Is(err, scaleset.RunnerNotFoundError), errors.Is(err, scaleset.NotFoundError):
|
||||
// The registration is gone, which is all this call wanted. Reached by
|
||||
// retrying after the removal landed but the deletion below did not,
|
||||
// and treating it as a failure would leave the runner stuck behind a
|
||||
// call that can never succeed again.
|
||||
log.Info("Runner is already removed from the service", "name", ephemeralRunner.Name, "runnerId", ephemeralRunner.Status.RunnerID)
|
||||
|
||||
default:
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
|
||||
// The registration is gone, so drop the finalizer that exists to remove it.
|
||||
// Otherwise deleting the runner below queues a second removal for a runner
|
||||
// the service has already forgotten, which is one wasted API call for every
|
||||
// runner a scale down takes.
|
||||
if controllerutil.ContainsFinalizer(ephemeralRunner, ephemeralRunnerActionsFinalizerName) {
|
||||
original := ephemeralRunner.DeepCopy()
|
||||
controllerutil.RemoveFinalizer(ephemeralRunner, ephemeralRunnerActionsFinalizerName)
|
||||
if err := r.Patch(ctx, ephemeralRunner, client.MergeFrom(original)); err != nil && !kerrors.IsNotFound(err) {
|
||||
return false, fmt.Errorf("failed to remove the runner registration finalizer: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
log.Info("Deleting ephemeral runner after removing from the service", "name", ephemeralRunner.Name, "runnerId", ephemeralRunner.Status.RunnerID)
|
||||
|
||||
@@ -78,6 +78,13 @@ func WithRemoveRunner(err error) ClientOption {
|
||||
}
|
||||
}
|
||||
|
||||
// WithRemoveRunnerFunc configures a function to handle RemoveRunner calls dynamically
|
||||
func WithRemoveRunnerFunc(fn func(context.Context, int64) error) ClientOption {
|
||||
return func(c *Client) {
|
||||
c.removeRunnerFunc = fn
|
||||
}
|
||||
}
|
||||
|
||||
// WithGenerateJitRunnerConfig configures the result of GenerateJitRunnerConfig
|
||||
func WithGenerateJitRunnerConfig(result *scaleset.RunnerScaleSetJitRunnerConfig, err error) ClientOption {
|
||||
return func(c *Client) {
|
||||
@@ -128,6 +135,7 @@ type Client struct {
|
||||
systemInfo scaleset.SystemInfo
|
||||
createRunnerScaleSetFunc func(context.Context, *scaleset.RunnerScaleSet) (*scaleset.RunnerScaleSet, error)
|
||||
updateRunnerScaleSetFunc func(context.Context, int, *scaleset.RunnerScaleSet) (*scaleset.RunnerScaleSet, error)
|
||||
removeRunnerFunc func(context.Context, int64) error
|
||||
|
||||
getRunnerScaleSetResult struct {
|
||||
*scaleset.RunnerScaleSet
|
||||
@@ -212,6 +220,9 @@ func (c *Client) GetRunnerByName(ctx context.Context, runnerName string) (*scale
|
||||
}
|
||||
|
||||
func (c *Client) RemoveRunner(ctx context.Context, runnerID int64) error {
|
||||
if c.removeRunnerFunc != nil {
|
||||
return c.removeRunnerFunc(ctx, runnerID)
|
||||
}
|
||||
return c.removeRunnerResult.err
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,457 @@
|
||||
/*
|
||||
Copyright 2020 The actions-runner-controller authors.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package actionsgithubcom
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
"github.com/actions/scaleset"
|
||||
"github.com/go-logr/logr"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"sigs.k8s.io/controller-runtime/pkg/manager"
|
||||
)
|
||||
|
||||
const (
|
||||
// unregistrationRetryDelay is how long a request waits before it is attempted
|
||||
// again after the service reported that the runner is still executing a job.
|
||||
unregistrationRetryDelay = 30 * time.Second
|
||||
|
||||
// unregistrationMinWorkers is the lower bound on the size of the worker pool.
|
||||
unregistrationMinWorkers = 4
|
||||
|
||||
// unregistrationMaxIdleWait bounds how long a worker sleeps when the queue
|
||||
// holds nothing it can act on yet. Pushes wake workers directly, so this is
|
||||
// only a backstop against a wake-up that was coalesced away.
|
||||
unregistrationMaxIdleWait = 30 * time.Second
|
||||
)
|
||||
|
||||
// runnerSelfDeregistered reports whether the runner removed its own
|
||||
// registration from the Actions service before its EphemeralRunner was deleted.
|
||||
//
|
||||
// It is an inference drawn from state the controller already has, not a lookup,
|
||||
// because the whole point is to avoid the API call. It is deliberately
|
||||
// conservative in the direction that costs an API call rather than the one that
|
||||
// leaks a registration.
|
||||
//
|
||||
// The Succeeded phase is set from a single observation: the runner container
|
||||
// terminated with exit code 0. Runners are configured as ephemeral, so a runner
|
||||
// that reaches a clean exit has already removed its own registration on the way
|
||||
// out, and asking the service to remove it again is a wasted round trip on the
|
||||
// hottest path the controller has.
|
||||
//
|
||||
// Every other phase is reachable with the registration still in place. A runner
|
||||
// that never started, was killed, exited non-zero, or was found to be outdated
|
||||
// did not get to deregister itself.
|
||||
func runnerSelfDeregistered(ephemeralRunner *v1alpha1.EphemeralRunner) bool {
|
||||
return ephemeralRunner.Status.Phase == v1alpha1.EphemeralRunnerPhaseSucceeded
|
||||
}
|
||||
|
||||
// runnerUnregistration is a single queued attempt to remove one runner from the
|
||||
// Actions service.
|
||||
type runnerUnregistration struct {
|
||||
// runner is a copy of the EphemeralRunner taken before it was deleted from
|
||||
// the cluster. The workers hold the whole object rather than just the runner
|
||||
// ID because resolving the credentials of the scale set the runner belongs
|
||||
// to needs its GitHub configuration, and by the time a worker gets to it the
|
||||
// object is usually gone from the API server.
|
||||
runner *v1alpha1.EphemeralRunner
|
||||
|
||||
// runnerID identifies the registration to remove. It is carried separately
|
||||
// because it is not always the one the runner's status reports: a runner
|
||||
// deleted before the controller recorded the ID has its registration
|
||||
// identified from the jitconfig secret instead.
|
||||
runnerID int
|
||||
|
||||
// readyAt is the earliest time this request may be attempted. The zero value
|
||||
// means it is ready immediately.
|
||||
readyAt time.Time
|
||||
}
|
||||
|
||||
// key identifies the registration this request is trying to remove.
|
||||
//
|
||||
// The runner ID alone would be the obvious key, but it is assigned by the
|
||||
// Actions service and only unique within the scale set's GitHub scope. One
|
||||
// controller can serve several of those, so the runner is named as well.
|
||||
func (r runnerUnregistration) key() unregistrationKey {
|
||||
return unregistrationKey{
|
||||
namespace: r.runner.Namespace,
|
||||
name: r.runner.Name,
|
||||
runnerID: r.runnerID,
|
||||
}
|
||||
}
|
||||
|
||||
type unregistrationKey struct {
|
||||
namespace string
|
||||
name string
|
||||
runnerID int
|
||||
}
|
||||
|
||||
// RunnerUnregistrationQueue removes runners from the Actions service outside of
|
||||
// the EphemeralRunner reconcile loop.
|
||||
//
|
||||
// Finalizing an EphemeralRunner is two unrelated pieces of work: deleting the
|
||||
// Kubernetes resources it owns, and removing its registration from the Actions
|
||||
// service. Only the first is local and fast. Doing them in sequence puts pod
|
||||
// garbage collection behind an external, eventually consistent API, so a burst
|
||||
// of completing jobs leaves completed pods standing around waiting on calls
|
||||
// they have no real dependency on.
|
||||
//
|
||||
// So the reconciler hands the registration over to this queue and moves on.
|
||||
// Push appends to a slice under a mutex and returns; a fixed pool of workers
|
||||
// issues the API calls. Nothing the Actions service does can extend a reconcile.
|
||||
//
|
||||
// The queue is in-memory only, and deliberately so. Requests still queued when
|
||||
// the controller restarts or loses leader election are dropped, and nothing
|
||||
// recovers them. That is acceptable because this is not the only thing that
|
||||
// removes a runner: the Actions service drops a registration on its own once
|
||||
// the runner stops reporting in, so a lost request delays that cleanup rather
|
||||
// than leaking the registration. Making the queue durable would mean writing to
|
||||
// the API server for every deletion, which is the cost this type exists to
|
||||
// avoid.
|
||||
//
|
||||
// A zero value is not usable. Use NewRunnerUnregistrationQueue.
|
||||
type RunnerUnregistrationQueue struct {
|
||||
log logr.Logger
|
||||
secretResolver SecretResolver
|
||||
workers int
|
||||
retryDelay time.Duration
|
||||
|
||||
mu sync.Mutex
|
||||
|
||||
// ready holds requests that can be attempted now, oldest first. It is
|
||||
// consumed through readyHead rather than by resliding from the front,
|
||||
// because a burst of deletions arrives all at once and shifting the tail
|
||||
// down on every take would make draining it quadratic, under the lock that
|
||||
// Push needs.
|
||||
ready []runnerUnregistration
|
||||
readyHead int
|
||||
|
||||
// delayed holds requests waiting out a retry delay. Only a runner the
|
||||
// service reported as still executing a job ends up here, so it stays short
|
||||
// enough to scan on every take.
|
||||
delayed []runnerUnregistration
|
||||
|
||||
// claimed holds the key of every request that has been taken on and not yet
|
||||
// finished with, whether it is waiting in the queue or in the middle of an
|
||||
// API call. Pushing a key that is already in here drops the request, so the
|
||||
// service is never asked twice to remove the same registration.
|
||||
//
|
||||
// A retry keeps its claim, which is what this mainly protects: two copies of
|
||||
// a request for a runner that is still executing a job would sit in the
|
||||
// delayed list retrying in lockstep for as long as the job runs.
|
||||
claimed map[unregistrationKey]struct{}
|
||||
|
||||
// notify carries a single wake-up token for a waiting worker. It is only
|
||||
// ever sent to without blocking, so that a push is never slowed down by the
|
||||
// state of the pool.
|
||||
notify chan struct{}
|
||||
}
|
||||
|
||||
// The manager owns the lifecycle of the pool.
|
||||
var _ manager.Runnable = (*RunnerUnregistrationQueue)(nil)
|
||||
|
||||
// NewRunnerUnregistrationQueue returns a queue drained by a pool of workers
|
||||
// goroutines, never fewer than unregistrationMinWorkers. The controller
|
||||
// defaults to a single concurrent reconcile, and sizing the pool off that alone
|
||||
// would leave one goroutine to absorb every burst.
|
||||
func NewRunnerUnregistrationQueue(log logr.Logger, secretResolver SecretResolver, workers int) *RunnerUnregistrationQueue {
|
||||
return &RunnerUnregistrationQueue{
|
||||
log: log,
|
||||
secretResolver: secretResolver,
|
||||
workers: max(workers, unregistrationMinWorkers),
|
||||
retryDelay: unregistrationRetryDelay,
|
||||
claimed: make(map[unregistrationKey]struct{}),
|
||||
notify: make(chan struct{}, 1),
|
||||
}
|
||||
}
|
||||
|
||||
// Push queues the removal of runnerID, the registration held by the given
|
||||
// runner, from the Actions service.
|
||||
//
|
||||
// It copies what the workers need and returns. It never blocks, never fails and
|
||||
// never touches the network, so no caller can be delayed by how far behind the
|
||||
// pool is or by how the service is behaving.
|
||||
//
|
||||
// A registration that is already queued, or is in the middle of being removed,
|
||||
// is not queued again, so the service is asked to remove it exactly once.
|
||||
//
|
||||
// Pushing to a nil queue drops the request. The only way to get one is to build
|
||||
// an EphemeralRunnerReconciler by hand and never call SetupWithManager, which
|
||||
// wires a queue up when the field is left unset.
|
||||
func (q *RunnerUnregistrationQueue) Push(ephemeralRunner *v1alpha1.EphemeralRunner, runnerID int) {
|
||||
if q == nil {
|
||||
return
|
||||
}
|
||||
request := runnerUnregistration{runner: ephemeralRunner.DeepCopy(), runnerID: runnerID}
|
||||
if !q.push(request) {
|
||||
q.log.V(1).Info(
|
||||
"Runner is already queued for removal from the service",
|
||||
"ephemeralRunner", types.NamespacedName{Namespace: ephemeralRunner.Namespace, Name: ephemeralRunner.Name},
|
||||
"runnerId", runnerID,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// Start drains the queue until ctx is cancelled.
|
||||
//
|
||||
// It implements manager.Runnable so the pool shares the manager's lifecycle,
|
||||
// starting alongside the controllers and stopping when the manager shuts down.
|
||||
// Whatever is still queued at that point is dropped; see the type
|
||||
// documentation for why that is safe.
|
||||
func (q *RunnerUnregistrationQueue) Start(ctx context.Context) error {
|
||||
q.log.Info("Starting runner unregistration workers", "workers", q.workers)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(q.workers)
|
||||
for range q.workers {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
q.work(ctx)
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
q.log.Info("Stopped runner unregistration workers", "dropped", q.len())
|
||||
return nil
|
||||
}
|
||||
|
||||
// work is the loop of a single worker. It takes whatever is ready, and when
|
||||
// nothing is, waits for the shorter of the next request becoming ready and the
|
||||
// next push.
|
||||
func (q *RunnerUnregistrationQueue) work(ctx context.Context) {
|
||||
for {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
|
||||
request, wait, ok := q.next(time.Now())
|
||||
if ok {
|
||||
q.unregister(ctx, request)
|
||||
continue
|
||||
}
|
||||
|
||||
// Nothing is ready, so park until something is. wait is how long the
|
||||
// earliest delayed request still has to go, and it is an upper bound
|
||||
// rather than a commitment: a push cuts it short. That is what keeps a
|
||||
// request needing no wait from queueing behind ones that do, even when
|
||||
// every worker is parked on a retry 30 seconds out.
|
||||
timer := time.NewTimer(wait)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
timer.Stop()
|
||||
return
|
||||
case <-q.notify:
|
||||
timer.Stop()
|
||||
case <-timer.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// unregister issues the removal for a single request.
|
||||
//
|
||||
// Failures are logged and dropped rather than retried, with one exception: the
|
||||
// service refuses to remove a runner it still considers to be executing a job,
|
||||
// and that request goes back on the queue. Everything else is left to the
|
||||
// service, which removes a registration that stops reporting in on its own.
|
||||
// There is no reconcile to fail here and no caller waiting on the result, so
|
||||
// retrying anything else would only build a backlog of calls that the service
|
||||
// is already going to make unnecessary.
|
||||
func (q *RunnerUnregistrationQueue) unregister(ctx context.Context, request runnerUnregistration) {
|
||||
runner := request.runner
|
||||
log := q.log.WithValues(
|
||||
"ephemeralRunner", types.NamespacedName{Namespace: runner.Namespace, Name: runner.Name},
|
||||
"runnerId", request.runnerID,
|
||||
)
|
||||
|
||||
// The claim taken when this was queued is held until the request is done
|
||||
// with, so nothing can queue the same registration alongside it. A retry is
|
||||
// not done with, and keeps the claim.
|
||||
retrying := false
|
||||
defer func() {
|
||||
if !retrying {
|
||||
q.release(request)
|
||||
}
|
||||
}()
|
||||
|
||||
actionsClient, err := q.secretResolver.GetActionsService(ctx, runner)
|
||||
if err != nil {
|
||||
log.Error(err, "Failed to get actions client to remove the runner from the service; leaving the runner for the service to clean up")
|
||||
return
|
||||
}
|
||||
|
||||
err = actionsClient.RemoveRunner(ctx, int64(request.runnerID))
|
||||
switch {
|
||||
case err == nil:
|
||||
log.Info("Removed runner from the service")
|
||||
|
||||
case errors.Is(err, scaleset.RunnerNotFoundError), errors.Is(err, scaleset.NotFoundError):
|
||||
// The registration is gone, which is the outcome this request wanted.
|
||||
log.Info("Runner is already removed from the service")
|
||||
|
||||
case errors.Is(err, scaleset.JobStillRunningError):
|
||||
log.Info("Runner is still executing a job, retrying the removal later", "retryAfter", q.retryDelay)
|
||||
retrying = true
|
||||
q.pushAfter(request, q.retryDelay)
|
||||
|
||||
case ctx.Err() != nil:
|
||||
log.Info("Shutting down before the runner could be removed from the service; leaving the runner for the service to clean up")
|
||||
|
||||
default:
|
||||
log.Error(err, "Failed to remove the runner from the service; leaving the runner for the service to clean up")
|
||||
}
|
||||
}
|
||||
|
||||
// next takes the oldest request that is ready to be attempted at now.
|
||||
//
|
||||
// When nothing is ready it reports how long to wait before asking again, which
|
||||
// is the time until the earliest delayed request, capped at
|
||||
// unregistrationMaxIdleWait.
|
||||
func (q *RunnerUnregistrationQueue) next(now time.Time) (runnerUnregistration, time.Duration, bool) {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
wait := q.promoteLocked(now)
|
||||
|
||||
if q.readyHead == len(q.ready) {
|
||||
return runnerUnregistration{}, wait, false
|
||||
}
|
||||
|
||||
request := q.ready[q.readyHead]
|
||||
// Drop the reference so that the consumed entry does not keep the copy of
|
||||
// the runner alive until the slice is reused.
|
||||
q.ready[q.readyHead] = runnerUnregistration{}
|
||||
q.readyHead++
|
||||
if q.readyHead == len(q.ready) {
|
||||
q.ready = q.ready[:0]
|
||||
q.readyHead = 0
|
||||
} else if q.readyHead >= len(q.ready)-q.readyHead {
|
||||
remaining := len(q.ready) - q.readyHead
|
||||
copy(q.ready, q.ready[q.readyHead:])
|
||||
clear(q.ready[remaining:])
|
||||
q.ready = q.ready[:remaining]
|
||||
q.readyHead = 0
|
||||
}
|
||||
|
||||
if q.readyHead < len(q.ready) || len(q.delayed) > 0 {
|
||||
// Hand the wake-up on so that a burst of pushes, which only ever leaves a
|
||||
// single token behind, spreads across the pool instead of being worked
|
||||
// through by whichever worker happened to take the first one. Waking for
|
||||
// a request that turns out to still be delayed costs one more pass
|
||||
// through this function.
|
||||
q.wake()
|
||||
}
|
||||
|
||||
return request, 0, true
|
||||
}
|
||||
|
||||
// promoteLocked moves every delayed request that has come due onto the ready
|
||||
// list and returns how long to wait for the earliest of the ones that have not,
|
||||
// capped at unregistrationMaxIdleWait.
|
||||
func (q *RunnerUnregistrationQueue) promoteLocked(now time.Time) time.Duration {
|
||||
if len(q.delayed) == 0 {
|
||||
return unregistrationMaxIdleWait
|
||||
}
|
||||
|
||||
// Partitioned in a single pass. A burst of runners refused together comes
|
||||
// due together, and removing them one at a time would shift the rest of the
|
||||
// list on every promotion, under the lock Push needs.
|
||||
wait := unregistrationMaxIdleWait
|
||||
kept := q.delayed[:0]
|
||||
for _, request := range q.delayed {
|
||||
if request.readyAt.After(now) {
|
||||
wait = min(wait, request.readyAt.Sub(now))
|
||||
kept = append(kept, request)
|
||||
continue
|
||||
}
|
||||
|
||||
request.readyAt = time.Time{}
|
||||
q.ready = append(q.ready, request)
|
||||
}
|
||||
|
||||
// Compacting leaves the promoted requests duplicated in the tail, where they
|
||||
// would keep their copy of the runner alive until the slice is reused.
|
||||
clear(q.delayed[len(kept):])
|
||||
q.delayed = kept
|
||||
|
||||
return wait
|
||||
}
|
||||
|
||||
// pushAfter queues request again, to be attempted no earlier than delay from
|
||||
// now.
|
||||
//
|
||||
// The request keeps the claim it already holds, so this cannot be turned away
|
||||
// as a duplicate of itself.
|
||||
func (q *RunnerUnregistrationQueue) pushAfter(request runnerUnregistration, delay time.Duration) {
|
||||
request.readyAt = time.Now().Add(delay)
|
||||
q.enqueue(request)
|
||||
}
|
||||
|
||||
// push queues a request for a registration that is not already queued, and
|
||||
// reports whether it took it on.
|
||||
func (q *RunnerUnregistrationQueue) push(request runnerUnregistration) bool {
|
||||
q.mu.Lock()
|
||||
if _, ok := q.claimed[request.key()]; ok {
|
||||
q.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
q.claimed[request.key()] = struct{}{}
|
||||
q.mu.Unlock()
|
||||
|
||||
q.enqueue(request)
|
||||
return true
|
||||
}
|
||||
|
||||
// release gives up the claim on a request that is done with, whatever the
|
||||
// outcome was. A later request to remove the same registration is then taken on
|
||||
// as a new one.
|
||||
func (q *RunnerUnregistrationQueue) release(request runnerUnregistration) {
|
||||
q.mu.Lock()
|
||||
delete(q.claimed, request.key())
|
||||
q.mu.Unlock()
|
||||
}
|
||||
|
||||
func (q *RunnerUnregistrationQueue) enqueue(request runnerUnregistration) {
|
||||
q.mu.Lock()
|
||||
if request.readyAt.IsZero() {
|
||||
q.ready = append(q.ready, request)
|
||||
} else {
|
||||
q.delayed = append(q.delayed, request)
|
||||
}
|
||||
q.mu.Unlock()
|
||||
|
||||
q.wake()
|
||||
}
|
||||
|
||||
// wake releases one worker from its wait. The send is non-blocking against a
|
||||
// single-slot buffer, so a push is never delayed by the pool, and a wake-up
|
||||
// that nobody is waiting for is kept for whoever waits next.
|
||||
func (q *RunnerUnregistrationQueue) wake() {
|
||||
select {
|
||||
case q.notify <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (q *RunnerUnregistrationQueue) len() int {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
return len(q.ready) - q.readyHead + len(q.delayed)
|
||||
}
|
||||
@@ -0,0 +1,258 @@
|
||||
package actionsgithubcom
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
"github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient"
|
||||
scalefake "github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake"
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/stretchr/testify/require"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
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"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||
)
|
||||
|
||||
// serviceLatencies are the round trip times the Actions service is stood up
|
||||
// with. Zero is not a latency anyone observes; it is there to separate the cost
|
||||
// of the call itself from the cost of waiting for it.
|
||||
var serviceLatencies = []time.Duration{0, 1 * time.Millisecond, 25 * time.Millisecond}
|
||||
|
||||
// Everything here is built with logr.Discard rather than the controller-runtime
|
||||
// global logger. That logger is a promise nobody fulfills in a benchmark
|
||||
// binary, and 30 seconds in it resolves itself to a fallback that writes a
|
||||
// stack trace to stderr. A run long enough to hit that gets its output
|
||||
// corrupted mid-line and its timings skewed by the writes.
|
||||
|
||||
// benchmarkActionsClient returns a client whose RemoveRunner takes latency to
|
||||
// answer.
|
||||
func benchmarkActionsClient(latency time.Duration) multiclient.Client {
|
||||
return scalefake.NewClient(scalefake.WithRemoveRunnerFunc(func(context.Context, int64) error {
|
||||
if latency > 0 {
|
||||
time.Sleep(latency)
|
||||
}
|
||||
return nil
|
||||
}))
|
||||
}
|
||||
|
||||
// newFinalizeBenchmarkReconciler builds a reconciler over an in-memory API
|
||||
// server holding nothing but one runner, its pod and its secret. queue may be
|
||||
// nil, which makes the push a no-op.
|
||||
//
|
||||
// One of these per iteration. The fake client scans every object it holds on
|
||||
// every read, so a store built up across iterations would charge the reconcile
|
||||
// for how many runners the benchmark happens to have finalized already.
|
||||
func newFinalizeBenchmarkReconciler(b *testing.B, scheme *runtime.Scheme, queue *RunnerUnregistrationQueue, phase v1alpha1.EphemeralRunnerPhase, n int) (*EphemeralRunnerReconciler, types.NamespacedName, int) {
|
||||
b.Helper()
|
||||
|
||||
c := fake.NewClientBuilder().
|
||||
WithScheme(scheme).
|
||||
WithStatusSubresource(&v1alpha1.EphemeralRunner{}).
|
||||
Build()
|
||||
|
||||
reconciler := &EphemeralRunnerReconciler{
|
||||
Client: c,
|
||||
Scheme: scheme,
|
||||
Log: logr.Discard(),
|
||||
UnregistrationQueue: queue,
|
||||
ResourceBuilder: ResourceBuilder{
|
||||
ResourceCache: newTestResourceCache(),
|
||||
},
|
||||
}
|
||||
|
||||
key, runnerID := createFinalizeBenchmarkRunner(b, c, phase, n)
|
||||
return reconciler, key, runnerID
|
||||
}
|
||||
|
||||
// benchmarkScheme registers only the types the finalizer path touches. The fake
|
||||
// client walks the scheme on every operation, and the full client-go scheme
|
||||
// puts more time into that bookkeeping than into the reconcile being measured.
|
||||
func benchmarkScheme(b *testing.B) *runtime.Scheme {
|
||||
b.Helper()
|
||||
|
||||
scheme := runtime.NewScheme()
|
||||
require.NoError(b, corev1.AddToScheme(scheme))
|
||||
require.NoError(b, v1alpha1.AddToScheme(scheme))
|
||||
return scheme
|
||||
}
|
||||
|
||||
// createFinalizeBenchmarkRunner puts a deleted EphemeralRunner, its pod and its
|
||||
// jitconfig secret in front of the reconciler, which is the state the finalizer
|
||||
// path runs against.
|
||||
func createFinalizeBenchmarkRunner(b *testing.B, c client.Client, phase v1alpha1.EphemeralRunnerPhase, n int) (types.NamespacedName, int) {
|
||||
b.Helper()
|
||||
|
||||
// A distinct runner every iteration, as in production. The API server is
|
||||
// rebuilt each time and would not care, but the queue is not: it holds a
|
||||
// claim on every registration handed to it, and reusing one would measure
|
||||
// the duplicate being turned away rather than the push.
|
||||
name := fmt.Sprintf("runner-%d", n)
|
||||
ctx := context.Background()
|
||||
ephemeralRunner := newExampleRunner(name, "default", "config-secret")
|
||||
ephemeralRunner.Finalizers = []string{ephemeralRunnerFinalizerName, ephemeralRunnerActionsFinalizerName}
|
||||
// Offset so the untimed warmup pass, which runs at -1, still gets a valid
|
||||
// registration and takes the same path as the measured ones.
|
||||
ephemeralRunner.Status.RunnerID = n + 2
|
||||
ephemeralRunner.Status.Phase = phase
|
||||
require.NoError(b, c.Create(ctx, ephemeralRunner))
|
||||
|
||||
require.NoError(b, c.Create(ctx, &corev1.Pod{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"},
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{{Name: v1alpha1.EphemeralRunnerContainerName, Image: "ghcr.io/actions/actions-runner"}},
|
||||
},
|
||||
}))
|
||||
require.NoError(b, c.Create(ctx, &corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"},
|
||||
Data: map[string][]byte{jitTokenKey: []byte("jit")},
|
||||
}))
|
||||
|
||||
require.NoError(b, c.Delete(ctx, ephemeralRunner))
|
||||
|
||||
return types.NamespacedName{Namespace: "default", Name: name}, ephemeralRunner.Status.RunnerID
|
||||
}
|
||||
|
||||
// BenchmarkEphemeralRunnerFinalize measures one pass of the finalizer path: the
|
||||
// work that stands between a completed job and its pod being collected.
|
||||
//
|
||||
// The variants are the three ways that pass can go:
|
||||
//
|
||||
// - skipped: the runner exited cleanly, so its registration is already gone
|
||||
// and nothing is queued. This is the path every completed job takes.
|
||||
// - queued: the registration is handed to the workers. Present behaviour for
|
||||
// a runner that may still be registered.
|
||||
// - synchronous: the removal is issued inline before the local cleanup, which
|
||||
// is the behaviour this replaced.
|
||||
//
|
||||
// The API server behind it is the controller runtime fake, which costs
|
||||
// milliseconds per reconcile and sets a floor well above what queueing or
|
||||
// skipping saves on a single deletion. So read this as the shape rather than
|
||||
// the size: queued is flat across service latencies and synchronous is not.
|
||||
// What a deletion pays to queue is BenchmarkRunnerUnregistrationQueuePush, and
|
||||
// what it saves by skipping is that same figure.
|
||||
func BenchmarkEphemeralRunnerFinalize(b *testing.B) {
|
||||
b.Run("skipped", func(b *testing.B) {
|
||||
queue := NewRunnerUnregistrationQueue(logr.Discard(), nil, 0)
|
||||
benchmarkFinalize(b, queue, v1alpha1.EphemeralRunnerPhaseSucceeded, nil)
|
||||
})
|
||||
|
||||
for _, latency := range serviceLatencies {
|
||||
b.Run(fmt.Sprintf("queued/latency=%s", latency), func(b *testing.B) {
|
||||
queue := NewRunnerUnregistrationQueue(logr.Discard(), &stubSecretResolver{client: benchmarkActionsClient(latency)}, 0)
|
||||
benchmarkFinalize(b, queue, v1alpha1.EphemeralRunnerPhaseRunning, nil)
|
||||
})
|
||||
}
|
||||
|
||||
for _, latency := range serviceLatencies {
|
||||
b.Run(fmt.Sprintf("synchronous/latency=%s", latency), func(b *testing.B) {
|
||||
actionsClient := benchmarkActionsClient(latency)
|
||||
// The removal ran before any of the local cleanup, so the reconcile
|
||||
// carried it. Issued here in front of a reconcile that queues nothing,
|
||||
// which is the same two pieces of work in the same order.
|
||||
benchmarkFinalize(b, nil, v1alpha1.EphemeralRunnerPhaseRunning, func(ctx context.Context, runnerID int) {
|
||||
_ = actionsClient.RemoveRunner(ctx, int64(runnerID))
|
||||
})
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func benchmarkFinalize(b *testing.B, queue *RunnerUnregistrationQueue, phase v1alpha1.EphemeralRunnerPhase, before func(context.Context, int)) {
|
||||
scheme := benchmarkScheme(b)
|
||||
ctx := context.Background()
|
||||
|
||||
// One untimed pass first. Everything here is measured in fractions of a
|
||||
// millisecond against a cold heap, and whichever variant runs first should
|
||||
// not be charged for warming the process up.
|
||||
warmup, key, _ := newFinalizeBenchmarkReconciler(b, scheme, queue, phase, -1)
|
||||
if _, err := warmup.Reconcile(ctx, ctrl.Request{NamespacedName: key}); err != nil {
|
||||
b.Fatalf("reconcile: %v", err)
|
||||
}
|
||||
|
||||
b.ResetTimer()
|
||||
b.StopTimer()
|
||||
for n := range b.N {
|
||||
reconciler, key, runnerID := newFinalizeBenchmarkReconciler(b, scheme, queue, phase, n)
|
||||
|
||||
b.StartTimer()
|
||||
if before != nil {
|
||||
before(ctx, runnerID)
|
||||
}
|
||||
_, err := reconciler.Reconcile(ctx, ctrl.Request{NamespacedName: key})
|
||||
b.StopTimer()
|
||||
|
||||
if err != nil {
|
||||
b.Fatalf("reconcile: %v", err)
|
||||
}
|
||||
if err := reconciler.Get(ctx, key, new(v1alpha1.EphemeralRunner)); err == nil {
|
||||
b.Fatal("runner was not finalized")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// BenchmarkRunnerUnregistrationQueuePush measures what the reconciler pays to
|
||||
// hand a removal over. It is the whole cost the Actions service imposes on the
|
||||
// finalizer path now, and it must stay flat: Push is called under no lock the
|
||||
// reconciler holds, but every deletion goes through it.
|
||||
func BenchmarkRunnerUnregistrationQueuePush(b *testing.B) {
|
||||
q := NewRunnerUnregistrationQueue(logr.Discard(), nil, 0)
|
||||
runner := newUnregistrationTestRunner("runner", 0, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
|
||||
b.ReportAllocs()
|
||||
runnerID := 0
|
||||
for b.Loop() {
|
||||
// A distinct registration every time, as in production. Pushing one
|
||||
// runner over and over would measure the duplicate being turned away
|
||||
// rather than the cost of queueing.
|
||||
runnerID++
|
||||
runner.Status.RunnerID = runnerID
|
||||
pushTestRunner(q, runner)
|
||||
}
|
||||
b.StopTimer()
|
||||
|
||||
// Nothing drains it here, so the queue holds every push and the claim on
|
||||
// every one of them is never released. Reported so the number above is read
|
||||
// as the cost of a push into a queue that deep.
|
||||
b.ReportMetric(float64(q.len()), "queued")
|
||||
}
|
||||
|
||||
// BenchmarkRunnerUnregistrationQueueDrain measures taking a burst back off the
|
||||
// queue, which is what the workers do when a scale set finishes.
|
||||
//
|
||||
// Reported per request. A burst arrives all at once and is drained in order, so
|
||||
// this has to stay flat as the burst grows: taking from the front by resliding
|
||||
// the tail down would make it climb with the size of the burst, and it climbs
|
||||
// while holding the lock that Push needs.
|
||||
func BenchmarkRunnerUnregistrationQueueDrain(b *testing.B) {
|
||||
for _, burst := range []int{100, 1_000, 10_000} {
|
||||
b.Run(fmt.Sprintf("burst=%d", burst), func(b *testing.B) {
|
||||
q := NewRunnerUnregistrationQueue(logr.Discard(), nil, 0)
|
||||
runner := newUnregistrationTestRunner("runner", 42, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
now := time.Now()
|
||||
|
||||
for b.Loop() {
|
||||
b.StopTimer()
|
||||
for range burst {
|
||||
// Queued directly, because what is being measured is taking
|
||||
// requests back off, and a claim is only given up by the
|
||||
// worker that finishes with one.
|
||||
q.enqueue(runnerUnregistration{runner: runner})
|
||||
}
|
||||
b.StartTimer()
|
||||
|
||||
for range burst {
|
||||
if _, _, ok := q.next(now); !ok {
|
||||
b.Fatal("queue ran dry before the burst was drained")
|
||||
}
|
||||
}
|
||||
}
|
||||
b.ReportMetric(float64(b.Elapsed().Nanoseconds())/float64(b.N*burst), "ns/request")
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,822 @@
|
||||
package actionsgithubcom
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1"
|
||||
"github.com/actions/actions-runner-controller/apis/actions.github.com/v1alpha1/appconfig"
|
||||
"github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient"
|
||||
"github.com/actions/actions-runner-controller/controllers/actions.github.com/multiclient/fake"
|
||||
"github.com/actions/actions-runner-controller/controllers/actions.github.com/object"
|
||||
"github.com/actions/scaleset"
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
ctrlfake "sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
|
||||
"sigs.k8s.io/controller-runtime/pkg/log"
|
||||
)
|
||||
|
||||
// stubSecretResolver hands every caller the same client, so a test can drive
|
||||
// the queue without a Kubernetes API server behind it.
|
||||
type stubSecretResolver struct {
|
||||
client multiclient.Client
|
||||
err error
|
||||
}
|
||||
|
||||
func (s *stubSecretResolver) GetAppConfig(context.Context, object.ActionsGitHubObject) (*appconfig.AppConfig, error) {
|
||||
return nil, s.err
|
||||
}
|
||||
|
||||
func (s *stubSecretResolver) GetActionsService(context.Context, object.ActionsGitHubObject) (multiclient.Client, error) {
|
||||
if s.err != nil {
|
||||
return nil, s.err
|
||||
}
|
||||
return s.client, nil
|
||||
}
|
||||
|
||||
// queued returns everything the queue is still holding, ready requests first.
|
||||
func (q *RunnerUnregistrationQueue) queued() []runnerUnregistration {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
queued := make([]runnerUnregistration, 0, len(q.ready)-q.readyHead+len(q.delayed))
|
||||
queued = append(queued, q.ready[q.readyHead:]...)
|
||||
return append(queued, q.delayed...)
|
||||
}
|
||||
|
||||
// pushTestRunner queues a runner under the ID its own status reports, which is
|
||||
// the case everywhere except a runner deleted before the controller recorded
|
||||
// one.
|
||||
func pushTestRunner(q *RunnerUnregistrationQueue, runner *v1alpha1.EphemeralRunner) {
|
||||
q.Push(runner, runner.Status.RunnerID)
|
||||
}
|
||||
|
||||
func newTestUnregistrationQueue(t *testing.T, client multiclient.Client) *RunnerUnregistrationQueue {
|
||||
t.Helper()
|
||||
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{client: client}, 0)
|
||||
// Retries are timed in tens of seconds in production. A test that waits for
|
||||
// one should not be.
|
||||
q.retryDelay = 5 * time.Millisecond
|
||||
return q
|
||||
}
|
||||
|
||||
// startTestUnregistrationQueue runs the pool for the duration of the test.
|
||||
func startTestUnregistrationQueue(t *testing.T, q *RunnerUnregistrationQueue) {
|
||||
t.Helper()
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
defer close(stopped)
|
||||
require.NoError(t, q.Start(ctx))
|
||||
}()
|
||||
|
||||
t.Cleanup(func() {
|
||||
cancel()
|
||||
select {
|
||||
case <-stopped:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Error("unregistration workers did not stop after the context was cancelled")
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func newUnregistrationTestRunner(name string, runnerID int, phase v1alpha1.EphemeralRunnerPhase) *v1alpha1.EphemeralRunner {
|
||||
return &v1alpha1.EphemeralRunner{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"},
|
||||
Spec: v1alpha1.EphemeralRunnerSpec{
|
||||
GitHubConfigURL: "https://github.com/owner/repo",
|
||||
GitHubConfigSecret: "config-secret",
|
||||
},
|
||||
Status: v1alpha1.EphemeralRunnerStatus{
|
||||
RunnerID: runnerID,
|
||||
Phase: phase,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// TestRunnerSelfDeregistered pins the decision the whole change rests on: a
|
||||
// runner that exited with code 0 is not asked to be removed from the service.
|
||||
func TestRunnerSelfDeregistered(t *testing.T) {
|
||||
tt := map[string]struct {
|
||||
phase v1alpha1.EphemeralRunnerPhase
|
||||
want bool
|
||||
}{
|
||||
"succeeded runner deregistered itself": {phase: v1alpha1.EphemeralRunnerPhaseSucceeded, want: true},
|
||||
"running runner never got to deregister": {phase: v1alpha1.EphemeralRunnerPhaseRunning, want: false},
|
||||
"pending runner is registered but idle": {phase: v1alpha1.EphemeralRunnerPhasePending, want: false},
|
||||
"failed runner never got to deregister": {phase: v1alpha1.EphemeralRunnerPhaseFailed, want: false},
|
||||
"outdated runner never got to deregister": {phase: v1alpha1.EphemeralRunnerPhaseOutdated, want: false},
|
||||
"runner with no phase yet": {want: false},
|
||||
}
|
||||
|
||||
for name, tc := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
runner := newUnregistrationTestRunner("test-runner", 1, tc.phase)
|
||||
assert.Equal(t, tc.want, runnerSelfDeregistered(runner))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegisteredRunnerID(t *testing.T) {
|
||||
tt := map[string]struct {
|
||||
runnerID int
|
||||
secret map[string][]byte
|
||||
secretErr error
|
||||
actionsClient multiclient.Client
|
||||
actionsErr error
|
||||
want int
|
||||
wantErr bool
|
||||
}{
|
||||
"runner reports its own ID": {
|
||||
runnerID: 1,
|
||||
want: 1,
|
||||
},
|
||||
"the status is preferred over the secret": {
|
||||
runnerID: 1,
|
||||
secret: map[string][]byte{"runnerId": []byte("7")},
|
||||
want: 1,
|
||||
},
|
||||
// Registration happens before the status records the ID, so a runner
|
||||
// deleted in between is registered under an ID only the secret knows.
|
||||
"runner deleted before its ID was recorded": {
|
||||
secret: map[string][]byte{"runnerId": []byte("7")},
|
||||
want: 7,
|
||||
},
|
||||
"runner deleted before its JIT secret was created": {
|
||||
actionsClient: fake.NewClient(fake.WithGetRunnerByName(
|
||||
&scaleset.RunnerReference{ID: 7, RunnerScaleSetID: 1},
|
||||
nil,
|
||||
)),
|
||||
want: 7,
|
||||
},
|
||||
"runner without an ID, a secret, or a matching registration was never registered": {
|
||||
actionsClient: fake.NewClient(fake.WithGetRunnerByName(nil, nil)),
|
||||
want: 0,
|
||||
},
|
||||
"runner whose matching registration cannot be looked up": {
|
||||
actionsClient: fake.NewClient(fake.WithGetRunnerByName(
|
||||
nil,
|
||||
errors.New("Actions service is unavailable"),
|
||||
)),
|
||||
wantErr: true,
|
||||
},
|
||||
"runner whose matching registration belongs to another scale set": {
|
||||
actionsClient: fake.NewClient(fake.WithGetRunnerByName(
|
||||
&scaleset.RunnerReference{ID: 7, RunnerScaleSetID: 2},
|
||||
nil,
|
||||
)),
|
||||
wantErr: true,
|
||||
},
|
||||
"runner whose Actions client cannot be resolved": {
|
||||
actionsErr: errors.New("configuration cannot be read"),
|
||||
wantErr: true,
|
||||
},
|
||||
"runner whose secret cannot name a registration": {
|
||||
secret: map[string][]byte{"runnerId": []byte("not-a-number")},
|
||||
want: 0,
|
||||
},
|
||||
// An unreadable secret is not an answer. Reporting 0 would drop the
|
||||
// finalizer and lose the last record of a registration that may exist.
|
||||
"runner whose secret cannot be read": {
|
||||
secret: map[string][]byte{"runnerId": []byte("7")},
|
||||
secretErr: apierrors.NewServiceUnavailable("etcd is unhappy"),
|
||||
wantErr: true,
|
||||
},
|
||||
"the status is answered without reading the secret at all": {
|
||||
runnerID: 1,
|
||||
secretErr: apierrors.NewServiceUnavailable("etcd is unhappy"),
|
||||
want: 1,
|
||||
},
|
||||
}
|
||||
|
||||
for name, tc := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
runner := newUnregistrationTestRunner("test-runner", tc.runnerID, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
runner.Spec.RunnerScaleSetID = 1
|
||||
|
||||
scheme := runtime.NewScheme()
|
||||
require.NoError(t, corev1.AddToScheme(scheme))
|
||||
require.NoError(t, v1alpha1.AddToScheme(scheme))
|
||||
|
||||
builder := ctrlfake.NewClientBuilder().WithScheme(scheme)
|
||||
if tc.secret != nil {
|
||||
builder = builder.WithObjects(&corev1.Secret{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: runner.Name, Namespace: runner.Namespace},
|
||||
Data: tc.secret,
|
||||
})
|
||||
}
|
||||
if tc.secretErr != nil {
|
||||
builder = builder.WithInterceptorFuncs(interceptor.Funcs{
|
||||
Get: func(context.Context, client.WithWatch, client.ObjectKey, client.Object, ...client.GetOption) error {
|
||||
return tc.secretErr
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
reconciler := &EphemeralRunnerReconciler{
|
||||
Client: builder.Build(),
|
||||
ResourceBuilder: ResourceBuilder{
|
||||
SecretResolver: &stubSecretResolver{client: tc.actionsClient, err: tc.actionsErr},
|
||||
},
|
||||
}
|
||||
runnerID, err := reconciler.registeredRunnerID(t.Context(), runner, logr.Discard())
|
||||
if tc.wantErr {
|
||||
require.Error(t, err)
|
||||
return
|
||||
}
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, tc.want, runnerID)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueueUnregistration(t *testing.T) {
|
||||
tt := map[string]struct {
|
||||
patchErr error
|
||||
wantErr bool
|
||||
wantQueued bool
|
||||
}{
|
||||
"finalizer was removed with the runner": {
|
||||
patchErr: apierrors.NewNotFound(schema.GroupResource{Resource: "ephemeralrunners"}, "test-runner"),
|
||||
wantQueued: true,
|
||||
},
|
||||
"finalizer patch failed": {
|
||||
patchErr: errors.New("etcd is unhappy"),
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
|
||||
for name, tc := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
runner := newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
runner.Finalizers = []string{ephemeralRunnerActionsFinalizerName}
|
||||
q := NewRunnerUnregistrationQueue(logr.Discard(), nil, 0)
|
||||
|
||||
c := ctrlfake.NewClientBuilder().
|
||||
WithScheme(runtime.NewScheme()).
|
||||
WithInterceptorFuncs(interceptor.Funcs{
|
||||
Patch: func(context.Context, client.WithWatch, client.Object, client.Patch, ...client.PatchOption) error {
|
||||
return tc.patchErr
|
||||
},
|
||||
}).
|
||||
Build()
|
||||
reconciler := &EphemeralRunnerReconciler{Client: c, UnregistrationQueue: q}
|
||||
|
||||
err := reconciler.queueUnregistration(t.Context(), runner, logr.Discard())
|
||||
if tc.wantErr {
|
||||
require.Error(t, err)
|
||||
} else {
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
queued := q.queued()
|
||||
if tc.wantQueued {
|
||||
require.Len(t, queued, 1)
|
||||
assert.Equal(t, 42, queued[0].runnerID)
|
||||
} else {
|
||||
assert.Empty(t, queued)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewRunnerUnregistrationQueueWorkerCount(t *testing.T) {
|
||||
tt := map[string]struct {
|
||||
workers int
|
||||
want int
|
||||
}{
|
||||
"unset falls back to the floor": {workers: 0, want: 4},
|
||||
"below the floor is raised": {workers: 2, want: 4},
|
||||
"at the floor is kept": {workers: 4, want: 4},
|
||||
"above the floor is kept": {workers: 100, want: 100},
|
||||
"negative falls back to the floor": {workers: -1, want: 4},
|
||||
}
|
||||
|
||||
for name, tc := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, tc.workers)
|
||||
assert.Equal(t, tc.want, q.workers)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueRemovesRunner(t *testing.T) {
|
||||
removed := make(chan int64, 1)
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, runnerID int64) error {
|
||||
removed <- runnerID
|
||||
return nil
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
select {
|
||||
case runnerID := <-removed:
|
||||
assert.Equal(t, int64(42), runnerID)
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("runner was not removed from the service")
|
||||
}
|
||||
|
||||
assert.Eventually(t, func() bool { return q.len() == 0 }, 10*time.Second, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueuePushIsIndependentOfTheService(t *testing.T) {
|
||||
// A service call that never returns must not be able to hold up a push, which
|
||||
// is the property the reconciler depends on.
|
||||
release := make(chan struct{})
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(ctx context.Context, _ int64) error {
|
||||
select {
|
||||
case <-release:
|
||||
case <-ctx.Done():
|
||||
}
|
||||
return nil
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
const runners = 200
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
for i := range runners {
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("test-runner-%d", i), i+1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
}
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("pushing was blocked by the workers")
|
||||
}
|
||||
|
||||
// Every worker is stuck in the call above, so everything but the requests
|
||||
// they took is still queued.
|
||||
assert.GreaterOrEqual(t, q.len(), runners-q.workers)
|
||||
}
|
||||
|
||||
// TestRunnerUnregistrationQueueDoesNotParkBehindDelayedWork pins the property
|
||||
// that makes one shared pool safe to use for both ready and delayed requests: a
|
||||
// worker parked on a retry is parked on an upper bound, not a commitment.
|
||||
//
|
||||
// Without it a handful of runners refused with JobStillRunning would hold the
|
||||
// whole pool for the length of the retry delay, and everything pushed behind
|
||||
// them would sit in the queue waiting on work that has nothing to do with it.
|
||||
func TestRunnerUnregistrationQueueDoesNotParkBehindDelayedWork(t *testing.T) {
|
||||
const ready = 1000
|
||||
|
||||
removed := make(chan int64, ready)
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, runnerID int64) error {
|
||||
if runnerID < ready {
|
||||
// Never succeeds, so these stay in the delayed list for the whole
|
||||
// test and every worker sees them.
|
||||
return fmt.Errorf("removing runner: %w", scaleset.JobStillRunningError)
|
||||
}
|
||||
removed <- runnerID
|
||||
return nil
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
// Long enough that a worker which committed to it would miss the deadline
|
||||
// below by two orders of magnitude.
|
||||
q.retryDelay = 30 * time.Second
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
// Enough refusals to park every worker twice over.
|
||||
for i := range q.workers * 2 {
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("still-running-%d", i), i+1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
}
|
||||
require.Eventually(t, func() bool {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
return len(q.delayed) == q.workers*2
|
||||
}, 10*time.Second, time.Millisecond, "the refused runners never settled into the delayed list")
|
||||
|
||||
for i := range ready {
|
||||
runnerID := ready + i
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("ready-%d", runnerID), runnerID, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
}
|
||||
|
||||
for range ready {
|
||||
select {
|
||||
case <-removed:
|
||||
case <-time.After(20 * time.Second):
|
||||
t.Fatal("ready removals were held up behind runners waiting out a retry")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueRetriesWhileTheJobIsStillRunning(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var calls int
|
||||
succeeded := make(chan struct{})
|
||||
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, _ int64) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
|
||||
calls++
|
||||
if calls < 3 {
|
||||
return fmt.Errorf("removing runner: %w", scaleset.JobStillRunningError)
|
||||
}
|
||||
close(succeeded)
|
||||
return nil
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
select {
|
||||
case <-succeeded:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("runner removal was not retried until it succeeded")
|
||||
}
|
||||
|
||||
assert.Eventually(t, func() bool { return q.len() == 0 }, 10*time.Second, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueDeduplicates(t *testing.T) {
|
||||
runner := newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
|
||||
t.Run("a registration already queued is not queued again", func(t *testing.T) {
|
||||
q := newTestUnregistrationQueue(t, fake.NewClient())
|
||||
|
||||
for range 5 {
|
||||
pushTestRunner(q, runner)
|
||||
}
|
||||
|
||||
assert.Equal(t, 1, q.len())
|
||||
})
|
||||
|
||||
t.Run("a registration being removed right now is not queued again", func(t *testing.T) {
|
||||
// The window the claim covers is wider than the queue itself: a request
|
||||
// a worker has already taken is no longer queued, but the service has
|
||||
// not answered yet, so asking it again is the duplicate call.
|
||||
started := make(chan struct{})
|
||||
release := make(chan struct{})
|
||||
var calls atomic.Int64
|
||||
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(ctx context.Context, _ int64) error {
|
||||
if calls.Add(1) == 1 {
|
||||
close(started)
|
||||
}
|
||||
select {
|
||||
case <-release:
|
||||
case <-ctx.Done():
|
||||
}
|
||||
return nil
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
pushTestRunner(q, runner)
|
||||
select {
|
||||
case <-started:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("the removal was never attempted")
|
||||
}
|
||||
|
||||
pushTestRunner(q, runner)
|
||||
assert.Zero(t, q.len(), "the runner was queued again while it was being removed")
|
||||
})
|
||||
|
||||
t.Run("a retrying registration is not queued alongside itself", func(t *testing.T) {
|
||||
// This is the case that costs something. Two copies of a request for a
|
||||
// runner that is still executing a job would sit in the delayed list
|
||||
// retrying in lockstep for as long as the job runs.
|
||||
var calls atomic.Int64
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, _ int64) error {
|
||||
calls.Add(1)
|
||||
return fmt.Errorf("removing runner: %w", scaleset.JobStillRunningError)
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, runner)
|
||||
assert.Eventually(t, func() bool { return calls.Load() > 0 }, 10*time.Second, time.Millisecond)
|
||||
|
||||
for range 5 {
|
||||
pushTestRunner(q, runner)
|
||||
}
|
||||
|
||||
assert.LessOrEqual(t, q.len(), 1, "the retrying runner was queued more than once")
|
||||
})
|
||||
|
||||
t.Run("the same runner is queued again once the removal is done with", func(t *testing.T) {
|
||||
// The claim is not a memory of everything ever removed. A runner that
|
||||
// comes back around, with a request that was dropped on a failure, is
|
||||
// taken on again.
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, _ int64) error {
|
||||
return errors.New("the service is unhappy")
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, runner)
|
||||
assert.Eventually(t, func() bool { return q.len() == 0 }, 10*time.Second, time.Millisecond)
|
||||
|
||||
// Asked of push rather than of the queue length, because a worker is
|
||||
// free to drain the second request before the check runs.
|
||||
assert.Eventually(t, func() bool {
|
||||
return q.push(runnerUnregistration{runner: runner, runnerID: runner.Status.RunnerID})
|
||||
}, 10*time.Second, time.Millisecond, "the runner could not be queued again after its removal failed")
|
||||
})
|
||||
|
||||
t.Run("a different registration of the same runner is queued", func(t *testing.T) {
|
||||
// A runner deleted before the controller recorded its ID is queued under
|
||||
// the ID recovered from its jitconfig secret, which is a different
|
||||
// registration from whatever its status reports.
|
||||
q := newTestUnregistrationQueue(t, fake.NewClient())
|
||||
|
||||
q.Push(runner, 42)
|
||||
q.Push(runner, 43)
|
||||
|
||||
assert.Equal(t, 2, q.len())
|
||||
})
|
||||
|
||||
t.Run("runner IDs are only unique within their own GitHub scope", func(t *testing.T) {
|
||||
// One controller serves scale sets in different orgs, and the service
|
||||
// hands out runner IDs per scope, so the same ID can name two unrelated
|
||||
// registrations. Dropping one of them would leak it.
|
||||
q := newTestUnregistrationQueue(t, fake.NewClient())
|
||||
|
||||
pushTestRunner(q, newUnregistrationTestRunner("runner-in-one-scale-set", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
pushTestRunner(q, newUnregistrationTestRunner("runner-in-another", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
assert.Equal(t, 2, q.len())
|
||||
})
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueDropsFailedRemovals(t *testing.T) {
|
||||
tt := map[string]error{
|
||||
"runner is already gone": fmt.Errorf("removing runner: %w", scaleset.RunnerNotFoundError),
|
||||
"scale set is gone": fmt.Errorf("removing runner: %w", scaleset.NotFoundError),
|
||||
"service is unavailable": errors.New("boom"),
|
||||
"credentials are invalid": fmt.Errorf("removing runner: %w", scaleset.UnauthorizedError),
|
||||
}
|
||||
|
||||
for name, removeErr := range tt {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
called := make(chan struct{}, 1)
|
||||
client := fake.NewClient(fake.WithRemoveRunnerFunc(func(_ context.Context, _ int64) error {
|
||||
select {
|
||||
case called <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
return removeErr
|
||||
}))
|
||||
|
||||
q := newTestUnregistrationQueue(t, client)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
select {
|
||||
case <-called:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("the service was never called")
|
||||
}
|
||||
|
||||
// Only a runner that is still executing a job is retried. Everything
|
||||
// else is left to the service.
|
||||
assert.Eventually(t, func() bool { return q.len() == 0 }, 10*time.Second, 10*time.Millisecond)
|
||||
assert.Never(t, func() bool { return q.len() > 0 }, 100*time.Millisecond, 10*time.Millisecond)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueSurvivesAnUnresolvableRunner(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{err: errors.New("no such secret")}, 0)
|
||||
startTestUnregistrationQueue(t, q)
|
||||
|
||||
pushTestRunner(q, newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
assert.Eventually(t, func() bool { return q.len() == 0 }, 10*time.Second, 10*time.Millisecond)
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueNext(t *testing.T) {
|
||||
now := time.Now()
|
||||
|
||||
t.Run("reports the idle wait when empty", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
|
||||
_, wait, ok := q.next(now)
|
||||
assert.False(t, ok)
|
||||
assert.Equal(t, unregistrationMaxIdleWait, wait)
|
||||
})
|
||||
|
||||
t.Run("takes requests in order", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
pushTestRunner(q, newUnregistrationTestRunner("first", 1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
pushTestRunner(q, newUnregistrationTestRunner("second", 2, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
request, _, ok := q.next(now)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "first", request.runner.Name)
|
||||
|
||||
request, _, ok = q.next(now)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "second", request.runner.Name)
|
||||
|
||||
_, _, ok = q.next(now)
|
||||
assert.False(t, ok)
|
||||
})
|
||||
|
||||
t.Run("drained requests are not kept alive by the queue", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
for i := range 10 {
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("runner-%d", i), i+1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
}
|
||||
for range 10 {
|
||||
_, _, ok := q.next(now)
|
||||
require.True(t, ok)
|
||||
}
|
||||
|
||||
assert.Equal(t, 0, q.len())
|
||||
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
assert.Equal(t, 0, q.readyHead, "the ready list is reset once it is drained")
|
||||
for _, request := range q.ready[:cap(q.ready)] {
|
||||
assert.Nil(t, request.runner, "a taken request still references its runner")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("interleaved requests compact consumed ready storage", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
const depth = 8
|
||||
for i := range depth {
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("runner-%d", i), i+1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
}
|
||||
|
||||
for i := range 1_000 {
|
||||
pushTestRunner(q, newUnregistrationTestRunner(fmt.Sprintf("new-runner-%d", i), depth+i+1, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
_, _, ok := q.next(now)
|
||||
require.True(t, ok)
|
||||
|
||||
q.mu.Lock()
|
||||
assert.LessOrEqual(t, len(q.ready), 2*depth, "ready storage grew beyond the live queue depth")
|
||||
assert.LessOrEqual(t, cap(q.ready), 3*depth, "ready backing storage grew beyond the live queue depth")
|
||||
assert.Equal(t, depth, len(q.ready)-q.readyHead)
|
||||
q.mu.Unlock()
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("a delayed request does not hold up the ones behind it", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("delayed", 1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, time.Hour)
|
||||
pushTestRunner(q, newUnregistrationTestRunner("ready", 2, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
request, _, ok := q.next(now)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "ready", request.runner.Name)
|
||||
|
||||
// The delayed one is still there, and reports how long it has left.
|
||||
_, wait, ok := q.next(now)
|
||||
assert.False(t, ok)
|
||||
assert.Equal(t, unregistrationMaxIdleWait, wait, "a wait longer than the cap is reported as the cap")
|
||||
assert.Equal(t, 1, q.len())
|
||||
})
|
||||
|
||||
t.Run("a promoted request goes behind the requests already ready", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("retried", 1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, time.Minute)
|
||||
pushTestRunner(q, newUnregistrationTestRunner("ready", 2, v1alpha1.EphemeralRunnerPhaseRunning))
|
||||
|
||||
later := now.Add(2 * time.Minute)
|
||||
for _, want := range []string{"ready", "retried"} {
|
||||
request, _, ok := q.next(later)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, want, request.runner.Name)
|
||||
}
|
||||
|
||||
// A promoted request loses its delay, so it is not held back again.
|
||||
assert.Equal(t, 0, q.len())
|
||||
})
|
||||
|
||||
t.Run("promoting a burst keeps the ones still waiting", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
// Interleaved, so promoting the due ones has to compact around the rest
|
||||
// rather than just cut a prefix off.
|
||||
for i := range 10 {
|
||||
delay := time.Minute
|
||||
if i%2 == 0 {
|
||||
delay = time.Hour
|
||||
}
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner(fmt.Sprintf("runner-%d", i), i+1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, delay)
|
||||
}
|
||||
|
||||
later := now.Add(2 * time.Minute)
|
||||
for i := 1; i < 10; i += 2 {
|
||||
request, _, ok := q.next(later)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, fmt.Sprintf("runner-%d", i), request.runner.Name)
|
||||
}
|
||||
|
||||
_, _, ok := q.next(later)
|
||||
assert.False(t, ok, "the requests due in an hour were promoted early")
|
||||
assert.Equal(t, 5, q.len())
|
||||
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
for _, request := range q.delayed[len(q.delayed):cap(q.delayed)] {
|
||||
assert.Nil(t, request.runner, "a promoted request is still referenced by the delayed list")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("waits only until the earliest request is ready", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("later", 1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, unregistrationMaxIdleWait/2)
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("sooner", 2, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, unregistrationMaxIdleWait/4)
|
||||
|
||||
_, wait, ok := q.next(time.Now())
|
||||
assert.False(t, ok)
|
||||
assert.Greater(t, wait, time.Duration(0))
|
||||
assert.LessOrEqual(t, wait, unregistrationMaxIdleWait/4)
|
||||
})
|
||||
|
||||
t.Run("takes a request once it is ready", func(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("delayed", 1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, time.Minute)
|
||||
|
||||
_, _, ok := q.next(now)
|
||||
assert.False(t, ok)
|
||||
|
||||
request, _, ok := q.next(now.Add(2 * time.Minute))
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "delayed", request.runner.Name)
|
||||
})
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueuePushCopiesTheRunner(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{}, 0)
|
||||
|
||||
runner := newUnregistrationTestRunner("test-runner", 42, v1alpha1.EphemeralRunnerPhaseRunning)
|
||||
pushTestRunner(q, runner)
|
||||
runner.Status.RunnerID = 0
|
||||
runner.Name = "mutated"
|
||||
|
||||
request, _, ok := q.next(time.Now())
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "test-runner", request.runner.Name)
|
||||
assert.Equal(t, 42, request.runner.Status.RunnerID)
|
||||
}
|
||||
|
||||
func TestRunnerUnregistrationQueueStartStopsWithItsContext(t *testing.T) {
|
||||
q := NewRunnerUnregistrationQueue(log.Log, &stubSecretResolver{client: fake.NewClient()}, 0)
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
stopped := make(chan error, 1)
|
||||
go func() { stopped <- q.Start(ctx) }()
|
||||
|
||||
// Anything left behind at shutdown is dropped rather than drained.
|
||||
q.pushAfter(runnerUnregistration{
|
||||
runner: newUnregistrationTestRunner("delayed", 1, v1alpha1.EphemeralRunnerPhaseRunning),
|
||||
}, time.Hour)
|
||||
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case err := <-stopped:
|
||||
assert.NoError(t, err)
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("workers did not stop after the context was cancelled")
|
||||
}
|
||||
}
|
||||
@@ -350,12 +350,28 @@ func main() {
|
||||
}
|
||||
|
||||
ephemeralRunnerOpts := append(controllerOpts, actionsgithubcom.WithMaxConcurrentReconciles(opts.EphemeralRunnerMaxConcurrentReconciles))
|
||||
|
||||
// Removing a runner registration from the Actions service is not part of
|
||||
// deleting the Kubernetes resources an EphemeralRunner owns, and running
|
||||
// it inline puts local pod cleanup behind an external API. These workers
|
||||
// take it off the reconcile path instead.
|
||||
runnerUnregistrationQueue := actionsgithubcom.NewRunnerUnregistrationQueue(
|
||||
log.WithName("RunnerUnregistration").WithValues("version", build.Version),
|
||||
secretResolver,
|
||||
opts.EphemeralRunnerMaxConcurrentReconciles,
|
||||
)
|
||||
if err := mgr.Add(runnerUnregistrationQueue); err != nil {
|
||||
log.Error(err, "unable to add the runner unregistration workers")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
if err = (&actionsgithubcom.EphemeralRunnerReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Log: log.WithName("EphemeralRunner").WithValues("version", build.Version),
|
||||
Scheme: mgr.GetScheme(),
|
||||
PublishMetrics: metricsAddr != "0",
|
||||
ResourceBuilder: rb,
|
||||
Client: mgr.GetClient(),
|
||||
Log: log.WithName("EphemeralRunner").WithValues("version", build.Version),
|
||||
Scheme: mgr.GetScheme(),
|
||||
PublishMetrics: metricsAddr != "0",
|
||||
UnregistrationQueue: runnerUnregistrationQueue,
|
||||
ResourceBuilder: rb,
|
||||
}).SetupWithManager(mgr, ephemeralRunnerOpts...); err != nil {
|
||||
log.Error(err, "unable to create controller", "controller", "EphemeralRunner")
|
||||
os.Exit(1)
|
||||
|
||||
Reference in New Issue
Block a user