diff --git a/config/calico/Clusterfile b/config/calico/Clusterfile index cbb4edd6d..f7cf805ee 100644 --- a/config/calico/Clusterfile +++ b/config/calico/Clusterfile @@ -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 diff --git a/config/registry/Clusterfile b/config/registry/Clusterfile index 73dd4f825..a523f68b2 100644 --- a/config/registry/Clusterfile +++ b/config/registry/Clusterfile @@ -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 diff --git a/pkg/apply/applydrivers/apply_drivers_default.go b/pkg/apply/applydrivers/apply_drivers_default.go index 9d11d33d5..018d67440 100644 --- a/pkg/apply/applydrivers/apply_drivers_default.go +++ b/pkg/apply/applydrivers/apply_drivers_default.go @@ -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 +} diff --git a/pkg/apply/processor/create.go b/pkg/apply/processor/create.go new file mode 100644 index 000000000..aecd770fa --- /dev/null +++ b/pkg/apply/processor/create.go @@ -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 +} diff --git a/pkg/apply/processor/delete.go b/pkg/apply/processor/delete.go new file mode 100644 index 000000000..91cf141b1 --- /dev/null +++ b/pkg/apply/processor/delete.go @@ -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 +} diff --git a/pkg/apply/processor/interface.go b/pkg/apply/processor/interface.go new file mode 100644 index 000000000..40487af87 --- /dev/null +++ b/pkg/apply/processor/interface.go @@ -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 +} diff --git a/pkg/apply/run.go b/pkg/apply/run.go index 9e4fa4277..e6a0927c7 100644 --- a/pkg/apply/run.go +++ b/pkg/apply/run.go @@ -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 } diff --git a/pkg/apply/utils.go b/pkg/apply/utils.go index 089188a38..da448b90c 100644 --- a/pkg/apply/utils.go +++ b/pkg/apply/utils.go @@ -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 } diff --git a/pkg/checker/host_checker.go b/pkg/checker/host_checker.go index a4daae1a5..fa46adae7 100644 --- a/pkg/checker/host_checker.go +++ b/pkg/checker/host_checker.go @@ -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 { diff --git a/pkg/config/config.go b/pkg/config/config.go index 6ebab729a..7c90b2117 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -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 - } -} diff --git a/pkg/env/env.go b/pkg/env/env.go index 2fc56ed98..9d3e8bdbc 100644 --- a/pkg/env/env.go +++ b/pkg/env/env.go @@ -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) } diff --git a/pkg/env/env_test.go b/pkg/env/env_test.go index fc7fac5ef..4e9fa39a1 100644 --- a/pkg/env/env_test.go +++ b/pkg/env/env_test.go @@ -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 { diff --git a/pkg/filesystem/filesystem.go b/pkg/filesystem/filesystem.go index e6209e38d..ea8364496 100644 --- a/pkg/filesystem/filesystem.go +++ b/pkg/filesystem/filesystem.go @@ -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) } diff --git a/pkg/filesystem/rootfs/rootfs_default.go b/pkg/filesystem/rootfs/rootfs_default.go index 6baf17d61..48c00da6f 100644 --- a/pkg/filesystem/rootfs/rootfs_default.go +++ b/pkg/filesystem/rootfs/rootfs_default.go @@ -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 } diff --git a/pkg/guest/guest.go b/pkg/guest/guest.go new file mode 100644 index 000000000..350b806fd --- /dev/null +++ b/pkg/guest/guest.go @@ -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] [ep-1 foo bar] +// [/ep-1] [foo bar] [/ep-2] [ep-2] +// [/ep-1] [foo bar] [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") +} diff --git a/pkg/image/cluster.go b/pkg/image/cluster.go index 2d238e51e..0dc54322d 100644 --- a/pkg/image/cluster.go +++ b/pkg/image/cluster.go @@ -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) { diff --git a/pkg/image/image.go b/pkg/image/image.go index 827edeac8..f53539e6c 100644 --- a/pkg/image/image.go +++ b/pkg/image/image.go @@ -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 { diff --git a/pkg/image/interface.go b/pkg/image/interface.go index b2159e260..f82c184d1 100644 --- a/pkg/image/interface.go +++ b/pkg/image/interface.go @@ -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 } diff --git a/pkg/image/registry.go b/pkg/image/registry.go index 85269149e..8528827f9 100644 --- a/pkg/image/registry.go +++ b/pkg/image/registry.go @@ -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") } diff --git a/pkg/infra/README.md b/pkg/infra/README.md index 33415ce3e..7c44df498 100644 --- a/pkg/infra/README.md +++ b/pkg/infra/README.md @@ -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 === 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: - 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: - 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 diff --git a/pkg/types/v1beta1/register.go b/pkg/types/v1beta1/register.go index b65345c8e..0f0b69f30 100644 --- a/pkg/types/v1beta1/register.go +++ b/pkg/types/v1beta1/register.go @@ -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"} diff --git a/pkg/types/v1beta1/utils.go b/pkg/types/v1beta1/utils.go deleted file mode 100644 index 92e5c3bd2..000000000 --- a/pkg/types/v1beta1/utils.go +++ /dev/null @@ -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 -} diff --git a/pkg/utils/contants/bash.go b/pkg/utils/contants/bash.go index 41d65d9d9..e939ae94e 100644 --- a/pkg/utils/contants/bash.go +++ b/pkg/utils/contants/bash.go @@ -20,6 +20,7 @@ import "fmt" const ( DefaultBashFmt = "cd %s && bash %s" + CdAndExecCmd = "cd %s && %s" renderInit = "init" renderClean = "clean" renderInitRegistry = "init-registry" diff --git a/pkg/utils/json/json.go b/pkg/utils/json/json.go index fae7257ae..e28f6fc21 100644 --- a/pkg/utils/json/json.go +++ b/pkg/utils/json/json.go @@ -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)) +} diff --git a/pkg/utils/maps/maps.go b/pkg/utils/maps/maps.go index 446d0d59e..dc551154e 100644 --- a/pkg/utils/maps/maps.go +++ b/pkg/utils/maps/maps.go @@ -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 + } +} diff --git a/pkg/utils/maps/maps_test.go b/pkg/utils/maps/maps_test.go new file mode 100644 index 000000000..1fa0dfd3e --- /dev/null +++ b/pkg/utils/maps/maps_test.go @@ -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) + } + }) + } +} diff --git a/pkg/utils/strings/strings.go b/pkg/utils/strings/strings.go index 50a7557b3..7ec2e9731 100644 --- a/pkg/utils/strings/strings.go +++ b/pkg/utils/strings/strings.go @@ -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) +}