summaryrefslogtreecommitdiff
path: root/cloud
diff options
context:
space:
mode:
authorzhengxinwei <zhengxinwei@huawei.com>2024-01-04 15:27:04 +0800
committerzhengxinwei-f <zhengxinwei@huawei.com>2024-01-16 09:35:44 +0800
commitb1c1c6a8e5442fff2b4d932060648e009f639df1 (patch)
tree575f35f9559894e995060f15b029ff73f2d23282 /cloud
parentnotify the task manager of the task status (diff)
downloadkubeedge-b1c1c6a8e5442fff2b4d932060648e009f639df1.tar.gz
Implement task manager to complete cloud edge task execution
Signed-off-by: zhengxinwei-f <zhengxinwei@huawei.com>
Diffstat (limited to 'cloud')
-rw-r--r--cloud/cmd/cloudcore/app/server.go4
-rw-r--r--cloud/pkg/cloudhub/dispatcher/message_dispatcher.go4
-rw-r--r--cloud/pkg/cloudhub/servers/httpserver/report_task_status.go28
-rw-r--r--cloud/pkg/common/messagelayer/context.go8
-rw-r--r--cloud/pkg/common/modules/modules.go3
-rw-r--r--cloud/pkg/taskmanager/config/config.go38
-rw-r--r--cloud/pkg/taskmanager/manager/downstream.go64
-rw-r--r--cloud/pkg/taskmanager/manager/executor.go415
-rw-r--r--cloud/pkg/taskmanager/manager/upstream.go144
-rw-r--r--cloud/pkg/taskmanager/task_manager.go144
-rw-r--r--cloud/pkg/taskmanager/util/controller/controller.go153
-rw-r--r--cloud/pkg/taskmanager/util/manager/common.go62
-rw-r--r--cloud/pkg/taskmanager/util/manager/task_cache.go52
-rw-r--r--cloud/pkg/taskmanager/util/util.go169
-rw-r--r--cloud/pkg/taskmanager/util/util_test.go171
15 files changed, 1446 insertions, 13 deletions
diff --git a/cloud/cmd/cloudcore/app/server.go b/cloud/cmd/cloudcore/app/server.go
index fd346cd72..3f4949729 100644
--- a/cloud/cmd/cloudcore/app/server.go
+++ b/cloud/cmd/cloudcore/app/server.go
@@ -49,10 +49,10 @@ import (
"github.com/kubeedge/kubeedge/cloud/pkg/dynamiccontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller"
- "github.com/kubeedge/kubeedge/cloud/pkg/nodeupgradejobcontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/policycontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/router"
"github.com/kubeedge/kubeedge/cloud/pkg/synccontroller"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager"
"github.com/kubeedge/kubeedge/common/constants"
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1/validation"
@@ -160,7 +160,7 @@ func registerModules(c *v1alpha1.CloudCoreConfig) {
edgecontroller.Register(c.Modules.EdgeController)
devicecontroller.Register(c.Modules.DeviceController)
imageprepullcontroller.Register(c.Modules.ImagePrePullController)
- nodeupgradejobcontroller.Register(c.Modules.NodeUpgradeJobController)
+ taskmanager.Register(c.Modules.TaskManager)
synccontroller.Register(c.Modules.SyncController)
cloudstream.Register(c.Modules.CloudStream, c.CommonConfig)
router.Register(c.Modules.Router)
diff --git a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
index a236d3ebd..940acb3a5 100644
--- a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
+++ b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
@@ -453,6 +453,10 @@ func noAckRequired(msg *beehivemodel.Message) bool {
return true
case msg.GetSource() == modules.NodeUpgradeJobControllerModuleName || msg.GetSource() == modules.ImagePrePullControllerModuleName:
return true
+ case msg.GetSource() == modules.TaskManagerModuleName:
+ return true
+ case msg.GetSource() == modules.NodeUpgradeJobControllerModuleName:
+ return true
case msg.GetOperation() == beehivemodel.ResponseOperation:
content, ok := msg.Content.(string)
if ok && content == commonconst.MessageSuccessfulContent {
diff --git a/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go b/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
index 70e53cc57..0518b82cc 100644
--- a/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
+++ b/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
@@ -13,8 +13,7 @@ import (
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
beehiveModel "github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
- "github.com/kubeedge/kubeedge/cloud/pkg/nodeupgradejobcontroller/controller"
- commontypes "github.com/kubeedge/kubeedge/common/types"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
)
const (
@@ -23,9 +22,9 @@ const (
// reportTaskStatus report the status of task
func reportTaskStatus(request *restful.Request, response *restful.Response) {
- resp := commontypes.TaskStatus{}
-
+ resp := v1alpha1.TaskStatus{}
taskID := request.PathParameter("taskID")
+ taskType := request.PathParameter("taskType")
nodeID := request.PathParameter("nodeID")
lr := &io.LimitedReader{
@@ -34,23 +33,30 @@ func reportTaskStatus(request *restful.Request, response *restful.Response) {
}
body, err := io.ReadAll(lr)
if err != nil {
- response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to get req body: %v", err))
+ err = response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to get req body: %v", err))
+ if err != nil {
+ klog.Warning(err.Error())
+ }
return
}
if lr.N <= 0 {
- response.WriteError(http.StatusBadRequest, errors.NewRequestEntityTooLargeError(fmt.Sprintf("the request body can only be up to 1MB in size")))
+ err = response.WriteError(http.StatusBadRequest, errors.NewRequestEntityTooLargeError("the request body can only be up to 1MB in size"))
+ if err != nil {
+ klog.Warning(err.Error())
+ }
return
}
if err = json.Unmarshal(body, &resp); err != nil {
- response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to marshal task info: %v", err))
+ err = response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to marshal task info: %v", err))
+ if err != nil {
+ klog.Warning(err.Error())
+ }
return
}
- //TODO The resource operation should be refactored, like: {type}/task/{taskID}/node/{nodeId}
msg := beehiveModel.NewMessage("").SetRoute(modules.CloudHubModuleName, modules.CloudHubModuleGroup).
- SetResourceOperation(fmt.Sprintf("%s/%s/node/%s", resp.Type, taskID, nodeID), controller.NodeUpgrade).FillBody(resp)
- //TODO The message should be sent to the task manager.
- beehiveContext.Send(modules.NodeUpgradeJobControllerModuleName, *msg)
+ SetResourceOperation(fmt.Sprintf("task/%s/node/%s", taskID, nodeID), taskType).FillBody(resp)
+ beehiveContext.Send(modules.TaskManagerModuleName, *msg)
if _, err = response.Write([]byte("ok")); err != nil {
klog.Errorf("failed to send the task resp to edge , err: %v", err)
diff --git a/cloud/pkg/common/messagelayer/context.go b/cloud/pkg/common/messagelayer/context.go
index 8a2e9d4bd..cf632345f 100644
--- a/cloud/pkg/common/messagelayer/context.go
+++ b/cloud/pkg/common/messagelayer/context.go
@@ -96,6 +96,14 @@ func DynamicControllerMessageLayer() MessageLayer {
}
}
+func TaskManagerMessageLayer() MessageLayer {
+ return &ContextMessageLayer{
+ SendModuleName: modules.CloudHubModuleName,
+ ReceiveModuleName: modules.TaskManagerModuleName,
+ ResponseModuleName: modules.CloudHubModuleName,
+ }
+}
+
func NodeUpgradeJobControllerMessageLayer() MessageLayer {
return &ContextMessageLayer{
SendModuleName: modules.CloudHubModuleName,
diff --git a/cloud/pkg/common/modules/modules.go b/cloud/pkg/common/modules/modules.go
index f458aee5e..98665ccab 100644
--- a/cloud/pkg/common/modules/modules.go
+++ b/cloud/pkg/common/modules/modules.go
@@ -19,6 +19,9 @@ const (
ImagePrePullControllerModuleName = "imageprepullcontroller"
ImagePrePullControllerModuleGroup = "imageprepullcontroller"
+ TaskManagerModuleName = "taskmanager"
+ TaskManagerModuleGroup = "taskmanager"
+
SyncControllerModuleName = "synccontroller"
SyncControllerModuleGroup = "synccontroller"
diff --git a/cloud/pkg/taskmanager/config/config.go b/cloud/pkg/taskmanager/config/config.go
new file mode 100644
index 000000000..1366a7904
--- /dev/null
+++ b/cloud/pkg/taskmanager/config/config.go
@@ -0,0 +1,38 @@
+/*
+Copyright 2022 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 config
+
+import (
+ "sync"
+
+ "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
+)
+
+var Config Configure
+var once sync.Once
+
+type Configure struct {
+ v1alpha1.TaskManager
+}
+
+func InitConfigure(tm *v1alpha1.TaskManager) {
+ once.Do(func() {
+ Config = Configure{
+ TaskManager: *tm,
+ }
+ })
+}
diff --git a/cloud/pkg/taskmanager/manager/downstream.go b/cloud/pkg/taskmanager/manager/downstream.go
new file mode 100644
index 000000000..dca8b5c9c
--- /dev/null
+++ b/cloud/pkg/taskmanager/manager/downstream.go
@@ -0,0 +1,64 @@
+/*
+Copyright 2022 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 manager
+
+import (
+ "k8s.io/klog/v2"
+
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "github.com/kubeedge/beehive/pkg/core/model"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/messagelayer"
+)
+
+type DownstreamController struct {
+ downStreamChan chan model.Message
+ messageLayer messagelayer.MessageLayer
+}
+
+// Start DownstreamController
+func (dc *DownstreamController) Start() error {
+ klog.Info("Start NodeUpgradeJob Downstream Controller")
+
+ go dc.syncTask()
+
+ return nil
+}
+
+// syncTask is used to get events from informer
+func (dc *DownstreamController) syncTask() {
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("stop sync tasks")
+ return
+ case msg := <-dc.downStreamChan:
+ err := dc.messageLayer.Send(msg)
+ if err != nil {
+ klog.Errorf("Failed to send upgrade message %v due to error %v", msg.GetID(), err)
+ return
+ }
+ }
+ }
+}
+
+func NewDownstreamController(messageChan chan model.Message) (*DownstreamController, error) {
+ dc := &DownstreamController{
+ downStreamChan: messageChan,
+ messageLayer: messagelayer.TaskManagerMessageLayer(),
+ }
+ return dc, nil
+}
diff --git a/cloud/pkg/taskmanager/manager/executor.go b/cloud/pkg/taskmanager/manager/executor.go
new file mode 100644
index 000000000..83005c075
--- /dev/null
+++ b/cloud/pkg/taskmanager/manager/executor.go
@@ -0,0 +1,415 @@
+package manager
+
+import (
+ "fmt"
+ "reflect"
+ "strings"
+ "sync"
+ "time"
+
+ "github.com/google/uuid"
+ "k8s.io/apimachinery/pkg/util/wait"
+ "k8s.io/client-go/kubernetes"
+ "k8s.io/klog/v2"
+
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "github.com/kubeedge/beehive/pkg/core/model"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/nodeupgradecontroller"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/controller"
+ "github.com/kubeedge/kubeedge/common/constants"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
+ api "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/util/fsm"
+)
+
+type Executor struct {
+ task util.TaskMessage
+ statusChan chan *v1alpha1.TaskStatus
+ nodes []v1alpha1.TaskStatus
+ controller controller.Controller
+}
+
+func NewExecutorMachine(messageChan chan util.TaskMessage, downStreamChan chan model.Message) (*ExecutorMachine, error) {
+ executorMachine = &ExecutorMachine{
+ kubeClient: client.GetKubeClient(),
+ executors: map[string]*Executor{},
+ messageChan: messageChan,
+ downStreamChan: downStreamChan,
+ }
+ return executorMachine, nil
+}
+
+func GetExecutorMachine() *ExecutorMachine {
+ return executorMachine
+}
+
+// Start ExecutorMachine
+func (em *ExecutorMachine) Start() error {
+ klog.Info("Start ExecutorMachine")
+
+ go em.syncTask()
+
+ return nil
+}
+
+// syncTask is used to get events from informer
+func (em *ExecutorMachine) syncTask() {
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("stop sync tasks")
+ return
+ case msg := <-em.messageChan:
+ if msg.ShutDown {
+ klog.Errorf("delete executor %s ", msg.Name)
+ DeleteExecutor(msg)
+ break
+ }
+ err := GetExecutor(msg).HandleMessage(msg.Status)
+ if err != nil {
+ klog.Errorf("Failed to handel upgrade message due to error %s", err.Error())
+ break
+ }
+ }
+ }
+}
+
+type ExecutorMachine struct {
+ kubeClient kubernetes.Interface
+ executors map[string]*Executor
+ messageChan chan util.TaskMessage
+ downStreamChan chan model.Message
+ sync.Mutex
+}
+
+var executorMachine *ExecutorMachine
+
+func GetExecutor(msg util.TaskMessage) *Executor {
+ executorMachine.Lock()
+ e, ok := executorMachine.executors[fmt.Sprintf("%s::%s", msg.Type, msg.Name)]
+ executorMachine.Unlock()
+ if ok && e != nil {
+ return e
+ }
+ e, err := initExecutor(msg)
+ if err != nil {
+ klog.Error("executor init failed, error: %s", err.Error())
+ return nil
+ }
+ return e
+}
+
+func DeleteExecutor(msg util.TaskMessage) {
+ executorMachine.Lock()
+ defer executorMachine.Unlock()
+ delete(executorMachine.executors, fmt.Sprintf("%s::%s", msg.Type, msg.Name))
+}
+
+func (e *Executor) HandleMessage(status v1alpha1.TaskStatus) error {
+ if e == nil {
+ return fmt.Errorf("executor is nil")
+ }
+ e.statusChan <- &status
+ return nil
+}
+
+func (e *Executor) initMessage(node v1alpha1.TaskStatus) *model.Message {
+ // delete it in 1.18
+ if e.task.Type == util.TaskUpgrade {
+ msg := e.initHistoryMessage(node)
+ if msg != nil {
+ klog.Warningf("send history message to node")
+ return msg
+ }
+ }
+
+ msg := model.NewMessage("")
+ resource := buildTaskResource(e.task.Type, e.task.Name, node.NodeName)
+
+ taskReq := commontypes.NodeTaskRequest{
+ TaskID: e.task.Name,
+ Type: e.task.Type,
+ State: string(node.State),
+ }
+ taskReq.Item = e.task.Msg
+ if node.State == api.TaskChecking {
+ taskReq.Item = commontypes.NodePreCheckRequest{
+ CheckItem: e.task.CheckItem,
+ }
+ }
+ msg.BuildRouter(modules.TaskManagerModuleName, modules.TaskManagerModuleGroup, resource, util.TaskUpgrade).
+ FillBody(taskReq)
+ return msg
+}
+
+func (e *Executor) initHistoryMessage(node v1alpha1.TaskStatus) *model.Message {
+ resource := buildUpgradeResource(e.task.Name, node.NodeName)
+ req := e.task.Msg.(commontypes.NodeUpgradeJobRequest)
+ upgradeController := e.controller.(*nodeupgradecontroller.NodeUpgradeController)
+ edgeVersion, err := upgradeController.GetNodeVersion(node.NodeName)
+ if err != nil {
+ klog.Errorf("get node version failed: %s", err.Error())
+ return nil
+ }
+ less, err := util.VersionLess(edgeVersion, "v1.16.0")
+ if err != nil {
+ klog.Errorf("version less failed: %s", err.Error())
+ return nil
+ }
+ if !less {
+ return nil
+ }
+ klog.Warningf("edge version is %s, is less than version %s", edgeVersion, "v1.16.0")
+ upgradeReq := commontypes.NodeUpgradeJobRequest{
+ UpgradeID: e.task.Name,
+ HistoryID: uuid.New().String(),
+ UpgradeTool: "keadm",
+ Version: req.Version,
+ Image: req.Image,
+ }
+ msg := model.NewMessage("")
+ msg.BuildRouter(modules.NodeUpgradeJobControllerModuleName, modules.NodeUpgradeJobControllerModuleGroup, resource, util.TaskUpgrade).
+ FillBody(upgradeReq)
+ return msg
+}
+
+func initExecutor(message util.TaskMessage) (*Executor, error) {
+ controller, err := controller.GetController(message.Type)
+ if err != nil {
+ return nil, err
+ }
+ nodeStatus, err := controller.GetNodeStatus(message.Name)
+ if err != nil {
+ return nil, err
+ }
+ if len(nodeStatus) == 0 {
+ nodeList := controller.ValidateNode(message)
+ if len(nodeList) == 0 {
+ return nil, fmt.Errorf("no node need to be upgrade")
+ }
+ nodeStatus = make([]v1alpha1.TaskStatus, len(nodeList))
+ for i, node := range nodeList {
+ nodeStatus[i] = v1alpha1.TaskStatus{NodeName: node.Name}
+ }
+ err = controller.UpdateNodeStatus(message.Name, nodeStatus)
+ if err != nil {
+ return nil, err
+ }
+ }
+ e := &Executor{
+ task: message,
+ statusChan: make(chan *v1alpha1.TaskStatus, 10),
+ controller: controller,
+ nodes: nodeStatus,
+ }
+ go e.start()
+ executorMachine.executors[fmt.Sprintf("%s::%s", message.Type, message.Name)] = e
+ return e, nil
+}
+
+func (e *Executor) start() {
+ maxFailedNodes := float64(len(e.nodes)) * (e.task.FailureTolerate)
+ failedNodes := map[string]bool{}
+ worker := workers{
+ number: int(e.task.Concurrency),
+ jobs: make(map[string]int),
+ shuttingDown: false,
+ Mutex: sync.Mutex{},
+ }
+ index := 0
+ dealCompletedNode := func(node v1alpha1.TaskStatus) error {
+ if node.State == api.TaskFailed {
+ failedNodes[node.NodeName] = true
+ }
+ if float64(len(failedNodes)) < maxFailedNodes {
+ return nil
+ }
+ worker.shuttingDown = true
+ if len(worker.jobs) > 0 {
+ klog.Warningf("wait for all workers(%d/%d) for task %s to finish running ", len(worker.jobs), worker.number, e.task.Name)
+ return nil
+ }
+
+ errMsg := fmt.Sprintf("the number of failed nodes is %d/%d, which exceeds the failure tolerance threshold.", len(failedNodes), len(e.nodes))
+ _, err := e.controller.ReportTaskStatus(e.task.Name, fsm.Event{
+ Type: node.Event,
+ Action: node.Action,
+ ErrorMsg: errMsg,
+ })
+ if err != nil {
+ return fmt.Errorf("%s, report status failed, %s", errMsg, err.Error())
+ }
+ return fmt.Errorf(errMsg)
+ }
+
+ index, err := e.initWorker(dealCompletedNode, &worker)
+ if err != nil {
+ klog.Errorf(err.Error())
+ return
+ }
+
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("stop sync tasks")
+ return
+ case status := <-e.statusChan:
+ if status == nil || reflect.DeepEqual(*status, v1alpha1.TaskStatus{}) {
+ break
+ }
+ if !e.controller.StageCompleted(e.task.Name, status.State) {
+ break
+ }
+ var endNode int
+ endNode, err = worker.endJob(status.NodeName)
+ if err != nil {
+ klog.Errorf(err.Error())
+ break
+ }
+
+ e.nodes[endNode] = *status
+ err = dealCompletedNode(*status)
+ if err != nil {
+ klog.Warning(err.Error())
+ break
+ }
+
+ if index >= len(e.nodes) {
+ if len(worker.jobs) != 0 {
+ break
+ }
+ var state api.State
+ state, err = e.completedTaskStage(*status)
+ if err != nil {
+ klog.Errorf(err.Error())
+ break
+ }
+ if fsm.TaskFinish(state) {
+ DeleteExecutor(e.task)
+ klog.Infof("task %s is finish", e.task.Name)
+ return
+ }
+ // next stage
+ index, err = e.initWorker(dealCompletedNode, &worker)
+ if err != nil {
+ klog.Errorf(err.Error())
+ }
+ break
+ }
+
+ nextNode := e.nodes[index]
+ err = worker.addJob(nextNode, index, e)
+ if err != nil {
+ klog.Errorf(err.Error())
+ break
+ }
+ index++
+ }
+ }
+}
+
+func (e *Executor) completedTaskStage(node v1alpha1.TaskStatus) (api.State, error) {
+ state, err := e.controller.ReportTaskStatus(e.task.Name, fsm.Event{
+ Type: node.Event,
+ Action: api.ActionSuccess,
+ })
+ if err != nil {
+ return "", err
+ }
+ return state, nil
+}
+
+func (e *Executor) initWorker(dealCompletedNode func(node v1alpha1.TaskStatus) error, worker *workers) (int, error) {
+ var index int
+ var node v1alpha1.TaskStatus
+ isEndNode := true
+ for index, node = range e.nodes {
+ if e.controller.StageCompleted(e.task.Name, node.State) {
+ err := dealCompletedNode(node)
+ if err != nil {
+ return 0, err
+ }
+ continue
+ }
+ err := worker.addJob(node, index, e)
+ if err != nil {
+ klog.Info(err.Error())
+ isEndNode = false
+ break
+ }
+ }
+ if isEndNode {
+ index++
+ }
+ return index, nil
+}
+
+type workers struct {
+ number int
+ jobs map[string]int
+ sync.Mutex
+ shuttingDown bool
+}
+
+func (w *workers) addJob(node v1alpha1.TaskStatus, index int, e *Executor) error {
+ if w.shuttingDown {
+ return fmt.Errorf("workers is stopped")
+ }
+ w.Lock()
+ if len(w.jobs) >= w.number {
+ w.Unlock()
+ return fmt.Errorf("workers are all running, %v/%v", len(w.jobs), w.number)
+ }
+ w.jobs[node.NodeName] = index
+ w.Unlock()
+ msg := e.initMessage(node)
+ go e.handelTimeOutJob(index)
+ executorMachine.downStreamChan <- *msg
+ return nil
+}
+
+func (e *Executor) handelTimeOutJob(index int) {
+ lastState := e.nodes[index].State
+ err := wait.Poll(1*time.Second, time.Duration(*e.task.TimeOutSeconds)*time.Second, func() (bool, error) {
+ if lastState != e.nodes[index].State || fsm.TaskFinish(e.nodes[index].State) {
+ return true, nil
+ }
+ klog.V(4).Infof("node %s stage is not completed", e.nodes[index].NodeName)
+ return false, nil
+ })
+ if err != nil {
+ _, err = e.controller.ReportNodeStatus(e.task.Name, e.nodes[index].NodeName, fsm.Event{
+ Type: "TimeOut",
+ Action: api.ActionFailure,
+ ErrorMsg: fmt.Sprintf("node task execution timeout, %s", err.Error()),
+ })
+ if err != nil {
+ klog.Warningf(err.Error())
+ }
+ }
+}
+
+func (w *workers) endJob(job string) (int, error) {
+ index, ok := w.jobs[job]
+ if !ok {
+ return index, fmt.Errorf("end job %s error, job not exist", job)
+ }
+ w.Lock()
+ delete(w.jobs, job)
+ w.Unlock()
+ return index, nil
+}
+
+func buildTaskResource(task, taskID, nodeID string) string {
+ resource := strings.Join([]string{task, taskID, "node", nodeID}, constants.ResourceSep)
+ return resource
+}
+
+func buildUpgradeResource(upgradeID, nodeID string) string {
+ resource := strings.Join([]string{util.TaskUpgrade, upgradeID, "node", nodeID}, constants.ResourceSep)
+ return resource
+}
diff --git a/cloud/pkg/taskmanager/manager/upstream.go b/cloud/pkg/taskmanager/manager/upstream.go
new file mode 100644
index 000000000..12cad9f46
--- /dev/null
+++ b/cloud/pkg/taskmanager/manager/upstream.go
@@ -0,0 +1,144 @@
+/*
+Copyright 2022 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 manager
+
+import (
+ "encoding/json"
+
+ k8sinformer "k8s.io/client-go/informers"
+ "k8s.io/client-go/kubernetes"
+ "k8s.io/klog/v2"
+
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "github.com/kubeedge/beehive/pkg/core/model"
+ keclient "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/informers"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/messagelayer"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/config"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/controller"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+ crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
+ "github.com/kubeedge/kubeedge/pkg/util/fsm"
+)
+
+// UpstreamController subscribe messages from edge and sync to k8s api server
+type UpstreamController struct {
+ // downstream controller to update NodeUpgradeJob status in cache
+ dc *DownstreamController
+
+ kubeClient kubernetes.Interface
+ informer k8sinformer.SharedInformerFactory
+ crdClient crdClientset.Interface
+ messageLayer messagelayer.MessageLayer
+ // message channel
+ taskStatusChan chan model.Message
+}
+
+// Start UpstreamController
+func (uc *UpstreamController) Start() error {
+ klog.Info("Start Task Upstream Controller")
+
+ uc.taskStatusChan = make(chan model.Message, config.Config.Buffer.TaskStatus)
+ go uc.dispatchMessage()
+
+ for i := 0; i < int(config.Config.Load.TaskWorkers); i++ {
+ go uc.updateTaskStatus()
+ }
+ return nil
+}
+
+// Start UpstreamController
+func (uc *UpstreamController) dispatchMessage() {
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("Stop dispatch task upstream message")
+ return
+ default:
+ }
+
+ msg, err := uc.messageLayer.Receive()
+ if err != nil {
+ klog.Warningf("Receive message failed, %v", err)
+ continue
+ }
+
+ klog.V(4).Infof("task upstream controller receive msg %#v", msg)
+
+ uc.taskStatusChan <- msg
+ }
+}
+
+// updateTaskStatus update NodeUpgradeJob status field
+func (uc *UpstreamController) updateTaskStatus() {
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("Stop update NodeUpgradeJob status")
+ return
+ case msg := <-uc.taskStatusChan:
+ klog.V(4).Infof("Message: %s, operation is: %s, and resource is: %s", msg.GetID(), msg.GetOperation(), msg.GetResource())
+
+ // get nodeID and upgradeID from Upgrade msg:
+ nodeID := util.GetNodeName(msg.GetResource())
+ taskID := util.GetTaskID(msg.GetResource())
+
+ data, err := msg.GetContentData()
+ if err != nil {
+ klog.Errorf("failed to get node upgrade content data: %v", err)
+ continue
+ }
+
+ c, err := controller.GetController(msg.GetOperation())
+ if err != nil {
+ klog.Errorf("Failed to get controller: %v", err)
+ continue
+ }
+
+ resp := v1alpha1.TaskStatus{}
+ err = json.Unmarshal(data, &resp)
+ if err != nil {
+ klog.Errorf("Failed to unmarshal node upgrade response: %v", err)
+ continue
+ }
+ event := fsm.Event{
+ Type: resp.Event,
+ Action: resp.Action,
+ ErrorMsg: resp.Reason,
+ }
+
+ _, err = c.ReportNodeStatus(taskID, nodeID, event)
+ if err != nil {
+ klog.Errorf("Failed to report status: %v", err)
+ continue
+ }
+ }
+ }
+}
+
+// NewUpstreamController create UpstreamController from config
+func NewUpstreamController(dc *DownstreamController) (*UpstreamController, error) {
+ uc := &UpstreamController{
+ kubeClient: keclient.GetKubeClient(),
+ informer: informers.GetInformersManager().GetKubeInformerFactory(),
+ crdClient: keclient.GetCRDClient(),
+ messageLayer: messagelayer.TaskManagerMessageLayer(),
+ dc: dc,
+ }
+ return uc, nil
+}
diff --git a/cloud/pkg/taskmanager/task_manager.go b/cloud/pkg/taskmanager/task_manager.go
new file mode 100644
index 000000000..167d65d8f
--- /dev/null
+++ b/cloud/pkg/taskmanager/task_manager.go
@@ -0,0 +1,144 @@
+/*
+Copyright 2022 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 taskmanager
+
+import (
+ "time"
+
+ "k8s.io/klog/v2"
+
+ "github.com/kubeedge/beehive/pkg/core"
+ "github.com/kubeedge/beehive/pkg/core/model"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/config"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/manager"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/nodeupgradecontroller"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/controller"
+ "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
+)
+
+//
+//// TaskManager is controller for processing upgrading edge node from cloud
+//type NodeUpgradeJobController struct {
+//}
+//
+//// Name of controller
+//func (uc *NodeUpgradeJobController) Name() string {
+// return modules.NodeUpgradeJobControllerModuleName
+//}
+//
+//// Group of controller
+//func (uc *NodeUpgradeJobController) Group() string {
+// return modules.NodeUpgradeJobControllerModuleGroup
+//}
+//
+//// Enable indicates whether enable this module
+//func (uc *NodeUpgradeJobController) Enable() bool {
+// return true
+//}
+//
+//// Start controller
+//func (uc *NodeUpgradeJobController) Start() {
+//}
+//
+//func newNodeUpgradeJobController() *NodeUpgradeJobController {
+// return &NodeUpgradeJobController{}
+//}
+
+type TaskManager struct {
+ downstream *manager.DownstreamController
+ executorMachine *manager.ExecutorMachine
+ upstream *manager.UpstreamController
+ enable bool
+}
+
+var _ core.Module = (*TaskManager)(nil)
+
+func newTaskManager(enable bool) *TaskManager {
+ if !enable {
+ return &TaskManager{enable: enable}
+ }
+ taskMessage := make(chan util.TaskMessage, 10)
+ downStreamMessage := make(chan model.Message, 10)
+ downstream, err := manager.NewDownstreamController(downStreamMessage)
+ if err != nil {
+ klog.Exitf("New task manager downstream failed with error: %s", err)
+ }
+ upstream, err := manager.NewUpstreamController(downstream)
+ if err != nil {
+ klog.Exitf("New task manager upstream failed with error: %s", err)
+ }
+ executorMachine, err := manager.NewExecutorMachine(taskMessage, downStreamMessage)
+ if err != nil {
+ klog.Exitf("New executor machine failed with error: %s", err)
+ }
+
+ upgradeNodeController, err := nodeupgradecontroller.NewNodeUpgradeController(taskMessage)
+ if err != nil {
+ klog.Exitf("New upgrade node controller failed with error: %s", err)
+ }
+ controller.Register(util.TaskUpgrade, upgradeNodeController)
+
+ return &TaskManager{
+ downstream: downstream,
+ executorMachine: executorMachine,
+ upstream: upstream,
+ enable: enable,
+ }
+}
+
+func Register(dc *v1alpha1.TaskManager) {
+ config.InitConfigure(dc)
+ core.Register(newTaskManager(dc.Enable))
+ //core.Register(newNodeUpgradeJobController())
+}
+
+// Name of controller
+func (uc *TaskManager) Name() string {
+ return modules.TaskManagerModuleName
+}
+
+// Group of controller
+func (uc *TaskManager) Group() string {
+ return modules.TaskManagerModuleGroup
+}
+
+// Enable indicates whether enable this module
+func (uc *TaskManager) Enable() bool {
+ return uc.enable
+}
+
+// Start controller
+func (uc *TaskManager) Start() {
+ if err := uc.downstream.Start(); err != nil {
+ klog.Exitf("start task manager downstream failed with error: %s", err)
+ }
+ // wait for downstream controller to start and load NodeUpgradeJob
+ // TODO think about sync
+ time.Sleep(1 * time.Second)
+ if err := uc.upstream.Start(); err != nil {
+ klog.Exitf("start task manager upstream failed with error: %s", err)
+ }
+ if err := uc.executorMachine.Start(); err != nil {
+ klog.Exitf("start task manager executorMachine failed with error: %s", err)
+ }
+
+ if err := controller.StartAllController(); err != nil {
+ klog.Exitf("start controller failed with error: %s", err)
+ }
+}
diff --git a/cloud/pkg/taskmanager/util/controller/controller.go b/cloud/pkg/taskmanager/util/controller/controller.go
new file mode 100644
index 000000000..55d43f92c
--- /dev/null
+++ b/cloud/pkg/taskmanager/util/controller/controller.go
@@ -0,0 +1,153 @@
+package controller
+
+import (
+ "fmt"
+
+ v1 "k8s.io/api/core/v1"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ k8sinformer "k8s.io/client-go/informers"
+ "k8s.io/client-go/kubernetes"
+ "k8s.io/klog/v2"
+
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/manager"
+ api "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+ crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
+ "github.com/kubeedge/kubeedge/pkg/util/fsm"
+)
+
+type Controller interface {
+ Name() string
+ Start() error
+ ReportNodeStatus(string, string, fsm.Event) (api.State, error)
+ ReportTaskStatus(string, fsm.Event) (api.State, error)
+ ValidateNode(util.TaskMessage) []v1.Node
+ GetNodeStatus(string) ([]v1alpha1.TaskStatus, error)
+ UpdateNodeStatus(string, []v1alpha1.TaskStatus) error
+ StageCompleted(taskID string, state api.State) bool
+}
+
+type BaseController struct {
+ name string
+ Informer k8sinformer.SharedInformerFactory
+ TaskManager *manager.TaskCache
+ MessageChan chan util.TaskMessage
+ KubeClient kubernetes.Interface
+ CrdClient crdClientset.Interface
+}
+
+func (bc *BaseController) Name() string {
+ return bc.name
+}
+
+func (bc *BaseController) Start() error {
+ return fmt.Errorf("controller not implemented")
+}
+
+func (bc *BaseController) StageCompleted(taskID string, state api.State) bool {
+ return false
+}
+
+func (bc *BaseController) ValidateNode(taskMessage util.TaskMessage) []v1.Node {
+ var validateNodes []v1.Node
+ nodes, err := bc.getNodeList(taskMessage.NodeNames, taskMessage.LabelSelector)
+ if err != nil {
+ klog.Warningf("get node list error: %s", err.Error())
+ return nil
+ }
+ for _, node := range nodes {
+ if !util.IsEdgeNode(node) {
+ klog.Warningf("Node(%s) is not edge node", node.Name)
+ continue
+ }
+ ready := isNodeReady(node)
+ if !ready {
+ continue
+ }
+ validateNodes = append(validateNodes, *node)
+ }
+ return validateNodes
+}
+
+func (bc *BaseController) GetNodeStatus(name string) ([]v1alpha1.TaskStatus, error) {
+ return nil, fmt.Errorf("function GetNodeStatus need to be init")
+}
+
+func (bc *BaseController) UpdateNodeStatus(name string, status []v1alpha1.TaskStatus) error {
+ return fmt.Errorf("function UpdateNodeStatus need to be init")
+}
+
+func isNodeReady(node *v1.Node) bool {
+ for _, condition := range node.Status.Conditions {
+ if condition.Type == v1.NodeReady && condition.Status != v1.ConditionTrue {
+ klog.Warningf("Node(%s) is in NotReady state", node.Name)
+ return false
+ }
+ }
+ return true
+}
+
+func (bc *BaseController) ReportNodeStatus(taskID string, nodeID string, event fsm.Event) (api.State, error) {
+ return "", fmt.Errorf("function ReportNodeStatus need to be init")
+}
+
+func (bc *BaseController) ReportTaskStatus(taskID string, event fsm.Event) (api.State, error) {
+ return "", fmt.Errorf("function ReportTaskStatus need to be init")
+}
+
+var (
+ controllers = map[string]Controller{}
+)
+
+func Register(name string, controller Controller) {
+ if _, ok := controllers[name]; ok {
+ klog.Warning("controller %s exists ", name)
+ }
+ controllers[name] = controller
+}
+
+func StartAllController() error {
+ for name, controller := range controllers {
+ err := controller.Start()
+ if err != nil {
+ return fmt.Errorf("start %s controller failed: %s", name, err.Error())
+ }
+ }
+ return nil
+}
+
+func GetController(name string) (Controller, error) {
+ controller, ok := controllers[name]
+ if !ok {
+ return nil, fmt.Errorf("controller %s is not registered", name)
+ }
+ return controller, nil
+}
+
+func (bc *BaseController) getNodeList(nodeNames []string, labelSelector *metav1.LabelSelector) ([]*v1.Node, error) {
+ var nodesToUpgrade []*v1.Node
+
+ if len(nodeNames) != 0 {
+ for _, name := range nodeNames {
+ node, err := bc.Informer.Core().V1().Nodes().Lister().Get(name)
+ if err != nil {
+ return nil, fmt.Errorf("failed to get node with name %s: %v", name, err)
+ }
+ nodesToUpgrade = append(nodesToUpgrade, node)
+ }
+ } else if labelSelector != nil {
+ selector, err := metav1.LabelSelectorAsSelector(labelSelector)
+ if err != nil {
+ return nil, fmt.Errorf("labelSelector(%s) is not valid: %v", labelSelector, err)
+ }
+
+ nodes, err := bc.Informer.Core().V1().Nodes().Lister().List(selector)
+ if err != nil {
+ return nil, fmt.Errorf("failed to get nodes with label %s: %v", selector.String(), err)
+ }
+ nodesToUpgrade = nodes
+ }
+
+ return nodesToUpgrade, nil
+}
diff --git a/cloud/pkg/taskmanager/util/manager/common.go b/cloud/pkg/taskmanager/util/manager/common.go
new file mode 100644
index 000000000..6e8008abf
--- /dev/null
+++ b/cloud/pkg/taskmanager/util/manager/common.go
@@ -0,0 +1,62 @@
+/*
+Copyright 2022 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 manager
+
+import (
+ "k8s.io/apimachinery/pkg/runtime"
+ "k8s.io/apimachinery/pkg/watch"
+ "k8s.io/klog/v2"
+)
+
+// Manager define the interface of a Manager, NodeUpgradeJob Manager implement it
+type Manager interface {
+ Events() chan watch.Event
+}
+
+// CommonResourceEventHandler can be used by NodeUpgradeJob Manager
+type CommonResourceEventHandler struct {
+ events chan watch.Event
+}
+
+func (c *CommonResourceEventHandler) obj2Event(t watch.EventType, obj interface{}) {
+ eventObj, ok := obj.(runtime.Object)
+ if !ok {
+ klog.Warningf("unknown type: %T, ignore", obj)
+ return
+ }
+ c.events <- watch.Event{Type: t, Object: eventObj}
+}
+
+// OnAdd handle Add event
+func (c *CommonResourceEventHandler) OnAdd(obj interface{}, isInInitialList bool) {
+ c.obj2Event(watch.Added, obj)
+}
+
+// OnUpdate handle Update event
+func (c *CommonResourceEventHandler) OnUpdate(oldObj, newObj interface{}) {
+ c.obj2Event(watch.Modified, newObj)
+}
+
+// OnDelete handle Delete event
+func (c *CommonResourceEventHandler) OnDelete(obj interface{}) {
+ c.obj2Event(watch.Deleted, obj)
+}
+
+// NewCommonResourceEventHandler create CommonResourceEventHandler used by NodeUpgradeJob Manager
+func NewCommonResourceEventHandler(events chan watch.Event) *CommonResourceEventHandler {
+ return &CommonResourceEventHandler{events: events}
+}
diff --git a/cloud/pkg/taskmanager/util/manager/task_cache.go b/cloud/pkg/taskmanager/util/manager/task_cache.go
new file mode 100644
index 000000000..98c793fb1
--- /dev/null
+++ b/cloud/pkg/taskmanager/util/manager/task_cache.go
@@ -0,0 +1,52 @@
+/*
+Copyright 2022 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 manager
+
+import (
+ "sync"
+
+ "k8s.io/apimachinery/pkg/watch"
+ "k8s.io/client-go/tools/cache"
+
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/config"
+)
+
+// TaskCache is a manager watch CRD change event
+type TaskCache struct {
+ // events from watch kubernetes api server
+ events chan watch.Event
+
+ // CacheMap, key is NodeUpgradeJob.Name, value is *v1alpha1.NodeUpgradeJob{}
+ CacheMap sync.Map
+}
+
+// Events return a channel, can receive all NodeUpgradeJob event
+func (dmm *TaskCache) Events() chan watch.Event {
+ return dmm.events
+}
+
+// NewTaskCache create TaskCache from config
+func NewTaskCache(si cache.SharedIndexInformer) (*TaskCache, error) {
+ events := make(chan watch.Event, config.Config.Buffer.TaskEvent)
+ rh := NewCommonResourceEventHandler(events)
+ _, err := si.AddEventHandler(rh)
+ if err != nil {
+ return nil, err
+ }
+
+ return &TaskCache{events: events}, nil
+}
diff --git a/cloud/pkg/taskmanager/util/util.go b/cloud/pkg/taskmanager/util/util.go
new file mode 100644
index 000000000..c68e8574a
--- /dev/null
+++ b/cloud/pkg/taskmanager/util/util.go
@@ -0,0 +1,169 @@
+/*
+Copyright 2022 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 util
+
+import (
+ "fmt"
+ "k8s.io/klog/v2"
+ "strings"
+
+ "github.com/distribution/distribution/v3/reference"
+ metav1 "k8s.io/api/core/v1"
+ "k8s.io/apimachinery/pkg/apis/meta/v1"
+ versionutil "k8s.io/apimachinery/pkg/util/version"
+
+ "github.com/kubeedge/kubeedge/common/constants"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+)
+
+const (
+ NodeUpgradeJobStatusKey = "nodeupgradejob.operations.kubeedge.io/status"
+ NodeUpgradeJobStatusValue = ""
+ NodeUpgradeHistoryKey = "nodeupgradejob.operations.kubeedge.io/history"
+)
+
+const (
+ TaskUpgrade = "upgrade"
+ TaskRollback = "rollback"
+ TaskBackup = "backup"
+ TaskDockerPrePull = "pre-pull"
+
+ ISO8601UTC = "2006-01-02T15:04:05Z"
+)
+
+type TaskMessage struct {
+ Type string
+ Name string
+ TimeOutSeconds *uint32
+ ShutDown bool
+ CheckItem []string
+ Concurrency int32
+ FailureTolerate float64
+ NodeNames []string
+ LabelSelector *v1.LabelSelector
+ Status v1alpha1.TaskStatus
+ Msg interface{}
+}
+
+// FilterVersion returns true only if the edge node version already on the upgrade req
+// version is like: v1.22.6-kubeedge-v1.10.0-beta.0.185+95378fb019912a, expected is like v1.10.0
+func FilterVersion(version string, expected string) bool {
+ // if not correct version format, also return true
+ strs := strings.Split(version, "-")
+ if len(strs) < 3 {
+ klog.Warningf("version format should be {k8s version}-kubeedge-{edgecore version}, but got : %s", version)
+ return true
+ }
+
+ // filter nodes that already in the required version
+ less, err := VersionLess(strs[2], expected)
+ if err != nil {
+ klog.Warningf("version filter failed: %s", err.Error())
+ less = false
+ }
+ return !less
+}
+
+// IsEdgeNode checks whether a node is an Edge Node
+// only if label {"node-role.kubernetes.io/edge": ""} exists, it is an edge node
+func IsEdgeNode(node *metav1.Node) bool {
+ if node.Labels == nil {
+ return false
+ }
+ if _, ok := node.Labels[constants.EdgeNodeRoleKey]; !ok {
+ return false
+ }
+
+ if node.Labels[constants.EdgeNodeRoleKey] != constants.EdgeNodeRoleValue {
+ return false
+ }
+
+ return true
+}
+
+// RemoveDuplicateElement deduplicate
+func RemoveDuplicateElement(s []string) []string {
+ result := make([]string, 0, len(s))
+ temp := make(map[string]struct{}, len(s))
+
+ for _, item := range s {
+ if _, ok := temp[item]; !ok {
+ temp[item] = struct{}{}
+ result = append(result, item)
+ }
+ }
+
+ return result
+}
+
+// MergeAnnotationUpgradeHistory constructs the new history based on the origin history
+// and we'll only keep 3 records
+func MergeAnnotationUpgradeHistory(origin, fromVersion, toVersion string) string {
+ newHistory := fmt.Sprintf("%s->%s", fromVersion, toVersion)
+ if origin == "" {
+ return newHistory
+ }
+
+ sets := strings.Split(origin, ";")
+ if len(sets) > 2 {
+ sets = sets[1:]
+ }
+
+ sets = append(sets, newHistory)
+ return strings.Join(sets, ";")
+}
+
+// GetImageRepo gets repo from a container image
+func GetImageRepo(image string) (string, error) {
+ named, err := reference.ParseNormalizedNamed(image)
+ if err != nil {
+ return "", fmt.Errorf("failed to parse image name: %v", err)
+ }
+
+ return named.Name(), nil
+}
+
+func GetNodeName(resource string) string {
+ // task/${TaskID}/node/${NodeID}
+ s := strings.Split(resource, "/")
+ return s[3]
+}
+func GetTaskID(resource string) string {
+ // task/${TaskID}/node/${NodeID}
+ s := strings.Split(resource, "/")
+ return s[1]
+}
+
+func VersionLess(version1, version2 string) (bool, error) {
+ less := false
+ ver1, err := versionutil.ParseGeneric(version1)
+ if err != nil {
+ return less, fmt.Errorf("version1 error: %v", err)
+ }
+ ver2, err := versionutil.ParseGeneric(version2)
+ if err != nil {
+ return less, fmt.Errorf("version2 error: %v", err)
+ }
+ // If the remote Major version is bigger or if the Major versions are the same,
+ // but the remote Minor is bigger use the client version release. This handles Major bumps too.
+ if ver1.Major() < ver2.Major() ||
+ (ver1.Major() == ver2.Major()) && ver1.Minor() < ver2.Minor() ||
+ (ver1.Major() == ver2.Major() && ver1.Minor() == ver2.Minor()) && ver1.Patch() < ver2.Patch() {
+ less = true
+ }
+ return less, nil
+}
diff --git a/cloud/pkg/taskmanager/util/util_test.go b/cloud/pkg/taskmanager/util/util_test.go
new file mode 100644
index 000000000..ced6de9ef
--- /dev/null
+++ b/cloud/pkg/taskmanager/util/util_test.go
@@ -0,0 +1,171 @@
+/*
+Copyright 2022 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 util
+
+import (
+ "reflect"
+ "testing"
+)
+
+func TestFilterVersion(t *testing.T) {
+ tests := []struct {
+ name string
+ version string
+ expected string
+ expectResult bool
+ }{
+ {
+ name: "not match expected version",
+ version: "v1.22.6-kubeedge-v1.9.0",
+ expected: "v1.10.0",
+ expectResult: false,
+ },
+ {
+ name: "not match expected version",
+ version: "v1.22.6-kubeedge-v1.10.0-beta.0.194+77ea462f402efb",
+ expected: "v1.10.0",
+ expectResult: false,
+ },
+ {
+ name: "no right format version",
+ version: "v1.22.6",
+ expected: "v1.10.0",
+ expectResult: false,
+ },
+ {
+ name: "match expected version",
+ version: "v1.22.6-kubeedge-v1.10.0",
+ expected: "v1.10.0",
+ expectResult: true,
+ },
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ result := FilterVersion(test.version, test.expected)
+ if result != test.expectResult {
+ t.Errorf("Got = %v, Want = %v", result, test.expectResult)
+ }
+ })
+ }
+}
+
+func TestRemoveDuplicateElement(t *testing.T) {
+ tests := []struct {
+ name string
+ input []string
+ expected []string
+ }{
+ {
+ name: "case 1",
+ input: []string{"a", "b", "c"},
+ expected: []string{"a", "b", "c"},
+ },
+ {
+ name: "case 2",
+ input: []string{"a", "a", "b", "c", "b", "a", "a"},
+ expected: []string{"a", "b", "c"},
+ },
+ {
+ name: "case 3",
+ input: []string{},
+ expected: []string{},
+ },
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ result := RemoveDuplicateElement(test.input)
+ if !reflect.DeepEqual(result, test.expected) {
+ t.Errorf("Got = %v, Want = %v", result, test.expected)
+ }
+ })
+ }
+}
+
+func TestMergeAnnotationUpgradeHistory(t *testing.T) {
+ tests := []struct {
+ name string
+ origin string
+ fromVersion string
+ toVersion string
+ expected string
+ }{
+ {
+ name: "case 1: no history record exist",
+ origin: "",
+ fromVersion: "v1.10.0",
+ toVersion: "v1.10.1",
+ expected: "v1.10.0->v1.10.1",
+ },
+ {
+ name: "case 2: 1 history record exist",
+ origin: "v1.10.0->v1.10.1",
+ fromVersion: "v1.10.1",
+ toVersion: "v1.10.2",
+ expected: "v1.10.0->v1.10.1;v1.10.1->v1.10.2",
+ },
+ {
+ name: "case 2: 3 history record exist",
+ origin: "1.10.0->v1.10.1;v1.10.1->v1.10.2;v1.10.2->v1.10.3",
+ fromVersion: "v1.10.3",
+ toVersion: "v1.10.4",
+ expected: "v1.10.1->v1.10.2;v1.10.2->v1.10.3;v1.10.3->v1.10.4",
+ },
+ }
+
+ for _, test := range tests {
+ t.Run(test.name, func(t *testing.T) {
+ result := MergeAnnotationUpgradeHistory(test.origin, test.fromVersion, test.toVersion)
+ if result != test.expected {
+ t.Errorf("Got = %v, Want = %v", result, test.expected)
+ }
+ })
+ }
+}
+
+func TestGetImageRepo(t *testing.T) {
+ tests := []struct {
+ Image string
+ ExpectRepo string
+ }{
+ {Image: "name", ExpectRepo: "docker.io/library/name"},
+ {Image: "name:tag", ExpectRepo: "docker.io/library/name"},
+ {Image: "name@sha256:59329e44d499406bd2e620473b0ba0b531abb7e326cef0156f33e5957cdfe259", ExpectRepo: "docker.io/library/name"},
+ {Image: "org/name", ExpectRepo: "docker.io/org/name"},
+ {Image: "org/name:tag", ExpectRepo: "docker.io/org/name"},
+ {Image: "org/name@sha256:59329e44d499406bd2e620473b0ba0b531abb7e326cef0156f33e5957cdfe259", ExpectRepo: "docker.io/org/name"},
+ {Image: "registry:8080/name", ExpectRepo: "registry:8080/name"},
+ {Image: "registry:8080/name:tag", ExpectRepo: "registry:8080/name"},
+ {Image: "registry:8080/name@sha256:59329e44d499406bd2e620473b0ba0b531abb7e326cef0156f33e5957cdfe259", ExpectRepo: "registry:8080/name"},
+ {Image: "registry:8080/org/name", ExpectRepo: "registry:8080/org/name"},
+ {Image: "registry:8080/org/name:tag", ExpectRepo: "registry:8080/org/name"},
+ {Image: "registry:8080/org/name@sha256:59329e44d499406bd2e620473b0ba0b531abb7e326cef0156f33e5957cdfe259", ExpectRepo: "registry:8080/org/name"},
+ }
+
+ for _, test := range tests {
+ t.Run(test.Image, func(t *testing.T) {
+ repo, err := GetImageRepo(test.Image)
+ if err != nil {
+ t.Errorf("error: %v", err)
+ }
+ if repo != test.ExpectRepo {
+ t.Errorf("Got = %v, Want = %v", repo, test.ExpectRepo)
+ }
+ })
+ }
+}