From 2bd17b8cece251196c6ca6aa4e563595ddbd9581 Mon Sep 17 00:00:00 2001 From: cuisongliu Date: Tue, 12 Jul 2022 16:00:33 +0800 Subject: [PATCH] refactor(main): sealos run other server (#1292) Signed-off-by: cuisongliu --- .../applydrivers/apply_drivers_default.go | 2 ++ pkg/apply/processor/interface.go | 17 +++++++++++ pkg/apply/processor/scale.go | 1 + pkg/constants/contants.go | 3 ++ pkg/constants/data.go | 4 +-- pkg/filesystem/rootfs/rootfs_default.go | 24 ++++++++-------- pkg/guest/guest.go | 6 ++-- pkg/remote/remote.go | 3 +- pkg/runtime/image_shim.go | 28 +++++++++++++------ pkg/runtime/registry.go | 28 +++++++++++++------ pkg/runtime/reset.go | 6 +++- pkg/runtime/runtime_getter.go | 2 +- pkg/ssh/ssh.go | 5 ++++ pkg/utils/file/file_v2.go | 13 --------- pkg/utils/iputils/iputils_v2.go | 3 ++ pkg/utils/yaml/yaml.go | 9 ++++++ .../github.com/labring/lvscare/care/care.go | 4 +-- 17 files changed, 105 insertions(+), 53 deletions(-) diff --git a/pkg/apply/applydrivers/apply_drivers_default.go b/pkg/apply/applydrivers/apply_drivers_default.go index f7131ee3a..ad743ce71 100644 --- a/pkg/apply/applydrivers/apply_drivers_default.go +++ b/pkg/apply/applydrivers/apply_drivers_default.go @@ -108,6 +108,8 @@ func (c *Applier) updateStatus(err error) { } func (c *Applier) reconcileCluster() error { + //sync newVersion pki and etc dir in `.sealos/default/pki` and `.sealos/default/etc` + processor.SyncNewVersionConfig(c.ClusterDesired) if err := c.installApp(c.RunNewImages); err != nil { return err } diff --git a/pkg/apply/processor/interface.go b/pkg/apply/processor/interface.go index 2110c674e..0a27f4d95 100644 --- a/pkg/apply/processor/interface.go +++ b/pkg/apply/processor/interface.go @@ -15,6 +15,9 @@ package processor import ( + "path" + + "github.com/labring/sealos/pkg/utils/file" "github.com/pkg/errors" "github.com/labring/sealos/pkg/constants" @@ -32,6 +35,20 @@ type Interface interface { Execute(cluster *v2.Cluster) error } +func SyncNewVersionConfig(cluster *v2.Cluster) { + d := constants.NewData(cluster.Name) + if !file.IsFile(d.PkiPath()) { + src, target := path.Join(d.Homedir(), constants.PkiDirName), d.PkiPath() + logger.Info("sync new version copy pki config: %s %s", src, target) + _ = file.RecursionCopy(src, target) + } + if !file.IsFile(d.EtcPath()) { + src, target := path.Join(d.Homedir(), constants.EtcDirName), d.EtcPath() + logger.Info("sync new version copy etc config: %s %s", src, target) + _ = file.RecursionCopy(src, target) + } +} + func SyncClusterStatus(cluster *v2.Cluster, service types.ClusterService, imgService types.ImageService, reset bool) error { if cluster.Status.Mounts == nil { containers, err := service.List() diff --git a/pkg/apply/processor/scale.go b/pkg/apply/processor/scale.go index 918d37dd7..2dfe216be 100644 --- a/pkg/apply/processor/scale.go +++ b/pkg/apply/processor/scale.go @@ -85,6 +85,7 @@ func (c *ScaleProcessor) GetPipeLine() ([]func(cluster *v2.Cluster) error, error ) return todoList, nil } + func (c *ScaleProcessor) Delete(cluster *v2.Cluster) error { logger.Info("Executing pipeline Delete in ScaleProcessor.") err := c.Runtime.DeleteMasters(c.MastersToDelete) diff --git a/pkg/constants/contants.go b/pkg/constants/contants.go index 0647e5874..69832e7c1 100644 --- a/pkg/constants/contants.go +++ b/pkg/constants/contants.go @@ -22,6 +22,9 @@ const ( LvsCareStaticPodName = "kube-sealos-lvscare" YamlFileSuffix = "yaml" DefaultRegistryDomain = "sealos.hub" + DefaultRegistryUsername = "admin" + DefaultRegistryPassword = "passw0rd" + DefaultRegistryData = "/var/lib/registry" DefaultLvscareDomain = "lvscare.node.ip" DefaultLvsCareImage = "sealos.hub:5000/sealos/lvscare:latest" ImageKubeVersionKey = "version" diff --git a/pkg/constants/data.go b/pkg/constants/data.go index 0706e140b..f3f271b5e 100644 --- a/pkg/constants/data.go +++ b/pkg/constants/data.go @@ -95,14 +95,14 @@ func (d *data) RootFSManifestsPath() string { } func (d *data) EtcPath() string { - return filepath.Join(d.Homedir(), EtcDirName) + return filepath.Join(ClusterDir(d.clusterName), EtcDirName) } func (d *data) AdminFile() string { return filepath.Join(d.EtcPath(), "admin.conf") } func (d *data) PkiPath() string { - return filepath.Join(d.Homedir(), PkiDirName) + return filepath.Join(ClusterDir(d.clusterName), PkiDirName) } func (d *data) PkiEtcdPath() string { diff --git a/pkg/filesystem/rootfs/rootfs_default.go b/pkg/filesystem/rootfs/rootfs_default.go index 1f27fb5cf..a8b60fe39 100644 --- a/pkg/filesystem/rootfs/rootfs_default.go +++ b/pkg/filesystem/rootfs/rootfs_default.go @@ -26,7 +26,7 @@ import ( "github.com/labring/sealos/pkg/constants" "github.com/labring/sealos/pkg/ssh" "github.com/labring/sealos/pkg/utils/exec" - file2 "github.com/labring/sealos/pkg/utils/file" + "github.com/labring/sealos/pkg/utils/file" "github.com/labring/sealos/pkg/utils/iputils" "github.com/labring/sealos/pkg/utils/logger" @@ -65,11 +65,10 @@ func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string, initFl target := constants.NewData(f.getClusterName(cluster)).RootFSPath() eg, _ := errgroup.WithContext(context.Background()) envProcessor := env.NewEnvProcessor(cluster, f.images) - for _, cInfo := range f.images { src := cInfo eg.Go(func() error { - if !file2.IsExist(src.MountPoint) { + if !file.IsExist(src.MountPoint) { logger.Debug("Image %s not exist,render env continue", src.ImageName) return nil } @@ -77,7 +76,7 @@ func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string, initFl if err != nil { return errors.Wrap(err, "render env to rootfs failed") } - dirs, err := file2.StatDir(src.MountPoint, true) + dirs, err := file.StatDir(src.MountPoint, true) if err != nil { return errors.Wrap(err, "get rootfs files failed") } @@ -95,6 +94,10 @@ func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string, initFl } check := constants.NewBash(f.getClusterName(cluster), cluster.GetImageLabels()) sshClient := f.getSSH(cluster) + shim := runtime.ImageShim{ + SSHInterface: sshClient, + IP: cluster.GetMaster0IPAndPort(), + } for _, IP := range ipList { ip := IP eg.Go(func() error { @@ -129,7 +132,7 @@ func (f *defaultRootfs) mountRootfs(cluster *v2.Cluster, ipList []string, initFl if checkBash == "" { return nil } - if err := f.getSSH(cluster).CmdAsync(ip, envProcessor.WrapperShell(ip, check.CheckBash()), runtime.ApplyImageShimCMD(target)); err != nil { + if err := f.getSSH(cluster).CmdAsync(ip, envProcessor.WrapperShell(ip, check.CheckBash()), shim.ApplyCMD(target)); err != nil { return err } if err := f.getSSH(cluster).CmdAsync(ip, envProcessor.WrapperShell(ip, check.InitBash())); err != nil { @@ -186,7 +189,7 @@ func renderENV(mountDir string, ipList []string, p env.Interface) error { for _, ip := range ipList { for _, dir := range []string{renderEtc, renderChart, renderManifests} { logger.Debug("render env dir: %s", dir) - if file2.IsExist(dir) { + if file.IsExist(dir) { err := p.RenderAll(ip, dir) if err != nil { return err @@ -201,19 +204,14 @@ func CopyFiles(sshEntry ssh.Interface, isRegistry, isApp bool, ip, src, target s if err != nil { return fmt.Errorf("failed to copy files %s", err) } - - if isRegistry { + if isRegistry || isApp { return sshEntry.Copy(ip, src, target) } - targetIP := ip - if isApp { - targetIP = "127.0.0.1" - } for _, f := range files { if f.Name() == constants.RegistryDirName { continue } - err = sshEntry.Copy(targetIP, filepath.Join(src, f.Name()), filepath.Join(target, f.Name())) + err = sshEntry.Copy(ip, filepath.Join(src, f.Name()), filepath.Join(target, f.Name())) if err != nil { return fmt.Errorf("failed to copy sub files %v", err) } diff --git a/pkg/guest/guest.go b/pkg/guest/guest.go index c4b4a2918..57b23387f 100644 --- a/pkg/guest/guest.go +++ b/pkg/guest/guest.go @@ -19,8 +19,9 @@ import ( "path/filepath" "strings" + "github.com/labring/sealos/pkg/ssh" + "github.com/labring/sealos/pkg/constants" - "github.com/labring/sealos/pkg/utils/exec" fileutil "github.com/labring/sealos/pkg/utils/file" "github.com/labring/sealos/pkg/utils/logger" "github.com/labring/sealos/pkg/utils/maps" @@ -75,13 +76,14 @@ func (d *Default) Apply(cluster *v2.Cluster, mounts []v2.MountImage) error { _ = fileutil.CleanFiles(kubeConfig) }() } + sshInterface := ssh.NewSSHClient(&cluster.Spec.SSH, true) for _, value := range guestCMD { if value == "" { continue } logger.Info("guest cmd is %s", value) - if err := exec.Cmd("bash", "-c", fmt.Sprintf(constants.CdAndExecCmd, clusterRootfs, value)); err != nil { + if err := sshInterface.CmdAsync(cluster.GetMaster0IPAndPort(), fmt.Sprintf(constants.CdAndExecCmd, clusterRootfs, value)); err != nil { return err } } diff --git a/pkg/remote/remote.go b/pkg/remote/remote.go index 5d50b98e9..a9af60596 100644 --- a/pkg/remote/remote.go +++ b/pkg/remote/remote.go @@ -82,10 +82,11 @@ func (s *remote) IPVS(ip, vip string, masters []string) error { } func (s *remote) IPVSClean(ip, vip string) error { var ipvsCommandTemplate = template.Must(template.New("ipvs").Parse(`` + - `ipvs --vs {{.vip}} -C`, + `ipvs --vs {{.vip}} --ip {{.ip}} -C`, )) data := map[string]interface{}{ "vip": vip, + "ip": ip, } out, err := renderTemplate(ipvsCommandTemplate, data) if err != nil { diff --git a/pkg/runtime/image_shim.go b/pkg/runtime/image_shim.go index ab98d0057..045892216 100644 --- a/pkg/runtime/image_shim.go +++ b/pkg/runtime/image_shim.go @@ -20,6 +20,8 @@ import ( "fmt" "path" + "github.com/labring/sealos/pkg/ssh" + "github.com/labring/sealos/pkg/constants" "github.com/labring/sealos/pkg/utils/logger" "github.com/labring/sealos/pkg/utils/yaml" @@ -29,25 +31,33 @@ import ( var defaultRootDirectory = "/var/lib/image-cri-shim" -//GetImageShim default dir is /var/lib/image-cri-shim -func GetImageShim(rootfs string) string { +type ImageShim struct { + SSHInterface ssh.Interface + IP string +} + +//GetInfo default dir is /var/lib/image-cri-shim +func (is *ImageShim) GetInfo(rootfs string) string { const imageCustomConfig = "image-cri-shim.yaml" + is.SSHInterface.SetStdout(false) + defer is.SSHInterface.SetStdout(true) etcPath := path.Join(rootfs, constants.EtcDirName, imageCustomConfig) - registryConfig, err := yaml.Unmarshal(etcPath) + data, _ := is.SSHInterface.Cmd(is.IP, fmt.Sprintf("cat %s", etcPath)) + shimConfig, err := yaml.UnmarshalData(data) if err != nil { - logger.Debug("use default registry config") + logger.Debug("use default image shim config") return defaultRootDirectory } - image, _, _ := unstructured.NestedString(registryConfig, "image") + image, _, _ := unstructured.NestedString(shimConfig, "image") logger.Debug("show image shim info, image dir : %s ", image) return image } -func ApplyImageShimCMD(rootfs string) string { - shimData := GetImageShim(rootfs) +func (is *ImageShim) ApplyCMD(rootfs string) string { + shimData := is.GetInfo(rootfs) return fmt.Sprintf(constants.DefaultCPFmt, shimData, path.Join(rootfs, constants.ImagesDirName, constants.ImageShimDirName), shimData) } -func DeleteImageShimCMD(rootfs string) string { - return fmt.Sprintf("rm -rf %s", GetImageShim(rootfs)) +func (is *ImageShim) DeleteCMD(rootfs string) string { + return fmt.Sprintf("rm -rf %s", is.GetInfo(rootfs)) } diff --git a/pkg/runtime/registry.go b/pkg/runtime/registry.go index 78c2397b0..3a4e936eb 100644 --- a/pkg/runtime/registry.go +++ b/pkg/runtime/registry.go @@ -30,15 +30,21 @@ import ( "github.com/labring/sealos/pkg/types/v1beta1" ) -func GetRegistry(rootfs, defaultRegistry string) *v1beta1.RegistryConfig { +func (k *KubeadmRuntime) GetRegistryInfo(rootfs, defaultRegistry string) *v1beta1.RegistryConfig { const registryCustomConfig = "registry.yml" var DefaultConfig = &v1beta1.RegistryConfig{ - IP: defaultRegistry, - Domain: constants.DefaultRegistryDomain, - Port: "5000", + IP: defaultRegistry, + Domain: constants.DefaultRegistryDomain, + Port: "5000", + Username: constants.DefaultRegistryUsername, + Password: constants.DefaultRegistryPassword, + Data: constants.DefaultRegistryData, } + k.getSSHInterface().SetStdout(false) + defer k.getSSHInterface().SetStdout(true) etcPath := path.Join(rootfs, constants.EtcDirName, registryCustomConfig) - registryConfig, err := yaml.Unmarshal(etcPath) + out, _ := k.getSSHInterface().Cmd(k.getMaster0IPAPIServer(), fmt.Sprintf("cat %s", etcPath)) + registryConfig, err := yaml.UnmarshalData(out) if err != nil { logger.Warn("read registry config path error: %+v", err) logger.Info("use default registry config") @@ -72,15 +78,15 @@ func GetRegistry(rootfs, defaultRegistry string) *v1beta1.RegistryConfig { return rConfig } -func (k *KubeadmRuntime) htpasswd() error { +func (k *KubeadmRuntime) htpasswd() (string, error) { htpasswdPath := path.Join(k.getContentData().RootFSEtcPath(), "registry_htpasswd") registry := k.getRegistry() if registry.Username == "" && registry.Password == "" { - return nil + return "", nil } data := passwd.Htpasswd(registry.Username, registry.Password) logger.Debug("write htpasswd file: %s,data: %s", htpasswdPath, data) - return file.WriteFile(htpasswdPath, []byte(data)) + return htpasswdPath, file.WriteFile(htpasswdPath, []byte(data)) } func (k *KubeadmRuntime) ApplyRegistry() error { @@ -91,10 +97,14 @@ func (k *KubeadmRuntime) ApplyRegistry() error { return fmt.Errorf("copy registry data failed %v", err) } ip := k.getMaster0IPAndPort() - err = k.htpasswd() + htpasswdPath, err := k.htpasswd() if err != nil { return fmt.Errorf("generator registry htpasswd failed %v", err) } + err = k.sshCopy(registry.IP, htpasswdPath, htpasswdPath) + if err != nil { + return fmt.Errorf("copy generator registry htpasswd failed %v", err) + } err = k.execInitRegistry(ip) if err != nil { return fmt.Errorf("exec registry.sh failed %v", err) diff --git a/pkg/runtime/reset.go b/pkg/runtime/reset.go index a82da81ac..c7676f40e 100644 --- a/pkg/runtime/reset.go +++ b/pkg/runtime/reset.go @@ -71,7 +71,11 @@ func (k *KubeadmRuntime) resetMasters(nodes []string) { func (k *KubeadmRuntime) resetNode(node string) error { logger.Info("start to reset node: %s", node) resetCmd := fmt.Sprintf(remoteCleanMasterOrNode, vlogToStr(k.vlog), k.getEtcdDataDir()) - deleteShimCmd := DeleteImageShimCMD(k.getContentData().RootFSPath()) + shim := &ImageShim{ + SSHInterface: nil, + IP: k.getMaster0IPAndPort(), + } + deleteShimCmd := shim.DeleteCMD(k.getContentData().RootFSPath()) if err := k.sshCmdAsync(node, resetCmd); err != nil { logger.Error("failed to clean node, exec command %s failed, %v", resetCmd, err) } diff --git a/pkg/runtime/runtime_getter.go b/pkg/runtime/runtime_getter.go index 930a4e7c1..d2d314b4f 100644 --- a/pkg/runtime/runtime_getter.go +++ b/pkg/runtime/runtime_getter.go @@ -34,7 +34,7 @@ import ( ) func (k *KubeadmRuntime) getRegistry() *v1beta1.RegistryConfig { - return GetRegistry(k.getContentData().RootFSPath(), k.getMaster0IPAndPort()) + return k.GetRegistryInfo(k.getContentData().RootFSPath(), k.getMaster0IPAndPort()) } func (k *KubeadmRuntime) getKubeVersion() string { diff --git a/pkg/ssh/ssh.go b/pkg/ssh/ssh.go index d5e64979c..014f2d0ce 100644 --- a/pkg/ssh/ssh.go +++ b/pkg/ssh/ssh.go @@ -37,6 +37,7 @@ type Interface interface { Cmd(host, cmd string) ([]byte, error) //CmdToString is exec command on remote host, and return spilt standard output and standard error CmdToString(host, cmd, spilt string) (string, error) + SetStdout(enable bool) Ping(host string) error } @@ -50,6 +51,10 @@ type SSH struct { LocalAddress *[]net.Addr } +func (s *SSH) SetStdout(enable bool) { + s.isStdout = enable +} + func NewSSHClient(ssh *v2.SSH, isStdout bool) Interface { if ssh.User == "" { ssh.User = v2.DefaultUserRoot diff --git a/pkg/utils/file/file_v2.go b/pkg/utils/file/file_v2.go index c27ad7701..e42723632 100644 --- a/pkg/utils/file/file_v2.go +++ b/pkg/utils/file/file_v2.go @@ -104,19 +104,6 @@ func ReadAll(fileName string) ([]byte, error) { return content, nil } -// file ./test/dir/xxx.txt if dir ./test/dir not exist, create it -func MkFileFullPathDir(fileName string) error { - localDir := filepath.Dir(fileName) - if err := Mkdir(localDir); err != nil { - return fmt.Errorf("failed to create local dir %s: %v", localDir, err) - } - return nil -} - -func Mkdir(dirName string) error { - return os.MkdirAll(dirName, 0755) -} - func MkDirs(dirs ...string) error { if len(dirs) == 0 { return nil diff --git a/pkg/utils/iputils/iputils_v2.go b/pkg/utils/iputils/iputils_v2.go index 6e73b6eba..19106b8bc 100644 --- a/pkg/utils/iputils/iputils_v2.go +++ b/pkg/utils/iputils/iputils_v2.go @@ -117,6 +117,9 @@ func IsLocalHostAddrs() (*[]net.Addr, error) { } func IsLocalIP(ip string, addrs *[]net.Addr) bool { + if defaultIP, _, err := net.SplitHostPort(ip); err == nil { + ip = defaultIP + } for _, address := range *addrs { if ipnet, ok := address.(*net.IPNet); ok && !ipnet.IP.IsLoopback() && ipnet.IP.To4() != nil && ipnet.IP.String() == ip { return true diff --git a/pkg/utils/yaml/yaml.go b/pkg/utils/yaml/yaml.go index f38626588..659859e63 100644 --- a/pkg/utils/yaml/yaml.go +++ b/pkg/utils/yaml/yaml.go @@ -43,6 +43,15 @@ func Unmarshal(path string) (map[string]interface{}, error) { return data, nil } +func UnmarshalData(metadata []byte) (map[string]interface{}, error) { + var data map[string]interface{} + err := yaml.Unmarshal(metadata, &data) + if err != nil { + return nil, err + } + return data, nil +} + func ToJSON(bs []byte) (jsons []string) { reader := bytes.NewReader(bs) ext := runtime.RawExtension{} diff --git a/staging/src/github.com/labring/lvscare/care/care.go b/staging/src/github.com/labring/lvscare/care/care.go index 854cc255e..6c1e2ab53 100644 --- a/staging/src/github.com/labring/lvscare/care/care.go +++ b/staging/src/github.com/labring/lvscare/care/care.go @@ -73,7 +73,7 @@ func (care *LvsCare) SyncRouter() error { if len(LVS.VirtualServer) == 0 { return errors.New("virtual server can't empty") } - if len(LVS.RealServer) > 0 { + if LVS.TargetIP != nil { var ipv4 bool vIP, _, err := net.SplitHostPort(LVS.VirtualServer) if err != nil { @@ -92,7 +92,7 @@ func (care *LvsCare) SyncRouter() error { LVS.Route = route.NewRoute(vIP, LVS.TargetIP.String()) return LVS.Route.SetRoute() } - return errors.New("real server can't empty") + return nil } func SetTargetIP() error {