diff --git a/main.go b/main.go index eaa18322..139c1ee9 100644 --- a/main.go +++ b/main.go @@ -117,10 +117,10 @@ func run(cmd *cobra.Command, args []string) error { return fmt.Errorf("failed to add the new cluster controller: %v", err) } - logger.Info("adding etcd pod controller") + logger.Info("adding statefulset controller") - if err := cluster.AddPodController(ctx, mgr, maxConcurrentReconciles); err != nil { - return fmt.Errorf("failed to add the new cluster controller: %v", err) + if err := cluster.AddStatefulSetController(ctx, mgr, maxConcurrentReconciles); err != nil { + return fmt.Errorf("failed to add the statefulset controller: %v", err) } logger.Info("adding service controller") diff --git a/pkg/controller/cluster/cluster.go b/pkg/controller/cluster/cluster.go index 136ea69a..36833447 100644 --- a/pkg/controller/cluster/cluster.go +++ b/pkg/controller/cluster/cluster.go @@ -46,7 +46,6 @@ const ( namePrefix = "k3k" clusterController = "k3k-cluster-controller" clusterFinalizerName = "cluster.k3k.io/finalizer" - etcdPodFinalizerName = "etcdpod.k3k.io/finalizer" ClusterInvalidName = "system" defaultVirtualClusterCIDR = "10.52.0.0/16" diff --git a/pkg/controller/cluster/pod.go b/pkg/controller/cluster/statefulset.go similarity index 85% rename from pkg/controller/cluster/pod.go rename to pkg/controller/cluster/statefulset.go index 47f28399..aa5e42f3 100644 --- a/pkg/controller/cluster/pod.go +++ b/pkg/controller/cluster/statefulset.go @@ -15,7 +15,6 @@ import ( "k8s.io/client-go/util/retry" "sigs.k8s.io/controller-runtime/pkg/controller" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" - "sigs.k8s.io/controller-runtime/pkg/handler" "sigs.k8s.io/controller-runtime/pkg/manager" "sigs.k8s.io/controller-runtime/pkg/reconcile" @@ -35,32 +34,34 @@ import ( ) const ( - podController = "k3k-pod-controller" + statefulsetController = "k3k-statefulset-controller" + etcdPodFinalizerName = "etcdpod.k3k.io/finalizer" ) -type PodReconciler struct { +type StatefulSetReconciler struct { Client ctrlruntimeclient.Client Scheme *runtime.Scheme } // Add adds a new controller to the manager -func AddPodController(ctx context.Context, mgr manager.Manager, maxConcurrentReconciles int) error { +func AddStatefulSetController(ctx context.Context, mgr manager.Manager, maxConcurrentReconciles int) error { // initialize a new Reconciler - reconciler := PodReconciler{ + reconciler := StatefulSetReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), } return ctrl.NewControllerManagedBy(mgr). - Watches(&v1.Pod{}, handler.EnqueueRequestForOwner(mgr.GetScheme(), mgr.GetRESTMapper(), &apps.StatefulSet{}, handler.OnlyControllerOwner())). - Named(podController). + For(&apps.StatefulSet{}). + Owns(&v1.Pod{}). + Named(statefulsetController). WithOptions(controller.Options{MaxConcurrentReconciles: maxConcurrentReconciles}). Complete(&reconciler) } -func (p *PodReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { - log := ctrl.LoggerFrom(ctx).WithValues("statefulset", req.NamespacedName) - ctx = ctrl.LoggerInto(ctx, log) // enrich the current logger +func (p *StatefulSetReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { + log := ctrl.LoggerFrom(ctx) + log.Info("reconciling statefulset") s := strings.Split(req.Name, "-") if len(s) < 1 { @@ -102,7 +103,7 @@ func (p *PodReconciler) Reconcile(ctx context.Context, req reconcile.Request) (r return reconcile.Result{}, nil } -func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cluster, pod *v1.Pod) error { +func (p *StatefulSetReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cluster, pod *v1.Pod) error { log := ctrl.LoggerFrom(ctx) log.Info("handling server pod") @@ -120,9 +121,7 @@ func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cl if !pod.DeletionTimestamp.IsZero() { // check if cluster is deleted then remove the finalizer from the pod if cluster.Name == "" { - if controllerutil.ContainsFinalizer(pod, etcdPodFinalizerName) { - controllerutil.RemoveFinalizer(pod, etcdPodFinalizerName) - + if controllerutil.RemoveFinalizer(pod, etcdPodFinalizerName) { if err := p.Client.Update(ctx, pod); err != nil { return err } @@ -166,7 +165,7 @@ func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cl return nil } -func (p *PodReconciler) getETCDTLS(ctx context.Context, cluster *v1alpha1.Cluster) (*tls.Config, error) { +func (p *StatefulSetReconciler) getETCDTLS(ctx context.Context, cluster *v1alpha1.Cluster) (*tls.Config, error) { log := ctrl.LoggerFrom(ctx) log.Info("generating etcd TLS client certificate", "cluster", cluster) @@ -253,7 +252,7 @@ func removePeer(ctx context.Context, client *clientv3.Client, name, address stri return nil } -func (p *PodReconciler) clusterToken(ctx context.Context, cluster *v1alpha1.Cluster) (string, error) { +func (p *StatefulSetReconciler) clusterToken(ctx context.Context, cluster *v1alpha1.Cluster) (string, error) { var tokenSecret v1.Secret nn := types.NamespacedName{ diff --git a/tests/cluster_persistence_test.go b/tests/cluster_persistence_test.go index 4a2cb50f..9bddf3b9 100644 --- a/tests/cluster_persistence_test.go +++ b/tests/cluster_persistence_test.go @@ -4,8 +4,12 @@ import ( "context" "crypto/x509" "errors" + "fmt" "time" + "k8s.io/utils/ptr" + + corev1 "k8s.io/api/core/v1" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" @@ -101,6 +105,72 @@ var _ = When("a dynamic cluster is installed", Label("e2e"), func() { _, _ = virtualCluster.NewNginxPod("") }) + It("can delete the cluster", func() { + ctx := context.Background() + + By("Deleting cluster") + + err := k8sClient.Delete(ctx, virtualCluster.Cluster) + Expect(err).To(Not(HaveOccurred())) + + Eventually(func() []corev1.Pod { + By("listing the pods in the namespace") + + podList, err := k8s.CoreV1().Pods(virtualCluster.Cluster.Namespace).List(ctx, v1.ListOptions{}) + Expect(err).To(Not(HaveOccurred())) + + GinkgoLogr.Info("podlist", "len", len(podList.Items)) + + return podList.Items + }). + WithTimeout(2 * time.Minute). + WithPolling(time.Second). + Should(BeEmpty()) + }) + + It("can delete a HA cluster", func() { + ctx := context.Background() + + namespace := NewNamespace() + + By(fmt.Sprintf("Creating new virtual cluster in namespace %s", namespace.Name)) + + cluster := NewCluster(namespace.Name) + cluster.Spec.Persistence.Type = v1alpha1.DynamicPersistenceMode + cluster.Spec.Servers = ptr.To[int32](2) + + CreateCluster(cluster) + + client, restConfig := NewVirtualK8sClientAndConfig(cluster) + + By(fmt.Sprintf("Created virtual cluster %s/%s", cluster.Namespace, cluster.Name)) + + virtualCluster := &VirtualCluster{ + Cluster: cluster, + RestConfig: restConfig, + Client: client, + } + + By("Deleting cluster") + + err := k8sClient.Delete(ctx, virtualCluster.Cluster) + Expect(err).To(Not(HaveOccurred())) + + Eventually(func() []corev1.Pod { + By("listing the pods in the namespace") + + podList, err := k8s.CoreV1().Pods(virtualCluster.Cluster.Namespace).List(ctx, v1.ListOptions{}) + Expect(err).To(Not(HaveOccurred())) + + GinkgoLogr.Info("podlist", "len", len(podList.Items)) + + return podList.Items + }). + WithTimeout(time.Minute * 3). + WithPolling(time.Second). + Should(BeEmpty()) + }) + It("uses the same bootstrap secret after a restart", func() { ctx := context.Background()