refactor(master): runtime fix code (#894)

* refactor(master): runtime fix code

* refactor(master): runtime fix code
This commit is contained in:
cuisongliu
2022-03-24 15:56:20 +08:00
committed by GitHub
parent 146fac739f
commit 9517a614b3
15 changed files with 157 additions and 146 deletions
+2 -2
View File
@@ -89,7 +89,7 @@ etcd:
listen-metrics-urls: http://0.0.0.0:2381
---
apiVersion: kubeproxy.config.k8s.io/v1alpha1
apiVersion: kubeproxy.Config.k8s.io/v1alpha1
kind: KubeProxyConfiguration
mode: "ipvs"
ipvs:
@@ -97,7 +97,7 @@ ipvs:
- "10.103.97.2/32"
---
apiVersion: kubelet.config.k8s.io/v1beta1
apiVersion: kubelet.Config.k8s.io/v1beta1
kind: KubeletConfiguration
authentication:
anonymous:
+4 -4
View File
@@ -53,20 +53,20 @@ func (k *KubeadmRuntime) BashInitOnMaster0() error {
}
func (k *KubeadmRuntime) ConfigInitKubeadmToMaster0() error {
logger.Info("start to copy kubeadm config to master0")
logger.Info("start to copy kubeadm Config to master0")
data, err := k.generateInitConfigs()
if err != nil {
return fmt.Errorf("generator config init kubeadm config error: %s", err.Error())
return fmt.Errorf("generator Config init kubeadm Config error: %s", err.Error())
}
initConfigPath := path.Join(k.data.TmpPath(), contants.DefaultInitKubeadmFileName)
outConfigPath := path.Join(k.data.EtcPath(), contants.DefaultInitKubeadmFileName)
err = file.WriteFile(initConfigPath, data)
if err != nil {
return fmt.Errorf("write config init kubeadm config error: %s", err.Error())
return fmt.Errorf("write Config init kubeadm Config error: %s", err.Error())
}
err = k.sshCopy(k.getMaster0IP(), initConfigPath, outConfigPath)
if err != nil {
return fmt.Errorf("copy config init kubeadm config error: %s", err.Error())
return fmt.Errorf("copy Config init kubeadm Config error: %s", err.Error())
}
return nil
}
+12 -12
View File
@@ -59,12 +59,12 @@ func (k *KubeadmRuntime) setAPIVersion(apiVersion string) {
//GetterKubeadmAPIVersion is covert version to kubeadmAPIServerVersion
// The support matrix will look something like this now and in the future:
// v1.10 and earlier: v1alpha1
// v1.11: v1alpha1 read-only, writes only v1alpha2 config
// v1.12: v1alpha2 read-only, writes only v1alpha3 config. Errors if the user tries to use v1alpha1
// v1.13: v1alpha3 read-only, writes only v1beta1 config. Errors if the user tries to use v1alpha1 or v1alpha2
// v1.14: v1alpha3 convert only, writes only v1beta1 config. Errors if the user tries to use v1alpha1 or v1alpha2
// v1.15: v1beta1 read-only, writes only v1beta2 config. Errors if the user tries to use v1alpha1, v1alpha2 or v1alpha3
// v1.22: v1beta2 read-only, writes only v1beta3 config. Errors if the user tries to use v1beta1 and older
// v1.11: v1alpha1 read-only, writes only v1alpha2 Config
// v1.12: v1alpha2 read-only, writes only v1alpha3 Config. Errors if the user tries to use v1alpha1
// v1.13: v1alpha3 read-only, writes only v1beta1 Config. Errors if the user tries to use v1alpha1 or v1alpha2
// v1.14: v1alpha3 convert only, writes only v1beta1 Config. Errors if the user tries to use v1alpha1 or v1alpha2
// v1.15: v1beta1 read-only, writes only v1beta2 Config. Errors if the user tries to use v1alpha1, v1alpha2 or v1alpha3
// v1.22: v1beta2 read-only, writes only v1beta3 Config. Errors if the user tries to use v1beta1 and older
func getterKubeadmAPIVersion(kubeVersion string) string {
var apiVersion string
switch {
@@ -84,7 +84,7 @@ func getterKubeadmAPIVersion(kubeVersion string) string {
}
func (k *KubeadmRuntime) getCGroupDriver(node string) (string, error) {
driver, err := k.ctlInterface.CGroup(node)
driver, err := k.getRemoteInterface().CGroup(node)
if err != nil {
return "", err
}
@@ -93,13 +93,13 @@ func (k *KubeadmRuntime) getCGroupDriver(node string) (string, error) {
}
func (k *KubeadmRuntime) MergeKubeadmConfig() error {
if k.config.ClusterFileKubeConfig != nil {
if err := k.LoadFromClusterfile(k.config.ClusterFileKubeConfig); err != nil {
return fmt.Errorf("failed to load kubeadm config from clusterfile: %v", err)
if k.Config.ClusterFileKubeConfig != nil {
if err := k.LoadFromClusterfile(k.Config.ClusterFileKubeConfig); err != nil {
return fmt.Errorf("failed to load kubeadm Config from clusterfile: %v", err)
}
}
if err := k.Merge(k.getDefaultKubeadmConfig()); err != nil {
return fmt.Errorf("failed to merge kubeadm config: %v", err)
return fmt.Errorf("failed to merge kubeadm Config: %v", err)
}
k.setKubeadmAPIVersion()
return nil
@@ -121,7 +121,7 @@ func (k *KubeadmRuntime) getVipAndPort() string {
return fmt.Sprintf("%s:6443", k.getVip())
}
func (k *KubeadmRuntime) getAPIServerDomain() string {
return k.config.apiServerDomain
return k.Config.apiServerDomain
}
func (k *KubeadmRuntime) getClusterAPIServer() string {
return fmt.Sprintf("https://%s:6443", k.getAPIServerDomain())
+10 -10
View File
@@ -35,9 +35,9 @@ import (
"k8s.io/kubelet/config/v1beta1"
)
// Read config from https://github.com/alibaba/sealer/blob/main/docs/design/clusterfile-v2.md and overwrite default kubeadm.yaml
// Use github.com/imdario/mergo to merge kubeadm config in Clusterfile and the default kubeadm config
// Using a config filter to handle some edge cases
// Read Config from https://github.com/alibaba/sealer/blob/main/docs/design/clusterfile-v2.md and overwrite default kubeadm.yaml
// Use github.com/imdario/mergo to merge kubeadm Config in Clusterfile and the default kubeadm Config
// Using a Config filter to handle some edge cases
// https://github.com/kubernetes/kubernetes/blob/master/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go
// Using map to overwrite Kubeadm configs
@@ -62,7 +62,7 @@ const (
// LoadFromClusterfile :Load KubeadmConfig from Clusterfile.
// If it has `KubeadmConfig` in Clusterfile, load every field to each configuration.
// If Kubeadm raw config in Clusterfile, just load it.
// If Kubeadm raw Config in Clusterfile, just load it.
func (k *KubeadmConfig) LoadFromClusterfile(kubeadmConfig *KubeadmConfig) error {
if kubeadmConfig == nil {
return nil
@@ -71,8 +71,8 @@ func (k *KubeadmConfig) LoadFromClusterfile(kubeadmConfig *KubeadmConfig) error
return mergo.Merge(k, kubeadmConfig)
}
// Merge Using github.com/imdario/mergo to merge KubeadmConfig to the CloudImage default kubeadm config, overwrite some field.
// if defaultKubeadmConfig file not exist, use default raw kubeadm config to merge k.KubeConfigSpec empty value
// Merge Using github.com/imdario/mergo to merge KubeadmConfig to the CloudImage default kubeadm Config, overwrite some field.
// if defaultKubeadmConfig file not exist, use default raw kubeadm Config to merge k.KubeConfigSpec empty value
func (k *KubeadmConfig) Merge(kubeadmYamlPath string) error {
var (
defaultKubeadmConfig *KubeadmConfig
@@ -87,12 +87,12 @@ func (k *KubeadmConfig) Merge(kubeadmYamlPath string) error {
}
defaultKubeadmConfig, err = LoadKubeadmConfigs(kubeadmYamlPath, DecodeCRDFromFile)
if err != nil {
return fmt.Errorf("failed to found kubeadm config from %s: %v", kubeadmYamlPath, err)
return fmt.Errorf("failed to found kubeadm Config from %s: %v", kubeadmYamlPath, err)
}
k.APIServer.CertSANs = append(k.APIServer.CertSANs, defaultKubeadmConfig.APIServer.CertSANs...)
err = mergo.Merge(k, defaultKubeadmConfig)
if err != nil {
return fmt.Errorf("failed to merge kubeadm config: %v", err)
return fmt.Errorf("failed to merge kubeadm Config: %v", err)
}
//using the DefaultKubeadmConfig configuration merge
return k.Merge("")
@@ -140,11 +140,11 @@ func NewKubeadmConfig() interface{} {
func DecodeCRDFromFile(filePath string, kind string) (interface{}, error) {
file, err := os.Open(filepath.Clean(filePath))
if err != nil {
return nil, fmt.Errorf("failed to dump config %v", err)
return nil, fmt.Errorf("failed to dump Config %v", err)
}
defer func() {
if err := file.Close(); err != nil {
logger.Warn("failed to dump config close clusterfile failed %v", err)
logger.Warn("failed to dump Config close clusterfile failed %v", err)
}
}()
return DecodeCRDFromReader(file, kind)
+2 -2
View File
@@ -18,11 +18,11 @@ package runtime
import "path"
const RemoteCopyKubeConfig = `rm -rf .kube/config && mkdir -p .kube && cp /etc/kubernetes/admin.conf .kube/config`
const RemoteCopyKubeConfig = `rm -rf .kube/Config && mkdir -p .kube && cp /etc/kubernetes/admin.conf .kube/Config`
func (k *KubeadmRuntime) copyNodeKubeConfig(hosts []string) error {
srcKubeFile := k.data.AdminFile()
desKubeFile := path.Join(".kube", "config")
desKubeFile := path.Join(".kube", "Config")
return k.sendFileToHosts(hosts, srcKubeFile, desKubeFile)
}
+5 -5
View File
@@ -74,20 +74,20 @@ func (k *KubeadmRuntime) sendJoinCPConfig(joinMaster []string) error {
}
func (k *KubeadmRuntime) ConfigJoinMasterKubeadmToMaster(master string) error {
logger.Info("start to copy kubeadm join config to master: %s", master)
logger.Info("start to copy kubeadm join Config to master: %s", master)
data, err := k.generateJoinMasterConfigs(master)
if err != nil {
return fmt.Errorf("generator config join master kubeadm config error: %s", err.Error())
return fmt.Errorf("generator Config join master kubeadm Config error: %s", err.Error())
}
joinConfigPath := path.Join(k.data.TmpPath(), contants.DefaultJoinMasterKubeadmFileName)
outConfigPath := path.Join(k.data.EtcPath(), contants.DefaultJoinMasterKubeadmFileName)
err = file.WriteFile(joinConfigPath, data)
if err != nil {
return fmt.Errorf("write config join master kubeadm config error: %s", err.Error())
return fmt.Errorf("write Config join master kubeadm Config error: %s", err.Error())
}
err = k.sshCopy(master, joinConfigPath, outConfigPath)
if err != nil {
return fmt.Errorf("copy config join master kubeadm config error: %s", err.Error())
return fmt.Errorf("copy Config join master kubeadm Config error: %s", err.Error())
}
return nil
}
@@ -102,7 +102,7 @@ func (k *KubeadmRuntime) joinMasters(masters []string) error {
return fmt.Errorf("filesystem init failed %v", err)
}
if err = ssh.WaitSSHReady(k.sshInterface, 6, masters...); err != nil {
if err = ssh.WaitSSHReady(k.getSSHInterface(), 6, masters...); err != nil {
return errors.Wrap(err, "join masters wait for ssh ready time out")
}
+6 -6
View File
@@ -35,7 +35,7 @@ func (k *KubeadmRuntime) joinNodes(newNodesIPList []string) error {
if err != nil {
return fmt.Errorf("filesystem init failed %v", err)
}
if err = ssh.WaitSSHReady(k.sshInterface, 6, newNodesIPList...); err != nil {
if err = ssh.WaitSSHReady(k.getSSHInterface(), 6, newNodesIPList...); err != nil {
return errors.Wrap(err, "join nodes wait for ssh ready time out")
}
@@ -47,7 +47,7 @@ func (k *KubeadmRuntime) joinNodes(newNodesIPList []string) error {
logger.Info("start to join %s as worker", node)
err = k.ConfigJoinNodeKubeadmToNode(node)
if err != nil {
return fmt.Errorf("failed to copy join node kubeadm config %s %v", node, err)
return fmt.Errorf("failed to copy join node kubeadm Config %s %v", node, err)
}
err = k.execHostsAppend(node, k.getVip(), k.getAPIServerDomain())
if err != nil {
@@ -83,20 +83,20 @@ func (k *KubeadmRuntime) joinNodes(newNodesIPList []string) error {
}
func (k *KubeadmRuntime) ConfigJoinNodeKubeadmToNode(node string) error {
logger.Info("start to copy kubeadm join config to node: %s", node)
logger.Info("start to copy kubeadm join Config to node: %s", node)
data, err := k.generateJoinNodeConfigs(node)
if err != nil {
return fmt.Errorf("generator config join kubeadm config error: %s", err.Error())
return fmt.Errorf("generator Config join kubeadm Config error: %s", err.Error())
}
joinConfigPath := path.Join(k.data.TmpPath(), contants.DefaultJoinNodeKubeadmFileName)
outConfigPath := path.Join(k.data.EtcPath(), contants.DefaultJoinNodeKubeadmFileName)
err = file.WriteFile(joinConfigPath, data)
if err != nil {
return fmt.Errorf("write config join kubeadm config error: %s", err.Error())
return fmt.Errorf("write Config join kubeadm Config error: %s", err.Error())
}
err = k.sshCopy(node, joinConfigPath, outConfigPath)
if err != nil {
return fmt.Errorf("copy config join kubeadm config error: %s", err.Error())
return fmt.Errorf("copy Config join kubeadm Config error: %s", err.Error())
}
return nil
}
+52 -4
View File
@@ -20,6 +20,10 @@ import (
"fmt"
"path"
"github.com/fanux/sealos/pkg/utils/contants"
"github.com/fanux/sealos/pkg/utils/yaml"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"github.com/fanux/sealos/pkg/passwd"
"github.com/fanux/sealos/pkg/utils/file"
@@ -28,18 +32,61 @@ import (
const DefaultCPFmt = "mkdir -p %s && cp -rf %s/* %s/"
func GetRegistry(rootfs, defaultRegistry string) *RegistryConfig {
const registryCustomConfig = "registry.yml"
var DefaultConfig = &RegistryConfig{
IP: defaultRegistry,
Domain: contants.DefaultRegistryDomain,
Port: "5000",
}
etcPath := path.Join(rootfs, contants.EtcDirName, registryCustomConfig)
registryConfig, err := yaml.Unmarshal(etcPath)
if err != nil {
logger.Debug("use default registry config")
return DefaultConfig
}
domain, _, _ := unstructured.NestedString(registryConfig, "domain")
port, _, _ := unstructured.NestedString(registryConfig, "port")
username, _, _ := unstructured.NestedString(registryConfig, "username")
password, _, _ := unstructured.NestedString(registryConfig, "password")
data, _, _ := unstructured.NestedString(registryConfig, "data")
ip, _, _ := unstructured.NestedString(registryConfig, "ip")
if ip == "" {
ip = defaultRegistry
}
if domain == "" {
domain = DefaultConfig.Domain
}
if port == "" {
domain = DefaultConfig.Port
}
rConfig := RegistryConfig{
IP: ip,
Domain: domain,
Port: port,
Username: username,
Password: password,
Data: data,
}
logger.Debug("show registry info, IP: %s, Domain: %s", rConfig.IP, rConfig.Domain)
return &rConfig
}
func (k *KubeadmRuntime) htpasswd() error {
htpasswdPath := path.Join(k.data.EtcPath(), "registry_htpasswd")
if k.registry.Username == "" && k.registry.Password == "" {
registry := k.getRegistry()
if registry.Username == "" && registry.Password == "" {
return nil
}
data := passwd.Htpasswd(k.registry.Username, k.registry.Password)
data := passwd.Htpasswd(registry.Username, registry.Password)
return file.WriteFile(htpasswdPath, []byte(data))
}
func (k *KubeadmRuntime) ApplyRegistry() error {
logger.Info("start to apply registry")
err := k.sshCmdAsync(k.registry.IP, fmt.Sprintf(DefaultCPFmt, k.registry.Data, k.data.RootFSRegistryPath(), k.registry.Data))
registry := k.getRegistry()
err := k.sshCmdAsync(registry.IP, fmt.Sprintf(DefaultCPFmt, registry.Data, k.data.RootFSRegistryPath(), registry.Data))
if err != nil {
return fmt.Errorf("copy registry data failed %v", err)
}
@@ -67,7 +114,8 @@ func (k *KubeadmRuntime) DeleteRegistry() error {
func (k *KubeadmRuntime) registryAuth(ip string) error {
logger.Info("registry auth in node %s", ip)
err := k.execHostsAppend(ip, k.registry.IP, k.registry.Domain)
registry := k.getRegistry()
err := k.execHostsAppend(ip, registry.IP, registry.Domain)
if err != nil {
return fmt.Errorf("add registry hosts failed %v", err)
}
+1 -1
View File
@@ -77,7 +77,7 @@ func (k *KubeadmRuntime) resetNode(node string) error {
if err != nil {
return fmt.Errorf("exec clean.sh failed %v", err)
}
err = k.execHostsDelete(node, k.registry.Domain)
err = k.execHostsDelete(node, k.getRegistry().Domain)
if err != nil {
return fmt.Errorf("delete registry hosts failed %v", err)
}
+8 -16
View File
@@ -23,26 +23,22 @@ import (
"github.com/fanux/sealos/pkg/utils/contants"
"github.com/fanux/sealos/pkg/env"
"github.com/fanux/sealos/pkg/remote"
v2 "github.com/fanux/sealos/pkg/types/v1beta1"
"github.com/fanux/sealos/pkg/utils/logger"
"github.com/fanux/sealos/pkg/utils/ssh"
)
type KubeadmRuntime struct {
*sync.Mutex
cluster *v2.Cluster
imageService image.Service
registry RegistryConfig
cluster *v2.Cluster
imageInfo *image.BuilderInfo
*KubeadmConfig
*config
*Config
*client
}
//nolint
type config struct {
// Clusterfile: the absolute path, we need to read kubeadm config from Clusterfile
type Config struct {
// Clusterfile: the absolute path, we need to read kubeadm Config from Clusterfile
ClusterFileKubeConfig *KubeadmConfig
apiServerDomain string
vlog int
@@ -50,11 +46,7 @@ type config struct {
//nolint
type client struct {
sshInterface ssh.Interface
envInterface env.Interface
ctlInterface remote.Interface
data contants.Data
bash contants.Bash
data contants.Data
}
type RegistryConfig struct {
@@ -138,11 +130,11 @@ func newKubeadmRuntime(clusterName string) (Interface, error) {
if err != nil {
return nil, err
}
k.imageService = imageService
if err = k.setData(clusterName); err != nil {
return nil, err
}
if err = k.setRegistry(); err != nil {
k.imageInfo, err = imageService.Inspect(k.cluster.Spec.Image)
if err != nil {
return nil, err
}
if err = k.setClient(); err != nil {
+44 -23
View File
@@ -20,17 +20,20 @@ import (
"context"
"fmt"
"github.com/fanux/sealos/pkg/env"
"github.com/fanux/sealos/pkg/remote"
"github.com/fanux/sealos/pkg/utils/contants"
"github.com/fanux/sealos/pkg/utils/logger"
"github.com/fanux/sealos/pkg/utils/ssh"
"golang.org/x/sync/errgroup"
)
func (k *KubeadmRuntime) getRegistry() *RegistryConfig {
return GetRegistry(k.data.RootFSPath(), k.getMaster0IP())
}
func (k *KubeadmRuntime) getKubeVersion() string {
labels, err := k.getImageLabels()
if err != nil {
logger.Painc("get kubernetes version error: %+v")
return ""
}
labels := k.getImageLabels()
image := labels["version"]
if image == "" {
logger.Painc("not fount kubernetes version")
@@ -64,10 +67,7 @@ func (k *KubeadmRuntime) getMaster0IPAPIServer() string {
}
func (k *KubeadmRuntime) getLvscareImage() (string, error) {
labels, err := k.getImageLabels()
if err != nil {
return "", err
}
labels := k.getImageLabels()
image := labels["image"]
if image == "" {
image = contants.DefaultLvsCareImage
@@ -76,7 +76,7 @@ func (k *KubeadmRuntime) getLvscareImage() (string, error) {
}
func (k *KubeadmRuntime) execIPVS(ip string, masters []string) error {
return k.ctlInterface.IPVS(ip, k.getVipAndPort(), masters)
return k.getRemoteInterface().IPVS(ip, k.getVipAndPort(), masters)
}
func (k *KubeadmRuntime) syncNodeIPVSYaml(masterIPs []string) error {
@@ -105,17 +105,17 @@ func (k *KubeadmRuntime) execIPVSPod(ip string, masters []string) error {
if err != nil {
return err
}
return k.ctlInterface.StaticPod(ip, k.getVipAndPort(), contants.LvsCareStaticPodName, image, masters)
return k.getRemoteInterface().StaticPod(ip, k.getVipAndPort(), contants.LvsCareStaticPodName, image, masters)
}
func (k *KubeadmRuntime) execToken(ip string) (string, error) {
return k.ctlInterface.Token(ip)
return k.getRemoteInterface().Token(ip)
}
func (k *KubeadmRuntime) execHostname(ip string) (string, error) {
return k.ctlInterface.Hostname(ip)
return k.getRemoteInterface().Hostname(ip)
}
func (k *KubeadmRuntime) execHostsAppend(ip, host, domain string) error {
return k.ctlInterface.HostsAdd(ip, host, domain)
return k.getRemoteInterface().HostsAdd(ip, host, domain)
}
func (k *KubeadmRuntime) execCert(ip string) error {
@@ -123,34 +123,55 @@ func (k *KubeadmRuntime) execCert(ip string) error {
if err != nil {
return err
}
return k.ctlInterface.Cert(ip, k.getCertSANS(), ip, hostname, k.getServiceCIDR(), k.getDNSDomain())
return k.getRemoteInterface().Cert(ip, k.getCertSANS(), ip, hostname, k.getServiceCIDR(), k.getDNSDomain())
}
func (k *KubeadmRuntime) execHostsDelete(ip, domain string) error {
return k.ctlInterface.HostsDelete(ip, domain)
return k.getRemoteInterface().HostsDelete(ip, domain)
}
func (k *KubeadmRuntime) execInit(ip string) error {
return k.sshInterface.CmdAsync(ip, k.envInterface.WrapperShell(ip, k.bash.InitBash()))
return k.getSSHInterface().CmdAsync(ip, k.getENVInterface().WrapperShell(ip, k.getScriptsBash().InitBash()))
}
func (k *KubeadmRuntime) execClean(ip string) error {
return k.sshInterface.CmdAsync(ip, k.envInterface.WrapperShell(ip, k.bash.CleanBash()))
return k.getSSHInterface().CmdAsync(ip, k.getENVInterface().WrapperShell(ip, k.getScriptsBash().CleanBash()))
}
func (k *KubeadmRuntime) execInitRegistry(ip string) error {
return k.sshInterface.CmdAsync(ip, k.envInterface.WrapperShell(ip, k.bash.InitRegistryBash()))
return k.getSSHInterface().CmdAsync(ip, k.getENVInterface().WrapperShell(ip, k.getScriptsBash().InitRegistryBash()))
}
func (k *KubeadmRuntime) execCleanRegistry(ip string) error {
return k.sshInterface.CmdAsync(ip, k.envInterface.WrapperShell(ip, k.bash.CleanRegistryBash()))
return k.getSSHInterface().CmdAsync(ip, k.getENVInterface().WrapperShell(ip, k.getScriptsBash().CleanRegistryBash()))
}
func (k *KubeadmRuntime) execAuth(ip string) error {
return k.sshInterface.CmdAsync(ip, k.envInterface.WrapperShell(ip, k.bash.AuthBash()))
return k.getSSHInterface().CmdAsync(ip, k.getENVInterface().WrapperShell(ip, k.getScriptsBash().AuthBash()))
}
func (k *KubeadmRuntime) sshCmdAsync(host string, cmd ...string) error {
return k.sshInterface.CmdAsync(host, cmd...)
return k.getSSHInterface().CmdAsync(host, cmd...)
}
func (k *KubeadmRuntime) sshCopy(host, srcFilePath, dstFilePath string) error {
return k.sshInterface.Copy(host, srcFilePath, dstFilePath)
return k.getSSHInterface().Copy(host, srcFilePath, dstFilePath)
}
func (k *KubeadmRuntime) getImageLabels() map[string]string {
return k.imageInfo.OCIv1.Config.Labels
}
func (k *KubeadmRuntime) getSSHInterface() ssh.Interface {
return ssh.NewSSHClient(&k.cluster.Spec.SSH, true)
}
func (k *KubeadmRuntime) getENVInterface() env.Interface {
return env.NewEnvProcessor(k.cluster)
}
func (k *KubeadmRuntime) getRemoteInterface() remote.Interface {
return remote.New(k.getClusterName(), k.getSSHInterface())
}
func (k *KubeadmRuntime) getScriptsBash() contants.Bash {
render := k.getImageLabels()
return contants.NewBash(k.getClusterName(), render)
}
+1 -52
View File
@@ -18,68 +18,17 @@ package runtime
import (
"fmt"
"path"
"github.com/fanux/sealos/pkg/env"
"github.com/fanux/sealos/pkg/remote"
"github.com/fanux/sealos/pkg/utils/contants"
"github.com/fanux/sealos/pkg/utils/decode"
fileutil "github.com/fanux/sealos/pkg/utils/file"
"github.com/fanux/sealos/pkg/utils/logger"
"github.com/fanux/sealos/pkg/utils/ssh"
"github.com/fanux/sealos/pkg/utils/yaml"
"github.com/pkg/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
)
func (k *KubeadmRuntime) setRegistry() error {
const registryCustomConfig = "registry.yml"
etcPath := path.Join(k.data.RootFSPath(), contants.EtcDirName, registryCustomConfig)
registryConfig, err := yaml.Unmarshal(etcPath)
if err != nil {
return err
}
domain, _, _ := unstructured.NestedString(registryConfig, "domain")
port, _, _ := unstructured.NestedFloat64(registryConfig, "port")
username, _, _ := unstructured.NestedString(registryConfig, "username")
password, _, _ := unstructured.NestedString(registryConfig, "password")
data, _, _ := unstructured.NestedString(registryConfig, "data")
rConfig := RegistryConfig{
IP: k.getMaster0IP(),
Domain: domain,
Port: fmt.Sprintf("%d", int(port)),
Username: username,
Password: password,
Data: data,
}
k.registry = rConfig
return nil
}
func (k *KubeadmRuntime) getImageLabels() (map[string]string, error) {
data, err := k.imageService.Inspect(k.getImageName())
if err != nil {
return nil, err
}
return data.OCIv1.Config.Labels, err
}
func (k *KubeadmRuntime) getImageName() string {
return k.cluster.Spec.Image
}
func (k *KubeadmRuntime) setClient() error {
sshInterface := ssh.NewSSHClient(&k.cluster.Spec.SSH, true)
k.client = &client{}
k.sshInterface = sshInterface
k.envInterface = env.NewEnvProcessor(k.cluster)
k.ctlInterface = remote.New(k.getClusterName(), sshInterface)
k.data = contants.NewData(k.getClusterName())
render, err := k.getImageLabels()
if err != nil {
return err
}
k.bash = contants.NewBash(k.getClusterName(), render)
return nil
}
@@ -102,7 +51,7 @@ func (k *KubeadmRuntime) setData(clusterName string) error {
}
k.cluster = &clusters[0]
k.KubeadmConfig = &KubeadmConfig{}
k.config = &config{
k.Config = &Config{
ClusterFileKubeConfig: kubeadmConfig,
apiServerDomain: DefaultAPIServerDomain,
}
+4 -4
View File
@@ -40,13 +40,13 @@ func vlogToStr(vlog int) string {
func (k *KubeadmRuntime) Command(version string, name CommandType) (cmd string) {
const (
InitMaster115Lower = `kubeadm init --config=%s --experimental-upload-certs`
InitMaster115Lower = `kubeadm init --Config=%s --experimental-upload-certs`
JoinMaster115Lower = "kubeadm join %s:6443 --token %s %s --experimental-control-plane --certificate-key %s"
JoinNode115Lower = "kubeadm join %s:6443 --token %s %s"
InitMaser115Upper = `kubeadm init --config=%s --upload-certs`
JoinMaster115Upper = "kubeadm join --config=%s"
JoinNode115Upper = "kubeadm join --config=%s"
InitMaser115Upper = `kubeadm init --Config=%s --upload-certs`
JoinMaster115Upper = "kubeadm join --Config=%s"
JoinNode115Upper = "kubeadm join --Config=%s"
)
initConfigPath := path.Join(k.data.EtcPath(), contants.DefaultInitKubeadmFileName)
+2 -2
View File
@@ -68,7 +68,7 @@ func (k *KubeadmRuntime) SendJoinMasterKubeConfigs(masters []string, files ...st
}
}
if k.ReplaceKubeConfigV1991V1992(masters) {
logger.Info("set kubernetes v1.19.1 v1.19.2 kube config")
logger.Info("set kubernetes v1.19.1 v1.19.2 kube Config")
}
return nil
}
@@ -83,7 +83,7 @@ func (k *KubeadmRuntime) ReplaceKubeConfigV1991V1992(masters []string) bool {
for _, v := range masters {
replaceCmd := fmt.Sprintf(RemoteReplaceKubeConfig, KUBESCHEDULERCONFIGFILE, v, KUBECONTROLLERCONFIGFILE, v, KUBESCHEDULERCONFIGFILE)
if err := k.sshCmdAsync(v, replaceCmd); err != nil {
logger.Info("failed to replace kube config on %s:%v ", v, err)
logger.Info("failed to replace kube Config on %s:%v ", v, err)
return false
}
}
+4 -3
View File
@@ -19,9 +19,10 @@ import (
)
const (
LvsCareStaticPodName = "kube-sealyun-lvscare"
YamlFileSuffix = "yaml"
DefaultLvsCareImage = "sealyun.hub:5000/sealyun/lvscare:latest"
LvsCareStaticPodName = "kube-sealyun-lvscare"
YamlFileSuffix = "yaml"
DefaultRegistryDomain = "selayun.hub"
DefaultLvsCareImage = "sealyun.hub:5000/sealyun/lvscare:latest"
)
//CRD kind