From c192e2ddbc126adad6fd01c9d4c0b509349e77ff Mon Sep 17 00:00:00 2001 From: cuisongliu Date: Thu, 28 Apr 2022 16:37:42 +0800 Subject: [PATCH] feature(main): no change images (#956) * feature(main): no change images * feature(main): delete node for sealos delete cmd * feature(main): auth for registry * feature(main): add node lic * feature(main): fix app install for add * feature(main): fix push install feature * feature(main): fix push install feature --- cmd/sealos/cmd/add.go | 61 ++++++ cmd/sealos/cmd/delete.go | 62 ++++++ .../applydrivers/apply_drivers_default.go | 63 +++++- pkg/apply/args.go | 14 ++ pkg/apply/processor/create.go | 4 +- pkg/apply/processor/delete.go | 4 +- pkg/apply/processor/install.go | 6 +- pkg/apply/processor/scale.go | 192 ++++++++++++++++++ pkg/apply/run.go | 18 +- pkg/apply/scale.go | 173 ++++++++++++++++ pkg/apply/utils.go | 16 ++ pkg/image/binary/cluster.go | 23 ++- pkg/registry/util.go | 31 +-- pkg/utils/iputils/iputils_v2.go | 23 ++- 14 files changed, 643 insertions(+), 47 deletions(-) create mode 100644 cmd/sealos/cmd/add.go create mode 100644 cmd/sealos/cmd/delete.go create mode 100644 pkg/apply/processor/scale.go create mode 100644 pkg/apply/scale.go diff --git a/cmd/sealos/cmd/add.go b/cmd/sealos/cmd/add.go new file mode 100644 index 000000000..48720c7c9 --- /dev/null +++ b/cmd/sealos/cmd/add.go @@ -0,0 +1,61 @@ +// Copyright © 2022 cuisongliu@qq.com. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package cmd + +import ( + "errors" + + "github.com/fanux/sealos/pkg/apply" + "github.com/spf13/cobra" +) + +// addCmd represents the delete command +func newAddCmd() *cobra.Command { + var deleteCmd = &cobra.Command{ + Use: "add", + Short: "add some node", + Args: cobra.NoArgs, + Example: ` +add to default cluster: + sealer add --masters x.x.x.x --nodes x.x.x.x + 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() + }, + PreRunE: func(cmd *cobra.Command, args []string) error { + if addArgs.Nodes == "" && addArgs.Masters == "" { + return errors.New("node and master not empty in same time") + } + return nil + }, + } + addArgs = &apply.ScaleArgs{} + deleteCmd.Flags().StringVarP(&addArgs.Masters, "masters", "m", "", "reduce Count or IPList to masters") + deleteCmd.Flags().StringVarP(&addArgs.Nodes, "nodes", "n", "", "reduce Count or IPList to nodes") + deleteCmd.Flags().StringVarP(&addArgs.ClusterName, "cluster", "c", "default", "delete a kubernetes cluster with cluster name") + return deleteCmd +} + +var addArgs *apply.ScaleArgs + +func init() { + rootCmd.AddCommand(newAddCmd()) +} diff --git a/cmd/sealos/cmd/delete.go b/cmd/sealos/cmd/delete.go new file mode 100644 index 000000000..eec2b476d --- /dev/null +++ b/cmd/sealos/cmd/delete.go @@ -0,0 +1,62 @@ +// Copyright © 2021 Alibaba Group Holding Ltd. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package cmd + +import ( + "errors" + + "github.com/fanux/sealos/pkg/apply" + "github.com/fanux/sealos/pkg/runtime" + "github.com/spf13/cobra" +) + +// deleteCmd represents the delete command +func newDeleteCmd() *cobra.Command { + var deleteCmd = &cobra.Command{ + Use: "delete", + Short: "delete some node", + Args: cobra.NoArgs, + Example: ` +delete to default cluster: + sealer delete --masters x.x.x.x --nodes x.x.x.x + sealer delete --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 { + applier, err := apply.NewScaleApplierFromArgs(deleteArgs, "delete") + if err != nil { + return err + } + return applier.Apply() + }, + PreRunE: func(cmd *cobra.Command, args []string) error { + if deleteArgs.Nodes == "" && deleteArgs.Masters == "" { + return errors.New("node and master not empty in same time") + } + return nil + }, + } + deleteArgs = &apply.ScaleArgs{} + deleteCmd.Flags().StringVarP(&deleteArgs.Masters, "masters", "m", "", "reduce Count or IPList to masters") + deleteCmd.Flags().StringVarP(&deleteArgs.Nodes, "nodes", "n", "", "reduce Count or IPList to nodes") + deleteCmd.Flags().StringVarP(&deleteArgs.ClusterName, "cluster", "c", "default", "delete a kubernetes cluster with cluster name") + deleteCmd.Flags().BoolVar(&runtime.ForceDelete, "force", false, "We also can input an --force flag to delete cluster by force") + return deleteCmd +} + +var deleteArgs *apply.ScaleArgs + +func init() { + rootCmd.AddCommand(newDeleteCmd()) +} diff --git a/pkg/apply/applydrivers/apply_drivers_default.go b/pkg/apply/applydrivers/apply_drivers_default.go index f47193bfb..2ba13d2c7 100644 --- a/pkg/apply/applydrivers/apply_drivers_default.go +++ b/pkg/apply/applydrivers/apply_drivers_default.go @@ -17,6 +17,7 @@ package applydrivers import ( "fmt" + "github.com/fanux/sealos/pkg/utils/iputils" "github.com/fanux/sealos/pkg/utils/strings" "github.com/fanux/sealos/pkg/apply/processor" @@ -39,6 +40,19 @@ func NewDefaultApplier(cluster *v2.Cluster) (Interface, error) { return &Applier{ ClusterDesired: cluster, ClusterFile: cFile, + ClusterCurrent: cFile.GetCluster(), + }, nil +} + +func NewDefaultScaleApplier(current, cluster *v2.Cluster) (Interface, error) { + if cluster.Name == "" { + return nil, fmt.Errorf("cluster name cannot be empty") + } + cFile := clusterfile.NewClusterFile(contants.Clusterfile(cluster.Name)) + return &Applier{ + ClusterDesired: cluster, + ClusterFile: cFile, + ClusterCurrent: current, }, nil } @@ -67,7 +81,15 @@ func (c *Applier) Apply() error { } func (c *Applier) reconcileCluster() error { - return c.installApp() + if err := c.installApp(); err != nil { + return err + } + mj, md := iputils.GetDiffHosts(c.ClusterCurrent.GetMasterIPList(), c.ClusterDesired.GetMasterIPList()) + nj, nd := iputils.GetDiffHosts(c.ClusterCurrent.GetNodeIPList(), c.ClusterDesired.GetNodeIPList()) + //if len(mj) == 0 && len(md) == 0 && len(nj) == 0 && len(nd) == 0 { + // return c.upgrade() + //} + return c.scaleCluster(mj, md, nj, nd) } func (c *Applier) initCluster() error { @@ -85,7 +107,7 @@ func (c *Applier) initCluster() error { return nil } -func diffImages(spec, curr *v2.Cluster) []string { +func diffImages(spec, curr *v2.Cluster) v2.ImageList { pullImages := make([]string, 0) for _, img := range spec.Spec.Image { if strings.NotIn(img, curr.Spec.Image) { @@ -102,15 +124,38 @@ func (c *Applier) installApp() error { } current := c.ClusterFile.GetCluster() pullImages := diffImages(c.ClusterDesired, current) - installProcessor, err := processor.NewInstallProcessor(c.ClusterFile, pullImages) - if err != nil { - return err - } - err = installProcessor.Execute(c.ClusterDesired) - if err != nil { - return err + if len(pullImages) != 0 { + installProcessor, err := processor.NewInstallProcessor(c.ClusterFile, pullImages) + if err != nil { + return err + } + err = installProcessor.Execute(c.ClusterDesired) + if err != nil { + return err + } } + logger.Info("no change exec install app images") + return nil +} +func (c *Applier) scaleCluster(mj, md, nj, nd []string) error { + logger.Info("Start to scale this cluster") + logger.Debug("current cluster: master %s, worker %s", c.ClusterCurrent.GetMasterIPList(), c.ClusterCurrent.GetNodeIPList()) + logger.Debug("desired cluster: master %s, worker %s", c.ClusterDesired.GetMasterIPList(), c.ClusterDesired.GetNodeIPList()) + if len(mj) == 0 && len(md) == 0 && len(nj) == 0 && len(nd) == 0 { + logger.Info("succeeded in scaling this cluster: no change nodes") + return nil + } + scaleProcessor, err := processor.NewScaleProcessor(c.ClusterFile, c.ClusterDesired.Spec.Image, mj, md, nj, nd) + if err != nil { + return err + } + cluster := c.ClusterDesired + err = scaleProcessor.Execute(cluster) + if err != nil { + return err + } + logger.Info("succeeded in scaling this cluster") return nil } diff --git a/pkg/apply/args.go b/pkg/apply/args.go index 2b6b20eaa..5d961b470 100644 --- a/pkg/apply/args.go +++ b/pkg/apply/args.go @@ -28,3 +28,17 @@ type RunArgs struct { CustomEnv []string CustomCMD []string } + +type ScaleArgs struct { + Masters string + Nodes string + ClusterName string +} + +func (a ScaleArgs) ToRunArgs() *RunArgs { + return &RunArgs{ + Masters: a.Masters, + Nodes: a.Nodes, + ClusterName: a.ClusterName, + } +} diff --git a/pkg/apply/processor/create.go b/pkg/apply/processor/create.go index 8b4389eae..7aece8117 100644 --- a/pkg/apply/processor/create.go +++ b/pkg/apply/processor/create.go @@ -60,7 +60,7 @@ func (c *CreateProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, erro var todoList []func(cluster *v2.Cluster) error todoList = append(todoList, //c.GetPhasePluginFunc(plugin.PhaseOriginally), - c.CreateCluster, + c.PreProcess, c.RunConfig, c.MountRootfs, //c.GetPhasePluginFunc(plugin.PhasePreInit), @@ -73,7 +73,7 @@ func (c *CreateProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, erro return todoList, nil } -func (c *CreateProcessor) CreateCluster(cluster *v2.Cluster) error { +func (c *CreateProcessor) PreProcess(cluster *v2.Cluster) error { err := c.RegistryManager.Pull(cluster.Spec.Image...) if err != nil { return err diff --git a/pkg/apply/processor/delete.go b/pkg/apply/processor/delete.go index e439ce700..98606891a 100644 --- a/pkg/apply/processor/delete.go +++ b/pkg/apply/processor/delete.go @@ -41,9 +41,9 @@ type DeleteProcessor struct { func (d DeleteProcessor) Execute(cluster *v2.Cluster) (err error) { d.cManifestList, err = d.ClusterManager.Inspect(cluster.Name, 0, len(cluster.Spec.Image)) if err != nil { - logger.Warn("delete process failed to inspect cluster, %v", err) + logger.Warn("delete process failed to inspect cluster,make sure you install k8s cluster.") + return err } - d.imgList, err = d.ImageManager.Inspect(cluster.Spec.Image...) if err != nil { return fmt.Errorf("failed to inspect image, %v", err) diff --git a/pkg/apply/processor/install.go b/pkg/apply/processor/install.go index 635e6f810..de878b95f 100644 --- a/pkg/apply/processor/install.go +++ b/pkg/apply/processor/install.go @@ -57,7 +57,7 @@ func (c *InstallProcessor) Execute(cluster *v2.Cluster) error { func (c *InstallProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, error) { var todoList []func(cluster *v2.Cluster) error todoList = append(todoList, - c.ChangeCluster, + c.PreProcess, c.RunConfig, c.MountRootfs, //i.GetPhasePluginFunc(plugin.PhasePreGuest), @@ -67,7 +67,7 @@ func (c *InstallProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, err return todoList, nil } -func (c *InstallProcessor) ChangeCluster(cluster *v2.Cluster) error { +func (c *InstallProcessor) PreProcess(cluster *v2.Cluster) error { err := c.ClusterFile.Process() if err != nil { return err @@ -114,7 +114,7 @@ func (c *InstallProcessor) RunGuest(cluster *v2.Cluster) error { return c.Guest.Apply(cluster, images) } -func NewInstallProcessor(clusterFile clusterfile.Interface, images []string) (Interface, error) { +func NewInstallProcessor(clusterFile clusterfile.Interface, images v2.ImageList) (Interface, error) { imgSvc, err := image.NewImageService() if err != nil { return nil, err diff --git a/pkg/apply/processor/scale.go b/pkg/apply/processor/scale.go new file mode 100644 index 000000000..74bea2b6e --- /dev/null +++ b/pkg/apply/processor/scale.go @@ -0,0 +1,192 @@ +// Copyright © 2021 Alibaba Group Holding Ltd. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package processor + +import ( + "context" + "fmt" + + "github.com/fanux/sealos/pkg/clusterfile" + "github.com/fanux/sealos/pkg/config" + "github.com/fanux/sealos/pkg/filesystem" + "github.com/fanux/sealos/pkg/image" + "github.com/fanux/sealos/pkg/image/types" + "github.com/fanux/sealos/pkg/runtime" + v2 "github.com/fanux/sealos/pkg/types/v1beta1" + "github.com/fanux/sealos/pkg/utils/contants" + "github.com/fanux/sealos/pkg/utils/logger" + "github.com/fanux/sealos/pkg/utils/yaml" + "golang.org/x/sync/errgroup" +) + +type ScaleProcessor struct { + ClusterFile clusterfile.Interface + Runtime runtime.Interface + ImageManager types.Service + ClusterManager types.ClusterService + pullImages []string + imageList types.ImageListOCIV1 + cManifestList types.ClusterManifestList + MastersToJoin []string + MastersToDelete []string + NodesToJoin []string + NodesToDelete []string + IsScaleUp bool +} + +func (c *ScaleProcessor) Execute(cluster *v2.Cluster) error { + pipLine, err := c.GetPipeLine() + if err != nil { + return err + } + + for _, f := range pipLine { + if err = f(cluster); err != nil { + return err + } + } + + return nil +} +func (c *ScaleProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, error) { + var todoList []func(cluster *v2.Cluster) error + if c.IsScaleUp { + todoList = append(todoList, + c.PreProcess, + c.RunConfig, + c.MountRootfs, + //s.GetPhasePluginFunc(plugin.PhasePreJoin), + c.Join, + //s.GetPhasePluginFunc(plugin.PhasePostJoin), + ) + return todoList, nil + } + + todoList = append(todoList, + c.PreProcess, + c.Delete, + //c.ApplyCleanPlugin, + c.UnMountRootfs, + ) + return todoList, nil +} +func (c *ScaleProcessor) Delete(cluster *v2.Cluster) error { + err := c.Runtime.DeleteMasters(c.MastersToDelete) + if err != nil { + return err + } + return c.Runtime.DeleteNodes(c.NodesToDelete) +} +func (c *ScaleProcessor) Join(cluster *v2.Cluster) error { + err := c.Runtime.JoinMasters(c.MastersToJoin) + if err != nil { + return err + } + return c.Runtime.JoinNodes(c.NodesToJoin) +} + +func (c ScaleProcessor) UnMountRootfs(cluster *v2.Cluster) error { + hosts := append(cluster.GetMasterIPAndPortList(), cluster.GetNodeIPAndPortList()...) + if c.cManifestList == nil { + logger.Warn("delete process unmount rootfs skip is cluster not mount rootfs") + return nil + } + fs, err := filesystem.NewRootfsMounter(c.cManifestList, c.imageList) + if err != nil { + return err + } + return fs.UnMountRootfs(cluster, hosts) +} + +func (c *ScaleProcessor) PreProcess(cluster *v2.Cluster) error { + err := c.ClusterFile.Process() + if err != nil { + return err + } + img, err := c.ImageManager.Inspect(c.pullImages...) + if err != nil { + return err + } + c.imageList = img + c.cManifestList, err = c.ClusterManager.Inspect(cluster.Name, 0, len(c.pullImages)) + if err != nil { + return err + } + if c.IsScaleUp { + clusterPath := contants.Clusterfile(cluster.Name) + if err = yaml.MarshalYamlToFile(clusterPath, cluster); err != nil { + return err + } + } + runTime, err := runtime.NewDefaultRuntime(cluster, c.ClusterFile.GetKubeadmConfig(), c.imageList) + if err != nil { + return fmt.Errorf("failed to init runtime, %v", err) + } + c.Runtime = runTime + + return err +} + +func (c *ScaleProcessor) RunConfig(cluster *v2.Cluster) error { + eg, _ := errgroup.WithContext(context.Background()) + for _, cManifest := range c.cManifestList { + manifest := cManifest + eg.Go(func() error { + cfg := config.NewConfiguration(manifest.MountPoint, c.ClusterFile.GetConfigs()) + return cfg.Dump(contants.Clusterfile(cluster.Name)) + }) + } + return eg.Wait() +} + +func (c *ScaleProcessor) MountRootfs(cluster *v2.Cluster) error { + hosts := append(cluster.GetMasterIPAndPortList(), cluster.GetNodeIPAndPortList()...) + fs, err := filesystem.NewRootfsMounter(c.cManifestList, c.imageList) + if err != nil { + return err + } + + return fs.MountRootfs(cluster, hosts, false) +} + +func NewScaleProcessor(clusterFile clusterfile.Interface, images v2.ImageList, masterToJoin, masterToDelete, nodeToJoin, nodeToDelete []string) (Interface, error) { + imgSvc, err := image.NewImageService() + if err != nil { + return nil, err + } + + clusterSvc, err := image.NewClusterService() + if err != nil { + return nil, err + } + + var up bool + // only scale up or scale down at a time + if len(masterToJoin) > 0 || len(nodeToJoin) > 0 { + up = true + } + + return &ScaleProcessor{ + MastersToDelete: masterToDelete, + MastersToJoin: masterToJoin, + NodesToDelete: nodeToDelete, + NodesToJoin: nodeToJoin, + ClusterFile: clusterFile, + ImageManager: imgSvc, + ClusterManager: clusterSvc, + pullImages: images, + IsScaleUp: up, + }, nil +} diff --git a/pkg/apply/run.go b/pkg/apply/run.go index 03d73a54d..61469f172 100644 --- a/pkg/apply/run.go +++ b/pkg/apply/run.go @@ -20,15 +20,11 @@ import ( "path/filepath" "strconv" - "github.com/fanux/sealos/pkg/checker" - "github.com/fanux/sealos/pkg/apply/applydrivers" "github.com/fanux/sealos/pkg/clusterfile" - fileutil "github.com/fanux/sealos/pkg/utils/file" - "github.com/fanux/sealos/pkg/utils/yaml" - v2 "github.com/fanux/sealos/pkg/types/v1beta1" "github.com/fanux/sealos/pkg/utils/contants" + fileutil "github.com/fanux/sealos/pkg/utils/file" "github.com/fanux/sealos/pkg/utils/iputils" "github.com/fanux/sealos/pkg/utils/logger" strings2 "github.com/fanux/sealos/pkg/utils/strings" @@ -64,7 +60,7 @@ func NewApplierFromArgs(imageName []string, args *RunArgs) (applydrivers.Interfa if err := c.SetClusterArgs(imageName, args); err != nil { return nil, err } - if err := c.Process(args); err != nil { + if err := Process(c.cluster); err != nil { return nil, err } return applydrivers.NewDefaultApplier(c.cluster) @@ -130,16 +126,6 @@ func (r *ClusterArgs) SetClusterArgs(imageList []string, args *RunArgs) error { return nil } -func (r *ClusterArgs) Process(args *RunArgs) error { - clusterPath := contants.Clusterfile(args.ClusterName) - err := checker.RunCheckList([]checker.Interface{checker.NewHostChecker()}, r.cluster, checker.PhasePre) - if err != nil { - return err - } - logger.Debug("write cluster file to local storage: %s", clusterPath) - return yaml.MarshalYamlToFile(clusterPath, r.cluster) -} - func (r *ClusterArgs) setHostWithIpsPort(ips []string, roles []string) { defaultPort := strconv.Itoa(int(r.cluster.Spec.SSH.Port)) hostMap := map[string]*v2.Host{} diff --git a/pkg/apply/scale.go b/pkg/apply/scale.go new file mode 100644 index 000000000..bd49548c9 --- /dev/null +++ b/pkg/apply/scale.go @@ -0,0 +1,173 @@ +// Copyright © 2021 Alibaba Group Holding Ltd. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package apply + +import ( + "fmt" + "strconv" + "strings" + + "github.com/fanux/sealos/pkg/checker" + + "github.com/fanux/sealos/pkg/apply/applydrivers" + "github.com/fanux/sealos/pkg/clusterfile" + v2 "github.com/fanux/sealos/pkg/types/v1beta1" + "github.com/fanux/sealos/pkg/utils/contants" + fileutil "github.com/fanux/sealos/pkg/utils/file" + "github.com/fanux/sealos/pkg/utils/iputils" + strings2 "github.com/fanux/sealos/pkg/utils/strings" +) + +// NewScaleApplierFromArgs will filter ip list from command parameters. +func NewScaleApplierFromArgs(scaleArgs *ScaleArgs, flag string) (applydrivers.Interface, error) { + var cluster *v2.Cluster + var curr *v2.Cluster + clusterPath := contants.Clusterfile(scaleArgs.ClusterName) + if !fileutil.IsExist(clusterPath) { + cluster = initCluster(scaleArgs.ClusterName) + curr = cluster + } else { + clusterFile := clusterfile.NewClusterFile(clusterPath) + err := clusterFile.Process() + if err != nil { + return nil, err + } + cluster = clusterFile.GetCluster() + curr = clusterFile.GetCluster().DeepCopy() + } + + if scaleArgs.Nodes == "" && scaleArgs.Masters == "" { + return nil, fmt.Errorf("the node or master parameter was not committed") + } + var err error + switch flag { + case "add": + err = Join(cluster, scaleArgs.ToRunArgs()) + if err != nil { + return nil, err + } + err = Process(cluster) + case "delete": + err = Delete(cluster, scaleArgs.ToRunArgs()) + } + if err != nil { + return nil, err + } + return applydrivers.NewDefaultScaleApplier(curr, cluster) +} + +func Process(cluster *v2.Cluster) error { + err := checker.RunCheckList([]checker.Interface{checker.NewHostChecker()}, cluster, checker.PhasePre) + if err != nil { + return err + } + return nil +} + +func Join(cluster *v2.Cluster, scalingArgs *RunArgs) error { + return joinNodes(cluster, scalingArgs) +} + +func joinNodes(cluster *v2.Cluster, scaleArgs *RunArgs) error { + if err := PreProcessIPList(scaleArgs); err != nil { + return err + } + if (!IsIPList(scaleArgs.Nodes) && scaleArgs.Nodes != "") || (!IsIPList(scaleArgs.Masters) && scaleArgs.Masters != "") { + return fmt.Errorf(" Parameter error: The current mode should submit iplist!") + } + + 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) + } + } + } + return nil +} + +func Delete(cluster *v2.Cluster, scaleArgs *RunArgs) error { + return deleteNodes(cluster, scaleArgs) +} + +func deleteNodes(cluster *v2.Cluster, scaleArgs *RunArgs) error { + if err := PreProcessIPList(scaleArgs); err != nil { + return err + } + if (!IsIPList(scaleArgs.Nodes) && scaleArgs.Nodes != "") || (!IsIPList(scaleArgs.Masters) && scaleArgs.Masters != "") { + return fmt.Errorf(" Parameter error: The current mode should submit iplist!") + } + //master0 machine cannot be deleted + 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 { + if strings2.InList(v2.MASTER, cluster.Spec.Hosts[i].Roles) { + cluster.Spec.Hosts[i].IPS = returnFilteredIPList(cluster.Spec.Hosts[i].IPS, strings.Split(scaleArgs.Masters, ","), defaultPort) + } + } + } + if scaleArgs.Nodes != "" && IsIPList(scaleArgs.Nodes) { + for i := range cluster.Spec.Hosts { + if strings2.InList(v2.NODE, cluster.Spec.Hosts[i].Roles) { + cluster.Spec.Hosts[i].IPS = returnFilteredIPList(cluster.Spec.Hosts[i].IPS, strings.Split(scaleArgs.Nodes, ","), defaultPort) + } + } + } + return nil +} + +func returnFilteredIPList(clusterIPList []string, toBeDeletedIPList []string, defaultPort string) (res []string) { + toBeDeletedIPList = fillIPAndPort(toBeDeletedIPList, defaultPort) + for _, ip := range clusterIPList { + if strings2.NotIn(ip, toBeDeletedIPList) { + res = append(res, fmt.Sprintf("%s:%s", ip, defaultPort)) + } + } + return +} + +func fillIPAndPort(ipList []string, defaultPort string) []string { + var ipAndPorts []string + for _, ip := range ipList { + targetIP, targetPort := iputils.GetHostIPAndPortOrDefault(ip, defaultPort) + ipAndPort := fmt.Sprintf("%s:%s", targetIP, targetPort) + ipAndPorts = append(ipAndPorts, ipAndPort) + } + return ipAndPorts +} diff --git a/pkg/apply/utils.go b/pkg/apply/utils.go index 763cc2014..77e3b7a3f 100644 --- a/pkg/apply/utils.go +++ b/pkg/apply/utils.go @@ -18,6 +18,8 @@ package apply import ( "fmt" + "net" + "strings" v2 "github.com/fanux/sealos/pkg/types/v1beta1" "github.com/fanux/sealos/pkg/utils/iputils" @@ -56,3 +58,17 @@ func PreProcessIPList(joinArgs *RunArgs) error { func removeIPListDuplicatesAndEmpty(ipList []string) []string { return strings2.RemoveDuplicate(strings2.RemoveStrSlice(ipList, []string{""})) } + +func IsIPList(args string) bool { + ipList := strings.Split(args, ",") + + for _, i := range ipList { + if !strings.Contains(i, ":") { + return net.ParseIP(i) != nil + } + if _, err := net.ResolveTCPAddr("tcp", i); err != nil { + return false + } + } + return true +} diff --git a/pkg/image/binary/cluster.go b/pkg/image/binary/cluster.go index c924ff64b..7b1ad89fd 100644 --- a/pkg/image/binary/cluster.go +++ b/pkg/image/binary/cluster.go @@ -36,14 +36,33 @@ type ClusterService struct { } func (d *ClusterService) Create(name string, index int, images ...string) (types.ClusterManifestList, error) { + data := exec.BashEval("buildah containers --json") + infos, err := listContainer(data) + if err != nil { + return nil, err + } + + for _, info := range infos { + for i := range images { + containerName := fmt.Sprintf("%s-%d", name, index+i) + if info.Containername == containerName { + cmd := fmt.Sprintf("buildah unmount %s && buildah rm %s", info.Containername, info.Containername) + if err = exec.Cmd("bash", "-c", cmd); err != nil { + return nil, err + } + } + } + } + var cmd strings.Builder for i, image := range images { - cmd.WriteString(fmt.Sprintf(" buildah from --pull=never --name %s-%d %s && buildah mount %s-%d ", name, index+i, image, name, index+i)) + containerName := fmt.Sprintf("%s-%d", name, index+i) + cmd.WriteString(fmt.Sprintf(" buildah from --pull=never --name %s %s && buildah mount %s ", containerName, image, containerName)) if i != len(images)-1 { cmd.WriteString(" && ") } } - err := exec.Cmd("bash", "-c", cmd.String()) + err = exec.Cmd("bash", "-c", cmd.String()) if err != nil { return nil, err } diff --git a/pkg/registry/util.go b/pkg/registry/util.go index 44d55c56c..38e475520 100644 --- a/pkg/registry/util.go +++ b/pkg/registry/util.go @@ -18,6 +18,8 @@ import ( "fmt" "strings" + "github.com/fanux/sealos/pkg/utils/logger" + "github.com/docker/docker/api/types" fileutil "github.com/fanux/sealos/pkg/utils/file" "k8s.io/apimachinery/pkg/util/json" @@ -94,17 +96,22 @@ func ParseNormalizedNamed(s string) (Named, error) { func GetAuthInfo() (map[string]types.AuthConfig, error) { authFile := "/run/user/0/containers/auth.json" - type auths struct { - Auths map[string]types.AuthConfig `json:"auths"` + if !fileutil.IsExist(authFile) { + logger.Warn("if you access private registry,you must be 'sealos login' or 'buildah login'") + } else { + type auths struct { + Auths map[string]types.AuthConfig `json:"auths"` + } + aus := &auths{} + data, err := fileutil.ReadAll(authFile) + if err != nil { + return nil, err + } + err = json.Unmarshal(data, aus) + if err != nil { + return nil, err + } + return aus.Auths, nil } - aus := &auths{} - data, err := fileutil.ReadAll(authFile) - if err != nil { - return nil, err - } - err = json.Unmarshal(data, aus) - if err != nil { - return nil, err - } - return aus.Auths, nil + return nil, nil } diff --git a/pkg/utils/iputils/iputils_v2.go b/pkg/utils/iputils/iputils_v2.go index dfcc8ea81..5d0bfb539 100644 --- a/pkg/utils/iputils/iputils_v2.go +++ b/pkg/utils/iputils/iputils_v2.go @@ -20,6 +20,8 @@ import ( "net" "strings" + "k8s.io/apimachinery/pkg/util/sets" + "github.com/fanux/sealos/pkg/utils/logger" ) @@ -30,6 +32,19 @@ func GetHostIP(host string) string { } return strings.Split(host, ":")[0] } +func GetDiffHosts(hostsOld, hostsNew []string) (add, sub []string) { + // Difference returns a set of objects that are not in s2 + // For example: + // s1 = {a1, a2, a3} + // s2 = {a1, a2, a4, a5} + // s1.Difference(s2) = {a3} + // s2.Difference(s1) = {a4, a5} + oldSet := sets.NewString(hostsOld...) + newSet := sets.NewString(hostsNew...) + add = newSet.Difference(oldSet).List() + sub = oldSet.Difference(newSet).List() + return +} func GetHostIPs(hosts []string) []string { var ips []string @@ -50,7 +65,13 @@ func GetHostIPAndPortOrDefault(host, Default string) (string, string) { func GetSSHHostIPAndPort(host string) (string, string) { return GetHostIPAndPortOrDefault(host, "22") } - +func GetHostIPAndPortSlice(hosts []string, Default string) (res []string) { + for _, ip := range hosts { + _ip, port := GetHostIPAndPortOrDefault(ip, Default) + res = append(res, fmt.Sprintf("%s:%s", _ip, port)) + } + return +} func GetHostIPSlice(hosts []string) (res []string) { for _, ip := range hosts { res = append(res, GetHostIP(ip))