diff --git a/k3k-kubelet/kubelet.go b/k3k-kubelet/kubelet.go index 5ae491b5..62f09a4f 100644 --- a/k3k-kubelet/kubelet.go +++ b/k3k-kubelet/kubelet.go @@ -211,9 +211,9 @@ func clusterIP(ctx context.Context, serviceName, clusterNamespace string, hostCl return service.Spec.ClusterIP, nil } -func (k *kubelet) registerNode(ctx context.Context, agentIP string, cfg config) error { +func (k *kubelet) registerNode(agentIP string, cfg config) error { providerFunc := k.newProviderFunc(cfg) - nodeOpts := k.nodeOpts(ctx, cfg.KubeletPort, cfg.ClusterNamespace, cfg.ClusterName, cfg.AgentHostname, agentIP) + nodeOpts := k.nodeOpts(cfg.KubeletPort, cfg.ClusterNamespace, cfg.ClusterName, cfg.AgentHostname, agentIP) var err error @@ -277,7 +277,7 @@ func (k *kubelet) newProviderFunc(cfg config) nodeutil.NewProviderFunc { } } -func (k *kubelet) nodeOpts(ctx context.Context, srvPort int, namespace, name, hostname, agentIP string) nodeutil.NodeOpt { +func (k *kubelet) nodeOpts(srvPort int, namespace, name, hostname, agentIP string) nodeutil.NodeOpt { return func(c *nodeutil.NodeConfig) error { c.HTTPListenAddr = fmt.Sprintf(":%d", srvPort) // set up the routes @@ -288,7 +288,7 @@ func (k *kubelet) nodeOpts(ctx context.Context, srvPort int, namespace, name, ho c.Handler = mux - tlsConfig, err := loadTLSConfig(ctx, k.hostClient, name, namespace, k.name, hostname, k.token, agentIP) + tlsConfig, err := loadTLSConfig(name, namespace, k.name, hostname, k.token, agentIP) if err != nil { return errors.New("unable to get tls config: " + err.Error()) } @@ -369,17 +369,10 @@ func kubeconfigBytes(url string, serverCA, clientCert, clientKey []byte) ([]byte return clientcmd.Write(*config) } -func loadTLSConfig(ctx context.Context, hostClient ctrlruntimeclient.Client, clusterName, clusterNamespace, nodeName, hostname, token, agentIP string) (*tls.Config, error) { - var ( - cluster v1alpha1.Cluster - b *bootstrap.ControlRuntimeBootstrap - ) +func loadTLSConfig(clusterName, clusterNamespace, nodeName, hostname, token, agentIP string) (*tls.Config, error) { + var b *bootstrap.ControlRuntimeBootstrap - if err := hostClient.Get(ctx, types.NamespacedName{Name: clusterName, Namespace: clusterNamespace}, &cluster); err != nil { - return nil, err - } - - endpoint := fmt.Sprintf("%s.%s", server.ServiceName(cluster.Name), cluster.Namespace) + endpoint := fmt.Sprintf("%s.%s", server.ServiceName(clusterName), clusterNamespace) if err := retry.OnError(controller.Backoff, func(err error) bool { return err != nil diff --git a/k3k-kubelet/main.go b/k3k-kubelet/main.go index fe5bc2e3..b79bf521 100644 --- a/k3k-kubelet/main.go +++ b/k3k-kubelet/main.go @@ -73,7 +73,7 @@ func run(cmd *cobra.Command, args []string) error { return fmt.Errorf("failed to create new virtual kubelet instance: %w", err) } - if err := k.registerNode(ctx, k.agentIP, cfg); err != nil { + if err := k.registerNode(k.agentIP, cfg); err != nil { return fmt.Errorf("failed to register new node: %w", err) } diff --git a/main.go b/main.go index f62a71de..eaa18322 100644 --- a/main.go +++ b/main.go @@ -123,6 +123,12 @@ func run(cmd *cobra.Command, args []string) error { return fmt.Errorf("failed to add the new cluster controller: %v", err) } + logger.Info("adding service controller") + + if err := cluster.AddServiceController(ctx, mgr, maxConcurrentReconciles); err != nil { + return fmt.Errorf("failed to add the new service controller: %v", err) + } + logger.Info("adding clusterpolicy controller") if err := policy.Add(mgr, config.ClusterCIDR, maxConcurrentReconciles); err != nil { diff --git a/pkg/controller/cluster/service.go b/pkg/controller/cluster/service.go new file mode 100644 index 00000000..ca4bafba --- /dev/null +++ b/pkg/controller/cluster/service.go @@ -0,0 +1,117 @@ +package cluster + +import ( + "context" + "fmt" + + "k8s.io/apimachinery/pkg/api/equality" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/tools/clientcmd" + "sigs.k8s.io/controller-runtime/pkg/manager" + "sigs.k8s.io/controller-runtime/pkg/predicate" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + v1 "k8s.io/api/core/v1" + ctrl "sigs.k8s.io/controller-runtime" + ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client" + + "github.com/rancher/k3k/k3k-kubelet/translate" + "github.com/rancher/k3k/pkg/controller" +) + +const ( + serviceController = "k3k-service-controller" +) + +type ServiceReconciler struct { + HostClient ctrlruntimeclient.Client +} + +// Add adds a new controller to the manager +func AddServiceController(ctx context.Context, mgr manager.Manager, maxConcurrentReconciles int) error { + reconciler := ServiceReconciler{ + HostClient: mgr.GetClient(), + } + + return ctrl.NewControllerManagedBy(mgr). + Named(serviceController). + For(&v1.Service{}).WithEventFilter(predicate.NewPredicateFuncs(reconciler.filterResources)). + Complete(&reconciler) +} + +func (r *ServiceReconciler) filterResources(object ctrlruntimeclient.Object) bool { + _, hasClusterNameLabel := object.GetLabels()[translate.ClusterNameLabel] + + _, hasNameAnnotation := object.GetAnnotations()[translate.ResourceNameAnnotation] + + _, hasNamespaceAnnotation := object.GetAnnotations()[translate.ResourceNamespaceAnnotation] + + // Return true only if all were found + return hasClusterNameLabel && hasNameAnnotation && hasNamespaceAnnotation +} + +func (r *ServiceReconciler) Reconcile(ctx context.Context, req reconcile.Request) (reconcile.Result, error) { + log := ctrl.LoggerFrom(ctx) + log.Info("ensuring service status to virtual cluster") + + var ( + virtServiceName, virtServiceNamespace string + hostService, virtService v1.Service + ) + + if err := r.HostClient.Get(ctx, req.NamespacedName, &hostService); err != nil { + return reconcile.Result{}, ctrlruntimeclient.IgnoreNotFound(err) + } + + // get cluster information from the object + labels := hostService.GetLabels() + + clusterName, ok := labels[translate.ClusterNameLabel] + if !ok { + return reconcile.Result{}, fmt.Errorf("cluster name label is not found") + } + + clusterNamespace := hostService.GetNamespace() + + virtClient, err := r.newVirtualClient(ctx, &hostService, clusterName, clusterNamespace) + if err != nil { + return reconcile.Result{}, fmt.Errorf("failed to get cluster info: %v", err) + } + + virtServiceName = hostService.Annotations[translate.ResourceNameAnnotation] + virtServiceNamespace = hostService.Annotations[translate.ResourceNamespaceAnnotation] + + if !hostService.DeletionTimestamp.IsZero() { + return reconcile.Result{}, nil + } + + if err := virtClient.Get(ctx, types.NamespacedName{Name: virtServiceName, Namespace: virtServiceNamespace}, &virtService); err != nil { + return reconcile.Result{}, fmt.Errorf("failed to get virt service: %v", err) + } + + if !equality.Semantic.DeepEqual(virtService.Status.LoadBalancer, hostService.Status.LoadBalancer) { + virtService.Status.LoadBalancer = hostService.Status.LoadBalancer + if err := virtClient.Status().Update(ctx, &virtService); err != nil { + return reconcile.Result{}, err + } + } + + return reconcile.Result{}, nil +} + +func (r *ServiceReconciler) newVirtualClient(ctx context.Context, obj ctrlruntimeclient.Object, clusterName, clusterNamespace string) (ctrlruntimeclient.Client, error) { + var clusterKubeConfig v1.Secret + + kubeconfigSecretName := controller.SafeConcatNameWithPrefix(clusterName, "kubeconfig") + + if err := r.HostClient.Get(ctx, types.NamespacedName{Name: kubeconfigSecretName, Namespace: clusterNamespace}, &clusterKubeConfig); err != nil { + return nil, err + } + + virtConfig, err := clientcmd.RESTConfigFromKubeConfig(clusterKubeConfig.Data["kubeconfig.yaml"]) + if err != nil { + return nil, fmt.Errorf("failed to create config from kubeconfig file: %v", err) + } + + return ctrlruntimeclient.New(virtConfig, ctrlruntimeclient.Options{}) +}