summaryrefslogtreecommitdiff
path: root/keadm
diff options
context:
space:
mode:
authorKubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com>2024-01-17 17:50:27 +0800
committerGitHub <noreply@github.com>2024-01-17 17:50:27 +0800
commitdadbeb0e72565822afca48cb3f66a17473d7d585 (patch)
tree253eded79262e2f0d83a3da66ae945c9985eec0f /keadm
parentMerge pull request #5321 from luomengY/ci-cri (diff)
parentsupport edge nodes upgrade and image pre pull (diff)
downloadkubeedge-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.go2
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/rollback.go113
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/upgrade.go173
-rw-r--r--keadm/cmd/keadm/app/cmd/util/common.go63
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
+}