diff --git a/cmd/sealos/cmd/add.go b/cmd/sealos/cmd/add.go index 48720c7c9..42213ca84 100644 --- a/cmd/sealos/cmd/add.go +++ b/cmd/sealos/cmd/add.go @@ -33,12 +33,12 @@ add to default cluster: sealer add --masters x.x.x.x-x.x.x.y --nodes x.x.x.x-x.x.x.y `, RunE: func(cmd *cobra.Command, args []string) error { - return errors.New("add feature no support") - //applier, err := apply.NewScaleApplierFromArgs(addArgs, "add") - //if err != nil { - // return err - //} - //return applier.Apply() + //return errors.New("add feature no support") + applier, err := apply.NewScaleApplierFromArgs(addArgs, "add") + if err != nil { + return err + } + return applier.Apply() }, PreRunE: func(cmd *cobra.Command, args []string) error { if addArgs.Nodes == "" && addArgs.Masters == "" { diff --git a/pkg/apply/processor/scale.go b/pkg/apply/processor/scale.go index 74bea2b6e..d71f4dd1c 100644 --- a/pkg/apply/processor/scale.go +++ b/pkg/apply/processor/scale.go @@ -94,7 +94,11 @@ func (c *ScaleProcessor) Join(cluster *v2.Cluster) error { if err != nil { return err } - return c.Runtime.JoinNodes(c.NodesToJoin) + err = c.Runtime.JoinNodes(c.NodesToJoin) + if err != nil { + return err + } + return c.Runtime.SyncNodeIPVS(cluster.GetMasterIPAndPortList(), cluster.GetNodeIPAndPortList()) } func (c ScaleProcessor) UnMountRootfs(cluster *v2.Cluster) error { @@ -152,13 +156,13 @@ func (c *ScaleProcessor) RunConfig(cluster *v2.Cluster) error { } func (c *ScaleProcessor) MountRootfs(cluster *v2.Cluster) error { - hosts := append(cluster.GetMasterIPAndPortList(), cluster.GetNodeIPAndPortList()...) + hosts := append(c.MastersToJoin, c.NodesToJoin...) fs, err := filesystem.NewRootfsMounter(c.cManifestList, c.imageList) if err != nil { return err } - return fs.MountRootfs(cluster, hosts, false) + return fs.MountRootfs(cluster, hosts, true) } func NewScaleProcessor(clusterFile clusterfile.Interface, images v2.ImageList, masterToJoin, masterToDelete, nodeToJoin, nodeToDelete []string) (Interface, error) { diff --git a/pkg/apply/scale.go b/pkg/apply/scale.go index c1f779691..7f61002d3 100644 --- a/pkg/apply/scale.go +++ b/pkg/apply/scale.go @@ -92,35 +92,68 @@ func joinNodes(cluster *v2.Cluster, scaleArgs *RunArgs) error { if (!IsIPList(scaleArgs.Nodes) && scaleArgs.Nodes != "") || (!IsIPList(scaleArgs.Masters) && scaleArgs.Masters != "") { return fmt.Errorf(" Parameter error: The current mode should submit iplist!") } + var hosts []v2.Host + var hasMaster bool + for i := 0; i < len(cluster.Spec.Hosts); i++ { + role := cluster.Spec.Hosts[i].Roles + if strings2.InList(v2.MASTER, role) { + hasMaster = true + res := iputils.GetHostIPAndPortSlice(cluster.Spec.Hosts[i].IPS, strconv.Itoa(int(cluster.Spec.SSH.Port))) + hosts = append(hosts, v2.Host{ + IPS: res, + Roles: role, + Env: cluster.Spec.Hosts[i].Env, + }) + } + } + if !hasMaster { + return fmt.Errorf("not found `master` role from file") + } + var ipAndPorts []string + waitAddMasters := strings.Split(scaleArgs.Masters, ",") + for _, j := range waitAddMasters { + if j == "" { + continue + } + _ip, port := iputils.GetHostIPAndPortOrDefault(j, strconv.Itoa(int(cluster.Spec.SSH.Port))) + ipAndPorts = append(ipAndPorts, fmt.Sprintf("%s:%s", _ip, port)) + } + if len(ipAndPorts) > 0 { + hosts = append(hosts, v2.Host{ + IPS: ipAndPorts, + Roles: []string{v2.MASTER, string(v2.AMD64)}, + }) + } - if scaleArgs.Masters != "" && IsIPList(scaleArgs.Masters) { - for i := 0; i < len(cluster.Spec.Hosts); i++ { - role := cluster.Spec.Hosts[i].Roles - if strings2.InList(v2.MASTER, role) { - ipset := iputils.GetHostIPAndPortSlice(strings.Split(scaleArgs.Masters, ","), strconv.Itoa(int(cluster.Spec.SSH.Port))) - cluster.Spec.Hosts[i].IPS = removeIPListDuplicatesAndEmpty(append(cluster.Spec.Hosts[i].IPS, ipset...)) - break - } - if i == len(cluster.Spec.Hosts)-1 { - return fmt.Errorf("not found `master` role from file") - } - } - } //add join node - if scaleArgs.Nodes != "" && IsIPList(scaleArgs.Nodes) { - for i := 0; i < len(cluster.Spec.Hosts); i++ { - role := cluster.Spec.Hosts[i].Roles - if strings2.InList(v2.NODE, role) { - ipset := iputils.GetHostIPAndPortSlice(strings.Split(scaleArgs.Nodes, ","), strconv.Itoa(int(cluster.Spec.SSH.Port))) - cluster.Spec.Hosts[i].IPS = removeIPListDuplicatesAndEmpty(append(cluster.Spec.Hosts[i].IPS, ipset...)) - break - } - if i == len(cluster.Spec.Hosts)-1 { - hosts := v2.Host{IPS: removeIPListDuplicatesAndEmpty(strings.Split(scaleArgs.Nodes, ",")), Roles: []string{v2.NODE, string(v2.AMD64)}} - cluster.Spec.Hosts = append(cluster.Spec.Hosts, hosts) - } + for i := 0; i < len(cluster.Spec.Hosts); i++ { + role := cluster.Spec.Hosts[i].Roles + if strings2.InList(v2.Node, role) { + res := iputils.GetHostIPAndPortSlice(cluster.Spec.Hosts[i].IPS, strconv.Itoa(int(cluster.Spec.SSH.Port))) + hosts = append(hosts, v2.Host{ + IPS: res, + Roles: role, + Env: cluster.Spec.Hosts[i].Env, + }) } } + ipAndPorts = []string{} + waitAddNodes := strings.Split(scaleArgs.Nodes, ",") + for _, j := range waitAddNodes { + if j == "" { + continue + } + _ip, port := iputils.GetHostIPAndPortOrDefault(j, strconv.Itoa(int(cluster.Spec.SSH.Port))) + ipAndPorts = append(ipAndPorts, fmt.Sprintf("%s:%s", _ip, port)) + } + if len(ipAndPorts) > 0 { + hosts = append(hosts, v2.Host{ + IPS: ipAndPorts, + Roles: []string{v2.Node, string(v2.AMD64)}, + }) + } + logger.Debug("des nodes: %v", hosts) + cluster.Spec.Hosts = hosts return nil } @@ -139,6 +172,7 @@ func deleteNodes(cluster *v2.Cluster, scaleArgs *RunArgs) error { if strings2.InList(cluster.GetMaster0IPAndPort(), strings.Split(scaleArgs.Masters, ",")) { return fmt.Errorf("master0 machine cannot be deleted") } + defaultPort := strconv.Itoa(int(cluster.Spec.SSH.Port)) if scaleArgs.Masters != "" && IsIPList(scaleArgs.Masters) { for i := range cluster.Spec.Hosts { @@ -154,6 +188,13 @@ func deleteNodes(cluster *v2.Cluster, scaleArgs *RunArgs) error { } } } + var hosts []v2.Host + for _, host := range cluster.Spec.Hosts { + if len(host.IPS) != 0 { + hosts = append(hosts, host) + } + } + cluster.Spec.Hosts = hosts return nil } diff --git a/pkg/checker/host_checker.go b/pkg/checker/host_checker.go index 3e895f287..87b11c10a 100644 --- a/pkg/checker/host_checker.go +++ b/pkg/checker/host_checker.go @@ -33,7 +33,6 @@ func (a HostChecker) Check(cluster *v2.Cluster, phase string) error { for _, hosts := range cluster.Spec.Hosts { ipList = append(ipList, hosts.IPS...) } - if err := checkHostnameUnique(cluster, ipList); err != nil { return err } diff --git a/pkg/runtime/kubeadm.go b/pkg/runtime/kubeadm.go index 588ff392f..d5570f128 100644 --- a/pkg/runtime/kubeadm.go +++ b/pkg/runtime/kubeadm.go @@ -206,10 +206,11 @@ func (k *KubeadmRuntime) setKubernetesToken() error { if err != nil { return err } + } else { + k.Token = &t } } } - k.setJoinToken(k.Token.JoinToken) k.setTokenCaCertHash(k.Token.DiscoveryTokenCaCertHash) k.setCertificateKey(k.Token.CertificateKey) @@ -246,6 +247,9 @@ func (k *KubeadmRuntime) getTokenCaCertHash() []string { func (k *KubeadmRuntime) setCertificateKey(certificateKey string) { k.InitConfiguration.CertificateKey = certificateKey + if k.JoinConfiguration.ControlPlane == nil { + k.JoinConfiguration.ControlPlane = &kubeadm.JoinControlPlane{} + } k.JoinConfiguration.ControlPlane.CertificateKey = certificateKey } @@ -270,6 +274,14 @@ func (k *KubeadmRuntime) setJoinAdvertiseAddress(advertiseAddress string) { } k.JoinConfiguration.ControlPlane.LocalAPIEndpoint.AdvertiseAddress = advertiseAddress } +func (k *KubeadmRuntime) setDefaultEtcdData(etcdData string) { + if k.ClusterConfiguration.Etcd.Local == nil { + k.ClusterConfiguration.Etcd.Local = &kubeadm.LocalEtcd{} + } + if k.ClusterConfiguration.Etcd.Local.DataDir == "" { + k.ClusterConfiguration.Etcd.Local.DataDir = etcdData + } +} func (k *KubeadmRuntime) cleanJoinLocalAPIEndPoint() { k.JoinConfiguration.ControlPlane = nil @@ -334,8 +346,7 @@ func (k *KubeadmRuntime) generateInitConfigs() ([]byte, error) { k.APIServer.ExtraArgs = make(map[string]string) } k.APIServer.ExtraArgs["etcd-servers"] = getEtcdEndpointsWithHTTPSPrefix(k.getMasterIPList()) - k.ClusterConfiguration.Etcd.Local = &kubeadm.LocalEtcd{DataDir: "/var/lib/etcd"} - + k.setDefaultEtcdData("/var/lib/etcd") k.IPVS.ExcludeCIDRs = append(k.KubeProxyConfiguration.IPVS.ExcludeCIDRs, fmt.Sprintf("%s/32", k.getVip())) k.IPVS.ExcludeCIDRs = strings2.RemoveDuplicate(k.IPVS.ExcludeCIDRs)