mirror of
https://github.com/labring/sealos.git
synced 2026-09-24 15:46:19 +08:00
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
This commit is contained in:
@@ -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())
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
+2
-16
@@ -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{}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
+19
-12
@@ -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
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
|
||||
Reference in New Issue
Block a user