diff --git a/conf/app.conf b/conf/app.conf index d146190..7b10408 100644 --- a/conf/app.conf +++ b/conf/app.conf @@ -19,6 +19,13 @@ casdoorApplication = app-casibase ; -- Outbound proxy (optional) ---- socks5Proxy = 127.0.0.1:10808 +; -- Built-in image registry mirror ----------------------------------------- +; auto = probe registry-1.docker.io at startup; enable mirrors only when the +; canonical registry is unreachable (default) +; always = always route built-in image pulls through the mirrors +; never = always pull from the canonical registries +imageRegistryMirror = auto + ; -- Helm images (optional) ------------------------------------------------- ; Set true only when implicit/latest chart images must use IfNotPresent. helmImplicitLatestPullPolicy = false @@ -27,3 +34,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 5e54286..ec82fcf 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/node.go b/controllers/node.go index c3bee89..f91290d 100644 --- a/controllers/node.go +++ b/controllers/node.go @@ -184,9 +184,16 @@ func (c *ApiController) GetWorkerKubeconfig() { c.ResponseError("generate worker kubeconfig: " + err.Error()) return } - c.ResponseOk(map[string]string{ + resp := map[string]string{ "nodeName": wk.NodeName, "kubeconfig": wk.Kubeconfig, - "containerdConfig": deploy.GenerateContainerdConfig(cfg.SandboxImage, cfg.Socks5Proxy), - }) + "containerdConfig": deploy.GenerateContainerdConfig(cfg.SandboxImage), + } + if cfg.UseRegistryMirror { + // Manually joined workers need the per-registry mirror files too; + // containerdConfig alone only points config_path at an empty certs.d. + resp["dockerHubHostsToml"] = deploy.GenerateDockerHubHostsToml() + resp["k8sRegistryHostsToml"] = deploy.GenerateK8sRegistryHostsToml() + } + c.ResponseOk(resp) } diff --git a/controllers/service.go b/controllers/service.go index 8d34820..bffdddc 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/deploy/containerd_config.go b/deploy/containerd_config.go index fbbb1fd..5864711 100644 --- a/deploy/containerd_config.go +++ b/deploy/containerd_config.go @@ -2,11 +2,13 @@ package deploy import "fmt" +const generatedRegistryHostsMarker = "# Generated by CasOS" + // GenerateContainerdConfig returns the content for /etc/containerd/config.toml. // For containerd 2.x the registry mirrors are configured via a hosts-dir; // use GenerateDockerHubHostsToml / GenerateK8sRegistryHostsToml for those files. -func GenerateContainerdConfig(sandboxImage, socks5Proxy string) string { - base := fmt.Sprintf(`# Generated by CasOS +func GenerateContainerdConfig(sandboxImage string) string { + return fmt.Sprintf(`# Generated by CasOS version = 2 [plugins.'io.containerd.cri.v1.images'] @@ -15,24 +17,18 @@ version = 2 [plugins.'io.containerd.cri.v1.runtime'.containerd.runtimes.runc.options] SystemdCgroup = true -`, sandboxImage) - - if socks5Proxy == "" { - return base - } - // In restricted areas: point the CRI image plugin at the hosts-dir so - // per-registry mirrors (hosts.toml files) are picked up automatically. - return base + ` [plugins.'io.containerd.cri.v1.images'.registry] config_path = '/etc/containerd/certs.d' -` +`, sandboxImage) } // GenerateDockerHubHostsToml returns the content for -// /etc/containerd/certs.d/docker.io/hosts.toml in restricted areas. +// /etc/containerd/certs.d/docker.io/hosts.toml. The canonical server remains +// the fallback when the mirror is unavailable. func GenerateDockerHubHostsToml() string { - return `server = "https://registry-1.docker.io" + return generatedRegistryHostsMarker + ` +server = "https://registry-1.docker.io" [host."https://docker.1ms.run"] capabilities = ["pull", "resolve"] @@ -40,9 +36,11 @@ func GenerateDockerHubHostsToml() string { } // GenerateK8sRegistryHostsToml returns the content for -// /etc/containerd/certs.d/registry.k8s.io/hosts.toml in restricted areas. +// /etc/containerd/certs.d/registry.k8s.io/hosts.toml. The canonical server +// remains the fallback when the mirror is unavailable. func GenerateK8sRegistryHostsToml() string { - return `server = "https://registry.k8s.io" + return generatedRegistryHostsMarker + ` +server = "https://registry.k8s.io" [host."https://registry.aliyuncs.com/google_containers"] capabilities = ["pull", "resolve"] diff --git a/deploy/installer.go b/deploy/installer.go index 48c5991..eefe683 100644 --- a/deploy/installer.go +++ b/deploy/installer.go @@ -7,6 +7,16 @@ import ( const nodeDeployResolverPath = "/etc/casos-resolv.conf" +const ( + dockerHubHostsPath = "/etc/containerd/certs.d/docker.io/hosts.toml" + k8sRegistryHostsPath = "/etc/containerd/certs.d/registry.k8s.io/hosts.toml" +) + +type registryMirrorFileRunner interface { + RunRootContext(ctx context.Context, command string) (string, error) + WriteFileContext(ctx context.Context, path, content, mode string) error +} + func (d *NodeDeployer) installNodeBinaries(ctx context.Context, runner *NodeDeploySSHRunner, arch, k8sVersion string) error { version := k8sVersion cniVersion := defaultNodeDeployCNIVersion @@ -47,16 +57,11 @@ test -f %[1]s`, nodeDeployResolverPath)); err != nil { } d.logStep(nodeDeployPhaseConfiguring, "Configuring containerd") - if err := runner.WriteFileContext(ctx, "/etc/containerd/config.toml", GenerateContainerdConfig(d.config.SandboxImage, d.config.Socks5Proxy), "0644"); err != nil { + if err := runner.WriteFileContext(ctx, "/etc/containerd/config.toml", GenerateContainerdConfig(d.config.SandboxImage), "0644"); err != nil { return fmt.Errorf("write /etc/containerd/config.toml: %w", err) } - if d.config.Socks5Proxy != "" { - if err := runner.WriteFileContext(ctx, "/etc/containerd/certs.d/docker.io/hosts.toml", GenerateDockerHubHostsToml(), "0644"); err != nil { - return fmt.Errorf("write /etc/containerd/certs.d/docker.io/hosts.toml: %w", err) - } - if err := runner.WriteFileContext(ctx, "/etc/containerd/certs.d/registry.k8s.io/hosts.toml", GenerateK8sRegistryHostsToml(), "0644"); err != nil { - return fmt.Errorf("write /etc/containerd/certs.d/registry.k8s.io/hosts.toml: %w", err) - } + if err := reconcileRegistryMirrorFiles(ctx, runner, d.config.UseRegistryMirror); err != nil { + return err } if _, err := runner.RunRootContext(ctx, "systemctl enable --now containerd && systemctl restart containerd"); err != nil { return fmt.Errorf("start containerd: %w", err) @@ -95,6 +100,29 @@ fi`, version, version, arch, version, arch, cniVersion, arch, cniVersion) return nil } +func reconcileRegistryMirrorFiles(ctx context.Context, runner registryMirrorFileRunner, enabled bool) error { + if enabled { + if err := runner.WriteFileContext(ctx, dockerHubHostsPath, GenerateDockerHubHostsToml(), "0644"); err != nil { + return fmt.Errorf("write %s: %w", dockerHubHostsPath, err) + } + if err := runner.WriteFileContext(ctx, k8sRegistryHostsPath, GenerateK8sRegistryHostsToml(), "0644"); err != nil { + return fmt.Errorf("write %s: %w", k8sRegistryHostsPath, err) + } + return nil + } + + cleanupCommand := fmt.Sprintf(`set -e +for path in %s %s; do + if [ -f "$path" ] && [ "$(sed -n '1p' "$path")" = %s ]; then + rm -f -- "$path" + fi +done`, shellSingleQuote(dockerHubHostsPath), shellSingleQuote(k8sRegistryHostsPath), shellSingleQuote(generatedRegistryHostsMarker)) + if _, err := runner.RunRootContext(ctx, cleanupCommand); err != nil { + return fmt.Errorf("remove managed containerd registry hosts: %w", err) + } + return nil +} + func (d *NodeDeployer) writeNodeFiles(ctx context.Context, runner *NodeDeploySSHRunner, nodeName, kubeconfig string) error { ca, err := extractCertificateAuthority(kubeconfig) if err != nil { diff --git a/deploy/types.go b/deploy/types.go index 1026fc8..4ae379a 100644 --- a/deploy/types.go +++ b/deploy/types.go @@ -117,17 +117,17 @@ type Config struct { ApiserverBind string ApiserverPort int SandboxImage string - Socks5Proxy string + UseRegistryMirror bool GenerateKubeconfig KubeconfigGenerator } func ConfigFromServerConfig(cfg server.Config) Config { return Config{ - AdvertiseAddress: cfg.AdvertiseAddress, - ApiserverBind: cfg.ApiserverBind, - ApiserverPort: cfg.ApiserverPort, - SandboxImage: cfg.SandboxImage, - Socks5Proxy: cfg.Socks5Proxy, + AdvertiseAddress: cfg.AdvertiseAddress, + ApiserverBind: cfg.ApiserverBind, + ApiserverPort: cfg.ApiserverPort, + SandboxImage: cfg.SandboxImage, + UseRegistryMirror: cfg.UseRegistryMirror, GenerateKubeconfig: func(nodeName, apiserverURL string) (*NodeKubeconfig, error) { wk, err := server.GenerateWorkerKubeconfigForServer(cfg, nodeName, apiserverURL) if err != nil { diff --git a/main.go b/main.go index aa96034..e442c42 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, srvCfg); 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/bootstrap.go b/server/bootstrap.go index 8697296..7045150 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 ad819ab..64c66b6 100644 --- a/server/config.go +++ b/server/config.go @@ -6,6 +6,7 @@ import ( "strings" "github.com/casosorg/casos/conf" + "github.com/sirupsen/logrus" ) // Config holds control-plane settings populated from app.conf. @@ -17,13 +18,17 @@ type Config struct { WebhookPort int // HTTPS port for the Casbin admission webhook server DSN string // MySQL DSN forwarded to kine SandboxImage string // containerd sandbox (pause) image, empty = upstream default - Socks5Proxy string // outbound socks5 proxy, e.g. 127.0.0.1:10808 CoreDNSImage string // CoreDNS image used by the built-in DNS bootstrap LocalPathProvisionerImage string // local-path-provisioner controller image 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 + ServiceLBImage string // hostPort proxy image used by the built-in ServiceLB 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 + UseRegistryMirror bool // route built-in image pulls through the registry mirrors (resolved from imageRegistryMirror) } // ConfigFromAppConf reads server config from the beego app.conf. @@ -48,25 +53,25 @@ func ConfigFromAppConf() (Config, error) { webhookPort := conf.GetConfigIntDefault("webhookPort", 9443) - socks5Proxy := conf.GetConfigString("socks5Proxy") - - sandboxImage := conf.GetConfigString("sandboxImage") - if sandboxImage == "" { - if socks5Proxy != "" { - sandboxImage = "registry.aliyuncs.com/google_containers/pause:3.10.1" - } else { - sandboxImage = "registry.k8s.io/pause:3.10.1" - } - } + sandboxImage := conf.GetConfigStringDefault("sandboxImage", "registry.k8s.io/pause:3.10.1") storageProvisionerEnabled := conf.GetConfigBoolDefault("storageProvisionerEnabled", true) - 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") + ingressControllerEnabled := conf.GetConfigBoolDefault("ingressControllerEnabled", false) + serviceLBEnabled := conf.GetConfigBoolDefault("serviceLBEnabled", false) + coreDNSImage := conf.GetConfigStringDefault("coreDNSImage", "docker.io/coredns/coredns:1.12.4") + localPathProvisionerImage := conf.GetConfigStringDefault("localPathProvisionerImage", "docker.io/rancher/local-path-provisioner:v0.0.32") + localPathHelperImage := conf.GetConfigStringDefault("localPathHelperImage", "docker.io/library/busybox:1.37.0") flannelImage := conf.GetConfigStringDefault("flannelImage", defaultFlannelImage) flannelCNIPluginImage := conf.GetConfigStringDefault("flannelCNIPluginImage", defaultFlannelCNIPluginImage) + ingressControllerImage := conf.GetConfigStringDefault("ingressControllerImage", defaultIngressControllerImage) + serviceLBImage := conf.GetConfigStringDefault("serviceLBImage", defaultServiceLBImage) + + useRegistryMirror, err := resolveRegistryMirror() + if err != nil { + return Config{}, err + } - return Config{ + config := Config{ DataDir: dataDir, ApiserverBind: bind, AdvertiseAddress: advertise, @@ -74,14 +79,27 @@ func ConfigFromAppConf() (Config, error) { WebhookPort: webhookPort, DSN: dsn, SandboxImage: sandboxImage, - Socks5Proxy: socks5Proxy, CoreDNSImage: coreDNSImage, LocalPathProvisionerImage: localPathProvisionerImage, LocalPathHelperImage: localPathHelperImage, FlannelImage: flannelImage, FlannelCNIPluginImage: flannelCNIPluginImage, + IngressControllerImage: ingressControllerImage, + ServiceLBImage: serviceLBImage, StorageProvisionerEnabled: storageProvisionerEnabled, - }, nil + IngressControllerEnabled: ingressControllerEnabled, + ServiceLBEnabled: serviceLBEnabled, + UseRegistryMirror: useRegistryMirror, + } + normalizeApplicationAccessConfig(&config) + return config, nil +} + +func normalizeApplicationAccessConfig(config *Config) { + if config != nil && config.IngressControllerEnabled && !config.ServiceLBEnabled { + logrus.Warn("ingressControllerEnabled requires serviceLBEnabled; disabling the built-in Ingress controller") + config.IngressControllerEnabled = false + } } // injectDBName inserts dbName into a MySQL DSN of the form diff --git a/server/flannel_bootstrap.go b/server/flannel_bootstrap.go index 5e1fe7e..97d0285 100644 --- a/server/flannel_bootstrap.go +++ b/server/flannel_bootstrap.go @@ -19,8 +19,8 @@ import ( ) const ( - defaultFlannelImage = "docker.1ms.run/flannel/flannel:v0.27.4" - defaultFlannelCNIPluginImage = "docker.1ms.run/flannel/flannel-cni-plugin:v1.8.0-flannel1" + defaultFlannelImage = "docker.io/flannel/flannel:v0.27.4" + defaultFlannelCNIPluginImage = "docker.io/flannel/flannel-cni-plugin:v1.8.0-flannel1" flannelNamespace = "kube-flannel" flannelServiceAccount = "flannel" flannelConfigMap = "kube-flannel-cfg" diff --git a/server/ingress_bootstrap.go b/server/ingress_bootstrap.go new file mode 100644 index 0000000..230905b --- /dev/null +++ b/server/ingress_bootstrap.go @@ -0,0 +1,434 @@ +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/library/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"}}, + }, + VolumeMounts: []corev1.VolumeMount{{Name: "tmp", MountPath: "/tmp"}}, + 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")}, + }, + }}, + Volumes: []corev1.Volume{{Name: "tmp", VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{}}}}, + }, + }, + }, + } + return deployment +} diff --git a/server/registry_mirror.go b/server/registry_mirror.go new file mode 100644 index 0000000..31ec335 --- /dev/null +++ b/server/registry_mirror.go @@ -0,0 +1,75 @@ +package server + +import ( + "context" + "fmt" + "net/http" + "strings" + "sync" + "time" + + "github.com/casosorg/casos/conf" +) + +// imageRegistryMirror modes accepted in app.conf. +const ( + registryMirrorModeAuto = "auto" + registryMirrorModeAlways = "always" + registryMirrorModeNever = "never" +) + +// registryProbeURL is probed once (mode "auto") to decide whether the canonical +// registries are reachable. Any HTTP response, including the 401 that Docker +// Hub returns for anonymous requests, counts as reachable; only a connection +// failure or timeout marks the environment as restricted. +const ( + registryProbeURL = "https://registry-1.docker.io/v2/" + registryProbeTimeout = 4 * time.Second +) + +var ( + registryProbeOnce sync.Once + registryProbeRestricted bool +) + +// resolveRegistryMirror reads imageRegistryMirror from app.conf and decides +// whether built-in image pulls should be routed through the configured +// registry mirrors. Mode "auto" (the default) probes the canonical registry +// directly instead of guessing from timezone, locale, or IP geolocation; the +// result is cached for the lifetime of the process. An explicit +// "always"/"never" skips the probe entirely. +func resolveRegistryMirror() (bool, error) { + mode := strings.ToLower(strings.TrimSpace(conf.GetConfigStringDefault("imageRegistryMirror", registryMirrorModeAuto))) + switch mode { + case registryMirrorModeAlways: + return true, nil + case registryMirrorModeNever: + return false, nil + case registryMirrorModeAuto: + return canonicalRegistryRestricted(), nil + default: + return false, fmt.Errorf("invalid imageRegistryMirror %q in app.conf: expected auto, always, or never", mode) + } +} + +// canonicalRegistryRestricted reports whether the canonical Docker registry is +// unreachable from this host. The probe runs at most once per process. +func canonicalRegistryRestricted() bool { + registryProbeOnce.Do(func() { + ctx, cancel := context.WithTimeout(context.Background(), registryProbeTimeout) + defer cancel() + req, err := http.NewRequestWithContext(ctx, http.MethodHead, registryProbeURL, nil) + if err != nil { + registryProbeRestricted = false + return + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + registryProbeRestricted = true + return + } + resp.Body.Close() + registryProbeRestricted = false + }) + return registryProbeRestricted +} diff --git a/server/servicelb.go b/server/servicelb.go new file mode 100644 index 0000000..6620f86 --- /dev/null +++ b/server/servicelb.go @@ -0,0 +1,727 @@ +package server + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "net" + "os" + "reflect" + "sort" + "strconv" + "strings" + "time" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + discoveryv1 "k8s.io/api/discovery/v1" + apiequality "k8s.io/apimachinery/pkg/api/equality" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/util/intstr" + "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" + appsinternal "k8s.io/kubernetes/pkg/apis/apps/v1" + + "github.com/sirupsen/logrus" +) + +const ( + serviceLBNamespace = "kube-system" + serviceLBManagedIPsAnnotation = "casos.io/service-lb-managed-ips" + serviceLBDisabledAnnotation = "casos.io/service-lb-disabled" + serviceLBManagedByLabel = "app.kubernetes.io/managed-by" + serviceLBComponentLabel = "app.kubernetes.io/component" + serviceLBServiceNameLabel = "casos.io/service-lb-service-name" + serviceLBServiceNamespaceLabel = "casos.io/service-lb-service-namespace" + serviceLBServiceUIDLabel = "casos.io/service-lb-service-uid" + serviceLBClass = "casos.io/service-lb" + serviceLBLeaderLease = "casos-service-lb" + defaultServiceLBImage = "docker.io/rancher/klipper-lb:v0.4.17" +) + +// StartServiceLB starts the built-in bare-metal LoadBalancer controller. Each +// managed Service gets a hostPort proxy DaemonSet; only nodes with Ready proxy +// Pods are published in status.loadBalancer. User-owned Service spec is not +// used as the steady-state data plane. +func StartServiceLB(ctx context.Context, cfg *rest.Config, srvCfg 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: serviceLBNamespace}, + 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, srvCfg.ServiceLBImage) + }, + 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, image string) { + const interval = 5 * time.Second + for { + if err := reconcileServiceLBWithImage(ctx, client, image); 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 reconcileServiceLBWithImage(ctx context.Context, client kubernetes.Interface, image string) error { + if strings.TrimSpace(image) == "" { + image = defaultServiceLBImage + } + 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) + } + existingDaemonSets, err := listServiceLBDaemonSets(ctx, client) + if err != nil { + return err + } + existingDaemonSetNames := make(map[string]struct{}, len(existingDaemonSets)) + for i := range existingDaemonSets { + existingDaemonSetNames[existingDaemonSets[i].Name] = struct{}{} + } + + desiredDaemonSets := make(map[string]struct{}) + reconcileErrors := make([]error, 0) + for i := range services.Items { + service := &services.Items[i] + if !serviceLBManages(service) { + _, hadDaemonSet := existingDaemonSetNames[serviceLBDaemonSetName(service)] + if hadDaemonSet || serviceLBWasExplicitlyOwned(service) || service.Annotations[serviceLBManagedIPsAnnotation] != "" { + 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 + } + + desiredDaemonSets[serviceLBDaemonSetName(service)] = struct{}{} + desired, err := buildServiceLBDaemonSet(service, image) + if err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("build ServiceLB DaemonSet for %s/%s: %w", service.Namespace, service.Name, err)) + continue + } + if err := createOrUpdateServiceLBDaemonSet(ctx, client, desired); err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("reconcile ServiceLB DaemonSet for %s/%s: %w", service.Namespace, service.Name, err)) + continue + } + nodeIPs, err := serviceLBReadyProxyNodeIPs(ctx, client, service, nodes.Items) + if err != nil { + reconcileErrors = append(reconcileErrors, fmt.Errorf("select ServiceLB 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 status for %s/%s: %w", service.Namespace, service.Name, err)) + } + } + if err := cleanupOrphanedServiceLBDaemonSets(ctx, client, desiredDaemonSets); err != nil { + reconcileErrors = append(reconcileErrors, err) + } + return errors.Join(reconcileErrors...) +} + +func cleanupServiceLB(ctx context.Context, client kubernetes.Interface) error { + daemonSets, err := listServiceLBDaemonSets(ctx, client) + if err != nil { + return err + } + daemonSetNames := make(map[string]struct{}, len(daemonSets)) + for i := range daemonSets { + daemonSetNames[daemonSets[i].Name] = struct{}{} + } + 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] + _, hadDaemonSet := daemonSetNames[serviceLBDaemonSetName(service)] + if !hadDaemonSet && !serviceLBWasExplicitlyOwned(service) && 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)) + } + } + for i := range daemonSets { + if err := client.AppsV1().DaemonSets(serviceLBNamespace).Delete(ctx, daemonSets[i].Name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + errs = append(errs, fmt.Errorf("delete ServiceLB DaemonSet %s/%s: %w", serviceLBNamespace, daemonSets[i].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 + } + if service.Spec.LoadBalancerClass == nil { + return true + } + return *service.Spec.LoadBalancerClass == serviceLBClass +} + +func serviceLBWasExplicitlyOwned(service *corev1.Service) bool { + return service != nil && service.Spec.LoadBalancerClass != nil && *service.Spec.LoadBalancerClass == serviceLBClass +} + +func serviceLBDaemonSetName(service *corev1.Service) string { + identity := string(service.UID) + if identity == "" { + identity = service.Namespace + "/" + service.Name + } + sum := sha256.Sum256([]byte(identity)) + suffix := hex.EncodeToString(sum[:6]) + name := strings.Trim(service.Name, "-") + const prefix = "svclb-" + maxNameLen := 63 - len(prefix) - 1 - len(suffix) + if len(name) > maxNameLen { + name = strings.TrimRight(name[:maxNameLen], "-") + } + if name == "" { + name = "service" + } + return prefix + name + "-" + suffix +} + +func serviceLBLabels(service *corev1.Service) map[string]string { + return map[string]string{ + serviceLBManagedByLabel: "casos", + serviceLBComponentLabel: "service-lb", + serviceLBServiceNameLabel: service.Name, + serviceLBServiceNamespaceLabel: service.Namespace, + serviceLBServiceUIDLabel: string(service.UID), + } +} + +func buildServiceLBDaemonSet(service *corev1.Service, image string) (*appsv1.DaemonSet, error) { + if service == nil { + return nil, fmt.Errorf("service is required") + } + if len(service.Spec.Ports) == 0 { + return nil, fmt.Errorf("service has no ports") + } + if strings.TrimSpace(image) == "" { + image = defaultServiceLBImage + } + name := serviceLBDaemonSetName(service) + labels := serviceLBLabels(service) + labels["app.kubernetes.io/name"] = name + selectorLabels := map[string]string{"app.kubernetes.io/name": name} + maxUnavailable := intstr.FromInt(1) + automountToken := false + sourceRanges := append([]string{}, service.Spec.LoadBalancerSourceRanges...) + if len(sourceRanges) == 0 { + sourceRanges = []string{"0.0.0.0/0"} + if serviceUsesIPFamily(service, corev1.IPv6Protocol) { + sourceRanges = append(sourceRanges, "::/0") + } + } + sort.Strings(sourceRanges) + + podSecurityContext := &corev1.PodSecurityContext{} + if serviceUsesIPFamily(service, corev1.IPv4Protocol) { + podSecurityContext.Sysctls = append(podSecurityContext.Sysctls, corev1.Sysctl{Name: "net.ipv4.ip_forward", Value: "1"}) + } + if serviceUsesIPFamily(service, corev1.IPv6Protocol) { + podSecurityContext.Sysctls = append(podSecurityContext.Sysctls, corev1.Sysctl{Name: "net.ipv6.conf.all.forwarding", Value: "1"}) + } + + containers := make([]corev1.Container, 0, len(service.Spec.Ports)) + for _, port := range service.Spec.Ports { + if port.Port <= 0 { + return nil, fmt.Errorf("service port %q has invalid port %d", port.Name, port.Port) + } + container := corev1.Container{ + Name: fmt.Sprintf("lb-%s-%d", strings.ToLower(string(port.Protocol)), port.Port), + Image: image, + ImagePullPolicy: corev1.PullIfNotPresent, + Ports: []corev1.ContainerPort{{ + Name: fmt.Sprintf("lb-%s-%d", strings.ToLower(string(port.Protocol)), port.Port), + ContainerPort: port.Port, + HostPort: port.Port, + Protocol: port.Protocol, + }}, + Env: []corev1.EnvVar{ + {Name: "SRC_PORT", Value: strconv.Itoa(int(port.Port))}, + {Name: "SRC_RANGES", Value: strings.Join(sourceRanges, ",")}, + {Name: "DEST_PROTO", Value: string(port.Protocol)}, + }, + SecurityContext: &corev1.SecurityContext{Capabilities: &corev1.Capabilities{Add: []corev1.Capability{"NET_ADMIN"}}}, + } + if service.Spec.ExternalTrafficPolicy == corev1.ServiceExternalTrafficPolicyTypeLocal { + if port.NodePort == 0 { + return nil, fmt.Errorf("service port %q requires a NodePort for externalTrafficPolicy=Local", port.Name) + } + container.Env = append( + container.Env, + corev1.EnvVar{Name: "DEST_PORT", Value: strconv.Itoa(int(port.NodePort))}, + corev1.EnvVar{Name: "DEST_IPS", ValueFrom: &corev1.EnvVarSource{FieldRef: &corev1.ObjectFieldSelector{FieldPath: "status.hostIPs"}}}, + ) + } else { + clusterIPs := serviceClusterIPs(service) + if len(clusterIPs) == 0 { + return nil, fmt.Errorf("service has no routable ClusterIP") + } + container.Env = append( + container.Env, + corev1.EnvVar{Name: "DEST_PORT", Value: strconv.Itoa(int(port.Port))}, + corev1.EnvVar{Name: "DEST_IPS", Value: strings.Join(clusterIPs, ",")}, + ) + } + containers = append(containers, container) + } + + return &appsv1.DaemonSet{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: serviceLBNamespace, Labels: labels}, + Spec: appsv1.DaemonSetSpec{ + Selector: &metav1.LabelSelector{MatchLabels: selectorLabels}, + UpdateStrategy: appsv1.DaemonSetUpdateStrategy{ + Type: appsv1.RollingUpdateDaemonSetStrategyType, + RollingUpdate: &appsv1.RollingUpdateDaemonSet{MaxUnavailable: &maxUnavailable}, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: mergeStringMap(labels, selectorLabels)}, + Spec: corev1.PodSpec{ + AutomountServiceAccountToken: &automountToken, + PriorityClassName: "system-node-critical", + SecurityContext: podSecurityContext, + Affinity: &corev1.Affinity{NodeAffinity: &corev1.NodeAffinity{RequiredDuringSchedulingIgnoredDuringExecution: &corev1.NodeSelector{ + NodeSelectorTerms: []corev1.NodeSelectorTerm{{MatchExpressions: []corev1.NodeSelectorRequirement{ + {Key: "node-role.kubernetes.io/control-plane", Operator: corev1.NodeSelectorOpDoesNotExist}, + {Key: "node-role.kubernetes.io/master", Operator: corev1.NodeSelectorOpDoesNotExist}, + }}}, + }}}, + Tolerations: []corev1.Toleration{ + {Key: "CriticalAddonsOnly", Operator: corev1.TolerationOpExists}, + {Key: "casos.io/bootstrap", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule}, + }, + Containers: containers, + }, + }, + }, + }, nil +} + +func createOrUpdateServiceLBDaemonSet(ctx context.Context, client kubernetes.Interface, desired *appsv1.DaemonSet) error { + daemonSets := client.AppsV1().DaemonSets(desired.Namespace) + current, err := daemonSets.Get(ctx, desired.Name, metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + _, err = daemonSets.Create(ctx, desired, metav1.CreateOptions{}) + return err + } + if err != nil { + return err + } + if current.Labels[serviceLBManagedByLabel] != "casos" || current.Labels[serviceLBComponentLabel] != "service-lb" { + return fmt.Errorf("DaemonSet %s/%s exists and is not managed by CasOS", desired.Namespace, desired.Name) + } + for _, key := range []string{serviceLBServiceNameLabel, serviceLBServiceNamespaceLabel, serviceLBServiceUIDLabel} { + if current.Labels[key] != desired.Labels[key] { + return fmt.Errorf("DaemonSet %s/%s belongs to a different Service", desired.Namespace, desired.Name) + } + } + desired.Labels = mergeStringMap(current.Labels, desired.Labels) + desired.Annotations = mergeStringMap(current.Annotations, desired.Annotations) + currentDefaulted := current.DeepCopy() + desiredDefaulted := desired.DeepCopy() + appsinternal.SetObjectDefaults_DaemonSet(currentDefaulted) + appsinternal.SetObjectDefaults_DaemonSet(desiredDefaulted) + if apiequality.Semantic.DeepEqual(currentDefaulted.Labels, desiredDefaulted.Labels) && + apiequality.Semantic.DeepEqual(currentDefaulted.Annotations, desiredDefaulted.Annotations) && + apiequality.Semantic.DeepEqual(currentDefaulted.Spec, desiredDefaulted.Spec) { + return nil + } + desired.ResourceVersion = current.ResourceVersion + _, err = daemonSets.Update(ctx, desired, metav1.UpdateOptions{}) + return err +} + +func cleanupOrphanedServiceLBDaemonSets(ctx context.Context, client kubernetes.Interface, desired map[string]struct{}) error { + daemonSets, err := listServiceLBDaemonSets(ctx, client) + if err != nil { + return err + } + errs := make([]error, 0) + for i := range daemonSets { + daemonSet := &daemonSets[i] + if _, ok := desired[daemonSet.Name]; ok { + continue + } + if err := client.AppsV1().DaemonSets(daemonSet.Namespace).Delete(ctx, daemonSet.Name, metav1.DeleteOptions{}); err != nil && !apierrors.IsNotFound(err) { + errs = append(errs, fmt.Errorf("delete orphaned ServiceLB DaemonSet %s/%s: %w", daemonSet.Namespace, daemonSet.Name, err)) + } + } + return errors.Join(errs...) +} + +func listServiceLBDaemonSets(ctx context.Context, client kubernetes.Interface) ([]appsv1.DaemonSet, error) { + daemonSets, err := client.AppsV1().DaemonSets(serviceLBNamespace).List(ctx, metav1.ListOptions{LabelSelector: labels.SelectorFromSet(labels.Set{ + serviceLBManagedByLabel: "casos", + serviceLBComponentLabel: "service-lb", + }).String()}) + if err != nil { + return nil, fmt.Errorf("list ServiceLB DaemonSets: %w", err) + } + return daemonSets.Items, nil +} + +func serviceLBReadyProxyNodeIPs(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodes []corev1.Node) ([]string, error) { + selector := labels.SelectorFromSet(labels.Set{ + serviceLBServiceNameLabel: service.Name, + serviceLBServiceNamespaceLabel: service.Namespace, + serviceLBServiceUIDLabel: string(service.UID), + }).String() + pods, err := client.CoreV1().Pods(serviceLBNamespace).List(ctx, metav1.ListOptions{LabelSelector: selector}) + if err != nil { + return nil, fmt.Errorf("list ServiceLB proxy Pods: %w", err) + } + readyProxyNodes := make(map[string]struct{}) + for i := range pods.Items { + pod := &pods.Items[i] + if pod.Spec.NodeName == "" || pod.Status.PodIP == "" || !isReadyPod(pod) { + continue + } + readyProxyNodes[pod.Spec.NodeName] = struct{}{} + } + if service.Spec.ExternalTrafficPolicy == corev1.ServiceExternalTrafficPolicyTypeLocal { + localNodes, err := serviceReadyEndpointNodes(ctx, client, service) + if err != nil { + return nil, err + } + for nodeName := range readyProxyNodes { + if _, ok := localNodes[nodeName]; !ok { + delete(readyProxyNodes, nodeName) + } + } + } + return readyProxyNodeAddresses(nodes, readyProxyNodes, service.Spec.IPFamilies), nil +} + +func serviceReadyEndpointNodes(ctx context.Context, client kubernetes.Interface, service *corev1.Service) (map[string]struct{}, 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 nil, fmt.Errorf("list EndpointSlices: %w", err) + } + nodes := make(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 + } + nodes[*endpoint.NodeName] = struct{}{} + } + } + return nodes, nil +} + +func readyProxyNodeAddresses(nodes []corev1.Node, readyProxyNodes map[string]struct{}, families []corev1.IPFamily) []string { + result := make([]string, 0) + for _, node := range nodes { + if _, ok := readyProxyNodes[node.Name]; !ok || !isReadyNode(node) || isControlPlaneNode(node) { + continue + } + external := make([]string, 0) + internal := make([]string, 0) + for _, address := range node.Status.Addresses { + if !serviceLBAddressMatchesFamilies(address.Address, families) { + continue + } + switch address.Type { + case corev1.NodeExternalIP: + external = append(external, address.Address) + case corev1.NodeInternalIP: + internal = append(internal, address.Address) + } + } + if len(external) > 0 { + result = append(result, external...) + } else { + result = append(result, internal...) + } + } + sort.Strings(result) + return uniqueStrings(result) +} + +func serviceLBAddressMatchesFamilies(address string, families []corev1.IPFamily) bool { + ip := net.ParseIP(address) + if ip == nil { + return false + } + if len(families) == 0 { + return true + } + want := corev1.IPv6Protocol + if ip.To4() != nil { + want = corev1.IPv4Protocol + } + for _, family := range families { + if family == want { + return true + } + } + return false +} + +func isReadyPod(pod *corev1.Pod) bool { + for _, condition := range pod.Status.Conditions { + if condition.Type == corev1.PodReady { + return condition.Status == corev1.ConditionTrue + } + } + return false +} + +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 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 { + previousManagedIPs, migrationErr := serviceLBManagedIPs(service.Annotations) + current := service + var err error + if migrationErr == nil { + current, err = removeLegacyServiceLBExternalIPs(ctx, client, service, previousManagedIPs) + if err != nil { + return err + } + } + desiredManagedIPs := uniqueStrings(append([]string{}, nodeIPs...)) + desiredStatus := current.Status.DeepCopy() + desiredStatus.LoadBalancer.Ingress = nil + for _, ip := range desiredManagedIPs { + desiredStatus.LoadBalancer.Ingress = append(desiredStatus.LoadBalancer.Ingress, corev1.LoadBalancerIngress{IP: ip}) + } + if !apiequality.Semantic.DeepEqual(current.Status.LoadBalancer, desiredStatus.LoadBalancer) { + statusUpdate := current.DeepCopy() + statusUpdate.Status = *desiredStatus + current, err = client.CoreV1().Services(current.Namespace).UpdateStatus(ctx, statusUpdate, metav1.UpdateOptions{}) + if err != nil { + return fmt.Errorf("update LoadBalancer service status: %w", err) + } + } + if migrationErr != nil { + return fmt.Errorf("preserve malformed legacy ServiceLB state: %w", migrationErr) + } + return removeLegacyServiceLBManagedIPsAnnotation(ctx, client, current.Namespace, current.Name) +} + +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 + } + status := current.Status.DeepCopy() + status.LoadBalancer = corev1.LoadBalancerStatus{} + if !apiequality.Semantic.DeepEqual(current.Status.LoadBalancer, status.LoadBalancer) { + statusUpdate := current.DeepCopy() + statusUpdate.Status = *status + current, err = client.CoreV1().Services(current.Namespace).UpdateStatus(ctx, statusUpdate, metav1.UpdateOptions{}) + if err != nil { + return err + } + } + managedIPs, err := serviceLBManagedIPs(current.Annotations) + if err != nil { + return fmt.Errorf("preserve malformed legacy ServiceLB state: %w", err) + } + current, err = removeLegacyServiceLBExternalIPs(ctx, client, current, managedIPs) + if err != nil { + return err + } + return removeLegacyServiceLBManagedIPsAnnotation(ctx, client, current.Namespace, current.Name) + }) +} + +func removeLegacyServiceLBExternalIPs(ctx context.Context, client kubernetes.Interface, service *corev1.Service, managedIPs []string) (*corev1.Service, error) { + if len(managedIPs) == 0 || len(service.Spec.ExternalIPs) == 0 { + return service, nil + } + managed := make(map[string]struct{}, len(managedIPs)) + for _, ip := range managedIPs { + managed[ip] = struct{}{} + } + desired := service.DeepCopy() + desired.Spec.ExternalIPs = desired.Spec.ExternalIPs[:0] + for _, ip := range service.Spec.ExternalIPs { + if _, ok := managed[ip]; !ok { + desired.Spec.ExternalIPs = append(desired.Spec.ExternalIPs, ip) + } + } + if reflect.DeepEqual(service.Spec.ExternalIPs, desired.Spec.ExternalIPs) { + return service, nil + } + updated, err := client.CoreV1().Services(service.Namespace).Update(ctx, desired, metav1.UpdateOptions{}) + if err != nil { + return nil, fmt.Errorf("remove legacy ServiceLB externalIPs: %w", err) + } + return updated, nil +} + +func removeLegacyServiceLBManagedIPsAnnotation(ctx context.Context, client kubernetes.Interface, namespace, name string) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + current, err := client.CoreV1().Services(namespace).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + return err + } + if _, ok := current.Annotations[serviceLBManagedIPsAnnotation]; !ok { + return nil + } + updated := current.DeepCopy() + delete(updated.Annotations, serviceLBManagedIPsAnnotation) + _, err = client.CoreV1().Services(namespace).Update(ctx, updated, metav1.UpdateOptions{}) + return err + }) +} + +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 serviceClusterIPs(service *corev1.Service) []string { + clusterIPs := append([]string{}, service.Spec.ClusterIPs...) + if len(clusterIPs) == 0 && service.Spec.ClusterIP != "" { + clusterIPs = append(clusterIPs, service.Spec.ClusterIP) + } + result := make([]string, 0, len(clusterIPs)) + for _, ip := range clusterIPs { + if ip != "" && !strings.EqualFold(ip, corev1.ClusterIPNone) { + result = append(result, ip) + } + } + return uniqueStrings(result) +} + +func serviceUsesIPFamily(service *corev1.Service, family corev1.IPFamily) bool { + if len(service.Spec.IPFamilies) == 0 { + return family == corev1.IPv4Protocol + } + for _, current := range service.Spec.IPFamilies { + if current == family { + return true + } + } + return false +} + +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/servicelb_test.go b/server/servicelb_test.go new file mode 100644 index 0000000..f9fbbab --- /dev/null +++ b/server/servicelb_test.go @@ -0,0 +1,80 @@ +package server + +import ( + "regexp" + "strings" + "testing" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" +) + +func TestServiceLBDaemonSetName(t *testing.T) { + tests := []struct { + name string + service *corev1.Service + want string + }{ + { + name: "uses uid for the stable suffix", + service: &corev1.Service{ObjectMeta: metav1.ObjectMeta{ + Name: "api", + UID: types.UID("uid-1"), + }}, + want: "svclb-api-4a49acf8a6bd", + }, + { + name: "uses namespace and name before the api assigns a uid", + service: &corev1.Service{ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "api", + }}, + want: "svclb-api-d53b356d3e1e", + }, + { + name: "trims boundary hyphens", + service: &corev1.Service{ObjectMeta: metav1.ObjectMeta{ + Name: "---", + }}, + want: "svclb-service-ff24f66688f3", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := serviceLBDaemonSetName(tt.service); got != tt.want { + t.Fatalf("serviceLBDaemonSetName() = %q, want %q", got, tt.want) + } + }) + } +} + +func TestServiceLBDaemonSetNameFitsDNSLabel(t *testing.T) { + service := &corev1.Service{ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: strings.Repeat("a", 43) + "-" + strings.Repeat("b", 19), + }} + + got := serviceLBDaemonSetName(service) + if len(got) > 63 { + t.Fatalf("serviceLBDaemonSetName() length = %d, want at most 63: %q", len(got), got) + } + if matched := regexp.MustCompile(`^[a-z0-9]([-a-z0-9]*[a-z0-9])?$`).MatchString(got); !matched { + t.Fatalf("serviceLBDaemonSetName() = %q, want a DNS label", got) + } +} + +func TestServiceLBDaemonSetNameSeparatesNamespacesWithoutUID(t *testing.T) { + first := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Namespace: "default", Name: "api"}} + second := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Namespace: "other", Name: "api"}} + + firstName := serviceLBDaemonSetName(first) + secondName := serviceLBDaemonSetName(second) + if firstName == secondName { + t.Fatalf("serviceLBDaemonSetName() reused %q across namespaces", firstName) + } + if secondName != "svclb-api-25dcc26e6d02" { + t.Fatalf("serviceLBDaemonSetName() = %q, want %q", secondName, "svclb-api-25dcc26e6d02") + } +} diff --git a/server/storage_bootstrap.go b/server/storage_bootstrap.go index df80673..aec9781 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 b75979a..ffa8d56 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 e7fe473..144c75e 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},