diff options
| author | zhengxinwei <zhengxinwei@huawei.com> | 2024-01-04 15:27:04 +0800 |
|---|---|---|
| committer | zhengxinwei-f <zhengxinwei@huawei.com> | 2024-01-16 09:35:44 +0800 |
| commit | b1c1c6a8e5442fff2b4d932060648e009f639df1 (patch) | |
| tree | 575f35f9559894e995060f15b029ff73f2d23282 /cloud | |
| parent | notify the task manager of the task status (diff) | |
| download | kubeedge-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.go | 4 | ||||
| -rw-r--r-- | cloud/pkg/cloudhub/dispatcher/message_dispatcher.go | 4 | ||||
| -rw-r--r-- | cloud/pkg/cloudhub/servers/httpserver/report_task_status.go | 28 | ||||
| -rw-r--r-- | cloud/pkg/common/messagelayer/context.go | 8 | ||||
| -rw-r--r-- | cloud/pkg/common/modules/modules.go | 3 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/config/config.go | 38 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/manager/downstream.go | 64 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/manager/executor.go | 415 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/manager/upstream.go | 144 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/task_manager.go | 144 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/util/controller/controller.go | 153 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/util/manager/common.go | 62 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/util/manager/task_cache.go | 52 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/util/util.go | 169 | ||||
| -rw-r--r-- | cloud/pkg/taskmanager/util/util_test.go | 171 |
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) + } + }) + } +} |
