diff --git a/README.md b/README.md index 7b1674ce..11900557 100644 --- a/README.md +++ b/README.md @@ -188,6 +188,18 @@ You can limit Wave to only watch certain namespaces: --namespaces=your-namespace,other-namespace ``` +#### Gated Pod Deletion + +When the [webhooks](#webhooks) are enabled, Wave deletes the Pods of a +StatefulSet that are stuck with its placeholder scheduler once the required +ConfigMaps/Secrets appear. To keep those Pods instead, set: + +``` +--disable-gated-pod-deletion=true +``` + +The StatefulSet then stays in `Pending` until the Pods are deleted manually. + ## Quick Start If you haven't yet got Wave running on your cluster, see @@ -317,6 +329,17 @@ Pods will stay in state `Pending` instead of `ContainerCreating`. When required Secrets/ConfigMaps have been created Wave will restore the scheduler and add the config hash without requiring any restarts. +A Pod's `spec.schedulerName` is immutable, so a Pod that was already created +while scheduling was disabled stays `Pending` forever. Deployments and +DaemonSets replace such Pods themselves, but a StatefulSet using +`podManagementPolicy: OrderedReady` waits for them to become ready and never +rolls out the restored pod template. Wave therefore deletes them once scheduling +is re-enabled, so that the StatefulSet controller recreates them from the +current revision. Only Pods that are owned by the StatefulSet, still use Wave's +placeholder scheduler and have not been scheduled to a node are deleted, and +Wave emits a `GatedPodDeleted` event for each of them. This can be turned off +with [`--disable-gated-pod-deletion`](#gated-pod-deletion). + ## Communication - Found a bug? Please open an issue. diff --git a/charts/wave/templates/clusterrole.yaml b/charts/wave/templates/clusterrole.yaml index 4858de6e..2a786825 100644 --- a/charts/wave/templates/clusterrole.yaml +++ b/charts/wave/templates/clusterrole.yaml @@ -25,6 +25,14 @@ rules: - create - update - patch + - apiGroups: + - "" + resources: + - pods + verbs: + - list + - watch + - delete - apiGroups: - apps resources: diff --git a/charts/wave/templates/deployment.yaml b/charts/wave/templates/deployment.yaml index e070c926..f722873d 100644 --- a/charts/wave/templates/deployment.yaml +++ b/charts/wave/templates/deployment.yaml @@ -44,6 +44,9 @@ spec: {{- if .Values.webhooks.enabled }} - --enable-webhooks=true {{- end }} + {{- if .Values.disableGatedPodDeletion }} + - --disable-gated-pod-deletion=true + {{- end }} volumeMounts: {{- if .Values.webhooks.enabled }} - mountPath: /tmp/k8s-webhook-server/serving-certs diff --git a/charts/wave/values.yaml b/charts/wave/values.yaml index e71ad848..78a23a23 100644 --- a/charts/wave/values.yaml +++ b/charts/wave/values.yaml @@ -50,6 +50,10 @@ serviceAccount: webhooks: enabled: false +# Keep pods of a StatefulSet which are stuck with the placeholder scheduler that +# the webhooks use to hold back pods with missing ConfigMaps/Secrets +disableGatedPodDeletion: false + # Period for reconciliation # syncPeriod: 5m diff --git a/cmd/manager/main.go b/cmd/manager/main.go index eb1c67fa..38ae7f01 100644 --- a/cmd/manager/main.go +++ b/cmd/manager/main.go @@ -52,6 +52,7 @@ var ( showVersion = flag.Bool("version", false, "Show version and exit") enableWebhooks = flag.Bool("enable-webhooks", false, "Enable webhooks") namespaces = flag.String("namespaces", "", "Comma-separated list of namespaces to watch. Defaults to all namespaces.") + disableGatedPodDeletion = flag.Bool("disable-gated-pod-deletion", false, "Do not delete pods of a StatefulSet which are stuck with the placeholder scheduler after scheduling has been re-enabled") setupLog = ctrl.Log.WithName("setup") ) @@ -112,8 +113,9 @@ func main() { // Setup all Controllers setupLog.Info("Setting up controller") controllerConfig := controller.Config{ - UpdateRate: *updateRate, - UpdateBurst: *updateBurst, + UpdateRate: *updateRate, + UpdateBurst: *updateBurst, + DisableGatedPodDeletion: *disableGatedPodDeletion, } if err := controller.AddToManager(mgr, controllerConfig); err != nil { setupLog.Error(err, "unable to register controllers to the manager") diff --git a/config/default/rbac/rbac_role.yaml b/config/default/rbac/rbac_role.yaml index db20be00..9b419b1d 100644 --- a/config/default/rbac/rbac_role.yaml +++ b/config/default/rbac/rbac_role.yaml @@ -155,3 +155,11 @@ rules: - update - patch - delete +- apiGroups: + - "" + resources: + - pods + verbs: + - list + - watch + - delete diff --git a/config/rbac/role.yaml b/config/rbac/role.yaml index d52f5c4e..c0217947 100644 --- a/config/rbac/role.yaml +++ b/config/rbac/role.yaml @@ -19,6 +19,12 @@ rules: - create - patch - update +- resources: + - pods + verbs: + - delete + - list + - watch - apiGroups: - apps resources: diff --git a/hack/run-test-in-minikube.sh b/hack/run-test-in-minikube.sh index bd3619a1..42d687c3 100755 --- a/hack/run-test-in-minikube.sh +++ b/hack/run-test-in-minikube.sh @@ -138,5 +138,61 @@ while ! kubectl get cm test-completed; do fi done +# Without the webhooks scheduling is never disabled, so there is nothing to recover from +if [ "$1" = "production" ]; then + echo Creating a StatefulSet whose Secret does not exist yet... + kubectl create -f - <<'EOF' +apiVersion: apps/v1 +kind: StatefulSet +metadata: + name: test-sts + annotations: + wave.pusher.com/update-on-config-change: "true" +spec: + serviceName: test-sts + podManagementPolicy: OrderedReady + replicas: 1 + selector: + matchLabels: + app: test-sts + template: + metadata: + labels: + app: test-sts + spec: + containers: + - name: test + image: nixery.dev/shell/kubectl + command: ["/bin/sh", "-ec", "sleep infinity"] + volumeMounts: + - name: secret + mountPath: /etc/secret + volumes: + - name: secret + secret: + secretName: test-sts +EOF + + # The webhook holds the pod back with an invalid scheduler. Its + # spec.schedulerName is immutable, so wave has to delete the pod once the + # Secret exists, otherwise the StatefulSet never becomes ready. + kubectl wait --for=create pod/test-sts-0 --timeout=60s + kubectl create secret generic test-sts --from-literal=test=init + + ctr=0 + while [ "$(kubectl get statefulset test-sts -o jsonpath='{.status.readyReplicas}')" != "1" ]; do + echo Waiting for the StatefulSet to become ready + sleep 10 + ctr=$((ctr+1)) + if [ "$ctr" -gt 30 ]; then + echo "StatefulSet did not recover after its Secret was created" + kubectl get pods -o wide + kubectl get pod test-sts-0 -o jsonpath='{.spec.schedulerName}' + kubectl describe statefulset test-sts + exit 1 + fi + done +fi + echo Test passed exit 0 diff --git a/pkg/controller/add_statefulset.go b/pkg/controller/add_statefulset.go index 8b06b3fa..6259d44b 100644 --- a/pkg/controller/add_statefulset.go +++ b/pkg/controller/add_statefulset.go @@ -24,6 +24,6 @@ import ( func init() { // AddToManagerFuncs is a list of functions to create controllers and add them to a manager. AddToManagerFuncs = append(AddToManagerFuncs, func(mgr manager.Manager, cfg Config) error { - return statefulset.Add(mgr, cfg.UpdateRate, cfg.UpdateBurst) + return statefulset.Add(mgr, cfg.UpdateRate, cfg.UpdateBurst, cfg.DisableGatedPodDeletion) }) } diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 97491cf5..a7ac9ee4 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -24,6 +24,10 @@ import ( type Config struct { UpdateRate float64 // updates per second UpdateBurst int // maximum burst size + + // DisableGatedPodDeletion keeps Pods which are stuck with Wave's placeholder + // scheduler instead of deleting them + DisableGatedPodDeletion bool } // AddToManagerFuncs is a list of functions to add all Controllers to the Manager diff --git a/pkg/controller/statefulset/statefulset_controller.go b/pkg/controller/statefulset/statefulset_controller.go index 4186532f..d10a11e1 100644 --- a/pkg/controller/statefulset/statefulset_controller.go +++ b/pkg/controller/statefulset/statefulset_controller.go @@ -29,20 +29,23 @@ import ( // +kubebuilder:rbac:groups=apps,resources=statefulsets,verbs=get;list;watch;update;patch // +kubebuilder:rbac:groups=,resources=configmaps,verbs=get;list;watch;update;patch // +kubebuilder:rbac:groups=,resources=secrets,verbs=get;list;watch;update;patch +// +kubebuilder:rbac:groups=,resources=pods,verbs=list;watch;delete // +kubebuilder:rbac:groups=,resources=events,verbs=create;update;patch // Add creates a new StatefulSet Controller and adds it to the Manager with default RBAC. The Manager will set fields on the Controller // and Start it when the Manager is Started. -func Add(mgr manager.Manager, updateRate float64, updateBurst int) error { - r := newReconciler(mgr, updateRate, updateBurst) +func Add(mgr manager.Manager, updateRate float64, updateBurst int, disableGatedPodDeletion bool) error { + r := newReconciler(mgr, updateRate, updateBurst, disableGatedPodDeletion) return add(mgr, r, r.handler) } // newReconciler returns a new reconcile.Reconciler -func newReconciler(mgr manager.Manager, updateRate float64, updateBurst int) *ReconcileStatefulSet { +func newReconciler(mgr manager.Manager, updateRate float64, updateBurst int, disableGatedPodDeletion bool) *ReconcileStatefulSet { + handler := core.NewHandler[*appsv1.StatefulSet](mgr.GetClient(), mgr.GetEventRecorderFor("wave"), updateRate, updateBurst) + handler.DisableGatedPodDeletion = disableGatedPodDeletion return &ReconcileStatefulSet{ scheme: mgr.GetScheme(), - handler: core.NewHandler[*appsv1.StatefulSet](mgr.GetClient(), mgr.GetEventRecorderFor("wave"), updateRate, updateBurst), + handler: handler, } } diff --git a/pkg/controller/statefulset/statefulset_controller_suite_test.go b/pkg/controller/statefulset/statefulset_controller_suite_test.go index 06b0e641..4413a14b 100644 --- a/pkg/controller/statefulset/statefulset_controller_suite_test.go +++ b/pkg/controller/statefulset/statefulset_controller_suite_test.go @@ -137,7 +137,7 @@ var _ = BeforeSuite(func() { m = utils.Matcher{Client: c} var recFn reconcile.Reconciler - r := newReconciler(mgr, math.Inf(1), 1) + r := newReconciler(mgr, math.Inf(1), 1, false) recFn, requestsStart, requests = core.SetupControllerTestReconcile(r) Expect(add(mgr, recFn, r.handler)).NotTo(HaveOccurred()) diff --git a/pkg/core/controller_suite.go b/pkg/core/controller_suite.go index 5e253c1a..d3b0352a 100644 --- a/pkg/core/controller_suite.go +++ b/pkg/core/controller_suite.go @@ -75,20 +75,6 @@ func ControllerTestSuite[I InstanceType]( const modified = "modified" - var waitForInstanceReconciled = func(obj Object, times int) { - request := reconcile.Request{ - NamespacedName: types.NamespacedName{ - Name: obj.GetName(), - Namespace: obj.GetNamespace(), - }, - } - for range times { - // wait for reconcile for creating the DaemonSet - Eventually(*requestsStart, timeout).Should(Receive(Equal(request))) - Eventually(*requests, timeout).Should(Receive(Equal(request))) - } - } - var consistentlyInstanceNotReconciled = func(obj Object) { request := reconcile.Request{ NamespacedName: types.NamespacedName{ @@ -156,6 +142,7 @@ func ControllerTestSuite[I InstanceType]( &appsv1.DaemonSetList{}, &appsv1.DeploymentList{}, &appsv1.StatefulSetList{}, + &corev1.PodList{}, &corev1.ConfigMapList{}, &corev1.SecretList{}, &corev1.EventList{}, @@ -175,7 +162,7 @@ func ControllerTestSuite[I InstanceType]( // Create a instance and wait for it to be reconciled expectNoReconciles() m.Create(instance).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) }) Context("And it has the required annotation", func() { @@ -191,7 +178,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(instance, addAnnotation).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -236,7 +223,7 @@ func ControllerTestSuite[I InstanceType]( podTemplate.Spec.Containers = []corev1.Container{podTemplate.Spec.Containers[0]} SetPodTemplate(instance, podTemplate) Expect(m.Client.Update(context.TODO(), instance)).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -279,7 +266,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(cm1, modifyCM).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Reschedules once since we update the hash + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Reschedules once since we update the hash // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -301,7 +288,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(cm2, modifyCM).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Reschedules once since we update the hash + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Reschedules once since we update the hash // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -326,7 +313,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(s1, modifyS).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Reschedules once since we update the hash + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Reschedules once since we update the hash // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -351,7 +338,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(s2, modifyS).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Reschedules once since we update the hash + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Reschedules once since we update the hash // Get the updated instance m.Get(instance, timeout).Should(Succeed()) @@ -373,7 +360,7 @@ func ControllerTestSuite[I InstanceType]( } expectNoReconciles() m.Update(instance, removeAnnotations).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) m.Get(instance).Should(Succeed()) Eventually(instance, timeout).ShouldNot(utils.WithAnnotations(HaveKey(RequiredAnnotation))) }) @@ -401,7 +388,7 @@ func ControllerTestSuite[I InstanceType]( Eventually(instance, timeout).Should(utils.WithPodTemplateAnnotations(HaveKey(ConfigHashAnnotation))) expectNoReconciles() m.Delete(instance).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) }) It("Not longer exists", func() { m.Get(instance).Should(MatchError(MatchRegexp(`not found`))) @@ -455,11 +442,14 @@ func ControllerTestSuite[I InstanceType]( } annotations[RequiredAnnotation] = "true" instance.SetAnnotations(annotations) + }) + // JustBeforeEach so that nested contexts can still adjust the instance + JustBeforeEach(func() { // Create a instance and wait for it to be reconciled expectNoReconciles() m.Create(instance).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) }) It("Has scheduling disabled", func() { @@ -469,11 +459,11 @@ func ControllerTestSuite[I InstanceType]( }) Context("And the missing child is created", func() { - BeforeEach(func() { + JustBeforeEach(func() { expectNoReconciles() cm1 = utils.ExampleConfigMap1.DeepCopy() m.Create(cm1).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Two since updating the scheduler self-triggers + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Two since updating the scheduler self-triggers }) It("Has scheduling renabled", func() { @@ -483,6 +473,67 @@ func ControllerTestSuite[I InstanceType]( }) }) + Context("And pods were created while scheduling was disabled", func() { + var stuckPod *corev1.Pod + var scheduledPod *corev1.Pod + var unownedPod *corev1.Pod + + JustBeforeEach(func() { + m.Get(instance, timeout).Should(Succeed()) + + // A pod as the controller creates it from the gated revision + stuckPod = utils.MakePod("stuck", instance, kindOf(instance), SchedulingDisabledSchedulerName) + // A gated pod which a scheduler has placed already + scheduledPod = utils.MakePod("stuck-scheduled", instance, kindOf(instance), SchedulingDisabledSchedulerName) + scheduledPod.Spec.NodeName = "node1" + // A gated pod which belongs to somebody else + unownedPod = utils.MakePod("stuck-unowned", instance, kindOf(instance), SchedulingDisabledSchedulerName) + unownedPod.OwnerReferences = nil + + for _, pod := range []*corev1.Pod{stuckPod, scheduledPod, unownedPod} { + m.Create(pod).Should(Succeed()) + m.Get(pod, timeout).Should(Succeed()) + } + + // Creating the missing child re-enables scheduling + expectNoReconciles() + cm1 = utils.ExampleConfigMap1.DeepCopy() + m.Create(cm1).Should(Succeed()) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Two since updating the scheduler self-triggers + }) + + It("Keeps pods which are scheduled or not owned by the instance", func() { + m.Get(scheduledPod).Should(Succeed()) + Expect(scheduledPod.GetDeletionTimestamp()).To(BeNil()) + m.Get(unownedPod).Should(Succeed()) + Expect(unownedPod.GetDeletionTimestamp()).To(BeNil()) + }) + + It("Deletes the stuck pod of a StatefulSet", func() { + if _, isStatefulSet := any(instance).(*appsv1.StatefulSet); !isStatefulSet { + // Deployments and DaemonSets replace the stuck pods themselves + m.Get(stuckPod).Should(Succeed()) + Expect(stuckPod.GetDeletionTimestamp()).To(BeNil()) + return + } + m.Get(stuckPod, timeout).Should(MatchError(MatchRegexp(`not found`))) + }) + + Context("And the controller replaces the stuck pods itself", func() { + BeforeEach(func() { + // Only a StatefulSet needs help, and only with OrderedReady + if statefulSet, ok := any(instance).(*appsv1.StatefulSet); ok { + statefulSet.Spec.PodManagementPolicy = appsv1.ParallelPodManagement + } + }) + + It("Keeps the stuck pod", func() { + m.Get(stuckPod).Should(Succeed()) + Expect(stuckPod.GetDeletionTimestamp()).To(BeNil()) + }) + }) + }) + }) Context("When a instance with missing children in projection is reconciled", func() { @@ -499,7 +550,7 @@ func ControllerTestSuite[I InstanceType]( // Create a instance and wait for it to be reconciled expectNoReconciles() m.Create(instance).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) }) It("Has scheduling disabled", func() { @@ -513,7 +564,7 @@ func ControllerTestSuite[I InstanceType]( expectNoReconciles() cm6Alt := utils.ExampleConfigMap6WithoutKey3.DeepCopy() m.Create(cm6Alt).Should(Succeed()) - waitForInstanceReconciled(instance, 1) + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 1) }) It("Has Scheduling still disabled", func() { m.Get(instance, timeout).Should(Succeed()) @@ -527,7 +578,7 @@ func ControllerTestSuite[I InstanceType]( cm6Copy := utils.ExampleConfigMap6.DeepCopy() cm6.Data = cm6Copy.Data Expect(m.Client.Update(context.TODO(), cm6)).Should(Succeed()) - waitForInstanceReconciled(instance, 2) // Reschedules once since we update the hash + reenables scheduling + utils.ConsumeReconciles(*requestsStart, *requests, instance, timeout, 2) // Reschedules once since we update the hash + reenables scheduling }) It("Has scheduling renabled", func() { diff --git a/pkg/core/gated_pods.go b/pkg/core/gated_pods.go new file mode 100644 index 00000000..c14a18ff --- /dev/null +++ b/pkg/core/gated_pods.go @@ -0,0 +1,86 @@ +/* +Copyright 2018 Pusher Ltd. and Wave Contributors + +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 core + +import ( + "context" + "fmt" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/client" + logf "sigs.k8s.io/controller-runtime/pkg/log" +) + +// deleteStuckUnschedulablePods deletes the Pods of the StatefulSet that were +// created while Wave had scheduling disabled. +// +// A Pod's spec.schedulerName is immutable, so restoring the scheduler in the +// pod template leaves those Pods waiting forever for a scheduler that does not +// exist. With OrderedReady the StatefulSet controller waits for them to become +// ready before it rolls out the restored revision, so the StatefulSet stays +// wedged until they are deleted and recreated from the current revision. +func (h *Handler[I]) deleteStuckUnschedulablePods(ctx context.Context, statefulSet *appsv1.StatefulSet) error { + log := logf.Log.WithName("wave").WithValues("namespace", statefulSet.GetNamespace(), "name", statefulSet.GetName()) + + // The cache only holds the namespaces given by --namespaces + pods := &corev1.PodList{} + if err := h.List(ctx, pods, client.InNamespace(statefulSet.GetNamespace())); err != nil { + return fmt.Errorf("error listing pods: %v", err) + } + + for i := range pods.Items { + pod := &pods.Items[i] + if !isGatedPodOf(pod, statefulSet) { + continue + } + + // The UID precondition prevents a stale cache from deleting a healthy + // Pod which has replaced the gated one + uid := pod.GetUID() + if err := h.Delete(ctx, pod, client.Preconditions{UID: &uid}); err != nil { + if errors.IsNotFound(err) || errors.IsConflict(err) { + continue + } + return fmt.Errorf("error deleting pod %s/%s: %v", pod.GetNamespace(), pod.GetName(), err) + } + + log.V(0).Info("Deleted pod which was created while scheduling was disabled", "pod", pod.GetName()) + h.recorder.Eventf(statefulSet, corev1.EventTypeNormal, "GatedPodDeleted", "Deleted Pod %s which was created while scheduling was disabled", pod.GetName()) + } + + return nil +} + +// isGatedPodOf returns whether the Pod belongs to the StatefulSet and is still +// held back by the scheduler Wave uses to disable scheduling +func isGatedPodOf(pod *corev1.Pod, statefulSet *appsv1.StatefulSet) bool { + if pod.Spec.SchedulerName != SchedulingDisabledSchedulerName { + return false + } + // Never touch a Pod that has been placed on a node or is already going away + if pod.Spec.NodeName != "" || pod.GetDeletionTimestamp() != nil { + return false + } + for _, ref := range pod.GetOwnerReferences() { + if ref.UID == statefulSet.GetUID() { + return true + } + } + return false +} diff --git a/pkg/core/handler.go b/pkg/core/handler.go index 143e2123..02c6dd46 100644 --- a/pkg/core/handler.go +++ b/pkg/core/handler.go @@ -22,6 +22,7 @@ import ( "sync" "golang.org/x/time/rate" + appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/types" @@ -38,6 +39,10 @@ type Handler[I InstanceType] struct { watchedConfigmaps WatcherList watchedSecrets WatcherList updateThrottler *UpdateThrottler + + // DisableGatedPodDeletion keeps Pods which are stuck with Wave's placeholder + // scheduler instead of deleting them + DisableGatedPodDeletion bool } // NewHandler constructs a new instance of Handler @@ -144,6 +149,23 @@ func (h *Handler[I]) handlePodController(ctx context.Context, instance I) (recon return reconcile.Result{}, fmt.Errorf("error updating instance %s/%s: %v", instance.GetNamespace(), instance.GetName(), err) } } + + // Restoring the scheduler above only changed the pod template. A Pod's + // spec.schedulerName is immutable, so Pods that were created from the + // disabled template stay unschedulable and have to be deleted to be + // recreated with the restored scheduler. + if schedulingChange && !h.DisableGatedPodDeletion { + // Deployments, DaemonSets and StatefulSets using podManagementPolicy + // Parallel replace those Pods themselves. An unset policy defaults to + // OrderedReady. + statefulSet, isStatefulSet := any(instance).(*appsv1.StatefulSet) + if isStatefulSet && statefulSet.Spec.PodManagementPolicy != appsv1.ParallelPodManagement { + if err := h.deleteStuckUnschedulablePods(ctx, statefulSet); err != nil { + return reconcile.Result{}, err + } + } + } + return reconcile.Result{}, nil } diff --git a/test/utils/reconciles.go b/test/utils/reconciles.go new file mode 100644 index 00000000..94582750 --- /dev/null +++ b/test/utils/reconciles.go @@ -0,0 +1,53 @@ +/* +Copyright 2018 Pusher Ltd. and Wave Contributors + +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 utils + +import ( + "time" + + "github.com/onsi/gomega" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/reconcile" +) + +// ConsumeReconciles waits for the given number of reconciles of the object to finish +func ConsumeReconciles(started <-chan reconcile.Request, finished <-chan reconcile.Request, obj client.Object, timeout time.Duration, times int) { + request := reconcile.Request{ + NamespacedName: types.NamespacedName{ + Name: obj.GetName(), + Namespace: obj.GetNamespace(), + }, + } + for range times { + gomega.Eventually(started, timeout).Should(gomega.Receive(gomega.Equal(request))) + gomega.Eventually(finished, timeout).Should(gomega.Receive(gomega.Equal(request))) + } +} + +// DrainReconciles consumes reconciles until none are triggered anymore. The +// channels are unbuffered, so the reconciler blocks until they are read. +func DrainReconciles(started <-chan reconcile.Request, finished <-chan reconcile.Request) { + for { + select { + case <-started: + <-finished + case <-time.After(time.Millisecond * 100): + return + } + } +} diff --git a/test/utils/test_objects.go b/test/utils/test_objects.go index 036af844..e1a6d765 100644 --- a/test/utils/test_objects.go +++ b/test/utils/test_objects.go @@ -604,3 +604,52 @@ var ExampleSecret6 = &corev1.Secret{ "key3": "example6:key3", }, } + +// PodSpecWithSecret1 returns a pod spec which references Secret example1 as its +// only child, so tests do not have to create all the example ConfigMaps and Secrets +func PodSpecWithSecret1() corev1.PodSpec { + return corev1.PodSpec{ + SchedulerName: "default-scheduler", + Volumes: []corev1.Volume{ + { + Name: "secret1", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{ + SecretName: "example1", + }, + }, + }, + }, + Containers: []corev1.Container{ + { + Name: "container1", + Image: "container1", + }, + }, + } +} + +// MakePod builds a Pod controlled by the given apps/v1 owner, as its controller +// would create it from a pod template using the given scheduler +func MakePod(name string, owner metav1.Object, ownerKind string, schedulerName string) *corev1.Pod { + pod := &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: owner.GetNamespace(), + Labels: labels, + OwnerReferences: []metav1.OwnerReference{ + { + APIVersion: "apps/v1", + Kind: ownerKind, + Name: owner.GetName(), + UID: owner.GetUID(), + Controller: &trueValue, + BlockOwnerDeletion: &trueValue, + }, + }, + }, + Spec: PodSpecWithSecret1(), + } + pod.Spec.SchedulerName = schedulerName + return pod +}