refactor(master): add init and reset logic (#904)

* refactor(master): add init and reset logic
This commit is contained in:
cuisongliu
2022-03-29 22:11:10 +08:00
committed by GitHub
parent f0f74dda5e
commit 34244d27ef
27 changed files with 846 additions and 236 deletions
+1 -1
View File
@@ -3,7 +3,7 @@ kind: ClusterConfiguration
networking:
podSubnet: 100.64.0.0/10
---
apiVersion: apps.sealyun.com/v1
apiVersion: apps.sealyun.com/v1beta1
kind: Config
metadata:
name: calico
+2 -2
View File
@@ -1,4 +1,4 @@
apiVersion: apps.sealyun.com/v1
apiVersion: apps.sealyun.com/v1beta1
kind: Cluster
metadata:
creationTimestamp: null
@@ -29,7 +29,7 @@ spec:
user: root
---
apiVersion: apps.sealyun.com/v1
apiVersion: apps.sealyun.com/v1beta1
kind: Config
metadata:
creationTimestamp: null
+63 -34
View File
@@ -19,11 +19,13 @@ package applydrivers
import (
"fmt"
"github.com/fanux/sealos/pkg/apply/processor"
"github.com/fanux/sealos/pkg/utils/logger"
"github.com/fanux/sealos/pkg/utils/yaml"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/fanux/sealos/pkg/client-go/kubernetes"
"github.com/fanux/sealos/pkg/clusterfile"
"github.com/fanux/sealos/pkg/filesystem"
"github.com/fanux/sealos/pkg/filesystem/rootfs"
"github.com/fanux/sealos/pkg/image"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
"github.com/fanux/sealos/pkg/utils/contants"
"k8s.io/apimachinery/pkg/version"
@@ -33,33 +35,10 @@ func NewDefaultApplier(cluster *v2.Cluster) (Interface, error) {
if cluster.Name == "" {
return nil, fmt.Errorf("cluster name cannot be empty")
}
imgSvc, err := image.NewImageService()
if err != nil {
return nil, err
}
clusterSvc, err := image.NewDefaultClusterService()
if err != nil {
return nil, err
}
registrySvc, err := image.NewDefaultRegistryService()
if err != nil {
return nil, err
}
mounter, err := filesystem.NewRootfsMounter(clusterSvc)
if err != nil {
return nil, err
}
cFile := clusterfile.NewClusterFile(contants.Clusterfile(cluster.Name))
return &Applier{
ClusterDesired: cluster,
ImageManager: imgSvc,
ClusterFile: cFile,
RegistryManager: registrySvc,
ClusterManager: clusterSvc,
RootfsFSystem: mounter,
ClusterDesired: cluster,
ClusterFile: cFile,
}, nil
}
@@ -67,13 +46,63 @@ type Applier struct {
ClusterDesired *v2.Cluster
ClusterCurrent *v2.Cluster
ClusterFile clusterfile.Interface
ImageManager image.Service
RegistryManager image.RegistryService
ClusterManager image.ClusterService
RootfsFSystem rootfs.Interface
Client kubernetes.Client
CurrentClusterInfo *version.Info
}
func (*Applier) Apply() error { return nil }
func (*Applier) Delete() error { return nil }
func (c *Applier) Apply() error {
clusterPath := contants.Clusterfile(c.ClusterDesired.Name)
if c.ClusterDesired.CreationTimestamp.IsZero() {
if err := c.initCluster(); err != nil {
return err
}
c.ClusterDesired.CreationTimestamp = metav1.Now()
} else {
if err := c.reconcileCluster(); err != nil {
return err
}
}
logger.Debug("write cluster file to local storage: %s", clusterPath)
return yaml.MarshalYamlToFile(clusterPath, c.ClusterDesired)
}
func (c *Applier) reconcileCluster() error {
return nil
}
func (c *Applier) initCluster() error {
logger.Info("Start to create a new cluster: master %s, worker %s", c.ClusterDesired.GetMasterIPList(), c.ClusterDesired.GetNodeIPList())
createProcessor, err := processor.NewCreateProcessor(c.ClusterFile)
if err != nil {
return err
}
if err := createProcessor.Execute(c.ClusterDesired); err != nil {
return err
}
logger.Info("succeeded in creating a new cluster, enjoy it!")
return nil
}
func (c *Applier) Delete() error {
t := metav1.Now()
c.ClusterDesired.DeletionTimestamp = &t
return c.deleteCluster()
}
func (c *Applier) deleteCluster() error {
deleteProcessor, err := processor.NewDeleteProcessor(c.ClusterFile)
if err != nil {
return err
}
if err := deleteProcessor.Execute(c.ClusterDesired); err != nil {
return err
}
logger.Info("succeeded in deleting current cluster")
return nil
}
+166
View File
@@ -0,0 +1,166 @@
// Copyright © 2022 sealos.
//
// 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 (
"fmt"
"github.com/fanux/sealos/pkg/clusterfile"
"github.com/fanux/sealos/pkg/config"
"github.com/fanux/sealos/pkg/filesystem"
"github.com/fanux/sealos/pkg/guest"
"github.com/fanux/sealos/pkg/image"
"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/yaml"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
)
type CreateProcessor struct {
ClusterFile clusterfile.Interface
ImageManager image.Service
ClusterManager image.ClusterService
RegistryManager image.RegistryService
Runtime runtime.Interface
Guest guest.Interface
Config config.Interface
img *v1.Image
cManifest *image.ClusterManifest
}
func (c *CreateProcessor) Execute(cluster *v2.Cluster) error {
err := yaml.MarshalYamlToFile(contants.Clusterfile(cluster.GetClusterName()), cluster)
if err != nil {
return err
}
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 *CreateProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, error) {
var todoList []func(cluster *v2.Cluster) error
todoList = append(todoList,
//c.GetPhasePluginFunc(plugin.PhaseOriginally),
c.CreateCluster,
c.RunConfig,
c.MountRootfs,
//c.GetPhasePluginFunc(plugin.PhasePreInit),
c.Init,
c.Join,
//c.GetPhasePluginFunc(plugin.PhasePreGuest),
c.RunGuest,
c.DeleteCluster,
//c.GetPhasePluginFunc(plugin.PhasePostInstall),
)
return todoList, nil
}
func (c *CreateProcessor) CreateCluster(cluster *v2.Cluster) error {
err := c.RegistryManager.Pull(cluster.Spec.Image)
if err != nil {
return err
}
img, err := c.ImageManager.Inspect(cluster.Spec.Image)
if err != nil {
return err
}
c.img = img
runTime, err := runtime.NewDefaultRuntime(cluster, c.ClusterFile.GetKubeadmConfig(), img)
if err != nil {
return fmt.Errorf("failed to init runtime, %v", err)
}
c.Runtime = runTime
c.cManifest, err = c.ClusterManager.Create(cluster.Name, cluster.Spec.Image)
return err
}
func (c *CreateProcessor) RunConfig(cluster *v2.Cluster) error {
c.Config = config.NewConfiguration(c.cManifest.MountPoint, c.ClusterFile.GetConfigs())
return c.Config.Dump(contants.Clusterfile(cluster.GetClusterName()))
}
func (c *CreateProcessor) MountRootfs(cluster *v2.Cluster) error {
hosts := append(cluster.GetMasterIPList(), cluster.GetNodeIPList()...)
fs, err := filesystem.NewRootfsMounter(c.cManifest, c.img)
if err != nil {
return err
}
return fs.MountRootfs(cluster, hosts)
}
func (c *CreateProcessor) Init(cluster *v2.Cluster) error {
return c.Runtime.Init()
}
func (c *CreateProcessor) Join(cluster *v2.Cluster) error {
err := c.Runtime.JoinMasters(cluster.GetMasterIPList()[1:])
if err != nil {
return err
}
err = c.Runtime.JoinNodes(cluster.GetNodeIPList())
if err != nil {
return err
}
return yaml.MarshalYamlToFile(contants.Clusterfile(cluster.GetClusterName()), cluster)
}
func (c *CreateProcessor) RunGuest(cluster *v2.Cluster) error {
return c.Guest.Apply(cluster)
}
func (c *CreateProcessor) DeleteCluster(cluster *v2.Cluster) error {
return c.ClusterManager.Delete(cluster.Name)
}
func NewCreateProcessor(clusterFile clusterfile.Interface) (Interface, error) {
imgSvc, err := image.NewImageService()
if err != nil {
return nil, err
}
clusterSvc, err := image.NewDefaultClusterService()
if err != nil {
return nil, err
}
registrySvc, err := image.NewDefaultRegistryService()
if err != nil {
return nil, err
}
gs, err := guest.NewGuestManager()
if err != nil {
return nil, err
}
return &CreateProcessor{
ClusterFile: clusterFile,
ImageManager: imgSvc,
ClusterManager: clusterSvc,
RegistryManager: registrySvc,
Guest: gs,
}, nil
}
+116
View File
@@ -0,0 +1,116 @@
// Copyright © 2022 sealos.
//
// 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 (
"fmt"
"github.com/fanux/sealos/pkg/clusterfile"
"github.com/fanux/sealos/pkg/filesystem"
"github.com/fanux/sealos/pkg/image"
"github.com/fanux/sealos/pkg/runtime"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
"github.com/fanux/sealos/pkg/utils/contants"
fileutil "github.com/fanux/sealos/pkg/utils/file"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
)
type DeleteProcessor struct {
ClusterManager image.ClusterService
ImageManager image.Service
ClusterFile clusterfile.Interface
img *v1.Image
cManifest *image.ClusterManifest
}
// Execute :according to the different of desired cluster to delete cluster.
func (d DeleteProcessor) Execute(cluster *v2.Cluster) (err error) {
d.cManifest, err = d.ClusterManager.Inspect(cluster.Name)
if err != nil {
return fmt.Errorf("failed to inspect cluster, %v", err)
}
d.img, err = d.ImageManager.Inspect(cluster.Spec.Image)
if err != nil {
return fmt.Errorf("failed to inspect image, %v", err)
}
runTime, err := runtime.NewDefaultRuntime(cluster, d.ClusterFile.GetKubeadmConfig(), d.img)
if err != nil {
return fmt.Errorf("failed to init runtime, %v", err)
}
err = runTime.Reset()
if err != nil {
return err
}
pipLine, err := d.GetPipeLine()
if err != nil {
return err
}
for _, f := range pipLine {
if err = f(cluster); err != nil {
return err
}
}
return nil
}
func (d DeleteProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, error) {
var todoList []func(cluster *v2.Cluster) error
todoList = append(todoList,
d.UnMountRootfs,
d.UnMountImage,
d.CleanFS,
)
return todoList, nil
}
func (d DeleteProcessor) UnMountRootfs(cluster *v2.Cluster) error {
hosts := append(cluster.GetMasterIPList(), cluster.GetNodeIPList()...)
fs, err := filesystem.NewRootfsMounter(d.cManifest, d.img)
if err != nil {
return err
}
return fs.UnMountRootfs(cluster, hosts)
}
func (d DeleteProcessor) UnMountImage(cluster *v2.Cluster) error {
return d.ClusterManager.Delete(cluster.Name)
}
func (d DeleteProcessor) CleanFS(cluster *v2.Cluster) error {
workDir := contants.ClusterDir(cluster.Name)
dataDir := contants.NewData(cluster.Name).Homedir()
return fileutil.CleanFiles(workDir, dataDir)
}
func NewDeleteProcessor(clusterFile clusterfile.Interface) (Interface, error) {
imgSvc, err := image.NewImageService()
if err != nil {
return nil, err
}
clusterSvc, err := image.NewDefaultClusterService()
if err != nil {
return nil, err
}
return DeleteProcessor{
ClusterFile: clusterFile,
ImageManager: imgSvc,
ClusterManager: clusterSvc,
}, nil
}
+22
View File
@@ -0,0 +1,22 @@
// Copyright © 2022 sealos.
//
// 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 v2 "github.com/fanux/sealos/pkg/types/v1beta1"
type Interface interface {
// Execute :according to the different of desired cluster to do cluster apply.
Execute(cluster *v2.Cluster) error
}
+7 -1
View File
@@ -20,6 +20,8 @@ import (
"fmt"
"os"
"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"
@@ -112,11 +114,15 @@ func (r *ClusterArgs) SetClusterArgs(imageName string, args *RunArgs) error {
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
}
if !args.DryRun {
logger.Debug("write cluster file to local storage: %s", clusterPath)
return yaml.MarshalYamlToFile(clusterPath, r.cluster)
}
data, err := fileutil.ReadAll(clusterPath)
data, err := yaml.MarshalYamlConfigs(r.cluster)
if err != nil {
return err
}
-2
View File
@@ -22,7 +22,6 @@ import (
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
"github.com/fanux/sealos/pkg/utils/iputils"
strings2 "github.com/fanux/sealos/pkg/utils/strings"
v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/sets"
)
@@ -32,7 +31,6 @@ func initCluster(clusterName string) *v2.Cluster {
cluster.Kind = "Cluster"
cluster.APIVersion = v2.SchemeGroupVersion.String()
cluster.Annotations = make(map[string]string)
cluster.CreationTimestamp = v1.Now()
return cluster
}
+4
View File
@@ -19,6 +19,8 @@ import (
"strconv"
"time"
"github.com/fanux/sealos/pkg/utils/logger"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
"github.com/fanux/sealos/pkg/utils/ssh"
)
@@ -43,6 +45,7 @@ func NewHostChecker() Interface {
}
func checkHostnameUnique(cluster *v2.Cluster, ipList []string) error {
logger.Info("checker:hostname %v", ipList)
hostnameList := map[string]bool{}
for _, ip := range ipList {
s, err := ssh.NewSSHByCluster(cluster, false)
@@ -63,6 +66,7 @@ func checkHostnameUnique(cluster *v2.Cluster, ipList []string) error {
//Check whether the node time is synchronized
func checkTimeSync(cluster *v2.Cluster, ipList []string) error {
logger.Info("checker:timeSync %v", ipList)
for _, ip := range ipList {
s, err := ssh.NewSSHByCluster(cluster, false)
if err != nil {
+3 -22
View File
@@ -20,6 +20,8 @@ import (
"io/ioutil"
"path/filepath"
"github.com/fanux/sealos/pkg/utils/maps"
"github.com/fanux/sealos/pkg/utils/contants"
"github.com/fanux/sealos/pkg/types/v1beta1"
@@ -158,7 +160,7 @@ func getMergeConfigData(path string, data []byte) ([]byte, error) {
if len(configMap) == 0 {
continue
}
deepMerge(&configMap, &mergeConfigMap)
maps.DeepMerge(&configMap, &mergeConfigMap)
cfg, err := yaml.Marshal(&configMap)
if err != nil {
@@ -168,24 +170,3 @@ func getMergeConfigData(path string, data []byte) ([]byte, error) {
}
return bytes.Join(configs, []byte("\n---\n")), nil
}
func deepMerge(dst, src *map[string]interface{}) {
for srcK, srcV := range *src {
dstV, ok := (*dst)[srcK]
if !ok {
continue
}
dV, ok := dstV.(map[string]interface{})
// dstV is string type
if !ok {
(*dst)[srcK] = srcV
continue
}
sV, ok := srcV.(map[string]interface{})
if !ok {
continue
}
deepMerge(&dV, &sV)
(*dst)[srcK] = dV
}
}
+12 -36
View File
@@ -21,8 +21,10 @@ import (
"path/filepath"
"strings"
"github.com/fanux/sealos/pkg/utils/maps"
strings2 "github.com/fanux/sealos/pkg/utils/strings"
"github.com/fanux/sealos/pkg/types/v1beta1"
strlib "github.com/fanux/sealos/pkg/utils/strings"
)
const templateSuffix = ".tmpl"
@@ -48,30 +50,15 @@ func NewEnvProcessor(cluster *v1beta1.Cluster) Interface {
}
func (p *processor) WrapperEnv(host string) map[string]string {
env := make(map[string]string)
for k, v := range p.getHostEnv(host) {
switch value := v.(type) {
case []string:
env[k] = strings.Join(value, " ")
case string:
env[k] = value
}
envs := p.getHostEnv(host)
for k, v := range envs {
env[k] = v
}
return env
}
func (p *processor) WrapperShell(host, shell string) string {
var env string
for k, v := range p.getHostEnv(host) {
switch value := v.(type) {
case []string:
env = fmt.Sprintf("%s%s=(%s) ", env, k, strings.Join(value, " "))
case string:
env = fmt.Sprintf("%s%s=\"%s\" ", env, k, value)
}
}
if env == "" {
return shell
}
return fmt.Sprintf("%s&& %s", env, shell)
envs := p.getHostEnv(host)
return strings2.EnvFromMap(shell, envs)
}
func (p *processor) RenderAll(host, dir string) error {
@@ -102,20 +89,9 @@ func (p *processor) RenderAll(host, dir string) error {
})
}
func mergeList(dst, src []string) []string {
for _, s := range src {
if strlib.InList(s, dst) {
continue
}
dst = append(dst, s)
}
return dst
}
// Merge the host ENV and global env, the host env will overwrite cluster.Spec.Env
func (p *processor) getHostEnv(hostIP string) (env map[string]interface{}) {
func (p *processor) getHostEnv(hostIP string) (env map[string]string) {
var hostEnv []string
for _, host := range p.Spec.Hosts {
for _, ip := range host.IPS {
if ip == hostIP {
@@ -123,8 +99,8 @@ func (p *processor) getHostEnv(hostIP string) (env map[string]interface{}) {
}
}
}
hostEnvMap := maps.ListToMap(hostEnv)
specEnvMap := maps.ListToMap(p.Spec.Env)
hostEnv = mergeList(hostEnv, p.Spec.Env)
return v1beta1.ConvertEnv(hostEnv)
return maps.MergeMap(specEnvMap, hostEnvMap)
}
+2 -27
View File
@@ -15,36 +15,11 @@
package env
import (
"reflect"
"testing"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
)
func Test_convertEnv(t *testing.T) {
type args struct {
envList []string
}
tests := []struct {
name string
args args
wantEnv map[string]interface{}
}{
{
"test convert env",
args{envList: []string{"IP=127.0.0.1", "IP=192.168.0.2", "key=value"}},
map[string]interface{}{"IP": []string{"127.0.0.1", "192.168.0.2"}, "key": "value"},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if gotEnv := v2.ConvertEnv(tt.args.envList); !reflect.DeepEqual(gotEnv, tt.wantEnv) {
t.Errorf("convertEnv() = %v, want %v", gotEnv, tt.wantEnv)
}
})
}
}
func getTestCluster() *v2.Cluster {
return &v2.Cluster{
Spec: v2.ClusterSpec{
@@ -53,7 +28,7 @@ func getTestCluster() *v2.Cluster {
{
IPS: []string{"192.168.0.2", "192.168.0.3", "192.168.0.4"},
Roles: []string{"master"},
Env: []string{"key=bar", "key=foo", "foo=bar", "IP=127.0.0.2"},
Env: []string{"key=bar", "foo=bar xxx ddd fffff", "IP=127.0.0.2"},
},
},
SSH: v2.ClusterSSH{},
@@ -82,7 +57,7 @@ func Test_processor_WrapperShell(t *testing.T) {
host: "192.168.0.2",
shell: "echo $foo ${IP[@]}",
},
"key=(bar foo value) foo=bar IP=(127.0.0.2 127.0.0.1) && echo $foo ${IP[@]}",
"IP=(127.0.0.2) key=(bar) foo=(bar xxx ddd fffff) && echo $foo ${IP[@]}",
},
}
for _, tt := range tests {
+3 -2
View File
@@ -19,9 +19,10 @@ package filesystem
import (
"github.com/fanux/sealos/pkg/filesystem/rootfs"
img "github.com/fanux/sealos/pkg/image"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
)
// NewRootfsMounter :according to the Metadata file content to determine what kind of Filesystem will be load.
func NewRootfsMounter(service img.ClusterService) (rootfs.Interface, error) {
return rootfs.NewDefaultRootfs(service)
func NewRootfsMounter(cluster *img.ClusterManifest, img *v1.Image) (rootfs.Interface, error) {
return rootfs.NewDefaultRootfs(cluster, img)
}
+14 -12
View File
@@ -23,6 +23,8 @@ import (
"path"
"path/filepath"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/fanux/sealos/pkg/env"
"github.com/fanux/sealos/pkg/image"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
@@ -35,7 +37,9 @@ import (
)
type defaultRootfs struct {
clusterService image.ClusterService
//clusterService image.ClusterService
img *v1.Image
cluster *image.ClusterManifest
}
func (f *defaultRootfs) MountRootfs(cluster *v2.Cluster, hosts []string) error {
@@ -56,13 +60,10 @@ func (f *defaultRootfs) getSSH(cluster *v2.Cluster) ssh.Interface {
func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string) error {
target := contants.NewData(f.getClusterName(cluster)).RootFSPath()
data, err := f.clusterService.Inspect(f.getClusterName(cluster))
if err != nil {
return errors.Wrap(err, fmt.Sprintf("inspect container %s data failed", f.getClusterName(cluster)))
}
src := data.MountPoint
src := f.cluster.MountPoint
envProcessor := env.NewEnvProcessor(cluster)
err = renderENV(src, ipList, envProcessor)
err := renderENV(src, ipList, envProcessor)
if err != nil {
return errors.Wrap(err, "render env to rootfs failed")
}
@@ -71,16 +72,17 @@ func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string) error
return errors.Wrap(err, "run chmod to rootfs failed")
}
check := contants.NewBash(f.getClusterName(cluster), f.img.Config.Labels)
eg, _ := errgroup.WithContext(context.Background())
for _, IP := range ipList {
ip := IP
eg.Go(func() error {
sshClient := f.getSSH(cluster)
err := CopyFiles(sshClient, ip == cluster.GetMaster0IP(), ip, src, target)
err = CopyFiles(sshClient, ip == cluster.GetMaster0IP(), ip, src, target)
if err != nil {
return fmt.Errorf("copy rootfs failed %v", err)
}
return err
return f.getSSH(cluster).CmdAsync(ip, envProcessor.WrapperShell(ip, check.CheckBash()))
})
}
return eg.Wait()
@@ -107,7 +109,7 @@ func (f *defaultRootfs) unmountRootfs(cluster *v2.Cluster, ipList []string) erro
if err = eg.Wait(); err != nil {
return err
}
return f.clusterService.Delete(f.getClusterName(cluster))
return nil
}
func renderENV(mountDir string, ipList []string, p env.Interface) error {
@@ -151,6 +153,6 @@ func CopyFiles(sshEntry ssh.Interface, isRegistry bool, ip, src, target string)
return nil
}
func NewDefaultRootfs(service image.ClusterService) (Interface, error) {
return &defaultRootfs{clusterService: service}, nil
func NewDefaultRootfs(cluster *image.ClusterManifest, img *v1.Image) (Interface, error) {
return &defaultRootfs{cluster: cluster, img: img}, nil
}
+154
View File
@@ -0,0 +1,154 @@
// Copyright © 2022 sealos.
//
// 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 guest
import (
"fmt"
strings2 "strings"
"github.com/fanux/sealos/pkg/env"
"github.com/fanux/sealos/pkg/image"
"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/fork/golang/expansion"
"github.com/fanux/sealos/pkg/utils/maps"
"github.com/fanux/sealos/pkg/utils/ssh"
"github.com/fanux/sealos/pkg/utils/strings"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
)
type Interface interface {
Apply(cluster *v2.Cluster) error
Delete(cluster *v2.Cluster) error
}
type Default struct {
imageService image.Service
}
func NewGuestManager() (Interface, error) {
is, err := image.NewImageService()
if err != nil {
return nil, err
}
return &Default{imageService: is}, nil
}
func (d *Default) Apply(cluster *v2.Cluster) error {
clusterRootfs := runtime.GetContantData(cluster.Name).RootFSPath()
img, err := d.imageService.Inspect(cluster.Spec.Image)
if err != nil {
return fmt.Errorf("get cluster image failed, %s", err)
}
sshClient := ssh.NewSSHClient(&cluster.Spec.SSH, true)
envInterface := env.NewEnvProcessor(cluster)
envs := envInterface.WrapperEnv(cluster.GetMaster0IP()) //clusterfile
guestCMD := strings2.Join(d.getGuestCmd(envs, cluster, img), " ")
for _, value := range []string{guestCMD} {
if value == "" {
continue
}
if err = sshClient.CmdAsync(cluster.GetMaster0IP(), fmt.Sprintf(contants.CdAndExecCmd, clusterRootfs, value)); err != nil {
return err
}
}
return nil
}
//Image Entrypoint Image Cmd Container command Container args Command run
// [/ep-1] [foo bar] <not set> <not set> [ep-1 foo bar]
// [/ep-1] [foo bar] [/ep-2] <not set> [ep-2]
// [/ep-1] [foo bar] <not set> [zoo boo] [ep-1 zoo boo]
// [/ep-1] [foo bar] [/ep-2] [zoo boo] [ep-2 zoo boo]
func (d *Default) getGuestCmd(envs map[string]string, cluster *v2.Cluster, image *v1.Image) []string {
if image.Config.Env != nil {
baseEnvs := maps.ListToMap(image.Config.Env)
envs = maps.MergeMap(baseEnvs, envs)
}
mapping := expansion.MappingFuncFor(envs)
//If you do not supply command or args for a Container, the defaults defined in the Docker image are used.
//If you supply a command but no args for a Container, only the supplied command is used. The default EntryPoint and the default Cmd defined in the Docker image are ignored.
//If you supply only args for a Container, the default Entrypoint defined in the Docker image is run with the args that you supplied.
//If you supply a command and args, the default Entrypoint and the default Cmd defined in the Docker image are ignored. Your command is run with your args.
if len(cluster.Spec.Command) != 0 && len(cluster.Spec.Args) == 0 {
image.Config.Cmd = []string{}
}
command := make([]string, 0)
args := make([]string, 0)
if len(cluster.Spec.Command) != 0 {
for _, cmd := range cluster.Spec.Command {
command = append(command, expansion.Expand(cmd, mapping))
}
}
if len(cluster.Spec.Args) != 0 {
for _, arg := range cluster.Spec.Args {
args = append(args, expansion.Expand(arg, mapping))
}
}
if len(command) == 0 {
command = image.Config.Entrypoint
}
if len(args) == 0 {
args = image.Config.Cmd
}
resCmd := append(command, args...)
resCmd = expandSh(resCmd)
resCmd = expandBash(resCmd)
return resCmd
}
func expandBash(resCmd []string) []string {
defaultBash := []string{"bash", "/bin/bash", "/usr/bin/bash"}
if len(resCmd) > 0 {
if strings.InList(resCmd[0], defaultBash) {
if len(resCmd) > 2 {
if strings.TrimSpaceWS(resCmd[1]) == "-c" {
resCmd = resCmd[2:]
return resCmd
}
}
resCmd = resCmd[1:]
return resCmd
}
}
return resCmd
}
func expandSh(resCmd []string) []string {
defaultSh := []string{"sh", "/bin/sh", "/usr/bin/sh"}
if len(resCmd) > 0 {
if strings.InList(resCmd[0], defaultSh) {
if len(resCmd) > 2 {
if strings.TrimSpaceWS(resCmd[1]) == "-c" {
resCmd = resCmd[2:]
return resCmd
}
}
resCmd = resCmd[1:]
return resCmd
}
}
return resCmd
}
func (d Default) Delete(cluster *v2.Cluster) error {
panic("implement me")
}
+70 -3
View File
@@ -16,17 +16,84 @@ limitations under the License.
package image
import (
"fmt"
"github.com/fanux/sealos/pkg/utils/exec"
"github.com/pkg/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/util/json"
)
type defaultClusterService struct {
}
func (*defaultClusterService) Create(name, image string) (*ClusterManifest, error) {
return &ClusterManifest{}, nil
func (d *defaultClusterService) Create(name, image string) (*ClusterManifest, error) {
cmd := fmt.Sprintf("buildah from --name %s %s && buildah mount %s", name, image, name)
err := exec.CmdForPipe("bash", "-c", cmd)
if err != nil {
return nil, err
}
return d.Inspect(name)
}
func (*defaultClusterService) Delete(name string) error {
data := exec.BashEval("buildah containers --json")
infos, err := listContainer(data)
if err != nil {
return err
}
for _, info := range infos {
if info.Containername == name {
cmd := fmt.Sprintf("buildah unmount %s && buildah rm %s", name, name)
return exec.CmdForPipe("bash", "-c", cmd)
}
}
return nil
}
func (*defaultClusterService) Inspect(name string) (*ClusterManifest, error) {
return &ClusterManifest{}, nil
data := exec.BashEval(fmt.Sprintf("buildah inspect %s", name))
return inspectContainer(data)
}
func inspectContainer(data string) (*ClusterManifest, error) {
if data != "" {
var outStruct map[string]interface{}
err := json.Unmarshal([]byte(data), &outStruct)
if err != nil {
return nil, errors.Wrap(err, "decode out json from container inspect failed")
}
container, _, _ := unstructured.NestedString(outStruct, "Container")
containerID, _, _ := unstructured.NestedString(outStruct, "ContainerID")
mountPoint, _, _ := unstructured.NestedString(outStruct, "MountPoint")
manifest := &ClusterManifest{
Container: container,
ContainerID: containerID,
MountPoint: mountPoint,
}
return manifest, nil
}
return nil, errors.New("inspect output is empty")
}
type ClusterInfo struct {
ID string `json:"id"`
Builder bool `json:"builder"`
Imageid string `json:"imageid"`
Imagename string `json:"imagename"`
Containername string `json:"containername"`
}
func listContainer(data string) ([]ClusterInfo, error) {
if data != "" {
var outStruct []ClusterInfo
err := json.Unmarshal([]byte(data), &outStruct)
if err != nil {
return nil, errors.Wrap(err, "decode out json from list container failed")
}
return outStruct, nil
}
return nil, errors.New("inspect output is empty")
}
func NewDefaultClusterService() (ClusterService, error) {
+30 -1
View File
@@ -17,7 +17,14 @@ limitations under the License.
package image
import (
"fmt"
"github.com/fanux/sealos/pkg/utils/exec"
json2 "github.com/fanux/sealos/pkg/utils/json"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
"github.com/pkg/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/util/json"
)
// defaultImageService is the default service, which is used for image pull/push
@@ -33,7 +40,29 @@ func (d *defaultImageService) Remove(images ...string) error {
}
func (d *defaultImageService) Inspect(image string) (*v1.Image, error) {
panic("implement me")
data := exec.BashEval(fmt.Sprintf("buildah inspect %s", image))
return inspectImage(data)
}
func inspectImage(data string) (*v1.Image, error) {
if data != "" {
var outStruct map[string]interface{}
err := json.Unmarshal([]byte(data), &outStruct)
if err != nil {
return nil, errors.Wrap(err, "decode out json from image inspect failed")
}
imageData, _, err := unstructured.NestedFieldCopy(outStruct, "OCIv1")
if err != nil {
return nil, errors.Wrap(err, "decode out json from OCIv1 object failed")
}
img := &v1.Image{}
err = json2.Convert(imageData, img)
if err != nil {
return nil, errors.Wrap(err, "decode OCIv1 to v1.Image failed")
}
return img, nil
}
return nil, errors.New("inspect output is empty")
}
func (d *defaultImageService) Build(options BuildOptions, contextDir, imageName string) error {
-1
View File
@@ -34,7 +34,6 @@ type RegistryService interface {
Login(domain, username, passwd string) error
Logout(domain, username string) error
Pull(image string) error
PullIfNotExist(image string) error
Push(image string) error
Sync(localDir, imageName string) error
}
+8 -4
View File
@@ -16,6 +16,12 @@ limitations under the License.
package image
import (
"fmt"
"github.com/fanux/sealos/pkg/utils/exec"
)
type defaultRegistryService struct {
}
@@ -26,11 +32,9 @@ func (*defaultRegistryService) Logout(domain, username string) error {
panic("implement me")
}
func (*defaultRegistryService) Pull(image string) error {
panic("implement me")
}
func (*defaultRegistryService) PullIfNotExist(image string) error {
panic("implement me")
return exec.CmdForPipe("bash", "-c", fmt.Sprintf("buildah pull %s", image))
}
func (*defaultRegistryService) Push(image string) error {
panic("implement me")
}
+3 -3
View File
@@ -85,7 +85,7 @@ GOPATH=/Users/cuisongliu/Workspaces/go #gosetup
2022-01-06 22:11:21 [INFO] [ali_ecs.go:216] reconcile {"roles":["master","ssdxxx"],"cpu":2,"memory":4,"count":1,"disks":[{"capacity":50,"category":""}],"arch":"amd64","ecsType":"ecs.c7a.large","os":{"name":"","version":"","id":"centos_8_0_x64_20G_alibase_20210712.vhd"}} instances success [172.16.0.140]
2022-01-06 22:11:23 [INFO] [ali_provider.go:70] create resource success www.sealyun.com/EipID: eip-uf6ptsp0s0uadt9pr7zcq
<nil>=== RUN TestAliApply/modify_instance_system_disk
infra_test.go:82: output yaml: apiVersion: apps.sealyun.com/v1
infra_test.go:82: output yaml: apiVersion: apps.sealyun.com/v1beta1
kind: Infra
metadata:
creationTimestamp: null
@@ -165,7 +165,7 @@ GOPATH=/Users/cuisongliu/Workspaces/go #gosetup
2022-01-06 22:11:28 [INFO] [ali_ecs.go:206] get up IP list [172.16.0.140], host count 1
2022-01-06 22:11:28 [INFO] [ali_ecs.go:216] reconcile {"roles":["master","ssdxxx"],"cpu":2,"memory":4,"count":1,"disks":[{"capacity":50,"category":""}],"arch":"amd64","ecsType":"ecs.c7a.large","os":{"name":"","version":"","id":"centos_8_0_x64_20G_alibase_20210712.vhd"}} instances success [172.16.0.140]
infra_test.go:107: add server:<nil>
infra_test.go:109: output yaml: apiVersion: apps.sealyun.com/v1
infra_test.go:109: output yaml: apiVersion: apps.sealyun.com/v1beta1
kind: Infra
metadata:
creationTimestamp: null
@@ -278,7 +278,7 @@ Recommend: https://error-center.aliyun.com/status/search?Keyword=ORDER.QUANTITY_
RequestId: 1F933962-8116-55C7-A4A1-00F50835C480
Message: User quota has exceeded the limit. ,skip it
infra_test.go:124: delete:<nil>
infra_test.go:126: output yaml: apiVersion: apps.sealyun.com/v1
infra_test.go:126: output yaml: apiVersion: apps.sealyun.com/v1beta1
kind: Infra
metadata:
creationTimestamp: null
+1 -1
View File
@@ -22,4 +22,4 @@ import (
const GroupName = "apps.sealyun.com"
// SchemeGroupVersion is group version used to register these objects
var SchemeGroupVersion = schema.GroupVersion{Group: GroupName, Version: "v1"}
var SchemeGroupVersion = schema.GroupVersion{Group: GroupName, Version: "v1beta1"}
-82
View File
@@ -1,82 +0,0 @@
/*
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 v1beta1
import (
"strings"
"github.com/fanux/sealos/pkg/utils/fork/golang/expansion"
v1 "github.com/opencontainers/image-spec/specs-go/v1"
)
// ConvertEnv []string to map[string]interface{}, example [IP=127.0.0.1,IP=192.160.0.2,Key=value] will convert to {IP:[127.0.0.1,192.168.0.2],key:value}
func ConvertEnv(envList []string) (env map[string]interface{}) {
temp := make(map[string][]string)
env = make(map[string]interface{})
for _, e := range envList {
var kv []string
if kv = strings.SplitN(e, "=", 2); len(kv) != 2 {
continue
}
temp[kv[0]] = append(temp[kv[0]], kv[1])
}
for k, v := range temp {
if len(v) > 1 {
env[k] = v
continue
}
if len(v) == 1 {
env[k] = v[0]
}
}
return
}
func ConvertCMDRun() {
}
// ExpandCommandAndArgs expands the given Cluster's command by replacing variable references `with the values of given EnvVar.
func (c *Cluster) ExpandCommandAndArgs(envs map[string]string, image v1.Image) (command []string, args []string) {
mapping := expansion.MappingFuncFor(envs)
if len(c.Spec.Command) != 0 {
for _, cmd := range c.Spec.Command {
command = append(command, expansion.Expand(cmd, mapping))
}
}
if len(c.Spec.Args) != 0 {
for _, arg := range c.Spec.Args {
args = append(args, expansion.Expand(arg, mapping))
}
}
if len(command) == 0 {
command = image.Config.Entrypoint
}
if len(args) == 0 {
args = image.Config.Cmd
}
return command, args
}
+1
View File
@@ -20,6 +20,7 @@ import "fmt"
const (
DefaultBashFmt = "cd %s && bash %s"
CdAndExecCmd = "cd %s && %s"
renderInit = "init"
renderClean = "clean"
renderInitRegistry = "init-registry"
+10
View File
@@ -18,6 +18,7 @@ package json
import (
"github.com/fanux/sealos/pkg/utils/file"
"github.com/pkg/errors"
"k8s.io/apimachinery/pkg/util/json"
)
@@ -33,3 +34,12 @@ func Unmarshal(path string) (map[string]interface{}, error) {
}
return data, nil
}
func Convert(from interface{}, to interface{}) error {
var data []byte
var err error
if data, err = json.Marshal(from); err != nil {
return errors.WithStack(err)
}
return errors.WithStack(json.Unmarshal(data, to))
}
+37 -2
View File
@@ -31,9 +31,13 @@ func MapToString(data map[string]string) string {
}
func StringToMap(data string) map[string]string {
m := make(map[string]string)
list := strings.Split(data, ",")
for _, l := range list {
return ListToMap(list)
}
func ListToMap(data []string) map[string]string {
m := make(map[string]string)
for _, l := range data {
if l != "" {
kv := strings.Split(l, "=")
if len(kv) == 2 {
@@ -43,3 +47,34 @@ func StringToMap(data string) map[string]string {
}
return m
}
func MergeMap(ms ...map[string]string) map[string]string {
res := map[string]string{}
for _, m := range ms {
for k, v := range m {
res[k] = v
}
}
return res
}
func DeepMerge(dst, src *map[string]interface{}) {
for srcK, srcV := range *src {
dstV, ok := (*dst)[srcK]
if !ok {
continue
}
dV, ok := dstV.(map[string]interface{})
// dstV is string type
if !ok {
(*dst)[srcK] = srcV
continue
}
sV, ok := srcV.(map[string]interface{})
if !ok {
continue
}
DeepMerge(&dV, &sV)
(*dst)[srcK] = dV
}
}
+106
View File
@@ -0,0 +1,106 @@
/*
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 maps
import (
"reflect"
"testing"
)
func TestMergeMap(t *testing.T) {
type args struct {
dst map[string]string
src map[string]string
}
tests := []struct {
name string
args args
want map[string]string
}{
{
name: "default",
args: args{
dst: map[string]string{
"aa": "cc",
},
src: map[string]string{
"aa": "bb",
},
},
want: map[string]string{
"aa": "bb",
},
},
{
name: "default-add",
args: args{
dst: map[string]string{
"aa": "cc",
},
src: map[string]string{
"aa": "bb",
"bb": "bb",
},
},
want: map[string]string{
"aa": "bb",
"bb": "bb",
},
},
{
name: "default-replace",
args: args{
dst: map[string]string{
"aa": "bb",
"bb": "bb",
},
src: map[string]string{
"bb": "dd",
},
},
want: map[string]string{
"aa": "bb",
"bb": "dd",
},
},
{
name: "default-delete",
args: args{
dst: map[string]string{
"aa": "bb",
"bb": "bb",
},
src: map[string]string{
"cc": "cc",
},
},
want: map[string]string{
"aa": "bb",
"bb": "bb",
"cc": "cc",
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
data := MergeMap(tt.args.dst, tt.args.src)
if !reflect.DeepEqual(data, tt.want) {
t.Errorf("MergeMap() = %v, want %v", data, tt.want)
}
})
}
}
+11
View File
@@ -221,3 +221,14 @@ func IsLetterOrNumber(k string) bool {
}
return true
}
func EnvFromMap(shell string, envs map[string]string) string {
var env string
for k, v := range envs {
env = fmt.Sprintf("%s%s=(%s) ", env, k, v)
}
if env == "" {
return shell
}
return fmt.Sprintf("%s&& %s", env, shell)
}