From 15108288b4a13796eea41a4b56f41ca3e6564489 Mon Sep 17 00:00:00 2001 From: bugkeep <1921817430@qq.com> Date: Sat, 25 Jul 2026 21:53:49 +0800 Subject: [PATCH] feat: provide default application access data plane --- conf/app.conf | 2 + controllers/ingress.go | 42 +-- controllers/service.go | 42 +-- main.go | 5 + server/apiserver.go | 2 +- server/bootstrap.go | 12 +- server/config.go | 24 +- server/ingress_bootstrap.go | 432 +++++++++++++++++++++++++++++++ server/servicelb.go | 468 ++++++++++++++++++++++++++++++++++ server/storage_bootstrap.go | 31 ++- web/src/DeploymentListPage.js | 37 ++- web/src/ServiceListPage.js | 21 +- 12 files changed, 1062 insertions(+), 56 deletions(-) create mode 100644 server/ingress_bootstrap.go create mode 100644 server/servicelb.go diff --git a/conf/app.conf b/conf/app.conf index d1461906..3812b562 100644 --- a/conf/app.conf +++ b/conf/app.conf @@ -27,3 +27,5 @@ helmImplicitLatestPullPolicy = false apiserverPort = 6443 apiserverBind = 127.0.0.1 dataDir = /var/lib/casos +ingressControllerEnabled = true +serviceLBEnabled = true diff --git a/controllers/ingress.go b/controllers/ingress.go index 5e542867..ec82fcf8 100644 --- a/controllers/ingress.go +++ b/controllers/ingress.go @@ -18,14 +18,15 @@ type ingressRule struct { } type ingressSummary struct { - Namespace string `json:"namespace"` - Name string `json:"name"` - IngressClass string `json:"ingressClass"` - Rules []ingressRule `json:"rules"` - TLSEnabled bool `json:"tlsEnabled"` - TLSSecretName string `json:"tlsSecretName"` - CreatedAt string `json:"createdAt"` - ResourceVersion string `json:"resourceVersion"` + Namespace string `json:"namespace"` + Name string `json:"name"` + IngressClass string `json:"ingressClass"` + Rules []ingressRule `json:"rules"` + TLSEnabled bool `json:"tlsEnabled"` + TLSSecretName string `json:"tlsSecretName"` + LoadBalancerAddresses []string `json:"loadBalancerAddresses,omitempty"` + CreatedAt string `json:"createdAt"` + ResourceVersion string `json:"resourceVersion"` } func toIngressSummary(ing networkingv1.Ingress) ingressSummary { @@ -65,15 +66,24 @@ func toIngressSummary(ing networkingv1.Ingress) ingressSummary { if tlsEnabled { tlsSecretName = ing.Spec.TLS[0].SecretName } + loadBalancerAddresses := make([]string, 0, len(ing.Status.LoadBalancer.Ingress)) + for _, address := range ing.Status.LoadBalancer.Ingress { + if address.IP != "" { + loadBalancerAddresses = append(loadBalancerAddresses, address.IP) + } else if address.Hostname != "" { + loadBalancerAddresses = append(loadBalancerAddresses, address.Hostname) + } + } return ingressSummary{ - Namespace: ing.Namespace, - Name: ing.Name, - IngressClass: cls, - Rules: rules, - TLSEnabled: tlsEnabled, - TLSSecretName: tlsSecretName, - CreatedAt: ing.CreationTimestamp.UTC().Format("2006-01-02 15:04:05"), - ResourceVersion: ing.ResourceVersion, + Namespace: ing.Namespace, + Name: ing.Name, + IngressClass: cls, + Rules: rules, + TLSEnabled: tlsEnabled, + TLSSecretName: tlsSecretName, + LoadBalancerAddresses: loadBalancerAddresses, + CreatedAt: ing.CreationTimestamp.UTC().Format("2006-01-02 15:04:05"), + ResourceVersion: ing.ResourceVersion, } } diff --git a/controllers/service.go b/controllers/service.go index 8d34820e..bffdddcd 100644 --- a/controllers/service.go +++ b/controllers/service.go @@ -20,14 +20,15 @@ type portSummary struct { } type serviceSummary struct { - Namespace string `json:"namespace"` - Name string `json:"name"` - Type string `json:"type"` - ClusterIP string `json:"clusterIP"` - Selector map[string]string `json:"selector"` - Ports []portSummary `json:"ports"` - CreatedAt string `json:"createdAt"` - ResourceVersion string `json:"resourceVersion"` + Namespace string `json:"namespace"` + Name string `json:"name"` + Type string `json:"type"` + ClusterIP string `json:"clusterIP"` + Selector map[string]string `json:"selector"` + Ports []portSummary `json:"ports"` + LoadBalancerAddresses []string `json:"loadBalancerAddresses,omitempty"` + CreatedAt string `json:"createdAt"` + ResourceVersion string `json:"resourceVersion"` } func toSvcSummary(svc corev1.Service) serviceSummary { @@ -41,15 +42,24 @@ func toSvcSummary(svc corev1.Service) serviceSummary { NodePort: p.NodePort, }) } + loadBalancerAddresses := make([]string, 0, len(svc.Status.LoadBalancer.Ingress)) + for _, ingress := range svc.Status.LoadBalancer.Ingress { + if ingress.IP != "" { + loadBalancerAddresses = append(loadBalancerAddresses, ingress.IP) + } else if ingress.Hostname != "" { + loadBalancerAddresses = append(loadBalancerAddresses, ingress.Hostname) + } + } return serviceSummary{ - Namespace: svc.Namespace, - Name: svc.Name, - Type: string(svc.Spec.Type), - ClusterIP: svc.Spec.ClusterIP, - Selector: svc.Spec.Selector, - Ports: ports, - CreatedAt: svc.CreationTimestamp.UTC().Format("2006-01-02 15:04:05"), - ResourceVersion: svc.ResourceVersion, + Namespace: svc.Namespace, + Name: svc.Name, + Type: string(svc.Spec.Type), + ClusterIP: svc.Spec.ClusterIP, + Selector: svc.Spec.Selector, + Ports: ports, + LoadBalancerAddresses: loadBalancerAddresses, + CreatedAt: svc.CreationTimestamp.UTC().Format("2006-01-02 15:04:05"), + ResourceVersion: svc.ResourceVersion, } } diff --git a/main.go b/main.go index aa960348..f626c4ef 100644 --- a/main.go +++ b/main.go @@ -67,6 +67,11 @@ func main() { if err := server.Bootstrap(ctx, adminCfg, srvCfg); err != nil { logs.Warning("bootstrap: %v", err) } + if srvCfg.ServiceLBEnabled { + if err := server.StartServiceLB(ctx, adminCfg); err != nil { + logs.Warning("start service load balancer: %v", err) + } + } if err := server.StartScheduler(ctx, srvCfg); err != nil { logs.Warning("start scheduler: %v", err) } diff --git a/server/apiserver.go b/server/apiserver.go index 0dbd0c8a..1aa3af5b 100644 --- a/server/apiserver.go +++ b/server/apiserver.go @@ -150,7 +150,7 @@ func buildApiserverArgs(cfg Config, certDir, etcdEndpoint, authzKubeconfig strin "--service-cluster-ip-range=" + serviceClusterIPRange, "--allow-privileged=true", "--authorization-mode=" + authzMode(authzKubeconfig), - "--enable-admission-plugins=NodeRestriction,ValidatingAdmissionWebhook", + "--enable-admission-plugins=NodeRestriction,DefaultIngressClass,ValidatingAdmissionWebhook", "--tls-cert-file=" + filepath.Join(certDir, "apiserver.crt"), "--tls-private-key-file=" + filepath.Join(certDir, "apiserver.key"), "--client-ca-file=" + filepath.Join(certDir, "ca.crt"), diff --git a/server/bootstrap.go b/server/bootstrap.go index 86972966..7045150e 100644 --- a/server/bootstrap.go +++ b/server/bootstrap.go @@ -16,8 +16,8 @@ import ( "k8s.io/client-go/rest" ) -// Bootstrap creates cluster-wide resources required for worker-node components -// to function correctly. It is idempotent — safe to call on every startup. +// Bootstrap creates CasOS-managed cluster add-ons. It is idempotent and safe +// to call on every startup; individual add-ons can be disabled in config. func Bootstrap(ctx context.Context, cfg *rest.Config, srvCfg Config) error { client, err := kubernetes.NewForConfig(cfg) if err != nil { @@ -32,6 +32,14 @@ func Bootstrap(ctx context.Context, cfg *rest.Config, srvCfg Config) error { errs = append(errs, ensureNodeProxierBinding(ctx, client)) errs = append(errs, ensureFlannel(ctx, client, srvCfg)) errs = append(errs, ensureClusterDNS(ctx, client, srvCfg)) + if srvCfg.IngressControllerEnabled { + errs = append(errs, ensureIngressController(ctx, client, srvCfg)) + } else { + errs = append(errs, cleanupIngressController(ctx, client)) + } + if !srvCfg.ServiceLBEnabled { + errs = append(errs, cleanupServiceLB(ctx, client)) + } if srvCfg.StorageProvisionerEnabled { errs = append(errs, ensureDefaultStorageProvisioner(ctx, client, srvCfg)) } diff --git a/server/config.go b/server/config.go index ad819abd..2f779778 100644 --- a/server/config.go +++ b/server/config.go @@ -23,7 +23,10 @@ type Config struct { LocalPathHelperImage string // helper pod image used by local-path-provisioner FlannelImage string // Flannel daemon image used by the built-in network bootstrap FlannelCNIPluginImage string // Flannel CNI plugin image installed on worker hosts + IngressControllerImage string // Traefik image used by the built-in Ingress controller StorageProvisionerEnabled bool // install the built-in local-path provisioner for local clusters + IngressControllerEnabled bool // install the built-in Traefik controller + ServiceLBEnabled bool // run the built-in bare-metal LoadBalancer controller } // ConfigFromAppConf reads server config from the beego app.conf. @@ -60,13 +63,16 @@ func ConfigFromAppConf() (Config, error) { } storageProvisionerEnabled := conf.GetConfigBoolDefault("storageProvisionerEnabled", true) + ingressControllerEnabled := conf.GetConfigBoolDefault("ingressControllerEnabled", false) + serviceLBEnabled := conf.GetConfigBoolDefault("serviceLBEnabled", false) coreDNSImage := conf.GetConfigStringDefault("coreDNSImage", "docker.1ms.run/coredns/coredns:1.12.4") localPathProvisionerImage := conf.GetConfigStringDefault("localPathProvisionerImage", "docker.1ms.run/rancher/local-path-provisioner:v0.0.32") localPathHelperImage := conf.GetConfigStringDefault("localPathHelperImage", "docker.1ms.run/library/busybox:1.37.0") flannelImage := conf.GetConfigStringDefault("flannelImage", defaultFlannelImage) flannelCNIPluginImage := conf.GetConfigStringDefault("flannelCNIPluginImage", defaultFlannelCNIPluginImage) + ingressControllerImage := conf.GetConfigStringDefault("ingressControllerImage", defaultIngressControllerImage) - return Config{ + config := Config{ DataDir: dataDir, ApiserverBind: bind, AdvertiseAddress: advertise, @@ -80,8 +86,22 @@ func ConfigFromAppConf() (Config, error) { LocalPathHelperImage: localPathHelperImage, FlannelImage: flannelImage, FlannelCNIPluginImage: flannelCNIPluginImage, + IngressControllerImage: ingressControllerImage, StorageProvisionerEnabled: storageProvisionerEnabled, - }, nil + IngressControllerEnabled: ingressControllerEnabled, + ServiceLBEnabled: serviceLBEnabled, + } + if err := validateApplicationAccessConfig(config); err != nil { + return Config{}, err + } + return config, nil +} + +func validateApplicationAccessConfig(config Config) error { + if config.IngressControllerEnabled && !config.ServiceLBEnabled { + return fmt.Errorf("ingressControllerEnabled requires serviceLBEnabled because the built-in Ingress Service uses LoadBalancerClass %q", serviceLBClass) + } + return nil } // injectDBName inserts dbName into a MySQL DSN of the form diff --git a/server/ingress_bootstrap.go b/server/ingress_bootstrap.go new file mode 100644 index 00000000..cbe9e99a --- /dev/null +++ b/server/ingress_bootstrap.go @@ -0,0 +1,432 @@ +package server + +import ( + "context" + "errors" + "fmt" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + networkingv1 "k8s.io/api/networking/v1" + rbacv1 "k8s.io/api/rbac/v1" + apiequality "k8s.io/apimachinery/pkg/api/equality" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/client-go/kubernetes" +) + +const ( + ingressControllerNamespace = "kube-system" + ingressControllerName = "traefik" + ingressControllerClass = "traefik" + ingressControllerID = "traefik.io/ingress-controller" + defaultIngressControllerImage = "docker.io/traefik:v3.3.4" +) + +func ingressControllerLabels() map[string]string { + return map[string]string{ + "app.kubernetes.io/name": ingressControllerName, + "app.kubernetes.io/managed-by": "casos", + } +} + +func validateIngressControllerOwnership(object metav1.Object) error { + if object.GetLabels()["app.kubernetes.io/managed-by"] != "casos" { + return fmt.Errorf("Ingress controller %T %s already exists and is not managed by CasOS", object, object.GetName()) + } + return nil +} + +func ensureIngressController(ctx context.Context, client kubernetes.Interface, cfg Config) error { + if err := ensureIngressControllerOwnership(ctx, client); err != nil { + return err + } + if err := ensureNamespace(ctx, client, ingressControllerNamespace); err != nil { + return err + } + if err := createOrUpdateServiceAccount(ctx, client, &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{ + Name: ingressControllerName, + Namespace: ingressControllerNamespace, + Labels: ingressControllerLabels(), + }, + }, validateIngressControllerOwnership); err != nil { + return err + } + if err := ensureIngressControllerClusterRole(ctx, client); err != nil { + return err + } + if err := ensureIngressControllerClusterRoleBinding(ctx, client); err != nil { + return err + } + if err := ensureIngressClass(ctx, client); err != nil { + return err + } + if err := ensureIngressControllerService(ctx, client); err != nil { + return err + } + return ensureIngressControllerDeployment(ctx, client, cfg) +} + +func cleanupIngressController(ctx context.Context, client kubernetes.Interface) error { + resources := []struct { + kind string + get func() (metav1.Object, error) + delete func() error + }{ + {kind: "Deployment", get: func() (metav1.Object, error) { + return client.AppsV1().Deployments(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }, delete: func() error { + return client.AppsV1().Deployments(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{}) + }}, + {kind: "Service", get: func() (metav1.Object, error) { + return client.CoreV1().Services(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }, delete: func() error { + return client.CoreV1().Services(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{}) + }}, + {kind: "IngressClass", get: func() (metav1.Object, error) { + return client.NetworkingV1().IngressClasses().Get(ctx, ingressControllerClass, metav1.GetOptions{}) + }, delete: func() error { + return client.NetworkingV1().IngressClasses().Delete(ctx, ingressControllerClass, metav1.DeleteOptions{}) + }}, + {kind: "ClusterRoleBinding", get: func() (metav1.Object, error) { + return client.RbacV1().ClusterRoleBindings().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{}) + }, delete: func() error { + return client.RbacV1().ClusterRoleBindings().Delete(ctx, "casos:traefik-ingress-controller", metav1.DeleteOptions{}) + }}, + {kind: "ClusterRole", get: func() (metav1.Object, error) { + return client.RbacV1().ClusterRoles().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{}) + }, delete: func() error { + return client.RbacV1().ClusterRoles().Delete(ctx, "casos:traefik-ingress-controller", metav1.DeleteOptions{}) + }}, + {kind: "ServiceAccount", get: func() (metav1.Object, error) { + return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }, delete: func() error { + return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{}) + }}, + } + errs := make([]error, 0) + for _, resource := range resources { + object, err := resource.get() + if apierrors.IsNotFound(err) { + continue + } + if err != nil { + errs = append(errs, fmt.Errorf("get Ingress controller %s for cleanup: %w", resource.kind, err)) + continue + } + if object.GetLabels()["app.kubernetes.io/managed-by"] != "casos" { + continue + } + if err := resource.delete(); err != nil && !apierrors.IsNotFound(err) { + errs = append(errs, fmt.Errorf("delete Ingress controller %s %s: %w", resource.kind, object.GetName(), err)) + } + } + return errors.Join(errs...) +} + +func ensureIngressControllerOwnership(ctx context.Context, client kubernetes.Interface) error { + checks := []struct { + kind string + get func() (metav1.Object, error) + }{ + {kind: "ServiceAccount", get: func() (metav1.Object, error) { + return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }}, + {kind: "Service", get: func() (metav1.Object, error) { + return client.CoreV1().Services(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }}, + {kind: "Deployment", get: func() (metav1.Object, error) { + return client.AppsV1().Deployments(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{}) + }}, + {kind: "ClusterRole", get: func() (metav1.Object, error) { + return client.RbacV1().ClusterRoles().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{}) + }}, + {kind: "ClusterRoleBinding", get: func() (metav1.Object, error) { + return client.RbacV1().ClusterRoleBindings().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{}) + }}, + {kind: "IngressClass", get: func() (metav1.Object, error) { + return client.NetworkingV1().IngressClasses().Get(ctx, ingressControllerClass, metav1.GetOptions{}) + }}, + } + for _, check := range checks { + object, err := check.get() + if apierrors.IsNotFound(err) { + continue + } + if err != nil { + return fmt.Errorf("check Ingress controller %s ownership: %w", check.kind, err) + } + if err := validateIngressControllerOwnership(object); err != nil { + return fmt.Errorf("check Ingress controller %s ownership: %w", check.kind, err) + } + } + return nil +} + +func ensureIngressControllerClusterRole(ctx context.Context, client kubernetes.Interface) error { + role := &rbacv1.ClusterRole{ + ObjectMeta: metav1.ObjectMeta{ + Name: "casos:traefik-ingress-controller", + Labels: ingressControllerLabels(), + }, + Rules: []rbacv1.PolicyRule{ + { + APIGroups: []string{""}, + Resources: []string{"services", "endpoints", "nodes", "secrets", "namespaces"}, + Verbs: []string{"get", "list", "watch"}, + }, + { + APIGroups: []string{"discovery.k8s.io"}, + Resources: []string{"endpointslices"}, + Verbs: []string{"get", "list", "watch"}, + }, + { + APIGroups: []string{"networking.k8s.io"}, + Resources: []string{"ingresses", "ingressclasses"}, + Verbs: []string{"get", "list", "watch"}, + }, + { + APIGroups: []string{""}, + Resources: []string{"events"}, + Verbs: []string{"create", "update", "patch"}, + }, + { + APIGroups: []string{"networking.k8s.io"}, + Resources: []string{"ingresses/status"}, + Verbs: []string{"update", "patch"}, + }, + }, + } + return createOrUpdateClusterRole(ctx, client, role, validateIngressControllerOwnership) +} + +func ensureIngressControllerClusterRoleBinding(ctx context.Context, client kubernetes.Interface) error { + binding := &rbacv1.ClusterRoleBinding{ + ObjectMeta: metav1.ObjectMeta{ + Name: "casos:traefik-ingress-controller", + Labels: ingressControllerLabels(), + }, + RoleRef: rbacv1.RoleRef{ + APIGroup: "rbac.authorization.k8s.io", + Kind: "ClusterRole", + Name: "casos:traefik-ingress-controller", + }, + Subjects: []rbacv1.Subject{{ + Kind: "ServiceAccount", + Name: ingressControllerName, + Namespace: ingressControllerNamespace, + }}, + } + return createOrUpdateClusterRoleBinding(ctx, client, binding, validateIngressControllerOwnership) +} + +func ensureIngressClass(ctx context.Context, client kubernetes.Interface) error { + classes := client.NetworkingV1().IngressClasses() + defaultClassExists := false + existingClasses, err := classes.List(ctx, metav1.ListOptions{}) + if err != nil { + return fmt.Errorf("list IngressClasses: %w", err) + } + for i := range existingClasses.Items { + class := &existingClasses.Items[i] + if class.Name != ingressControllerClass && class.Annotations["ingressclass.kubernetes.io/is-default-class"] == "true" { + defaultClassExists = true + break + } + } + desired := &networkingv1.IngressClass{ + ObjectMeta: metav1.ObjectMeta{ + Name: ingressControllerClass, + Labels: ingressControllerLabels(), + Annotations: map[string]string{}, + }, + Spec: networkingv1.IngressClassSpec{Controller: ingressControllerID}, + } + if !defaultClassExists { + desired.Annotations["ingressclass.kubernetes.io/is-default-class"] = "true" + } + current, err := classes.Get(ctx, desired.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + if _, err := classes.Create(ctx, desired, metav1.CreateOptions{}); err != nil { + return fmt.Errorf("create IngressClass %s: %w", desired.Name, err) + } + return nil + } + if err != nil { + return fmt.Errorf("get IngressClass %s: %w", desired.Name, err) + } + if current.Spec.Controller != desired.Spec.Controller { + return fmt.Errorf("IngressClass %s is owned by controller %s", desired.Name, current.Spec.Controller) + } + if current.Labels["app.kubernetes.io/managed-by"] != "casos" { + return fmt.Errorf("IngressClass %s already exists and is not managed by CasOS", desired.Name) + } + desired.Labels = mergeStringMap(current.Labels, desired.Labels) + desired.Annotations = mergeStringMap(current.Annotations, desired.Annotations) + if defaultClassExists { + delete(desired.Annotations, "ingressclass.kubernetes.io/is-default-class") + } + if apiequality.Semantic.DeepEqual(current.Labels, desired.Labels) && + apiequality.Semantic.DeepEqual(current.Annotations, desired.Annotations) && + apiequality.Semantic.DeepEqual(current.Spec, desired.Spec) { + return nil + } + desired.ResourceVersion = current.ResourceVersion + if _, err := classes.Update(ctx, desired, metav1.UpdateOptions{}); err != nil { + return fmt.Errorf("update IngressClass %s: %w", desired.Name, err) + } + return nil +} + +func ensureIngressControllerService(ctx context.Context, client kubernetes.Interface) error { + desired := buildIngressControllerService() + services := client.CoreV1().Services(ingressControllerNamespace) + current, err := services.Get(ctx, desired.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + if _, err := services.Create(ctx, desired, metav1.CreateOptions{}); err != nil { + return fmt.Errorf("create Ingress controller Service: %w", err) + } + return nil + } + if err != nil { + return fmt.Errorf("get Ingress controller Service: %w", err) + } + if current.Labels["app.kubernetes.io/managed-by"] != "casos" { + return fmt.Errorf("Ingress controller Service %s/%s already exists and is not managed by CasOS", ingressControllerNamespace, ingressControllerName) + } + desired.Labels = mergeStringMap(current.Labels, desired.Labels) + desired.Annotations = mergeStringMap(current.Annotations, desired.Annotations) + desired.Spec.ClusterIP = current.Spec.ClusterIP + desired.Spec.ClusterIPs = current.Spec.ClusterIPs + desired.Spec.IPFamilies = current.Spec.IPFamilies + desired.Spec.IPFamilyPolicy = current.Spec.IPFamilyPolicy + desired.Spec.ExternalIPs = current.Spec.ExternalIPs + desired.Spec.SessionAffinity = current.Spec.SessionAffinity + desired.Spec.InternalTrafficPolicy = current.Spec.InternalTrafficPolicy + desired.Spec.LoadBalancerIP = current.Spec.LoadBalancerIP + desired.Spec.LoadBalancerClass = current.Spec.LoadBalancerClass + desired.Spec.LoadBalancerSourceRanges = current.Spec.LoadBalancerSourceRanges + desired.Spec.AllocateLoadBalancerNodePorts = current.Spec.AllocateLoadBalancerNodePorts + desired.Spec.HealthCheckNodePort = current.Spec.HealthCheckNodePort + for i := range desired.Spec.Ports { + for _, currentPort := range current.Spec.Ports { + if desired.Spec.Ports[i].Name == currentPort.Name { + desired.Spec.Ports[i].NodePort = currentPort.NodePort + break + } + } + } + if apiequality.Semantic.DeepEqual(current.Labels, desired.Labels) && + apiequality.Semantic.DeepEqual(current.Annotations, desired.Annotations) && + apiequality.Semantic.DeepEqual(current.Spec, desired.Spec) { + return nil + } + desired.ResourceVersion = current.ResourceVersion + if _, err := services.Update(ctx, desired, metav1.UpdateOptions{}); err != nil { + return fmt.Errorf("update Ingress controller Service: %w", err) + } + return nil +} + +func buildIngressControllerService() *corev1.Service { + return &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: ingressControllerName, + Namespace: ingressControllerNamespace, + Labels: ingressControllerLabels(), + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeLoadBalancer, + Selector: ingressControllerLabels(), + ExternalTrafficPolicy: corev1.ServiceExternalTrafficPolicyTypeCluster, + LoadBalancerClass: ptr(serviceLBClass), + Ports: []corev1.ServicePort{ + {Name: "web", Port: 80, TargetPort: intstr.FromInt(8000), Protocol: corev1.ProtocolTCP}, + {Name: "websecure", Port: 443, TargetPort: intstr.FromInt(8443), Protocol: corev1.ProtocolTCP}, + }, + }, + } +} + +func ensureIngressControllerDeployment(ctx context.Context, client kubernetes.Interface, cfg Config) error { + return createOrUpdateDeployment(ctx, client, buildIngressControllerDeployment(cfg), validateIngressControllerOwnership) +} + +func buildIngressControllerDeployment(cfg Config) *appsv1.Deployment { + replicas := int32(1) + image := cfg.IngressControllerImage + if image == "" { + image = defaultIngressControllerImage + } + labels := ingressControllerLabels() + deployment := &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: ingressControllerName, + Namespace: ingressControllerNamespace, + Labels: labels, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: &replicas, + Selector: &metav1.LabelSelector{MatchLabels: labels}, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: labels}, + Spec: corev1.PodSpec{ + ServiceAccountName: ingressControllerName, + SecurityContext: &corev1.PodSecurityContext{ + RunAsNonRoot: ptr(true), + RunAsUser: ptr(int64(65532)), + RunAsGroup: ptr(int64(65532)), + SeccompProfile: &corev1.SeccompProfile{Type: corev1.SeccompProfileTypeRuntimeDefault}, + }, + Tolerations: []corev1.Toleration{ + {Key: "CriticalAddonsOnly", Operator: corev1.TolerationOpExists}, + {Key: "node-role.kubernetes.io/control-plane", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule}, + {Key: "node-role.kubernetes.io/master", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule}, + {Key: "casos.io/bootstrap", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule}, + }, + Containers: []corev1.Container{{ + Name: ingressControllerName, + Image: image, + ImagePullPolicy: corev1.PullIfNotPresent, + Args: []string{ + "--providers.kubernetesingress=true", + "--providers.kubernetesingress.ingressclass=" + ingressControllerClass, + "--providers.kubernetesingress.ingressendpoint.publishedservice=" + ingressControllerNamespace + "/" + ingressControllerName, + "--entrypoints.web.address=:8000", + "--entrypoints.websecure.address=:8443", + "--ping=true", + }, + Ports: []corev1.ContainerPort{ + {Name: "web", ContainerPort: 8000, Protocol: corev1.ProtocolTCP}, + {Name: "websecure", ContainerPort: 8443, Protocol: corev1.ProtocolTCP}, + {Name: "ping", ContainerPort: 8080, Protocol: corev1.ProtocolTCP}, + }, + ReadinessProbe: &corev1.Probe{ + ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{ + Path: "/ping", Port: intstr.FromInt(8080), Scheme: corev1.URISchemeHTTP, + }}, + InitialDelaySeconds: 5, PeriodSeconds: 10, TimeoutSeconds: 5, FailureThreshold: 6, + }, + LivenessProbe: &corev1.Probe{ + ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{Path: "/ping", Port: intstr.FromInt(8080), Scheme: corev1.URISchemeHTTP}}, + InitialDelaySeconds: 15, PeriodSeconds: 10, TimeoutSeconds: 5, FailureThreshold: 6, + }, + SecurityContext: &corev1.SecurityContext{ + AllowPrivilegeEscalation: ptr(false), + ReadOnlyRootFilesystem: ptr(true), + Capabilities: &corev1.Capabilities{Drop: []corev1.Capability{"ALL"}}, + }, + Resources: corev1.ResourceRequirements{ + Requests: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("50m"), corev1.ResourceMemory: resource.MustParse("64Mi")}, + Limits: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("500m"), corev1.ResourceMemory: resource.MustParse("256Mi")}, + }, + }}, + }, + }, + }, + } + return deployment +} diff --git a/server/servicelb.go b/server/servicelb.go new file mode 100644 index 00000000..a12ff6a7 --- /dev/null +++ b/server/servicelb.go @@ -0,0 +1,468 @@ +package server + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net" + "os" + "reflect" + "time" + + corev1 "k8s.io/api/core/v1" + discoveryv1 "k8s.io/api/discovery/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/client-go/kubernetes" + coordinationclient "k8s.io/client-go/kubernetes/typed/coordination/v1" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/leaderelection" + "k8s.io/client-go/tools/leaderelection/resourcelock" + "k8s.io/client-go/util/retry" + + "github.com/sirupsen/logrus" +) + +const ( + serviceLBManagedIPsAnnotation = "casos.io/service-lb-managed-ips" + serviceLBDisabledAnnotation = "casos.io/service-lb-disabled" + serviceLBEnabledAnnotation = "casos.io/service-lb-enabled" + serviceLBClass = "casos.io/service-lb" + serviceLBLeaderLease = "casos-service-lb" +) + +// StartServiceLB starts the built-in bare-metal LoadBalancer reconciler. It +// publishes Ready Worker addresses and uses Service externalIPs so kube-proxy +// can route LoadBalancer traffic without a cloud provider. +func StartServiceLB(ctx context.Context, cfg *rest.Config) error { + if cfg == nil { + return fmt.Errorf("apiserver rest config is required") + } + client, err := kubernetes.NewForConfig(cfg) + if err != nil { + return fmt.Errorf("service load balancer client: %w", err) + } + if ctx == nil { + ctx = context.Background() + } + coordination, err := coordinationclient.NewForConfig(cfg) + if err != nil { + return fmt.Errorf("service load balancer coordination client: %w", err) + } + hostname, err := os.Hostname() + if err != nil { + return fmt.Errorf("service load balancer identity: %w", err) + } + identity := fmt.Sprintf("%s-%d", hostname, os.Getpid()) + elector, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ + Lock: &resourcelock.LeaseLock{ + LeaseMeta: metav1.ObjectMeta{Name: serviceLBLeaderLease, Namespace: ingressControllerNamespace}, + Client: coordination, + LockConfig: resourcelock.ResourceLockConfig{Identity: identity}, + }, + LeaseDuration: 15 * time.Second, + RenewDeadline: 10 * time.Second, + RetryPeriod: 2 * time.Second, + ReleaseOnCancel: true, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(leaderCtx context.Context) { runServiceLB(leaderCtx, client) }, + OnStoppedLeading: func() { logrus.Warn("service load balancer leadership lost") }, + }, + }) + if err != nil { + return fmt.Errorf("service load balancer leader election: %w", err) + } + go elector.Run(ctx) + return nil +} + +func runServiceLB(ctx context.Context, client kubernetes.Interface) { + const interval = 5 * time.Second + for { + if err := reconcileServiceLB(ctx, client); err != nil && ctx.Err() == nil { + logrus.Warnf("service load balancer reconciliation failed: %v", err) + } + timer := time.NewTimer(interval) + select { + case <-ctx.Done(): + timer.Stop() + return + case <-timer.C: + } + } +} + +func reconcileServiceLB(ctx context.Context, client kubernetes.Interface) error { + nodes, err := client.CoreV1().Nodes().List(ctx, metav1.ListOptions{}) + if err != nil { + return fmt.Errorf("list nodes: %w", err) + } + services, err := client.CoreV1().Services(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) + if err != nil { + return fmt.Errorf("list LoadBalancer services: %w", err) + } + reconcileErrors := make([]error, 0) + for i := range services.Items { + service := &services.Items[i] + if service.Spec.Type != corev1.ServiceTypeLoadBalancer { + continue + } + if !serviceLBManages(service) { + if _, managed := service.Annotations[serviceLBManagedIPsAnnotation]; managed { + if err := cleanupLoadBalancerService(ctx, client, service); err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("clean up LoadBalancer service %s/%s: %w", service.Namespace, service.Name, err)) + } + } + continue + } + nodeIPs, err := serviceLBNodeIPs(ctx, client, service, nodes.Items) + if err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("select LoadBalancer nodes for %s/%s: %w", service.Namespace, service.Name, err)) + continue + } + if err := reconcileLoadBalancerService(ctx, client, service, nodeIPs); err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("reconcile LoadBalancer service %s/%s: %w", service.Namespace, service.Name, err)) + } + } + return errors.Join(reconcileErrors...) +} + +func cleanupServiceLB(ctx context.Context, client kubernetes.Interface) error { + services, err := client.CoreV1().Services(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) + if err != nil { + return fmt.Errorf("list services for ServiceLB cleanup: %w", err) + } + errs := make([]error, 0) + for i := range services.Items { + service := &services.Items[i] + if service.Annotations[serviceLBManagedIPsAnnotation] == "" { + continue + } + if err := cleanupLoadBalancerService(ctx, client, service); err != nil { + errs = append(errs, fmt.Errorf("clean up ServiceLB state for %s/%s: %w", service.Namespace, service.Name, err)) + } + } + return errors.Join(errs...) +} + +func serviceLBManages(service *corev1.Service) bool { + if service == nil || service.Spec.Type != corev1.ServiceTypeLoadBalancer || service.Annotations[serviceLBDisabledAnnotation] == "true" { + return false + } + return (service.Spec.LoadBalancerClass != nil && *service.Spec.LoadBalancerClass == serviceLBClass) || service.Annotations[serviceLBEnabledAnnotation] == "true" +} + +type serviceLBNode struct { + name string + ips []string +} + +func readyServiceLBNodes(nodes []corev1.Node) []serviceLBNode { + return readyNodeAddresses(nodes, false) +} + +func readyNodeAddresses(nodes []corev1.Node, controlPlaneOnly bool) []serviceLBNode { + seen := map[string]struct{}{} + result := make([]serviceLBNode, 0) + for _, node := range nodes { + if !isReadyNode(node) || (controlPlaneOnly != isControlPlaneNode(node)) { + continue + } + externalAddresses := make([]string, 0) + internalAddresses := make([]string, 0) + for _, address := range node.Status.Addresses { + if address.Type != corev1.NodeExternalIP && address.Type != corev1.NodeInternalIP { + continue + } + ip := net.ParseIP(address.Address) + if ip == nil { + continue + } + value := ip.String() + if _, ok := seen[value]; ok { + continue + } + if address.Type == corev1.NodeExternalIP { + externalAddresses = append(externalAddresses, value) + } else { + internalAddresses = append(internalAddresses, value) + } + } + addresses := externalAddresses + if len(addresses) == 0 { + addresses = internalAddresses + } + if len(addresses) > 0 { + for _, address := range addresses { + seen[address] = struct{}{} + } + result = append(result, serviceLBNode{name: node.Name, ips: addresses}) + } + } + return result +} + +func serviceLBNodeIPs(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodes []corev1.Node) ([]string, error) { + candidates := readyServiceLBNodes(nodes) + if service.Spec.ExternalTrafficPolicy != corev1.ServiceExternalTrafficPolicyTypeLocal { + ready, err := serviceHasReadyEndpoint(ctx, client, service) + if err != nil { + return nil, err + } + if !ready { + return []string{}, nil + } + ips := make([]string, 0) + for _, node := range candidates { + ips = append(ips, node.ips...) + } + return uniqueStrings(ips), nil + } + selector := labels.SelectorFromSet(labels.Set{discoveryv1.LabelServiceName: service.Name}).String() + endpointSlices, err := client.DiscoveryV1().EndpointSlices(service.Namespace).List(ctx, metav1.ListOptions{ + LabelSelector: selector, + }) + if err != nil { + return nil, fmt.Errorf("list EndpointSlices: %w", err) + } + localNodes := map[string]struct{}{} + for _, endpointSlice := range endpointSlices.Items { + for _, endpoint := range endpointSlice.Endpoints { + if endpoint.NodeName == nil || (endpoint.Conditions.Ready != nil && !*endpoint.Conditions.Ready) { + continue + } + localNodes[*endpoint.NodeName] = struct{}{} + } + } + matchingIPs := func(candidates []serviceLBNode) []string { + ips := make([]string, 0) + for _, node := range candidates { + if _, ok := localNodes[node.name]; !ok { + continue + } + ips = append(ips, node.ips...) + } + return uniqueStrings(ips) + } + if workerIPs := matchingIPs(readyNodeAddresses(nodes, false)); len(workerIPs) > 0 { + return workerIPs, nil + } + return nil, nil +} + +func serviceHasReadyEndpoint(ctx context.Context, client kubernetes.Interface, service *corev1.Service) (bool, error) { + selector := labels.SelectorFromSet(labels.Set{discoveryv1.LabelServiceName: service.Name}).String() + endpointSlices, err := client.DiscoveryV1().EndpointSlices(service.Namespace).List(ctx, metav1.ListOptions{ + LabelSelector: selector, + }) + if err != nil { + return false, fmt.Errorf("list EndpointSlices: %w", err) + } + for _, endpointSlice := range endpointSlices.Items { + for _, endpoint := range endpointSlice.Endpoints { + if endpoint.Conditions.Ready == nil || *endpoint.Conditions.Ready { + return true, nil + } + } + } + return false, nil +} + +func isReadyNode(node corev1.Node) bool { + if node.Spec.Unschedulable { + return false + } + for _, condition := range node.Status.Conditions { + if condition.Type == corev1.NodeReady { + return condition.Status == corev1.ConditionTrue + } + } + return false +} + +func cleanupLoadBalancerService(ctx context.Context, client kubernetes.Interface, service *corev1.Service) error { + if service == nil || service.DeletionTimestamp != nil { + return nil + } + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + current, err := client.CoreV1().Services(service.Namespace).Get(ctx, service.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + return nil + } + if err != nil { + return err + } + if current.DeletionTimestamp != nil { + return nil + } + managedIPs, err := serviceLBManagedIPs(current.Annotations) + if err != nil { + return err + } + managedSet := make(map[string]struct{}, len(managedIPs)) + for _, ip := range managedIPs { + managedSet[ip] = struct{}{} + } + status := current.Status.DeepCopy() + status.LoadBalancer.Ingress = status.LoadBalancer.Ingress[:0] + for _, ingress := range current.Status.LoadBalancer.Ingress { + if _, managed := managedSet[ingress.IP]; !managed { + status.LoadBalancer.Ingress = append(status.LoadBalancer.Ingress, ingress) + } + } + if !reflect.DeepEqual(current.Status, *status) { + statusUpdate := current.DeepCopy() + statusUpdate.Status = *status + current, err = client.CoreV1().Services(current.Namespace).UpdateStatus(ctx, statusUpdate, metav1.UpdateOptions{}) + if err != nil { + return err + } + } + updated := current.DeepCopy() + updated.Spec.ExternalIPs = updated.Spec.ExternalIPs[:0] + for _, ip := range current.Spec.ExternalIPs { + if _, managed := managedSet[ip]; !managed { + updated.Spec.ExternalIPs = append(updated.Spec.ExternalIPs, ip) + } + } + delete(updated.Annotations, serviceLBManagedIPsAnnotation) + if reflect.DeepEqual(current.Spec.ExternalIPs, updated.Spec.ExternalIPs) && current.Annotations[serviceLBManagedIPsAnnotation] == "" { + return nil + } + _, err = client.CoreV1().Services(current.Namespace).Update(ctx, updated, metav1.UpdateOptions{}) + return err + }) +} + +func isControlPlaneNode(node corev1.Node) bool { + _, controlPlane := node.Labels["node-role.kubernetes.io/control-plane"] + _, master := node.Labels["node-role.kubernetes.io/master"] + return controlPlane || master +} + +func reconcileLoadBalancerService(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodeIPs []string) error { + if service == nil || service.DeletionTimestamp != nil { + return nil + } + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + current, err := client.CoreV1().Services(service.Namespace).Get(ctx, service.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + return nil + } + if err != nil { + return err + } + if current.DeletionTimestamp != nil || !serviceLBManages(current) { + return nil + } + return reconcileLoadBalancerServiceOnce(ctx, client, current, nodeIPs) + }) +} + +func reconcileLoadBalancerServiceOnce(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodeIPs []string) error { + managedIPs, err := serviceLBManagedIPs(service.Annotations) + if err != nil { + return err + } + managedSet := make(map[string]struct{}, len(managedIPs)) + for _, ip := range managedIPs { + managedSet[ip] = struct{}{} + } + userIPs := make([]string, 0, len(service.Spec.ExternalIPs)) + for _, ip := range service.Spec.ExternalIPs { + if _, ok := managedSet[ip]; !ok && !containsString(userIPs, ip) { + userIPs = append(userIPs, ip) + } + } + desiredManagedIPs := make([]string, 0, len(nodeIPs)) + for _, ip := range nodeIPs { + if !containsString(userIPs, ip) { + desiredManagedIPs = append(desiredManagedIPs, ip) + } + } + desiredExternalIPs := append(append([]string{}, userIPs...), nodeIPs...) + desiredExternalIPs = uniqueStrings(desiredExternalIPs) + encodedManagedIPs, err := json.Marshal(desiredManagedIPs) + if err != nil { + return fmt.Errorf("encode managed LoadBalancer IPs: %w", err) + } + desiredAnnotation := string(encodedManagedIPs) + + updated := service.DeepCopy() + if !reflect.DeepEqual(updated.Spec.ExternalIPs, desiredExternalIPs) || updated.Annotations[serviceLBManagedIPsAnnotation] != desiredAnnotation { + if updated.Annotations == nil { + updated.Annotations = map[string]string{} + } + updated.Spec.ExternalIPs = desiredExternalIPs + updated.Annotations[serviceLBManagedIPsAnnotation] = desiredAnnotation + current, err := client.CoreV1().Services(service.Namespace).Update(ctx, updated, metav1.UpdateOptions{}) + if err != nil { + return fmt.Errorf("update LoadBalancer service spec: %w", err) + } + updated = current + } + + desiredStatus := updated.Status.DeepCopy() + unmanagedIngress := make([]corev1.LoadBalancerIngress, 0, len(desiredStatus.LoadBalancer.Ingress)) + for _, ingress := range desiredStatus.LoadBalancer.Ingress { + if _, managed := managedSet[ingress.IP]; !managed { + unmanagedIngress = append(unmanagedIngress, ingress) + } + } + desiredStatus.LoadBalancer.Ingress = unmanagedIngress + for _, ip := range nodeIPs { + if !containsLoadBalancerIngress(desiredStatus.LoadBalancer.Ingress, ip) { + desiredStatus.LoadBalancer.Ingress = append(desiredStatus.LoadBalancer.Ingress, corev1.LoadBalancerIngress{IP: ip}) + } + } + if reflect.DeepEqual(updated.Status, *desiredStatus) { + return nil + } + updated.Status = *desiredStatus + if _, err := client.CoreV1().Services(service.Namespace).UpdateStatus(ctx, updated, metav1.UpdateOptions{}); err != nil && !apierrors.IsNotFound(err) { + return fmt.Errorf("update LoadBalancer service status: %w", err) + } + return nil +} + +func containsLoadBalancerIngress(ingresses []corev1.LoadBalancerIngress, ip string) bool { + for _, ingress := range ingresses { + if ingress.IP == ip { + return true + } + } + return false +} + +func serviceLBManagedIPs(annotations map[string]string) ([]string, error) { + value := annotations[serviceLBManagedIPsAnnotation] + if value == "" { + return nil, nil + } + var ips []string + if err := json.Unmarshal([]byte(value), &ips); err != nil { + return nil, fmt.Errorf("decode managed LoadBalancer IPs: %w", err) + } + return uniqueStrings(ips), nil +} + +func containsString(values []string, want string) bool { + for _, value := range values { + if value == want { + return true + } + } + return false +} + +func uniqueStrings(values []string) []string { + result := make([]string, 0, len(values)) + for _, value := range values { + if value != "" && !containsString(result, value) { + result = append(result, value) + } + } + return result +} diff --git a/server/storage_bootstrap.go b/server/storage_bootstrap.go index df806732..aec9781e 100644 --- a/server/storage_bootstrap.go +++ b/server/storage_bootstrap.go @@ -391,7 +391,18 @@ func hashConfigData(data map[string]string) string { return fmt.Sprintf("%x", sum[:]) } -func createOrUpdateServiceAccount(ctx context.Context, client kubernetes.Interface, sa *corev1.ServiceAccount) error { +type existingObjectValidator func(metav1.Object) error + +func validateExistingObject(object metav1.Object, validators []existingObjectValidator) error { + for _, validator := range validators { + if err := validator(object); err != nil { + return err + } + } + return nil +} + +func createOrUpdateServiceAccount(ctx context.Context, client kubernetes.Interface, sa *corev1.ServiceAccount, validators ...existingObjectValidator) error { current, err := client.CoreV1().ServiceAccounts(sa.Namespace).Get(ctx, sa.Name, metav1.GetOptions{}) if apierrors.IsNotFound(err) { _, err = client.CoreV1().ServiceAccounts(sa.Namespace).Create(ctx, sa, metav1.CreateOptions{}) @@ -403,6 +414,9 @@ func createOrUpdateServiceAccount(ctx context.Context, client kubernetes.Interfa if err != nil { return fmt.Errorf("get serviceaccount %s/%s: %w", sa.Namespace, sa.Name, err) } + if err := validateExistingObject(current, validators); err != nil { + return fmt.Errorf("validate serviceaccount %s/%s: %w", sa.Namespace, sa.Name, err) + } sa.Labels = mergeStringMap(current.Labels, sa.Labels) sa.Annotations = mergeStringMap(current.Annotations, sa.Annotations) sa.Secrets = current.Secrets @@ -415,7 +429,7 @@ func createOrUpdateServiceAccount(ctx context.Context, client kubernetes.Interfa return nil } -func createOrUpdateClusterRole(ctx context.Context, client kubernetes.Interface, role *rbacv1.ClusterRole) error { +func createOrUpdateClusterRole(ctx context.Context, client kubernetes.Interface, role *rbacv1.ClusterRole, validators ...existingObjectValidator) error { current, err := client.RbacV1().ClusterRoles().Get(ctx, role.Name, metav1.GetOptions{}) if apierrors.IsNotFound(err) { _, err = client.RbacV1().ClusterRoles().Create(ctx, role, metav1.CreateOptions{}) @@ -427,6 +441,9 @@ func createOrUpdateClusterRole(ctx context.Context, client kubernetes.Interface, if err != nil { return fmt.Errorf("get clusterrole %s: %w", role.Name, err) } + if err := validateExistingObject(current, validators); err != nil { + return fmt.Errorf("validate clusterrole %s: %w", role.Name, err) + } role.Labels = mergeStringMap(current.Labels, role.Labels) role.Annotations = mergeStringMap(current.Annotations, role.Annotations) role.ResourceVersion = current.ResourceVersion @@ -436,7 +453,7 @@ func createOrUpdateClusterRole(ctx context.Context, client kubernetes.Interface, return nil } -func createOrUpdateClusterRoleBinding(ctx context.Context, client kubernetes.Interface, binding *rbacv1.ClusterRoleBinding) error { +func createOrUpdateClusterRoleBinding(ctx context.Context, client kubernetes.Interface, binding *rbacv1.ClusterRoleBinding, validators ...existingObjectValidator) error { current, err := client.RbacV1().ClusterRoleBindings().Get(ctx, binding.Name, metav1.GetOptions{}) if apierrors.IsNotFound(err) { _, err = client.RbacV1().ClusterRoleBindings().Create(ctx, binding, metav1.CreateOptions{}) @@ -448,6 +465,9 @@ func createOrUpdateClusterRoleBinding(ctx context.Context, client kubernetes.Int if err != nil { return fmt.Errorf("get clusterrolebinding %s: %w", binding.Name, err) } + if err := validateExistingObject(current, validators); err != nil { + return fmt.Errorf("validate clusterrolebinding %s: %w", binding.Name, err) + } binding.Labels = mergeStringMap(current.Labels, binding.Labels) binding.Annotations = mergeStringMap(current.Annotations, binding.Annotations) binding.ResourceVersion = current.ResourceVersion @@ -478,7 +498,7 @@ func createOrUpdateConfigMap(ctx context.Context, client kubernetes.Interface, c return nil } -func createOrUpdateDeployment(ctx context.Context, client kubernetes.Interface, deployment *appsv1.Deployment) error { +func createOrUpdateDeployment(ctx context.Context, client kubernetes.Interface, deployment *appsv1.Deployment, validators ...existingObjectValidator) error { current, err := client.AppsV1().Deployments(deployment.Namespace).Get(ctx, deployment.Name, metav1.GetOptions{}) if apierrors.IsNotFound(err) { _, err = client.AppsV1().Deployments(deployment.Namespace).Create(ctx, deployment, metav1.CreateOptions{}) @@ -490,6 +510,9 @@ func createOrUpdateDeployment(ctx context.Context, client kubernetes.Interface, if err != nil { return fmt.Errorf("get deployment %s/%s: %w", deployment.Namespace, deployment.Name, err) } + if err := validateExistingObject(current, validators); err != nil { + return fmt.Errorf("validate deployment %s/%s: %w", deployment.Namespace, deployment.Name, err) + } deployment.Labels = mergeStringMap(current.Labels, deployment.Labels) deployment.Annotations = mergeStringMap(current.Annotations, deployment.Annotations) deployment.ResourceVersion = current.ResourceVersion diff --git a/web/src/DeploymentListPage.js b/web/src/DeploymentListPage.js index b75979a7..ffa8d56a 100644 --- a/web/src/DeploymentListPage.js +++ b/web/src/DeploymentListPage.js @@ -153,13 +153,22 @@ class DeploymentListPage extends React.Component { const {services, ingresses, nodeIP} = this.state; const urls = []; - if (nodeIP) { - const svc = services.find(s => s.name === deploy.name && s.namespace === deploy.namespace && s.type === "NodePort"); - if (svc) { - (svc.ports ?? []).filter(p => p.nodePort).forEach(p => { - urls.push({url: `http://${nodeIP}:${p.nodePort}`, type: "nodeport"}); - }); + const svc = services.find(s => s.name === deploy.name && s.namespace === deploy.namespace && (s.type === "NodePort" || s.type === "LoadBalancer")); + if (svc) { + let addresses; + if (svc.type === "LoadBalancer") { + addresses = svc.loadBalancerAddresses ?? []; + } else { + addresses = nodeIP ? [nodeIP] : []; } + const ports = svc.type === "LoadBalancer" + ? (svc.ports ?? []).filter(p => p.port) + : (svc.ports ?? []).filter(p => p.nodePort); + addresses.forEach(address => ports.forEach(p => { + const port = svc.type === "LoadBalancer" ? p.port : p.nodePort; + const host = address.includes(":") ? `[${address}]` : address; + urls.push({url: `http://${host}:${port}`, type: svc.type === "LoadBalancer" ? "loadbalancer" : "nodeport"}); + })); } const deployServiceNames = new Set( @@ -173,9 +182,19 @@ class DeploymentListPage extends React.Component { .filter(ing => ing.namespace === deploy.namespace) .forEach(ing => { (ing.rules ?? []).forEach(rule => { - if (deployServiceNames.has(rule.serviceName) && rule.host) { - const path = rule.path && rule.path !== "/" ? rule.path : ""; - urls.push({url: `http://${rule.host}${path}`, type: "domain"}); + if (!deployServiceNames.has(rule.serviceName)) { + return; + } + const path = rule.path && rule.path !== "/" ? rule.path : ""; + const scheme = ing.tlsEnabled ? "https" : "http"; + if (rule.host) { + urls.push({url: `${scheme}://${rule.host}${path}`, type: "domain"}); + } else { + const publishedAddresses = ing.loadBalancerAddresses ?? []; + publishedAddresses.forEach(address => { + const host = address.includes(":") ? `[${address}]` : address; + urls.push({url: `${scheme}://${host}${path}`, type: "ingress"}); + }); } }); }); diff --git a/web/src/ServiceListPage.js b/web/src/ServiceListPage.js index e7fe4735..144c75ed 100644 --- a/web/src/ServiceListPage.js +++ b/web/src/ServiceListPage.js @@ -238,17 +238,26 @@ class ServiceListPage extends React.Component { title: "Access URL", key: "accessUrl", render: (_, record) => { - if (record.type !== "NodePort" || !nodeIP) { + if (record.type !== "NodePort" && record.type !== "LoadBalancer") { return null; } - return (record.ports ?? []).filter(p => p.nodePort).map((p, i) => { - const url = `http://${nodeIP}:${p.nodePort}`; - return ( - + const loadBalancerAddresses = record.loadBalancerAddresses ?? []; + const ports = record.type === "LoadBalancer" + ? (record.ports ?? []).filter(p => p.port) + : (record.ports ?? []).filter(p => p.nodePort); + const addresses = record.type === "LoadBalancer" ? loadBalancerAddresses : (nodeIP ? [nodeIP] : []); + const links = []; + addresses.forEach(address => ports.forEach((p, i) => { + const port = record.type === "LoadBalancer" ? p.port : p.nodePort; + const host = address.includes(":") ? `[${address}]` : address; + const url = `http://${host}:${port}`; + links.push( + {url} ); - }); + })); + return links; }, }, {title: "Created", dataIndex: "createdAt", key: "createdAt", width: 180},