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},