feature(main): fix add feature and delete node feature (#962)

* feature(main): fix add feature and delete node feature

* feature(main): fix add feature and delete node feature
This commit is contained in:
cuisongliu
2022-04-29 10:26:11 +08:00
committed by GitHub
parent 89b5cea8ba
commit 15d1932f91
5 changed files with 93 additions and 38 deletions
+6 -6
View File
@@ -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 == "" {
+7 -3
View File
@@ -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) {
+66 -25
View File
@@ -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
}
-1
View File
@@ -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
}
+14 -3
View File
@@ -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)