diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2024-01-17 17:50:27 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2024-01-17 17:50:27 +0800 |
| commit | dadbeb0e72565822afca48cb3f66a17473d7d585 (patch) | |
| tree | 253eded79262e2f0d83a3da66ae945c9985eec0f /keadm | |
| parent | Merge pull request #5321 from luomengY/ci-cri (diff) | |
| parent | support edge nodes upgrade and image pre pull (diff) | |
| download | kubeedge-dadbeb0e72565822afca48cb3f66a17473d7d585.tar.gz | |
Merge pull request #5330 from ZhengXinwei-F/task-manager
Implement task manager to complete cloud edge task execution
Diffstat (limited to 'keadm')
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/cmd_others.go | 2 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/rollback.go | 113 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/upgrade.go | 173 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/util/common.go | 63 |
4 files changed, 243 insertions, 108 deletions
diff --git a/keadm/cmd/keadm/app/cmd/cmd_others.go b/keadm/cmd/keadm/app/cmd/cmd_others.go index 71c699d52..56c394608 100644 --- a/keadm/cmd/keadm/app/cmd/cmd_others.go +++ b/keadm/cmd/keadm/app/cmd/cmd_others.go @@ -89,5 +89,7 @@ func NewKubeedgeCommand() *cobra.Command { cmds.AddCommand(beta.NewBeta()) cmds.AddCommand(edge.NewEdgeUpgrade()) + cmds.AddCommand(edge.NewEdgeRollback()) + return cmds } diff --git a/keadm/cmd/keadm/app/cmd/edge/rollback.go b/keadm/cmd/keadm/app/cmd/edge/rollback.go new file mode 100644 index 000000000..cacd279a3 --- /dev/null +++ b/keadm/cmd/keadm/app/cmd/edge/rollback.go @@ -0,0 +1,113 @@ +/* +Copyright 2023 The KubeEdge Authors. + +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 edge + +import ( + "fmt" + "os" + + "github.com/spf13/cobra" + "k8s.io/klog/v2" + "sigs.k8s.io/yaml" + + "github.com/kubeedge/kubeedge/common/constants" + "github.com/kubeedge/kubeedge/keadm/cmd/keadm/app/cmd/util" + "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha2" + api "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/util/fsm" + "github.com/kubeedge/kubeedge/pkg/version" +) + +// NewEdgeUpgrade returns KubeEdge edge upgrade command. +func NewEdgeRollback() *cobra.Command { + rollbackOptions := newRollbackOptions() + + cmd := &cobra.Command{ + Use: "rollback", + Short: "rollback edge component. Rollback the edge node to the desired version.", + Long: "Rollback edge component. Rollback the edge node to the desired version.", + RunE: func(cmd *cobra.Command, args []string) error { + // rollback edge core + return rollbackEdgeCore(rollbackOptions) + }, + } + + addRollbackFlags(cmd, rollbackOptions) + return cmd +} + +// newJoinOptions returns a struct ready for being used for creating cmd join flags. +func newRollbackOptions() *RollbackOptions { + opts := &RollbackOptions{} + opts.HistoryVersion = version.Get().String() + opts.Config = constants.DefaultConfigDir + "edgecore.yaml" + + return opts +} + +func rollbackEdgeCore(ro *RollbackOptions) error { + // get EdgeCore configuration from edgecore.yaml config file + data, err := os.ReadFile(ro.Config) + if err != nil { + return fmt.Errorf("failed to read config file %s: %v", ro.Config, err) + } + + configure := &v1alpha2.EdgeCoreConfig{} + err = yaml.Unmarshal(data, configure) + if err != nil { + return fmt.Errorf("failed to unmarshal config file %s: %v", ro.Config, err) + } + event := &fsm.Event{ + Type: "RollBack", + Action: api.ActionSuccess, + } + defer func() { + // report upgrade result to cloudhub + if err = util.ReportTaskResult(configure, ro.TaskType, ro.TaskName, *event); err != nil { + klog.Warning("failed to report upgrade result to cloud: %v", err) + } + }() + + rbErr := rollback(ro.HistoryVersion, configure.DataBase.DataSource, ro.Config) + if rbErr != nil { + event.Action = api.ActionFailure + event.Msg = fmt.Sprintf("upgrade error: %v, rollback error: %v", err, rbErr) + } + + return nil +} + +type RollbackOptions struct { + HistoryVersion string + TaskType string + TaskName string + Config string +} + +func addRollbackFlags(cmd *cobra.Command, rollbackOptions *RollbackOptions) { + cmd.Flags().StringVar(&rollbackOptions.HistoryVersion, "history", rollbackOptions.HistoryVersion, + "Use this key to specify the origin version before upgrade") + + cmd.Flags().StringVar(&rollbackOptions.Config, "config", rollbackOptions.Config, + "Use this key to specify the path to the edgecore configuration file.") + + cmd.Flags().StringVar(&rollbackOptions.TaskType, "type", "rollback", + "Use this key to specify the task type for reporting status.") + + cmd.Flags().StringVar(&rollbackOptions.TaskName, "name", rollbackOptions.TaskName, + "Use this key to specify the task name for reporting status.") +} diff --git a/keadm/cmd/keadm/app/cmd/edge/upgrade.go b/keadm/cmd/keadm/app/cmd/edge/upgrade.go index 001323679..2d5f934d8 100644 --- a/keadm/cmd/keadm/app/cmd/edge/upgrade.go +++ b/keadm/cmd/keadm/app/cmd/edge/upgrade.go @@ -17,28 +17,21 @@ limitations under the License. package edge import ( - "bytes" - "crypto/tls" - "crypto/x509" - "encoding/json" "fmt" "io" - "net" - "net/http" "os" "path/filepath" - "time" "github.com/spf13/cobra" "k8s.io/klog/v2" "sigs.k8s.io/yaml" "github.com/kubeedge/kubeedge/common/constants" - commontypes "github.com/kubeedge/kubeedge/common/types" "github.com/kubeedge/kubeedge/keadm/cmd/keadm/app/cmd/common" "github.com/kubeedge/kubeedge/keadm/cmd/keadm/app/cmd/util" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha2" - upgradev1alpha1 "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1" + api "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1" + "github.com/kubeedge/kubeedge/pkg/util/fsm" ) var ( @@ -93,91 +86,96 @@ func (up *UpgradeOptions) upgrade() error { HistoryID: up.HistoryID, FromVersion: up.FromVersion, ToVersion: up.ToVersion, + TaskType: up.TaskType, Image: up.Image, + DisableBackup: up.DisableBackup, ConfigFilePath: up.Config, EdgeCoreConfig: configure, } + event := &fsm.Event{ + Type: "Upgrade", + Action: api.ActionSuccess, + } defer func() { // report upgrade result to cloudhub - if err := upgrade.reportUpgradeResult(); err != nil { + if err = util.ReportTaskResult(configure, upgrade.TaskType, upgrade.UpgradeID, *event); err != nil { klog.Errorf("failed to report upgrade result to cloud: %v", err) } // cleanup idempotency record - if err := os.Remove(idempotencyRecord); err != nil { + if err = os.Remove(idempotencyRecord); err != nil { klog.Errorf("failed to remove idempotency_record file(%s): %v", idempotencyRecord, err) } }() // only allow upgrade when last upgrade finished if util.FileExists(idempotencyRecord) { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackSuccess)) - upgrade.UpdateFailureReason("last upgrade not finished, not allowed upgrade again") + event.Action = api.ActionFailure + event.Msg = "last upgrade not finished, not allowed upgrade again" return fmt.Errorf("last upgrade not finished, not allowed upgrade again") } // create idempotency_record file if err := os.MkdirAll(filepath.Dir(idempotencyRecord), 0750); err != nil { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackSuccess)) reason := fmt.Sprintf("failed to create idempotency_record dir: %v", err) - upgrade.UpdateFailureReason(reason) + event.Action = api.ActionFailure + event.Msg = reason return fmt.Errorf(reason) } if _, err := os.Create(idempotencyRecord); err != nil { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackSuccess)) reason := fmt.Sprintf("failed to create idempotency_record file: %v", err) - upgrade.UpdateFailureReason(reason) + event.Action = api.ActionFailure + event.Msg = reason return fmt.Errorf(reason) } // run script to do upgrade operation err = upgrade.PreProcess() if err != nil { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackSuccess)) - upgrade.UpdateFailureReason(fmt.Sprintf("upgrade error: %v", err)) + event.Action = api.ActionFailure + event.Msg = fmt.Sprintf("upgrade pre process failed: %v", err) return fmt.Errorf("upgrade pre process failed: %v", err) } err = upgrade.Process() if err != nil { + event.Type = "Rollback" rbErr := upgrade.Rollback() if rbErr != nil { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackFailed)) - upgrade.UpdateFailureReason(fmt.Sprintf("upgrade error: %v, rollback error: %v", err, rbErr)) + event.Action = api.ActionFailure + event.Msg = rbErr.Error() } else { - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeFailedRollbackSuccess)) - upgrade.UpdateFailureReason(fmt.Sprintf("upgrade error: %v", err)) + event.Msg = err.Error() } return fmt.Errorf("upgrade process failed: %v", err) } - upgrade.UpdateStatus(string(upgradev1alpha1.UpgradeSuccess)) - return nil } func (up *Upgrade) PreProcess() error { - klog.Infof("upgrade preprocess start") - backupPath := filepath.Join(util.KubeEdgeBackupPath, up.FromVersion) - if err := os.MkdirAll(backupPath, 0750); err != nil { - return fmt.Errorf("mkdirall failed: %v", err) - } + // download the request version edgecore + klog.Infof("Begin to download version %s edgecore", up.ToVersion) + if !up.DisableBackup { + backupPath := filepath.Join(util.KubeEdgeBackupPath, up.FromVersion) + if err := os.MkdirAll(backupPath, 0750); err != nil { + return fmt.Errorf("mkdirall failed: %v", err) + } - // backup edgecore.db: copy from origin path to backup path - if err := copy(up.EdgeCoreConfig.DataBase.DataSource, filepath.Join(backupPath, "edgecore.db")); err != nil { - return fmt.Errorf("failed to backup db: %v", err) - } - // backup edgecore.yaml: copy from origin path to backup path - if err := copy(up.ConfigFilePath, filepath.Join(backupPath, "edgecore.yaml")); err != nil { - return fmt.Errorf("failed to back config: %v", err) - } - // backup edgecore: copy from origin path to backup path - if err := copy(filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), filepath.Join(backupPath, util.KubeEdgeBinaryName)); err != nil { - return fmt.Errorf("failed to backup edgecore: %v", err) + // backup edgecore.db: copy from origin path to backup path + if err := copy(up.EdgeCoreConfig.DataBase.DataSource, filepath.Join(backupPath, "edgecore.db")); err != nil { + return fmt.Errorf("failed to backup db: %v", err) + } + // backup edgecore.yaml: copy from origin path to backup path + if err := copy(up.ConfigFilePath, filepath.Join(backupPath, "edgecore.yaml")); err != nil { + return fmt.Errorf("failed to back config: %v", err) + } + // backup edgecore: copy from origin path to backup path + if err := copy(filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), filepath.Join(backupPath, util.KubeEdgeBinaryName)); err != nil { + return fmt.Errorf("failed to backup edgecore: %v", err) + } } - // download the request version edgecore - klog.Infof("Begin to download version %s edgecore", up.ToVersion) upgradePath := filepath.Join(util.KubeEdgeUpgradePath, up.ToVersion) if up.EdgeCoreConfig.Modules.Edged.TailoredKubeletConfig.ContainerRuntimeEndpoint == "" { up.EdgeCoreConfig.Modules.Edged.TailoredKubeletConfig.ContainerRuntimeEndpoint = up.EdgeCoreConfig.Modules.Edged.RemoteRuntimeEndpoint @@ -266,6 +264,10 @@ func (up *Upgrade) Process() error { } func (up *Upgrade) Rollback() error { + return rollback(up.FromVersion, up.EdgeCoreConfig.DataBase.DataSource, up.ConfigFilePath) +} + +func rollback(HistoryVersion, dataSource, configFilePath string) error { klog.Infof("upgrade rollback process start") // stop edgecore @@ -277,12 +279,12 @@ func (up *Upgrade) Rollback() error { // rollback origin config/db/binary // backup edgecore.db: copy from backup path to origin path - backupPath := filepath.Join(util.KubeEdgeBackupPath, up.FromVersion) - if err := copy(filepath.Join(backupPath, "edgecore.db"), up.EdgeCoreConfig.DataBase.DataSource); err != nil { + backupPath := filepath.Join(util.KubeEdgeBackupPath, HistoryVersion) + if err := copy(filepath.Join(backupPath, "edgecore.db"), dataSource); err != nil { return fmt.Errorf("failed to rollback db: %v", err) } // backup edgecore.yaml: copy from backup path to origin path - if err := copy(filepath.Join(backupPath, "edgecore.yaml"), up.ConfigFilePath); err != nil { + if err := copy(filepath.Join(backupPath, "edgecore.yaml"), configFilePath); err != nil { return fmt.Errorf("failed to back config: %v", err) } // backup edgecore: copy from backup path to origin path @@ -292,7 +294,7 @@ func (up *Upgrade) Rollback() error { // generate edgecore.service if util.HasSystemd() { - err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), up.ConfigFilePath), false) + err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), configFilePath), false) if err != nil { return fmt.Errorf("failed to create edgecore.service file: %v", err) } @@ -303,7 +305,6 @@ func (up *Upgrade) Rollback() error { if err != nil { return fmt.Errorf("failed to start origin edgecore: %v", err) } - return nil } @@ -315,67 +316,15 @@ func (up *Upgrade) UpdateFailureReason(reason string) { up.Reason = reason } -func (up *Upgrade) reportUpgradeResult() error { - resp := &commontypes.NodeUpgradeJobResponse{ - UpgradeID: up.UpgradeID, - HistoryID: up.HistoryID, - NodeName: up.EdgeCoreConfig.Modules.Edged.HostnameOverride, - FromVersion: up.FromVersion, - ToVersion: up.ToVersion, - Status: up.Status, - Reason: up.Reason, - } - - var caCrt []byte - caCertPath := up.EdgeCoreConfig.Modules.EdgeHub.TLSCAFile - caCrt, err := os.ReadFile(caCertPath) - if err != nil { - return fmt.Errorf("failed to read ca: %v", err) - } - - rootCAs := x509.NewCertPool() - rootCAs.AppendCertsFromPEM(caCrt) - - certFile := up.EdgeCoreConfig.Modules.EdgeHub.TLSCertFile - keyFile := up.EdgeCoreConfig.Modules.EdgeHub.TLSPrivateKeyFile - cliCrt, err := tls.LoadX509KeyPair(certFile, keyFile) - - transport := &http.Transport{ - DialContext: (&net.Dialer{ - Timeout: 30 * time.Second, - KeepAlive: 30 * time.Second, - }).DialContext, - // use TLS configuration - TLSClientConfig: &tls.Config{ - RootCAs: rootCAs, - InsecureSkipVerify: false, - Certificates: []tls.Certificate{cliCrt}, - }, - } - - client := &http.Client{Transport: transport, Timeout: 30 * time.Second} - - respData, err := json.Marshal(resp) - if err != nil { - return fmt.Errorf("marshal failed: %v", err) - } - url := up.EdgeCoreConfig.Modules.EdgeHub.HTTPServer + constants.DefaultNodeUpgradeURL - result, err := client.Post(url, "application/json", bytes.NewReader(respData)) - if err != nil { - return fmt.Errorf("post http request failed: %v", err) - } - defer result.Body.Close() - - return nil -} - type UpgradeOptions struct { - UpgradeID string - HistoryID string - FromVersion string - ToVersion string - Config string - Image string + UpgradeID string + HistoryID string + FromVersion string + ToVersion string + Config string + Image string + DisableBackup bool + TaskType string } type Upgrade struct { @@ -384,7 +333,9 @@ type Upgrade struct { FromVersion string ToVersion string Image string + DisableBackup bool ConfigFilePath string + TaskType string EdgeCoreConfig *v1alpha2.EdgeCoreConfig Status string @@ -409,4 +360,10 @@ func addUpgradeFlags(cmd *cobra.Command, upgradeOptions *UpgradeOptions) { cmd.Flags().StringVar(&upgradeOptions.Image, "image", upgradeOptions.Image, "Use this key to specify installation image to download.") + + cmd.Flags().StringVar(&upgradeOptions.TaskType, "type", "upgrade", + "Use this key to specify the task type for reporting status.") + + cmd.Flags().BoolVar(&upgradeOptions.DisableBackup, "disable-backup", upgradeOptions.DisableBackup, + "Use this key to specify the backup enable for upgrade.") } diff --git a/keadm/cmd/keadm/app/cmd/util/common.go b/keadm/cmd/keadm/app/cmd/util/common.go index 4cbc2c243..a5065484e 100644 --- a/keadm/cmd/keadm/app/cmd/util/common.go +++ b/keadm/cmd/keadm/app/cmd/util/common.go @@ -18,16 +18,22 @@ package util import ( "archive/tar" + "bytes" "compress/gzip" "crypto/sha512" + "crypto/tls" + "crypto/x509" + "encoding/json" "fmt" "io" + "net" "net/http" "os" "path/filepath" "regexp" "strconv" "strings" + "time" "github.com/blang/semver" "github.com/spf13/pflag" @@ -39,8 +45,11 @@ import ( "k8s.io/klog/v2" "github.com/kubeedge/kubeedge/common/constants" + commontypes "github.com/kubeedge/kubeedge/common/types" types "github.com/kubeedge/kubeedge/keadm/cmd/keadm/app/cmd/common" + "github.com/kubeedge/kubeedge/pkg/apis" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha2" + "github.com/kubeedge/kubeedge/pkg/util/fsm" pkgversion "github.com/kubeedge/kubeedge/pkg/version" ) @@ -639,3 +648,57 @@ func downloadServiceFile(componentType types.ComponentType, version semver.Versi } return nil } + +func ReportTaskResult(config *v1alpha2.EdgeCoreConfig, taskType, taskID string, event fsm.Event) error { + resp := &commontypes.NodeTaskResponse{ + NodeName: config.Modules.Edged.HostnameOverride, + Event: event.Type, + Action: event.Action, + Time: time.Now().Format(apis.ISO8601UTC), + Reason: event.Msg, + } + edgeHub := config.Modules.EdgeHub + var caCrt []byte + caCertPath := edgeHub.TLSCAFile + caCrt, err := os.ReadFile(caCertPath) + if err != nil { + return fmt.Errorf("failed to read ca: %v", err) + } + + rootCAs := x509.NewCertPool() + rootCAs.AppendCertsFromPEM(caCrt) + + certFile := edgeHub.TLSCertFile + keyFile := edgeHub.TLSPrivateKeyFile + cliCrt, err := tls.LoadX509KeyPair(certFile, keyFile) + + transport := &http.Transport{ + DialContext: (&net.Dialer{ + Timeout: 30 * time.Second, + KeepAlive: 30 * time.Second, + }).DialContext, + // use TLS configuration + TLSClientConfig: &tls.Config{ + RootCAs: rootCAs, + InsecureSkipVerify: false, + Certificates: []tls.Certificate{cliCrt}, + }, + } + + client := &http.Client{Transport: transport, Timeout: 30 * time.Second} + + respData, err := json.Marshal(resp) + if err != nil { + return fmt.Errorf("marshal failed: %v", err) + } + url := edgeHub.HTTPServer + fmt.Sprintf("/task/%s/name/%s/node/%s/status", taskType, taskID, config.Modules.Edged.HostnameOverride) + result, err := client.Post(url, "application/json", bytes.NewReader(respData)) + + if err != nil { + return fmt.Errorf("post http request failed: %v", err) + } + klog.Error("report result ", result) + defer result.Body.Close() + + return nil +} |
