refactor(main): sealos run other server (#1292)

Signed-off-by: cuisongliu <cuisongliu@qq.com>
This commit is contained in:
cuisongliu
2022-07-12 16:00:33 +08:00
committed by GitHub
parent 074d3b3562
commit 2bd17b8cec
17 changed files with 105 additions and 53 deletions
@@ -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
}
+17
View File
@@ -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()
+1
View File
@@ -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)
+3
View File
@@ -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"
+2 -2
View File
@@ -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 {
+11 -13
View File
@@ -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)
}
+4 -2
View File
@@ -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
}
}
+2 -1
View File
@@ -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 {
+19 -9
View File
@@ -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))
}
+19 -9
View File
@@ -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)
+5 -1
View File
@@ -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)
}
+1 -1
View File
@@ -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 {
+5
View File
@@ -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
-13
View File
@@ -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
+3
View File
@@ -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
+9
View File
@@ -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{}
@@ -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 {