diff --git a/main.go b/main.go index 3b721953..1580ffaf 100644 --- a/main.go +++ b/main.go @@ -109,24 +109,25 @@ func run(clx *cli.Context) error { } ctrlruntimelog.SetLogger(zapr.NewLogger(logger.Desugar().WithOptions(zap.AddCallerSkip(1)))) + logger.Info("adding cluster controller") - if err := cluster.Add(ctx, mgr, sharedAgentImage, sharedAgentImagePullPolicy, logger); err != nil { + if err := cluster.Add(ctx, mgr, sharedAgentImage, sharedAgentImagePullPolicy); err != nil { return fmt.Errorf("failed to add the new cluster controller: %v", err) } logger.Info("adding etcd pod controller") - if err := cluster.AddPodController(ctx, mgr, logger); err != nil { + if err := cluster.AddPodController(ctx, mgr); err != nil { return fmt.Errorf("failed to add the new cluster controller: %v", err) } logger.Info("adding clusterset controller") - if err := clusterset.Add(ctx, mgr, clusterCIDR, logger); err != nil { + if err := clusterset.Add(ctx, mgr, clusterCIDR); err != nil { return fmt.Errorf("failed to add the clusterset controller: %v", err) } if clusterCIDR == "" { logger.Info("adding networkpolicy node controller") - if err := clusterset.AddNodeController(ctx, mgr, logger); err != nil { + if err := clusterset.AddNodeController(ctx, mgr); err != nil { return fmt.Errorf("failed to add the clusterset node controller: %v", err) } } diff --git a/pkg/controller/cluster/cluster.go b/pkg/controller/cluster/cluster.go index a5fcae29..9a122914 100644 --- a/pkg/controller/cluster/cluster.go +++ b/pkg/controller/cluster/cluster.go @@ -13,8 +13,6 @@ import ( "github.com/rancher/k3k/pkg/controller/cluster/agent" "github.com/rancher/k3k/pkg/controller/cluster/server" "github.com/rancher/k3k/pkg/controller/cluster/server/bootstrap" - "github.com/rancher/k3k/pkg/log" - "go.uber.org/zap" v1 "k8s.io/api/core/v1" rbacv1 "k8s.io/api/rbac/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -51,12 +49,10 @@ type ClusterReconciler struct { Scheme *runtime.Scheme SharedAgentImage string SharedAgentImagePullPolicy string - logger *log.Logger } // Add adds a new controller to the manager -func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgentImagePullPolicy string, logger *log.Logger) error { - +func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgentImagePullPolicy string) error { discoveryClient, err := discovery.NewDiscoveryClientForConfig(mgr.GetConfig()) if err != nil { return err @@ -69,7 +65,6 @@ func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgent Scheme: mgr.GetScheme(), SharedAgentImage: sharedAgentImage, SharedAgentImagePullPolicy: sharedAgentImagePullPolicy, - logger: logger.Named(clusterController), } return ctrl.NewControllerManagedBy(mgr). @@ -81,78 +76,65 @@ func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgent } func (c *ClusterReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { - var ( - cluster v1alpha1.Cluster - podList v1.PodList - ) - log := c.logger.With("Cluster", req.NamespacedName) + log := ctrl.LoggerFrom(ctx).WithValues("cluster", req.NamespacedName) + ctx = ctrl.LoggerInto(ctx, log) // enrich the current logger + + log.Info("reconciling cluster") + + var cluster v1alpha1.Cluster if err := c.Client.Get(ctx, req.NamespacedName, &cluster); err != nil { - return reconcile.Result{}, ctrlruntimeclient.IgnoreNotFound(err) + return reconcile.Result{}, err } - // if the Version is not specified we will try to use the same Kubernetes version of the host. - // This version is stored in the Status object, and it will not be updated if already set. - if cluster.Spec.Version == "" && cluster.Status.HostVersion == "" { - hostVersion, err := c.DiscoveryClient.ServerVersion() - if err != nil { + // if DeletionTimestamp is not Zero -> finalize the object + if !cluster.DeletionTimestamp.IsZero() { + return c.finalizeCluster(ctx, cluster) + } + + // add finalizers + if !controllerutil.AddFinalizer(&cluster, clusterFinalizerName) { + if err := c.Client.Update(ctx, &cluster); err != nil { return reconcile.Result{}, err } + } - // update Status HostVersion - k8sVersion := strings.Split(hostVersion.GitVersion, "+")[0] - cluster.Status.HostVersion = k8sVersion + "-k3s1" + orig := cluster.DeepCopy() + + reconcilerErr := c.reconcileCluster(ctx, &cluster) + + // update Status if needed + if !reflect.DeepEqual(orig.Status, cluster.Status) { if err := c.Client.Status().Update(ctx, &cluster); err != nil { return reconcile.Result{}, err } } - if cluster.DeletionTimestamp.IsZero() { - if !controllerutil.ContainsFinalizer(&cluster, clusterFinalizerName) { - controllerutil.AddFinalizer(&cluster, clusterFinalizerName) - if err := c.Client.Update(ctx, &cluster); err != nil { - return reconcile.Result{}, err - } - } - log.Info("enqueue cluster") - return reconcile.Result{}, c.createCluster(ctx, &cluster, log) - } - - // remove finalizer from the server pods and update them. - matchingLabels := ctrlruntimeclient.MatchingLabels(map[string]string{"role": "server"}) - listOpts := &ctrlruntimeclient.ListOptions{Namespace: cluster.Namespace} - matchingLabels.ApplyToList(listOpts) - if err := c.Client.List(ctx, &podList, listOpts); err != nil { - return reconcile.Result{}, ctrlruntimeclient.IgnoreNotFound(err) - } - for _, pod := range podList.Items { - if controllerutil.ContainsFinalizer(&pod, etcdPodFinalizerName) { - controllerutil.RemoveFinalizer(&pod, etcdPodFinalizerName) - if err := c.Client.Update(ctx, &pod); err != nil { - return reconcile.Result{}, err - } - } - } - - if err := c.unbindNodeProxyClusterRole(ctx, &cluster); err != nil { - return reconcile.Result{}, err - } - - if controllerutil.ContainsFinalizer(&cluster, clusterFinalizerName) { - // remove finalizer from the cluster and update it. - controllerutil.RemoveFinalizer(&cluster, clusterFinalizerName) - if err := c.Client.Update(ctx, &cluster); err != nil { - return reconcile.Result{}, err - } - } - log.Info("deleting cluster") - return reconcile.Result{}, nil + return reconcile.Result{}, reconcilerErr } -func (c *ClusterReconciler) createCluster(ctx context.Context, cluster *v1alpha1.Cluster, log *zap.SugaredLogger) error { +func (c *ClusterReconciler) reconcileCluster(ctx context.Context, cluster *v1alpha1.Cluster) error { + log := ctrl.LoggerFrom(ctx) + + // if the Version is not specified we will try to use the same Kubernetes version of the host. + // This version is stored in the Status object, and it will not be updated if already set. + if cluster.Spec.Version == "" && cluster.Status.HostVersion == "" { + log.V(1).Info("cluster version not set") + + hostVersion, err := c.DiscoveryClient.ServerVersion() + if err != nil { + return err + } + + // update Status HostVersion + k8sVersion := strings.Split(hostVersion.GitVersion, "+")[0] + cluster.Status.HostVersion = k8sVersion + "-k3s1" + } + if err := c.validate(cluster); err != nil { - log.Errorw("invalid change", zap.Error(err)) + log.Error(err, "invalid change") return nil } + token, err := c.token(ctx, cluster) if err != nil { return err @@ -343,30 +325,6 @@ func (c *ClusterReconciler) bindNodeProxyClusterRole(ctx context.Context, cluste return c.Client.Update(ctx, clusterRoleBinding) } -func (c *ClusterReconciler) unbindNodeProxyClusterRole(ctx context.Context, cluster *v1alpha1.Cluster) error { - clusterRoleBinding := &rbacv1.ClusterRoleBinding{} - if err := c.Client.Get(ctx, types.NamespacedName{Name: "k3k-node-proxy"}, clusterRoleBinding); err != nil { - return fmt.Errorf("failed to get or find k3k-node-proxy ClusterRoleBinding: %w", err) - } - - subjectName := controller.SafeConcatNameWithPrefix(cluster.Name, agent.SharedNodeAgentName) - - var cleanedSubjects []rbacv1.Subject - for _, subject := range clusterRoleBinding.Subjects { - if subject.Name != subjectName || subject.Namespace != cluster.Namespace { - cleanedSubjects = append(cleanedSubjects, subject) - } - } - - // if no subject was removed, all good - if reflect.DeepEqual(clusterRoleBinding.Subjects, cleanedSubjects) { - return nil - } - - clusterRoleBinding.Subjects = cleanedSubjects - return c.Client.Update(ctx, clusterRoleBinding) -} - func (c *ClusterReconciler) agent(ctx context.Context, cluster *v1alpha1.Cluster, serviceIP, token string) error { agent := agent.New(cluster, serviceIP, c.SharedAgentImage, c.SharedAgentImagePullPolicy, token) agentsConfig := agent.Config() diff --git a/pkg/controller/cluster/cluster_finalize.go b/pkg/controller/cluster/cluster_finalize.go new file mode 100644 index 00000000..99c1c651 --- /dev/null +++ b/pkg/controller/cluster/cluster_finalize.go @@ -0,0 +1,79 @@ +package cluster + +import ( + "context" + "fmt" + "reflect" + + "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" + "github.com/rancher/k3k/pkg/controller" + "github.com/rancher/k3k/pkg/controller/cluster/agent" + v1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/reconcile" +) + +func (c *ClusterReconciler) finalizeCluster(ctx context.Context, cluster v1alpha1.Cluster) (reconcile.Result, error) { + log := ctrl.LoggerFrom(ctx) + log.Info("finalizing Cluster") + + // remove finalizer from the server pods and update them. + matchingLabels := ctrlruntimeclient.MatchingLabels(map[string]string{"role": "server"}) + listOpts := &ctrlruntimeclient.ListOptions{Namespace: cluster.Namespace} + matchingLabels.ApplyToList(listOpts) + + var podList v1.PodList + if err := c.Client.List(ctx, &podList, listOpts); err != nil { + return reconcile.Result{}, ctrlruntimeclient.IgnoreNotFound(err) + } + + for _, pod := range podList.Items { + if controllerutil.ContainsFinalizer(&pod, etcdPodFinalizerName) { + controllerutil.RemoveFinalizer(&pod, etcdPodFinalizerName) + if err := c.Client.Update(ctx, &pod); err != nil { + return reconcile.Result{}, err + } + } + } + + if err := c.unbindNodeProxyClusterRole(ctx, &cluster); err != nil { + return reconcile.Result{}, err + } + + if controllerutil.ContainsFinalizer(&cluster, clusterFinalizerName) { + // remove finalizer from the cluster and update it. + controllerutil.RemoveFinalizer(&cluster, clusterFinalizerName) + if err := c.Client.Update(ctx, &cluster); err != nil { + return reconcile.Result{}, err + } + } + return reconcile.Result{}, nil +} + +func (c *ClusterReconciler) unbindNodeProxyClusterRole(ctx context.Context, cluster *v1alpha1.Cluster) error { + clusterRoleBinding := &rbacv1.ClusterRoleBinding{} + if err := c.Client.Get(ctx, types.NamespacedName{Name: "k3k-node-proxy"}, clusterRoleBinding); err != nil { + return fmt.Errorf("failed to get or find k3k-node-proxy ClusterRoleBinding: %w", err) + } + + subjectName := controller.SafeConcatNameWithPrefix(cluster.Name, agent.SharedNodeAgentName) + + var cleanedSubjects []rbacv1.Subject + for _, subject := range clusterRoleBinding.Subjects { + if subject.Name != subjectName || subject.Namespace != cluster.Namespace { + cleanedSubjects = append(cleanedSubjects, subject) + } + } + + // if no subject was removed, all good + if reflect.DeepEqual(clusterRoleBinding.Subjects, cleanedSubjects) { + return nil + } + + clusterRoleBinding.Subjects = cleanedSubjects + return c.Client.Update(ctx, clusterRoleBinding) +} diff --git a/pkg/controller/cluster/cluster_suite_test.go b/pkg/controller/cluster/cluster_suite_test.go index 076ed7e8..2e512211 100644 --- a/pkg/controller/cluster/cluster_suite_test.go +++ b/pkg/controller/cluster/cluster_suite_test.go @@ -5,9 +5,9 @@ import ( "path/filepath" "testing" + "github.com/go-logr/zapr" "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" "github.com/rancher/k3k/pkg/controller/cluster" - "github.com/rancher/k3k/pkg/log" "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" @@ -53,11 +53,13 @@ var _ = BeforeSuite(func() { k8sClient, err = client.New(cfg, client.Options{Scheme: scheme}) Expect(err).NotTo(HaveOccurred()) + ctrl.SetLogger(zapr.NewLogger(zap.NewNop())) + mgr, err := ctrl.NewManager(cfg, ctrl.Options{Scheme: scheme}) Expect(err).NotTo(HaveOccurred()) ctx, cancel = context.WithCancel(context.Background()) - err = cluster.Add(ctx, mgr, "", "", &log.Logger{SugaredLogger: zap.NewNop().Sugar()}) + err = cluster.Add(ctx, mgr, "", "") Expect(err).NotTo(HaveOccurred()) go func() { diff --git a/pkg/controller/cluster/pod.go b/pkg/controller/cluster/pod.go index c997845c..6ec3d39e 100644 --- a/pkg/controller/cluster/pod.go +++ b/pkg/controller/cluster/pod.go @@ -15,10 +15,8 @@ import ( "github.com/rancher/k3k/pkg/controller/certs" "github.com/rancher/k3k/pkg/controller/cluster/server" "github.com/rancher/k3k/pkg/controller/cluster/server/bootstrap" - "github.com/rancher/k3k/pkg/log" "go.etcd.io/etcd/api/v3/v3rpc/rpctypes" clientv3 "go.etcd.io/etcd/client/v3" - "go.uber.org/zap" apps "k8s.io/api/apps/v1" v1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -41,16 +39,14 @@ const ( type PodReconciler struct { Client ctrlruntimeclient.Client Scheme *runtime.Scheme - logger *log.Logger } // Add adds a new controller to the manager -func AddPodController(ctx context.Context, mgr manager.Manager, logger *log.Logger) error { +func AddPodController(ctx context.Context, mgr manager.Manager) error { // initialize a new Reconciler reconciler := PodReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), - logger: logger.Named(podController), } return ctrl.NewControllerManagedBy(mgr). @@ -63,7 +59,8 @@ func AddPodController(ctx context.Context, mgr manager.Manager, logger *log.Logg } func (p *PodReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { - log := p.logger.With("Pod", req.NamespacedName) + log := ctrl.LoggerFrom(ctx).WithValues("pod", req.NamespacedName) + ctx = ctrl.LoggerInto(ctx, log) // enrich the current logger s := strings.Split(req.Name, "-") if len(s) < 1 { @@ -88,22 +85,28 @@ func (p *PodReconciler) Reconcile(ctx context.Context, req reconcile.Request) (r return reconcile.Result{}, ctrlruntimeclient.IgnoreNotFound(err) } for _, pod := range podList.Items { - log.Info("Handle etcd server pod") - if err := p.handleServerPod(ctx, cluster, &pod, log); err != nil { + + if err := p.handleServerPod(ctx, cluster, &pod); err != nil { return reconcile.Result{}, err } } return reconcile.Result{}, nil } -func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cluster, pod *v1.Pod, log *zap.SugaredLogger) error { - if _, ok := pod.Labels["role"]; ok { - if pod.Labels["role"] != "server" { - return nil - } - } else { +func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cluster, pod *v1.Pod) error { + log := ctrl.LoggerFrom(ctx) + log.Info("handling server pod") + + role, found := pod.Labels["role"] + if !found { return fmt.Errorf("server pod has no role label") } + + if role != "server" { + log.V(1).Info("pod has a different role: " + role) + return nil + } + // if etcd pod is marked for deletion then we need to remove it from the etcd member list before deletion if !pod.DeletionTimestamp.IsZero() { // check if cluster is deleted then remove the finalizer from the pod @@ -116,7 +119,7 @@ func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cl } return nil } - tlsConfig, err := p.getETCDTLS(ctx, &cluster, log) + tlsConfig, err := p.getETCDTLS(ctx, &cluster) if err != nil { return err } @@ -131,9 +134,10 @@ func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cl return err } - if err := removePeer(ctx, client, pod.Name, pod.Status.PodIP, log); err != nil { + if err := removePeer(ctx, client, pod.Name, pod.Status.PodIP); err != nil { return err } + // remove our finalizer from the list and update it. if controllerutil.ContainsFinalizer(pod, etcdPodFinalizerName) { controllerutil.RemoveFinalizer(pod, etcdPodFinalizerName) @@ -150,13 +154,16 @@ func (p *PodReconciler) handleServerPod(ctx context.Context, cluster v1alpha1.Cl return nil } -func (p *PodReconciler) getETCDTLS(ctx context.Context, cluster *v1alpha1.Cluster, log *zap.SugaredLogger) (*tls.Config, error) { - log.Infow("generating etcd TLS client certificate", "Cluster", cluster.Name, "Namespace", cluster.Namespace) +func (p *PodReconciler) getETCDTLS(ctx context.Context, cluster *v1alpha1.Cluster) (*tls.Config, error) { + log := ctrl.LoggerFrom(ctx) + log.Info("generating etcd TLS client certificate", "cluster", cluster) + token, err := p.clusterToken(ctx, cluster) if err != nil { return nil, err } endpoint := server.ServiceName(cluster.Name) + "." + cluster.Namespace + var b *bootstrap.ControlRuntimeBootstrap if err := retry.OnError(k3kcontroller.Backoff, func(err error) bool { return true @@ -191,7 +198,10 @@ func (p *PodReconciler) getETCDTLS(ctx context.Context, cluster *v1alpha1.Cluste } // removePeer removes a peer from the cluster. The peer name and IP address must both match. -func removePeer(ctx context.Context, client *clientv3.Client, name, address string, log *zap.SugaredLogger) error { +func removePeer(ctx context.Context, client *clientv3.Client, name, address string) error { + log := ctrl.LoggerFrom(ctx) + log.Info("removing peer from cluster", "name", name, "address", address) + ctx, cancel := context.WithTimeout(ctx, memberRemovalTimeout) defer cancel() members, err := client.MemberList(ctx) @@ -208,8 +218,9 @@ func removePeer(ctx context.Context, client *clientv3.Client, name, address stri if err != nil { return err } + if u.Hostname() == address { - log.Infow("Removing member from etcd", "name", member.Name, "id", member.ID, "address", address) + log.Info("removing member from etcd", "name", member.Name, "id", member.ID, "address", address) _, err := client.MemberRemove(ctx, member.ID) if errors.Is(err, rpctypes.ErrGRPCMemberNotFound) { return nil diff --git a/pkg/controller/cluster/token.go b/pkg/controller/cluster/token.go index cfc0be42..503a2df1 100644 --- a/pkg/controller/cluster/token.go +++ b/pkg/controller/cluster/token.go @@ -12,6 +12,7 @@ import ( apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" ) @@ -35,15 +36,16 @@ func (c *ClusterReconciler) token(ctx context.Context, cluster *v1alpha1.Cluster } func (c *ClusterReconciler) ensureTokenSecret(ctx context.Context, cluster *v1alpha1.Cluster) (string, error) { + log := ctrl.LoggerFrom(ctx) + // check if the secret is already created - var ( - tokenSecret v1.Secret - nn = types.NamespacedName{ - Name: TokenSecretName(cluster.Name), - Namespace: cluster.Namespace, - } - ) - if err := c.Client.Get(ctx, nn, &tokenSecret); err != nil { + key := types.NamespacedName{ + Name: TokenSecretName(cluster.Name), + Namespace: cluster.Namespace, + } + + var tokenSecret v1.Secret + if err := c.Client.Get(ctx, key, &tokenSecret); err != nil { if !apierrors.IsNotFound(err) { return "", err } @@ -52,15 +54,19 @@ func (c *ClusterReconciler) ensureTokenSecret(ctx context.Context, cluster *v1al if tokenSecret.Data != nil { return string(tokenSecret.Data["token"]), nil } - c.logger.Info("Token secret is not specified, creating a random token") + + log.Info("Token secret is not specified, creating a random token") + token, err := random(16) if err != nil { return "", err } + tokenSecret = TokenSecretObj(token, cluster.Name, cluster.Namespace) if err := controllerutil.SetControllerReference(cluster, &tokenSecret, c.Scheme); err != nil { return "", err } + if err := c.ensure(ctx, &tokenSecret, false); err != nil { return "", err } diff --git a/pkg/controller/clusterset/clusterset.go b/pkg/controller/clusterset/clusterset.go index 49f7ee2b..047349b2 100644 --- a/pkg/controller/clusterset/clusterset.go +++ b/pkg/controller/clusterset/clusterset.go @@ -7,8 +7,6 @@ import ( "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" k3kcontroller "github.com/rancher/k3k/pkg/controller" - "github.com/rancher/k3k/pkg/log" - "go.uber.org/zap" v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -37,17 +35,15 @@ type ClusterSetReconciler struct { Client ctrlruntimeclient.Client Scheme *runtime.Scheme ClusterCIDR string - logger *log.Logger } // Add adds a new controller to the manager -func Add(ctx context.Context, mgr manager.Manager, clusterCIDR string, logger *log.Logger) error { +func Add(ctx context.Context, mgr manager.Manager, clusterCIDR string) error { // initialize a new Reconciler reconciler := ClusterSetReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), ClusterCIDR: clusterCIDR, - logger: logger.Named(clusterSetController), } return ctrl.NewControllerManagedBy(mgr). @@ -121,22 +117,23 @@ func namespaceLabelsPredicate() predicate.Predicate { } func (c *ClusterSetReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { - log := c.logger.With("ClusterSet", req.NamespacedName) + log := ctrl.LoggerFrom(ctx).WithValues("clusterset", req.NamespacedName) + ctx = ctrl.LoggerInto(ctx, log) // enrich the current logger var clusterSet v1alpha1.ClusterSet if err := c.Client.Get(ctx, req.NamespacedName, &clusterSet); err != nil { return reconcile.Result{}, client.IgnoreNotFound(err) } - if err := c.reconcileNetworkPolicy(ctx, log, &clusterSet); err != nil { + if err := c.reconcileNetworkPolicy(ctx, &clusterSet); err != nil { return reconcile.Result{}, err } - if err := c.reconcileNamespacePodSecurityLabels(ctx, log, &clusterSet); err != nil { + if err := c.reconcileNamespacePodSecurityLabels(ctx, &clusterSet); err != nil { return reconcile.Result{}, err } - if err := c.reconcileClusters(ctx, log, &clusterSet); err != nil { + if err := c.reconcileClusters(ctx, &clusterSet); err != nil { return reconcile.Result{}, err } @@ -164,7 +161,8 @@ func (c *ClusterSetReconciler) Reconcile(ctx context.Context, req reconcile.Requ return reconcile.Result{}, nil } -func (c *ClusterSetReconciler) reconcileNetworkPolicy(ctx context.Context, log *zap.SugaredLogger, clusterSet *v1alpha1.ClusterSet) error { +func (c *ClusterSetReconciler) reconcileNetworkPolicy(ctx context.Context, clusterSet *v1alpha1.ClusterSet) error { + log := ctrl.LoggerFrom(ctx) log.Info("reconciling NetworkPolicy") networkPolicy, err := netpol(ctx, c.ClusterCIDR, clusterSet, c.Client) @@ -258,7 +256,8 @@ func netpol(ctx context.Context, clusterCIDR string, clusterSet *v1alpha1.Cluste }, nil } -func (c *ClusterSetReconciler) reconcileNamespacePodSecurityLabels(ctx context.Context, log *zap.SugaredLogger, clusterSet *v1alpha1.ClusterSet) error { +func (c *ClusterSetReconciler) reconcileNamespacePodSecurityLabels(ctx context.Context, clusterSet *v1alpha1.ClusterSet) error { + log := ctrl.LoggerFrom(ctx) log.Info("reconciling Namespace") var ns v1.Namespace @@ -293,7 +292,7 @@ func (c *ClusterSetReconciler) reconcileNamespacePodSecurityLabels(ctx context.C } if !reflect.DeepEqual(ns.Labels, newLabels) { - log.Debug("labels changed, updating namespace") + log.V(1).Info("labels changed, updating namespace") ns.Labels = newLabels return c.Client.Update(ctx, &ns) @@ -301,7 +300,8 @@ func (c *ClusterSetReconciler) reconcileNamespacePodSecurityLabels(ctx context.C return nil } -func (c *ClusterSetReconciler) reconcileClusters(ctx context.Context, log *zap.SugaredLogger, clusterSet *v1alpha1.ClusterSet) error { +func (c *ClusterSetReconciler) reconcileClusters(ctx context.Context, clusterSet *v1alpha1.ClusterSet) error { + log := ctrl.LoggerFrom(ctx) log.Info("reconciling Clusters") var clusters v1alpha1.ClusterList diff --git a/pkg/controller/clusterset/clusterset_suite_test.go b/pkg/controller/clusterset/clusterset_suite_test.go index 3afec6f3..466c6c9c 100644 --- a/pkg/controller/clusterset/clusterset_suite_test.go +++ b/pkg/controller/clusterset/clusterset_suite_test.go @@ -5,9 +5,9 @@ import ( "path/filepath" "testing" + "github.com/go-logr/zapr" "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" "github.com/rancher/k3k/pkg/controller/clusterset" - "github.com/rancher/k3k/pkg/log" "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" @@ -51,10 +51,10 @@ var _ = BeforeSuite(func() { mgr, err := ctrl.NewManager(cfg, ctrl.Options{Scheme: scheme}) Expect(err).NotTo(HaveOccurred()) - ctx, cancel = context.WithCancel(context.Background()) - nopLogger := &log.Logger{SugaredLogger: zap.NewNop().Sugar()} + ctrl.SetLogger(zapr.NewLogger(zap.NewNop())) - err = clusterset.Add(ctx, mgr, "", nopLogger) + ctx, cancel = context.WithCancel(context.Background()) + err = clusterset.Add(ctx, mgr, "") Expect(err).NotTo(HaveOccurred()) go func() { diff --git a/pkg/controller/clusterset/node.go b/pkg/controller/clusterset/node.go index 9fc64fcd..05c4a232 100644 --- a/pkg/controller/clusterset/node.go +++ b/pkg/controller/clusterset/node.go @@ -4,8 +4,6 @@ import ( "context" "github.com/rancher/k3k/pkg/apis/k3k.io/v1alpha1" - "github.com/rancher/k3k/pkg/log" - "go.uber.org/zap" v1 "k8s.io/api/core/v1" networkingv1 "k8s.io/api/networking/v1" "k8s.io/apimachinery/pkg/runtime" @@ -24,16 +22,14 @@ type NodeReconciler struct { Client ctrlruntimeclient.Client Scheme *runtime.Scheme ClusterCIDR string - logger *log.Logger } // AddNodeController adds a new controller to the manager -func AddNodeController(ctx context.Context, mgr manager.Manager, logger *log.Logger) error { +func AddNodeController(ctx context.Context, mgr manager.Manager) error { // initialize a new Reconciler reconciler := NodeReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), - logger: logger.Named(nodeController), } return ctrl.NewControllerManagedBy(mgr). @@ -46,7 +42,11 @@ func AddNodeController(ctx context.Context, mgr manager.Manager, logger *log.Log } func (n *NodeReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { - log := n.logger.With("Node", req.NamespacedName) + log := ctrl.LoggerFrom(ctx).WithValues("node", req.NamespacedName) + ctx = ctrl.LoggerInto(ctx, log) // enrich the current logger + + log.Info("reconciling node") + var clusterSetList v1alpha1.ClusterSetList if err := n.Client.List(ctx, &clusterSetList); err != nil { return reconcile.Result{}, err @@ -56,26 +56,33 @@ func (n *NodeReconciler) Reconcile(ctx context.Context, req reconcile.Request) ( return reconcile.Result{}, nil } - if err := n.ensureNetworkPolicies(ctx, clusterSetList, log); err != nil { + if err := n.ensureNetworkPolicies(ctx, clusterSetList); err != nil { return reconcile.Result{}, err } return reconcile.Result{}, nil } -func (n *NodeReconciler) ensureNetworkPolicies(ctx context.Context, clusterSetList v1alpha1.ClusterSetList, log *zap.SugaredLogger) error { +func (n *NodeReconciler) ensureNetworkPolicies(ctx context.Context, clusterSetList v1alpha1.ClusterSetList) error { + log := ctrl.LoggerFrom(ctx) + log.Info("ensuring network policies") + var setNetworkPolicy *networkingv1.NetworkPolicy for _, cs := range clusterSetList.Items { if cs.Spec.DisableNetworkPolicy { continue } + + log = log.WithValues("clusterset", cs.Namespace+"/"+cs.Name) + log.Info("updating NetworkPolicy for ClusterSet") + var err error - log.Infow("Updating NetworkPolicy for ClusterSet", "name", cs.Name, "namespace", cs.Namespace) setNetworkPolicy, err = netpol(ctx, "", &cs, n.Client) if err != nil { return err } - log.Debugw("New NetworkPolicy for clusterset", "name", cs.Name, "namespace", cs.Namespace) + + log.Info("new NetworkPolicy for clusterset") if err := n.Client.Update(ctx, setNetworkPolicy); err != nil { return err }