diff --git a/conf/app.conf b/conf/app.conf
index d146190..3812b56 100644
--- a/conf/app.conf
+++ b/conf/app.conf
@@ -27,3 +27,5 @@ helmImplicitLatestPullPolicy = false
apiserverPort = 6443
apiserverBind = 127.0.0.1
dataDir = /var/lib/casos
+ingressControllerEnabled = true
+serviceLBEnabled = true
diff --git a/controllers/ingress.go b/controllers/ingress.go
index 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/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/main.go b/main.go
index aa96034..f626c4e 100644
--- a/main.go
+++ b/main.go
@@ -67,6 +67,11 @@ func main() {
if err := server.Bootstrap(ctx, adminCfg, srvCfg); err != nil {
logs.Warning("bootstrap: %v", err)
}
+ if srvCfg.ServiceLBEnabled {
+ if err := server.StartServiceLB(ctx, adminCfg); err != nil {
+ logs.Warning("start service load balancer: %v", err)
+ }
+ }
if err := server.StartScheduler(ctx, srvCfg); err != nil {
logs.Warning("start scheduler: %v", err)
}
diff --git a/server/apiserver.go b/server/apiserver.go
index 0dbd0c8..1aa3af5 100644
--- a/server/apiserver.go
+++ b/server/apiserver.go
@@ -150,7 +150,7 @@ func buildApiserverArgs(cfg Config, certDir, etcdEndpoint, authzKubeconfig strin
"--service-cluster-ip-range=" + serviceClusterIPRange,
"--allow-privileged=true",
"--authorization-mode=" + authzMode(authzKubeconfig),
- "--enable-admission-plugins=NodeRestriction,ValidatingAdmissionWebhook",
+ "--enable-admission-plugins=NodeRestriction,DefaultIngressClass,ValidatingAdmissionWebhook",
"--tls-cert-file=" + filepath.Join(certDir, "apiserver.crt"),
"--tls-private-key-file=" + filepath.Join(certDir, "apiserver.key"),
"--client-ca-file=" + filepath.Join(certDir, "ca.crt"),
diff --git a/server/bootstrap.go b/server/bootstrap.go
index 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..2f77977 100644
--- a/server/config.go
+++ b/server/config.go
@@ -23,7 +23,10 @@ type Config struct {
LocalPathHelperImage string // helper pod image used by local-path-provisioner
FlannelImage string // Flannel daemon image used by the built-in network bootstrap
FlannelCNIPluginImage string // Flannel CNI plugin image installed on worker hosts
+ IngressControllerImage string // Traefik image used by the built-in Ingress controller
StorageProvisionerEnabled bool // install the built-in local-path provisioner for local clusters
+ IngressControllerEnabled bool // install the built-in Traefik controller
+ ServiceLBEnabled bool // run the built-in bare-metal LoadBalancer controller
}
// ConfigFromAppConf reads server config from the beego app.conf.
@@ -60,13 +63,16 @@ func ConfigFromAppConf() (Config, error) {
}
storageProvisionerEnabled := conf.GetConfigBoolDefault("storageProvisionerEnabled", true)
+ ingressControllerEnabled := conf.GetConfigBoolDefault("ingressControllerEnabled", false)
+ serviceLBEnabled := conf.GetConfigBoolDefault("serviceLBEnabled", false)
coreDNSImage := conf.GetConfigStringDefault("coreDNSImage", "docker.1ms.run/coredns/coredns:1.12.4")
localPathProvisionerImage := conf.GetConfigStringDefault("localPathProvisionerImage", "docker.1ms.run/rancher/local-path-provisioner:v0.0.32")
localPathHelperImage := conf.GetConfigStringDefault("localPathHelperImage", "docker.1ms.run/library/busybox:1.37.0")
flannelImage := conf.GetConfigStringDefault("flannelImage", defaultFlannelImage)
flannelCNIPluginImage := conf.GetConfigStringDefault("flannelCNIPluginImage", defaultFlannelCNIPluginImage)
+ ingressControllerImage := conf.GetConfigStringDefault("ingressControllerImage", defaultIngressControllerImage)
- return Config{
+ config := Config{
DataDir: dataDir,
ApiserverBind: bind,
AdvertiseAddress: advertise,
@@ -80,8 +86,22 @@ func ConfigFromAppConf() (Config, error) {
LocalPathHelperImage: localPathHelperImage,
FlannelImage: flannelImage,
FlannelCNIPluginImage: flannelCNIPluginImage,
+ IngressControllerImage: ingressControllerImage,
StorageProvisionerEnabled: storageProvisionerEnabled,
- }, nil
+ IngressControllerEnabled: ingressControllerEnabled,
+ ServiceLBEnabled: serviceLBEnabled,
+ }
+ if err := validateApplicationAccessConfig(config); err != nil {
+ return Config{}, err
+ }
+ return config, nil
+}
+
+func validateApplicationAccessConfig(config Config) error {
+ if config.IngressControllerEnabled && !config.ServiceLBEnabled {
+ return fmt.Errorf("ingressControllerEnabled requires serviceLBEnabled because the built-in Ingress Service uses LoadBalancerClass %q", serviceLBClass)
+ }
+ return nil
}
// injectDBName inserts dbName into a MySQL DSN of the form
diff --git a/server/ingress_bootstrap.go b/server/ingress_bootstrap.go
new file mode 100644
index 0000000..cbe9e99
--- /dev/null
+++ b/server/ingress_bootstrap.go
@@ -0,0 +1,432 @@
+package server
+
+import (
+ "context"
+ "errors"
+ "fmt"
+
+ appsv1 "k8s.io/api/apps/v1"
+ corev1 "k8s.io/api/core/v1"
+ networkingv1 "k8s.io/api/networking/v1"
+ rbacv1 "k8s.io/api/rbac/v1"
+ apiequality "k8s.io/apimachinery/pkg/api/equality"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ "k8s.io/apimachinery/pkg/api/resource"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ "k8s.io/apimachinery/pkg/util/intstr"
+ "k8s.io/client-go/kubernetes"
+)
+
+const (
+ ingressControllerNamespace = "kube-system"
+ ingressControllerName = "traefik"
+ ingressControllerClass = "traefik"
+ ingressControllerID = "traefik.io/ingress-controller"
+ defaultIngressControllerImage = "docker.io/traefik:v3.3.4"
+)
+
+func ingressControllerLabels() map[string]string {
+ return map[string]string{
+ "app.kubernetes.io/name": ingressControllerName,
+ "app.kubernetes.io/managed-by": "casos",
+ }
+}
+
+func validateIngressControllerOwnership(object metav1.Object) error {
+ if object.GetLabels()["app.kubernetes.io/managed-by"] != "casos" {
+ return fmt.Errorf("Ingress controller %T %s already exists and is not managed by CasOS", object, object.GetName())
+ }
+ return nil
+}
+
+func ensureIngressController(ctx context.Context, client kubernetes.Interface, cfg Config) error {
+ if err := ensureIngressControllerOwnership(ctx, client); err != nil {
+ return err
+ }
+ if err := ensureNamespace(ctx, client, ingressControllerNamespace); err != nil {
+ return err
+ }
+ if err := createOrUpdateServiceAccount(ctx, client, &corev1.ServiceAccount{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: ingressControllerName,
+ Namespace: ingressControllerNamespace,
+ Labels: ingressControllerLabels(),
+ },
+ }, validateIngressControllerOwnership); err != nil {
+ return err
+ }
+ if err := ensureIngressControllerClusterRole(ctx, client); err != nil {
+ return err
+ }
+ if err := ensureIngressControllerClusterRoleBinding(ctx, client); err != nil {
+ return err
+ }
+ if err := ensureIngressClass(ctx, client); err != nil {
+ return err
+ }
+ if err := ensureIngressControllerService(ctx, client); err != nil {
+ return err
+ }
+ return ensureIngressControllerDeployment(ctx, client, cfg)
+}
+
+func cleanupIngressController(ctx context.Context, client kubernetes.Interface) error {
+ resources := []struct {
+ kind string
+ get func() (metav1.Object, error)
+ delete func() error
+ }{
+ {kind: "Deployment", get: func() (metav1.Object, error) {
+ return client.AppsV1().Deployments(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }, delete: func() error {
+ return client.AppsV1().Deployments(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{})
+ }},
+ {kind: "Service", get: func() (metav1.Object, error) {
+ return client.CoreV1().Services(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }, delete: func() error {
+ return client.CoreV1().Services(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{})
+ }},
+ {kind: "IngressClass", get: func() (metav1.Object, error) {
+ return client.NetworkingV1().IngressClasses().Get(ctx, ingressControllerClass, metav1.GetOptions{})
+ }, delete: func() error {
+ return client.NetworkingV1().IngressClasses().Delete(ctx, ingressControllerClass, metav1.DeleteOptions{})
+ }},
+ {kind: "ClusterRoleBinding", get: func() (metav1.Object, error) {
+ return client.RbacV1().ClusterRoleBindings().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{})
+ }, delete: func() error {
+ return client.RbacV1().ClusterRoleBindings().Delete(ctx, "casos:traefik-ingress-controller", metav1.DeleteOptions{})
+ }},
+ {kind: "ClusterRole", get: func() (metav1.Object, error) {
+ return client.RbacV1().ClusterRoles().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{})
+ }, delete: func() error {
+ return client.RbacV1().ClusterRoles().Delete(ctx, "casos:traefik-ingress-controller", metav1.DeleteOptions{})
+ }},
+ {kind: "ServiceAccount", get: func() (metav1.Object, error) {
+ return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }, delete: func() error {
+ return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Delete(ctx, ingressControllerName, metav1.DeleteOptions{})
+ }},
+ }
+ errs := make([]error, 0)
+ for _, resource := range resources {
+ object, err := resource.get()
+ if apierrors.IsNotFound(err) {
+ continue
+ }
+ if err != nil {
+ errs = append(errs, fmt.Errorf("get Ingress controller %s for cleanup: %w", resource.kind, err))
+ continue
+ }
+ if object.GetLabels()["app.kubernetes.io/managed-by"] != "casos" {
+ continue
+ }
+ if err := resource.delete(); err != nil && !apierrors.IsNotFound(err) {
+ errs = append(errs, fmt.Errorf("delete Ingress controller %s %s: %w", resource.kind, object.GetName(), err))
+ }
+ }
+ return errors.Join(errs...)
+}
+
+func ensureIngressControllerOwnership(ctx context.Context, client kubernetes.Interface) error {
+ checks := []struct {
+ kind string
+ get func() (metav1.Object, error)
+ }{
+ {kind: "ServiceAccount", get: func() (metav1.Object, error) {
+ return client.CoreV1().ServiceAccounts(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }},
+ {kind: "Service", get: func() (metav1.Object, error) {
+ return client.CoreV1().Services(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }},
+ {kind: "Deployment", get: func() (metav1.Object, error) {
+ return client.AppsV1().Deployments(ingressControllerNamespace).Get(ctx, ingressControllerName, metav1.GetOptions{})
+ }},
+ {kind: "ClusterRole", get: func() (metav1.Object, error) {
+ return client.RbacV1().ClusterRoles().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{})
+ }},
+ {kind: "ClusterRoleBinding", get: func() (metav1.Object, error) {
+ return client.RbacV1().ClusterRoleBindings().Get(ctx, "casos:traefik-ingress-controller", metav1.GetOptions{})
+ }},
+ {kind: "IngressClass", get: func() (metav1.Object, error) {
+ return client.NetworkingV1().IngressClasses().Get(ctx, ingressControllerClass, metav1.GetOptions{})
+ }},
+ }
+ for _, check := range checks {
+ object, err := check.get()
+ if apierrors.IsNotFound(err) {
+ continue
+ }
+ if err != nil {
+ return fmt.Errorf("check Ingress controller %s ownership: %w", check.kind, err)
+ }
+ if err := validateIngressControllerOwnership(object); err != nil {
+ return fmt.Errorf("check Ingress controller %s ownership: %w", check.kind, err)
+ }
+ }
+ return nil
+}
+
+func ensureIngressControllerClusterRole(ctx context.Context, client kubernetes.Interface) error {
+ role := &rbacv1.ClusterRole{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: "casos:traefik-ingress-controller",
+ Labels: ingressControllerLabels(),
+ },
+ Rules: []rbacv1.PolicyRule{
+ {
+ APIGroups: []string{""},
+ Resources: []string{"services", "endpoints", "nodes", "secrets", "namespaces"},
+ Verbs: []string{"get", "list", "watch"},
+ },
+ {
+ APIGroups: []string{"discovery.k8s.io"},
+ Resources: []string{"endpointslices"},
+ Verbs: []string{"get", "list", "watch"},
+ },
+ {
+ APIGroups: []string{"networking.k8s.io"},
+ Resources: []string{"ingresses", "ingressclasses"},
+ Verbs: []string{"get", "list", "watch"},
+ },
+ {
+ APIGroups: []string{""},
+ Resources: []string{"events"},
+ Verbs: []string{"create", "update", "patch"},
+ },
+ {
+ APIGroups: []string{"networking.k8s.io"},
+ Resources: []string{"ingresses/status"},
+ Verbs: []string{"update", "patch"},
+ },
+ },
+ }
+ return createOrUpdateClusterRole(ctx, client, role, validateIngressControllerOwnership)
+}
+
+func ensureIngressControllerClusterRoleBinding(ctx context.Context, client kubernetes.Interface) error {
+ binding := &rbacv1.ClusterRoleBinding{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: "casos:traefik-ingress-controller",
+ Labels: ingressControllerLabels(),
+ },
+ RoleRef: rbacv1.RoleRef{
+ APIGroup: "rbac.authorization.k8s.io",
+ Kind: "ClusterRole",
+ Name: "casos:traefik-ingress-controller",
+ },
+ Subjects: []rbacv1.Subject{{
+ Kind: "ServiceAccount",
+ Name: ingressControllerName,
+ Namespace: ingressControllerNamespace,
+ }},
+ }
+ return createOrUpdateClusterRoleBinding(ctx, client, binding, validateIngressControllerOwnership)
+}
+
+func ensureIngressClass(ctx context.Context, client kubernetes.Interface) error {
+ classes := client.NetworkingV1().IngressClasses()
+ defaultClassExists := false
+ existingClasses, err := classes.List(ctx, metav1.ListOptions{})
+ if err != nil {
+ return fmt.Errorf("list IngressClasses: %w", err)
+ }
+ for i := range existingClasses.Items {
+ class := &existingClasses.Items[i]
+ if class.Name != ingressControllerClass && class.Annotations["ingressclass.kubernetes.io/is-default-class"] == "true" {
+ defaultClassExists = true
+ break
+ }
+ }
+ desired := &networkingv1.IngressClass{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: ingressControllerClass,
+ Labels: ingressControllerLabels(),
+ Annotations: map[string]string{},
+ },
+ Spec: networkingv1.IngressClassSpec{Controller: ingressControllerID},
+ }
+ if !defaultClassExists {
+ desired.Annotations["ingressclass.kubernetes.io/is-default-class"] = "true"
+ }
+ current, err := classes.Get(ctx, desired.Name, metav1.GetOptions{})
+ if apierrors.IsNotFound(err) {
+ if _, err := classes.Create(ctx, desired, metav1.CreateOptions{}); err != nil {
+ return fmt.Errorf("create IngressClass %s: %w", desired.Name, err)
+ }
+ return nil
+ }
+ if err != nil {
+ return fmt.Errorf("get IngressClass %s: %w", desired.Name, err)
+ }
+ if current.Spec.Controller != desired.Spec.Controller {
+ return fmt.Errorf("IngressClass %s is owned by controller %s", desired.Name, current.Spec.Controller)
+ }
+ if current.Labels["app.kubernetes.io/managed-by"] != "casos" {
+ return fmt.Errorf("IngressClass %s already exists and is not managed by CasOS", desired.Name)
+ }
+ desired.Labels = mergeStringMap(current.Labels, desired.Labels)
+ desired.Annotations = mergeStringMap(current.Annotations, desired.Annotations)
+ if defaultClassExists {
+ delete(desired.Annotations, "ingressclass.kubernetes.io/is-default-class")
+ }
+ if apiequality.Semantic.DeepEqual(current.Labels, desired.Labels) &&
+ apiequality.Semantic.DeepEqual(current.Annotations, desired.Annotations) &&
+ apiequality.Semantic.DeepEqual(current.Spec, desired.Spec) {
+ return nil
+ }
+ desired.ResourceVersion = current.ResourceVersion
+ if _, err := classes.Update(ctx, desired, metav1.UpdateOptions{}); err != nil {
+ return fmt.Errorf("update IngressClass %s: %w", desired.Name, err)
+ }
+ return nil
+}
+
+func ensureIngressControllerService(ctx context.Context, client kubernetes.Interface) error {
+ desired := buildIngressControllerService()
+ services := client.CoreV1().Services(ingressControllerNamespace)
+ current, err := services.Get(ctx, desired.Name, metav1.GetOptions{})
+ if apierrors.IsNotFound(err) {
+ if _, err := services.Create(ctx, desired, metav1.CreateOptions{}); err != nil {
+ return fmt.Errorf("create Ingress controller Service: %w", err)
+ }
+ return nil
+ }
+ if err != nil {
+ return fmt.Errorf("get Ingress controller Service: %w", err)
+ }
+ if current.Labels["app.kubernetes.io/managed-by"] != "casos" {
+ return fmt.Errorf("Ingress controller Service %s/%s already exists and is not managed by CasOS", ingressControllerNamespace, ingressControllerName)
+ }
+ desired.Labels = mergeStringMap(current.Labels, desired.Labels)
+ desired.Annotations = mergeStringMap(current.Annotations, desired.Annotations)
+ desired.Spec.ClusterIP = current.Spec.ClusterIP
+ desired.Spec.ClusterIPs = current.Spec.ClusterIPs
+ desired.Spec.IPFamilies = current.Spec.IPFamilies
+ desired.Spec.IPFamilyPolicy = current.Spec.IPFamilyPolicy
+ desired.Spec.ExternalIPs = current.Spec.ExternalIPs
+ desired.Spec.SessionAffinity = current.Spec.SessionAffinity
+ desired.Spec.InternalTrafficPolicy = current.Spec.InternalTrafficPolicy
+ desired.Spec.LoadBalancerIP = current.Spec.LoadBalancerIP
+ desired.Spec.LoadBalancerClass = current.Spec.LoadBalancerClass
+ desired.Spec.LoadBalancerSourceRanges = current.Spec.LoadBalancerSourceRanges
+ desired.Spec.AllocateLoadBalancerNodePorts = current.Spec.AllocateLoadBalancerNodePorts
+ desired.Spec.HealthCheckNodePort = current.Spec.HealthCheckNodePort
+ for i := range desired.Spec.Ports {
+ for _, currentPort := range current.Spec.Ports {
+ if desired.Spec.Ports[i].Name == currentPort.Name {
+ desired.Spec.Ports[i].NodePort = currentPort.NodePort
+ break
+ }
+ }
+ }
+ if apiequality.Semantic.DeepEqual(current.Labels, desired.Labels) &&
+ apiequality.Semantic.DeepEqual(current.Annotations, desired.Annotations) &&
+ apiequality.Semantic.DeepEqual(current.Spec, desired.Spec) {
+ return nil
+ }
+ desired.ResourceVersion = current.ResourceVersion
+ if _, err := services.Update(ctx, desired, metav1.UpdateOptions{}); err != nil {
+ return fmt.Errorf("update Ingress controller Service: %w", err)
+ }
+ return nil
+}
+
+func buildIngressControllerService() *corev1.Service {
+ return &corev1.Service{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: ingressControllerName,
+ Namespace: ingressControllerNamespace,
+ Labels: ingressControllerLabels(),
+ },
+ Spec: corev1.ServiceSpec{
+ Type: corev1.ServiceTypeLoadBalancer,
+ Selector: ingressControllerLabels(),
+ ExternalTrafficPolicy: corev1.ServiceExternalTrafficPolicyTypeCluster,
+ LoadBalancerClass: ptr(serviceLBClass),
+ Ports: []corev1.ServicePort{
+ {Name: "web", Port: 80, TargetPort: intstr.FromInt(8000), Protocol: corev1.ProtocolTCP},
+ {Name: "websecure", Port: 443, TargetPort: intstr.FromInt(8443), Protocol: corev1.ProtocolTCP},
+ },
+ },
+ }
+}
+
+func ensureIngressControllerDeployment(ctx context.Context, client kubernetes.Interface, cfg Config) error {
+ return createOrUpdateDeployment(ctx, client, buildIngressControllerDeployment(cfg), validateIngressControllerOwnership)
+}
+
+func buildIngressControllerDeployment(cfg Config) *appsv1.Deployment {
+ replicas := int32(1)
+ image := cfg.IngressControllerImage
+ if image == "" {
+ image = defaultIngressControllerImage
+ }
+ labels := ingressControllerLabels()
+ deployment := &appsv1.Deployment{
+ ObjectMeta: metav1.ObjectMeta{
+ Name: ingressControllerName,
+ Namespace: ingressControllerNamespace,
+ Labels: labels,
+ },
+ Spec: appsv1.DeploymentSpec{
+ Replicas: &replicas,
+ Selector: &metav1.LabelSelector{MatchLabels: labels},
+ Template: corev1.PodTemplateSpec{
+ ObjectMeta: metav1.ObjectMeta{Labels: labels},
+ Spec: corev1.PodSpec{
+ ServiceAccountName: ingressControllerName,
+ SecurityContext: &corev1.PodSecurityContext{
+ RunAsNonRoot: ptr(true),
+ RunAsUser: ptr(int64(65532)),
+ RunAsGroup: ptr(int64(65532)),
+ SeccompProfile: &corev1.SeccompProfile{Type: corev1.SeccompProfileTypeRuntimeDefault},
+ },
+ Tolerations: []corev1.Toleration{
+ {Key: "CriticalAddonsOnly", Operator: corev1.TolerationOpExists},
+ {Key: "node-role.kubernetes.io/control-plane", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule},
+ {Key: "node-role.kubernetes.io/master", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule},
+ {Key: "casos.io/bootstrap", Operator: corev1.TolerationOpExists, Effect: corev1.TaintEffectNoSchedule},
+ },
+ Containers: []corev1.Container{{
+ Name: ingressControllerName,
+ Image: image,
+ ImagePullPolicy: corev1.PullIfNotPresent,
+ Args: []string{
+ "--providers.kubernetesingress=true",
+ "--providers.kubernetesingress.ingressclass=" + ingressControllerClass,
+ "--providers.kubernetesingress.ingressendpoint.publishedservice=" + ingressControllerNamespace + "/" + ingressControllerName,
+ "--entrypoints.web.address=:8000",
+ "--entrypoints.websecure.address=:8443",
+ "--ping=true",
+ },
+ Ports: []corev1.ContainerPort{
+ {Name: "web", ContainerPort: 8000, Protocol: corev1.ProtocolTCP},
+ {Name: "websecure", ContainerPort: 8443, Protocol: corev1.ProtocolTCP},
+ {Name: "ping", ContainerPort: 8080, Protocol: corev1.ProtocolTCP},
+ },
+ ReadinessProbe: &corev1.Probe{
+ ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{
+ Path: "/ping", Port: intstr.FromInt(8080), Scheme: corev1.URISchemeHTTP,
+ }},
+ InitialDelaySeconds: 5, PeriodSeconds: 10, TimeoutSeconds: 5, FailureThreshold: 6,
+ },
+ LivenessProbe: &corev1.Probe{
+ ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{Path: "/ping", Port: intstr.FromInt(8080), Scheme: corev1.URISchemeHTTP}},
+ InitialDelaySeconds: 15, PeriodSeconds: 10, TimeoutSeconds: 5, FailureThreshold: 6,
+ },
+ SecurityContext: &corev1.SecurityContext{
+ AllowPrivilegeEscalation: ptr(false),
+ ReadOnlyRootFilesystem: ptr(true),
+ Capabilities: &corev1.Capabilities{Drop: []corev1.Capability{"ALL"}},
+ },
+ Resources: corev1.ResourceRequirements{
+ Requests: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("50m"), corev1.ResourceMemory: resource.MustParse("64Mi")},
+ Limits: corev1.ResourceList{corev1.ResourceCPU: resource.MustParse("500m"), corev1.ResourceMemory: resource.MustParse("256Mi")},
+ },
+ }},
+ },
+ },
+ },
+ }
+ return deployment
+}
diff --git a/server/servicelb.go b/server/servicelb.go
new file mode 100644
index 0000000..a12ff6a
--- /dev/null
+++ b/server/servicelb.go
@@ -0,0 +1,468 @@
+package server
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "fmt"
+ "net"
+ "os"
+ "reflect"
+ "time"
+
+ corev1 "k8s.io/api/core/v1"
+ discoveryv1 "k8s.io/api/discovery/v1"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ "k8s.io/apimachinery/pkg/labels"
+ "k8s.io/client-go/kubernetes"
+ coordinationclient "k8s.io/client-go/kubernetes/typed/coordination/v1"
+ "k8s.io/client-go/rest"
+ "k8s.io/client-go/tools/leaderelection"
+ "k8s.io/client-go/tools/leaderelection/resourcelock"
+ "k8s.io/client-go/util/retry"
+
+ "github.com/sirupsen/logrus"
+)
+
+const (
+ serviceLBManagedIPsAnnotation = "casos.io/service-lb-managed-ips"
+ serviceLBDisabledAnnotation = "casos.io/service-lb-disabled"
+ serviceLBEnabledAnnotation = "casos.io/service-lb-enabled"
+ serviceLBClass = "casos.io/service-lb"
+ serviceLBLeaderLease = "casos-service-lb"
+)
+
+// StartServiceLB starts the built-in bare-metal LoadBalancer reconciler. It
+// publishes Ready Worker addresses and uses Service externalIPs so kube-proxy
+// can route LoadBalancer traffic without a cloud provider.
+func StartServiceLB(ctx context.Context, cfg *rest.Config) error {
+ if cfg == nil {
+ return fmt.Errorf("apiserver rest config is required")
+ }
+ client, err := kubernetes.NewForConfig(cfg)
+ if err != nil {
+ return fmt.Errorf("service load balancer client: %w", err)
+ }
+ if ctx == nil {
+ ctx = context.Background()
+ }
+ coordination, err := coordinationclient.NewForConfig(cfg)
+ if err != nil {
+ return fmt.Errorf("service load balancer coordination client: %w", err)
+ }
+ hostname, err := os.Hostname()
+ if err != nil {
+ return fmt.Errorf("service load balancer identity: %w", err)
+ }
+ identity := fmt.Sprintf("%s-%d", hostname, os.Getpid())
+ elector, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{
+ Lock: &resourcelock.LeaseLock{
+ LeaseMeta: metav1.ObjectMeta{Name: serviceLBLeaderLease, Namespace: ingressControllerNamespace},
+ Client: coordination,
+ LockConfig: resourcelock.ResourceLockConfig{Identity: identity},
+ },
+ LeaseDuration: 15 * time.Second,
+ RenewDeadline: 10 * time.Second,
+ RetryPeriod: 2 * time.Second,
+ ReleaseOnCancel: true,
+ Callbacks: leaderelection.LeaderCallbacks{
+ OnStartedLeading: func(leaderCtx context.Context) { runServiceLB(leaderCtx, client) },
+ OnStoppedLeading: func() { logrus.Warn("service load balancer leadership lost") },
+ },
+ })
+ if err != nil {
+ return fmt.Errorf("service load balancer leader election: %w", err)
+ }
+ go elector.Run(ctx)
+ return nil
+}
+
+func runServiceLB(ctx context.Context, client kubernetes.Interface) {
+ const interval = 5 * time.Second
+ for {
+ if err := reconcileServiceLB(ctx, client); err != nil && ctx.Err() == nil {
+ logrus.Warnf("service load balancer reconciliation failed: %v", err)
+ }
+ timer := time.NewTimer(interval)
+ select {
+ case <-ctx.Done():
+ timer.Stop()
+ return
+ case <-timer.C:
+ }
+ }
+}
+
+func reconcileServiceLB(ctx context.Context, client kubernetes.Interface) error {
+ nodes, err := client.CoreV1().Nodes().List(ctx, metav1.ListOptions{})
+ if err != nil {
+ return fmt.Errorf("list nodes: %w", err)
+ }
+ services, err := client.CoreV1().Services(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
+ if err != nil {
+ return fmt.Errorf("list LoadBalancer services: %w", err)
+ }
+ reconcileErrors := make([]error, 0)
+ for i := range services.Items {
+ service := &services.Items[i]
+ if service.Spec.Type != corev1.ServiceTypeLoadBalancer {
+ continue
+ }
+ if !serviceLBManages(service) {
+ if _, managed := service.Annotations[serviceLBManagedIPsAnnotation]; managed {
+ if err := cleanupLoadBalancerService(ctx, client, service); err != nil {
+ reconcileErrors = append(reconcileErrors, fmt.Errorf("clean up LoadBalancer service %s/%s: %w", service.Namespace, service.Name, err))
+ }
+ }
+ continue
+ }
+ nodeIPs, err := serviceLBNodeIPs(ctx, client, service, nodes.Items)
+ if err != nil {
+ reconcileErrors = append(reconcileErrors, fmt.Errorf("select LoadBalancer nodes for %s/%s: %w", service.Namespace, service.Name, err))
+ continue
+ }
+ if err := reconcileLoadBalancerService(ctx, client, service, nodeIPs); err != nil {
+ reconcileErrors = append(reconcileErrors, fmt.Errorf("reconcile LoadBalancer service %s/%s: %w", service.Namespace, service.Name, err))
+ }
+ }
+ return errors.Join(reconcileErrors...)
+}
+
+func cleanupServiceLB(ctx context.Context, client kubernetes.Interface) error {
+ services, err := client.CoreV1().Services(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
+ if err != nil {
+ return fmt.Errorf("list services for ServiceLB cleanup: %w", err)
+ }
+ errs := make([]error, 0)
+ for i := range services.Items {
+ service := &services.Items[i]
+ if service.Annotations[serviceLBManagedIPsAnnotation] == "" {
+ continue
+ }
+ if err := cleanupLoadBalancerService(ctx, client, service); err != nil {
+ errs = append(errs, fmt.Errorf("clean up ServiceLB state for %s/%s: %w", service.Namespace, service.Name, err))
+ }
+ }
+ return errors.Join(errs...)
+}
+
+func serviceLBManages(service *corev1.Service) bool {
+ if service == nil || service.Spec.Type != corev1.ServiceTypeLoadBalancer || service.Annotations[serviceLBDisabledAnnotation] == "true" {
+ return false
+ }
+ return (service.Spec.LoadBalancerClass != nil && *service.Spec.LoadBalancerClass == serviceLBClass) || service.Annotations[serviceLBEnabledAnnotation] == "true"
+}
+
+type serviceLBNode struct {
+ name string
+ ips []string
+}
+
+func readyServiceLBNodes(nodes []corev1.Node) []serviceLBNode {
+ return readyNodeAddresses(nodes, false)
+}
+
+func readyNodeAddresses(nodes []corev1.Node, controlPlaneOnly bool) []serviceLBNode {
+ seen := map[string]struct{}{}
+ result := make([]serviceLBNode, 0)
+ for _, node := range nodes {
+ if !isReadyNode(node) || (controlPlaneOnly != isControlPlaneNode(node)) {
+ continue
+ }
+ externalAddresses := make([]string, 0)
+ internalAddresses := make([]string, 0)
+ for _, address := range node.Status.Addresses {
+ if address.Type != corev1.NodeExternalIP && address.Type != corev1.NodeInternalIP {
+ continue
+ }
+ ip := net.ParseIP(address.Address)
+ if ip == nil {
+ continue
+ }
+ value := ip.String()
+ if _, ok := seen[value]; ok {
+ continue
+ }
+ if address.Type == corev1.NodeExternalIP {
+ externalAddresses = append(externalAddresses, value)
+ } else {
+ internalAddresses = append(internalAddresses, value)
+ }
+ }
+ addresses := externalAddresses
+ if len(addresses) == 0 {
+ addresses = internalAddresses
+ }
+ if len(addresses) > 0 {
+ for _, address := range addresses {
+ seen[address] = struct{}{}
+ }
+ result = append(result, serviceLBNode{name: node.Name, ips: addresses})
+ }
+ }
+ return result
+}
+
+func serviceLBNodeIPs(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodes []corev1.Node) ([]string, error) {
+ candidates := readyServiceLBNodes(nodes)
+ if service.Spec.ExternalTrafficPolicy != corev1.ServiceExternalTrafficPolicyTypeLocal {
+ ready, err := serviceHasReadyEndpoint(ctx, client, service)
+ if err != nil {
+ return nil, err
+ }
+ if !ready {
+ return []string{}, nil
+ }
+ ips := make([]string, 0)
+ for _, node := range candidates {
+ ips = append(ips, node.ips...)
+ }
+ return uniqueStrings(ips), nil
+ }
+ selector := labels.SelectorFromSet(labels.Set{discoveryv1.LabelServiceName: service.Name}).String()
+ endpointSlices, err := client.DiscoveryV1().EndpointSlices(service.Namespace).List(ctx, metav1.ListOptions{
+ LabelSelector: selector,
+ })
+ if err != nil {
+ return nil, fmt.Errorf("list EndpointSlices: %w", err)
+ }
+ localNodes := map[string]struct{}{}
+ for _, endpointSlice := range endpointSlices.Items {
+ for _, endpoint := range endpointSlice.Endpoints {
+ if endpoint.NodeName == nil || (endpoint.Conditions.Ready != nil && !*endpoint.Conditions.Ready) {
+ continue
+ }
+ localNodes[*endpoint.NodeName] = struct{}{}
+ }
+ }
+ matchingIPs := func(candidates []serviceLBNode) []string {
+ ips := make([]string, 0)
+ for _, node := range candidates {
+ if _, ok := localNodes[node.name]; !ok {
+ continue
+ }
+ ips = append(ips, node.ips...)
+ }
+ return uniqueStrings(ips)
+ }
+ if workerIPs := matchingIPs(readyNodeAddresses(nodes, false)); len(workerIPs) > 0 {
+ return workerIPs, nil
+ }
+ return nil, nil
+}
+
+func serviceHasReadyEndpoint(ctx context.Context, client kubernetes.Interface, service *corev1.Service) (bool, error) {
+ selector := labels.SelectorFromSet(labels.Set{discoveryv1.LabelServiceName: service.Name}).String()
+ endpointSlices, err := client.DiscoveryV1().EndpointSlices(service.Namespace).List(ctx, metav1.ListOptions{
+ LabelSelector: selector,
+ })
+ if err != nil {
+ return false, fmt.Errorf("list EndpointSlices: %w", err)
+ }
+ for _, endpointSlice := range endpointSlices.Items {
+ for _, endpoint := range endpointSlice.Endpoints {
+ if endpoint.Conditions.Ready == nil || *endpoint.Conditions.Ready {
+ return true, nil
+ }
+ }
+ }
+ return false, nil
+}
+
+func isReadyNode(node corev1.Node) bool {
+ if node.Spec.Unschedulable {
+ return false
+ }
+ for _, condition := range node.Status.Conditions {
+ if condition.Type == corev1.NodeReady {
+ return condition.Status == corev1.ConditionTrue
+ }
+ }
+ return false
+}
+
+func cleanupLoadBalancerService(ctx context.Context, client kubernetes.Interface, service *corev1.Service) error {
+ if service == nil || service.DeletionTimestamp != nil {
+ return nil
+ }
+ return retry.RetryOnConflict(retry.DefaultRetry, func() error {
+ current, err := client.CoreV1().Services(service.Namespace).Get(ctx, service.Name, metav1.GetOptions{})
+ if apierrors.IsNotFound(err) {
+ return nil
+ }
+ if err != nil {
+ return err
+ }
+ if current.DeletionTimestamp != nil {
+ return nil
+ }
+ managedIPs, err := serviceLBManagedIPs(current.Annotations)
+ if err != nil {
+ return err
+ }
+ managedSet := make(map[string]struct{}, len(managedIPs))
+ for _, ip := range managedIPs {
+ managedSet[ip] = struct{}{}
+ }
+ status := current.Status.DeepCopy()
+ status.LoadBalancer.Ingress = status.LoadBalancer.Ingress[:0]
+ for _, ingress := range current.Status.LoadBalancer.Ingress {
+ if _, managed := managedSet[ingress.IP]; !managed {
+ status.LoadBalancer.Ingress = append(status.LoadBalancer.Ingress, ingress)
+ }
+ }
+ if !reflect.DeepEqual(current.Status, *status) {
+ statusUpdate := current.DeepCopy()
+ statusUpdate.Status = *status
+ current, err = client.CoreV1().Services(current.Namespace).UpdateStatus(ctx, statusUpdate, metav1.UpdateOptions{})
+ if err != nil {
+ return err
+ }
+ }
+ updated := current.DeepCopy()
+ updated.Spec.ExternalIPs = updated.Spec.ExternalIPs[:0]
+ for _, ip := range current.Spec.ExternalIPs {
+ if _, managed := managedSet[ip]; !managed {
+ updated.Spec.ExternalIPs = append(updated.Spec.ExternalIPs, ip)
+ }
+ }
+ delete(updated.Annotations, serviceLBManagedIPsAnnotation)
+ if reflect.DeepEqual(current.Spec.ExternalIPs, updated.Spec.ExternalIPs) && current.Annotations[serviceLBManagedIPsAnnotation] == "" {
+ return nil
+ }
+ _, err = client.CoreV1().Services(current.Namespace).Update(ctx, updated, metav1.UpdateOptions{})
+ return err
+ })
+}
+
+func isControlPlaneNode(node corev1.Node) bool {
+ _, controlPlane := node.Labels["node-role.kubernetes.io/control-plane"]
+ _, master := node.Labels["node-role.kubernetes.io/master"]
+ return controlPlane || master
+}
+
+func reconcileLoadBalancerService(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodeIPs []string) error {
+ if service == nil || service.DeletionTimestamp != nil {
+ return nil
+ }
+ return retry.RetryOnConflict(retry.DefaultRetry, func() error {
+ current, err := client.CoreV1().Services(service.Namespace).Get(ctx, service.Name, metav1.GetOptions{})
+ if apierrors.IsNotFound(err) {
+ return nil
+ }
+ if err != nil {
+ return err
+ }
+ if current.DeletionTimestamp != nil || !serviceLBManages(current) {
+ return nil
+ }
+ return reconcileLoadBalancerServiceOnce(ctx, client, current, nodeIPs)
+ })
+}
+
+func reconcileLoadBalancerServiceOnce(ctx context.Context, client kubernetes.Interface, service *corev1.Service, nodeIPs []string) error {
+ managedIPs, err := serviceLBManagedIPs(service.Annotations)
+ if err != nil {
+ return err
+ }
+ managedSet := make(map[string]struct{}, len(managedIPs))
+ for _, ip := range managedIPs {
+ managedSet[ip] = struct{}{}
+ }
+ userIPs := make([]string, 0, len(service.Spec.ExternalIPs))
+ for _, ip := range service.Spec.ExternalIPs {
+ if _, ok := managedSet[ip]; !ok && !containsString(userIPs, ip) {
+ userIPs = append(userIPs, ip)
+ }
+ }
+ desiredManagedIPs := make([]string, 0, len(nodeIPs))
+ for _, ip := range nodeIPs {
+ if !containsString(userIPs, ip) {
+ desiredManagedIPs = append(desiredManagedIPs, ip)
+ }
+ }
+ desiredExternalIPs := append(append([]string{}, userIPs...), nodeIPs...)
+ desiredExternalIPs = uniqueStrings(desiredExternalIPs)
+ encodedManagedIPs, err := json.Marshal(desiredManagedIPs)
+ if err != nil {
+ return fmt.Errorf("encode managed LoadBalancer IPs: %w", err)
+ }
+ desiredAnnotation := string(encodedManagedIPs)
+
+ updated := service.DeepCopy()
+ if !reflect.DeepEqual(updated.Spec.ExternalIPs, desiredExternalIPs) || updated.Annotations[serviceLBManagedIPsAnnotation] != desiredAnnotation {
+ if updated.Annotations == nil {
+ updated.Annotations = map[string]string{}
+ }
+ updated.Spec.ExternalIPs = desiredExternalIPs
+ updated.Annotations[serviceLBManagedIPsAnnotation] = desiredAnnotation
+ current, err := client.CoreV1().Services(service.Namespace).Update(ctx, updated, metav1.UpdateOptions{})
+ if err != nil {
+ return fmt.Errorf("update LoadBalancer service spec: %w", err)
+ }
+ updated = current
+ }
+
+ desiredStatus := updated.Status.DeepCopy()
+ unmanagedIngress := make([]corev1.LoadBalancerIngress, 0, len(desiredStatus.LoadBalancer.Ingress))
+ for _, ingress := range desiredStatus.LoadBalancer.Ingress {
+ if _, managed := managedSet[ingress.IP]; !managed {
+ unmanagedIngress = append(unmanagedIngress, ingress)
+ }
+ }
+ desiredStatus.LoadBalancer.Ingress = unmanagedIngress
+ for _, ip := range nodeIPs {
+ if !containsLoadBalancerIngress(desiredStatus.LoadBalancer.Ingress, ip) {
+ desiredStatus.LoadBalancer.Ingress = append(desiredStatus.LoadBalancer.Ingress, corev1.LoadBalancerIngress{IP: ip})
+ }
+ }
+ if reflect.DeepEqual(updated.Status, *desiredStatus) {
+ return nil
+ }
+ updated.Status = *desiredStatus
+ if _, err := client.CoreV1().Services(service.Namespace).UpdateStatus(ctx, updated, metav1.UpdateOptions{}); err != nil && !apierrors.IsNotFound(err) {
+ return fmt.Errorf("update LoadBalancer service status: %w", err)
+ }
+ return nil
+}
+
+func containsLoadBalancerIngress(ingresses []corev1.LoadBalancerIngress, ip string) bool {
+ for _, ingress := range ingresses {
+ if ingress.IP == ip {
+ return true
+ }
+ }
+ return false
+}
+
+func serviceLBManagedIPs(annotations map[string]string) ([]string, error) {
+ value := annotations[serviceLBManagedIPsAnnotation]
+ if value == "" {
+ return nil, nil
+ }
+ var ips []string
+ if err := json.Unmarshal([]byte(value), &ips); err != nil {
+ return nil, fmt.Errorf("decode managed LoadBalancer IPs: %w", err)
+ }
+ return uniqueStrings(ips), nil
+}
+
+func containsString(values []string, want string) bool {
+ for _, value := range values {
+ if value == want {
+ return true
+ }
+ }
+ return false
+}
+
+func uniqueStrings(values []string) []string {
+ result := make([]string, 0, len(values))
+ for _, value := range values {
+ if value != "" && !containsString(result, value) {
+ result = append(result, value)
+ }
+ }
+ return result
+}
diff --git a/server/storage_bootstrap.go b/server/storage_bootstrap.go
index 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},