summaryrefslogtreecommitdiff
path: root/cloud
diff options
context:
space:
mode:
authorzhengxinwei <zhengxinwei@huawei.com>2024-01-04 15:28:35 +0800
committerzhengxinwei-f <zhengxinwei@huawei.com>2024-01-16 20:18:43 +0800
commitd091c086a5c8620bf43c90a64835d55253227005 (patch)
treefaa19b8f9a64684541631e1580e9dc7d78a63f48 /cloud
parentImplement task manager to complete cloud edge task execution (diff)
downloadkubeedge-d091c086a5c8620bf43c90a64835d55253227005.tar.gz
support edge nodes upgrade and image pre pull
Signed-off-by: zhengxinwei-f <zhengxinwei@huawei.com>
Diffstat (limited to 'cloud')
-rw-r--r--cloud/cmd/cloudcore/app/server.go2
-rw-r--r--cloud/pkg/cloudhub/dispatcher/message_dispatcher.go9
-rw-r--r--cloud/pkg/cloudhub/servers/httpserver/report_task_status.go22
-rw-r--r--cloud/pkg/cloudhub/servers/httpserver/upgrade.go64
-rw-r--r--cloud/pkg/common/messagelayer/context.go16
-rw-r--r--cloud/pkg/common/modules/modules.go3
-rw-r--r--cloud/pkg/imageprepullcontroller/config/config.go38
-rw-r--r--cloud/pkg/imageprepullcontroller/controller/downstream.go268
-rw-r--r--cloud/pkg/imageprepullcontroller/controller/upstream.go147
-rw-r--r--cloud/pkg/imageprepullcontroller/controller/util.go131
-rw-r--r--cloud/pkg/imageprepullcontroller/imageprepullcontroller.go91
-rw-r--r--cloud/pkg/imageprepullcontroller/manager/common.go62
-rw-r--r--cloud/pkg/imageprepullcontroller/manager/imageprepull.go52
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/config/config.go38
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/controller/downstream.go447
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/controller/upstream.go244
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/controller/util.go124
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/controller/util_test.go223
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/manager/common.go62
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/manager/nodeupgrade.go52
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/nodeupgradejobcontroller.go91
-rw-r--r--cloud/pkg/taskmanager/config/config.go2
-rw-r--r--cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_controller.go338
-rw-r--r--cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_task.go133
-rw-r--r--cloud/pkg/taskmanager/manager/downstream.go4
-rw-r--r--cloud/pkg/taskmanager/manager/executor.go176
-rw-r--r--cloud/pkg/taskmanager/manager/upstream.go13
-rw-r--r--cloud/pkg/taskmanager/nodeupgradecontroller/node_upgrade_controller.go386
-rw-r--r--cloud/pkg/taskmanager/nodeupgradecontroller/upgrade_task.go122
-rw-r--r--cloud/pkg/taskmanager/task_manager.go37
-rw-r--r--cloud/pkg/taskmanager/util/controller/controller.go16
-rw-r--r--cloud/pkg/taskmanager/util/manager/common.go2
-rw-r--r--cloud/pkg/taskmanager/util/manager/task_cache.go2
-rw-r--r--cloud/pkg/taskmanager/util/util.go24
-rw-r--r--cloud/pkg/taskmanager/util/util_test.go6
35 files changed, 1205 insertions, 2242 deletions
diff --git a/cloud/cmd/cloudcore/app/server.go b/cloud/cmd/cloudcore/app/server.go
index 3f4949729..c4a7a2388 100644
--- a/cloud/cmd/cloudcore/app/server.go
+++ b/cloud/cmd/cloudcore/app/server.go
@@ -48,7 +48,6 @@ import (
"github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller"
"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/policycontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/router"
"github.com/kubeedge/kubeedge/cloud/pkg/synccontroller"
@@ -159,7 +158,6 @@ func registerModules(c *v1alpha1.CloudCoreConfig) {
cloudhub.Register(c.Modules.CloudHub)
edgecontroller.Register(c.Modules.EdgeController)
devicecontroller.Register(c.Modules.DeviceController)
- imageprepullcontroller.Register(c.Modules.ImagePrePullController)
taskmanager.Register(c.Modules.TaskManager)
synccontroller.Register(c.Modules.SyncController)
cloudstream.Register(c.Modules.CloudStream, c.CommonConfig)
diff --git a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
index 940acb3a5..18382f13b 100644
--- a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
+++ b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go
@@ -33,8 +33,8 @@ import (
"github.com/kubeedge/kubeedge/cloud/pkg/cloudhub/session"
"github.com/kubeedge/kubeedge/cloud/pkg/common/messagelayer"
"github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
- imageprepullcontroller "github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller/controller"
"github.com/kubeedge/kubeedge/cloud/pkg/synccontroller"
+ taskutil "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
commonconst "github.com/kubeedge/kubeedge/common/constants"
v2 "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao/v2"
"github.com/kubeedge/kubeedge/pkg/apis/reliablesyncs/v1alpha1"
@@ -184,8 +184,9 @@ func (md *messageDispatcher) DispatchUpstream(message *beehivemodel.Message, inf
message.Router.Resource = fmt.Sprintf("node/%s/%s", info.NodeID, message.Router.Resource)
beehivecontext.Send(modules.RouterModuleName, *message)
- case message.GetOperation() == imageprepullcontroller.ImagePrePull:
- beehivecontext.SendToGroup(modules.ImagePrePullControllerModuleGroup, *message)
+ case message.GetOperation() == taskutil.TaskPrePull ||
+ message.GetOperation() == taskutil.TaskUpgrade:
+ beehivecontext.SendToGroup(modules.TaskManagerModuleGroup, *message)
default:
err := md.PubToController(info, message)
@@ -451,8 +452,6 @@ func noAckRequired(msg *beehivemodel.Message) bool {
return true
case msg.GetGroup() == modules.UserGroup:
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:
diff --git a/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go b/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
index 0518b82cc..784470502 100644
--- a/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
+++ b/cloud/pkg/cloudhub/servers/httpserver/report_task_status.go
@@ -1,3 +1,19 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
package httpserver
import (
@@ -13,7 +29,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/pkg/apis/operations/v1alpha1"
+ "github.com/kubeedge/kubeedge/common/types"
)
const (
@@ -22,7 +38,7 @@ const (
// reportTaskStatus report the status of task
func reportTaskStatus(request *restful.Request, response *restful.Response) {
- resp := v1alpha1.TaskStatus{}
+ resp := types.NodeTaskResponse{}
taskID := request.PathParameter("taskID")
taskType := request.PathParameter("taskType")
nodeID := request.PathParameter("nodeID")
@@ -47,7 +63,7 @@ func reportTaskStatus(request *restful.Request, response *restful.Response) {
return
}
if err = json.Unmarshal(body, &resp); err != nil {
- err = response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to marshal task info: %v", err))
+ err = response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to unmarshal task info: %v", err))
if err != nil {
klog.Warning(err.Error())
}
diff --git a/cloud/pkg/cloudhub/servers/httpserver/upgrade.go b/cloud/pkg/cloudhub/servers/httpserver/upgrade.go
index c8b8eb65f..2dc343714 100644
--- a/cloud/pkg/cloudhub/servers/httpserver/upgrade.go
+++ b/cloud/pkg/cloudhub/servers/httpserver/upgrade.go
@@ -17,6 +17,7 @@ import (
"encoding/json"
"fmt"
"io"
+ "net/http"
"github.com/emicklei/go-restful"
"k8s.io/apimachinery/pkg/api/errors"
@@ -25,40 +26,73 @@ 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"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
commontypes "github.com/kubeedge/kubeedge/common/types"
+ api "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+)
+
+const (
+ UpgradeSuccess = "upgrade_success"
+ UpgradeFailedRollbackSuccess = "upgrade_failed_rollback_success"
+ UpgradeFailedRollbackFailed = "upgrade_failed_rollback_failed"
)
// upgradeEdge upgrade the edgecore version
func upgradeEdge(request *restful.Request, response *restful.Response) {
resp := commontypes.NodeUpgradeJobResponse{}
- defer func() {
- if _, err := response.Write([]byte("ok")); err != nil {
- klog.Errorf("failed to send upgrade edge resp, err: %v", err)
- }
- }()
+ taskID := resp.UpgradeID
+ taskType := util.TaskUpgrade
+ nodeID := resp.NodeName
- limit := int64(3 * 1024 * 1024)
lr := &io.LimitedReader{
R: request.Request.Body,
- N: limit + 1,
+ N: millionByte + 1,
}
body, err := io.ReadAll(lr)
if err != nil {
- klog.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 {
- klog.Errorf("%v", errors.NewRequestEntityTooLargeError(fmt.Sprintf("limit is %d", limit)))
+ 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 {
- klog.Errorf("failed to marshal upgrade info: %v", err)
+ if err = json.Unmarshal(body, &resp); err != nil {
+ err = response.WriteError(http.StatusBadRequest, fmt.Errorf("failed to marshal task info: %v", err))
+ if err != nil {
+ klog.Warning(err.Error())
+ }
return
}
+ newResp := commontypes.NodeTaskResponse{
+ NodeName: resp.NodeName,
+ Event: "Upgrade",
+ Action: api.ActionSuccess,
+ Reason: resp.Reason,
+ }
- msg := beehiveModel.NewMessage("").SetRoute(modules.CloudHubModuleName, modules.CloudHubModuleName).
- SetResourceOperation(fmt.Sprintf("%s/%s/node/%s", controller.NodeUpgrade, resp.UpgradeID, resp.NodeName), controller.NodeUpgrade).FillBody(resp)
- beehiveContext.Send(modules.NodeUpgradeJobControllerModuleName, *msg)
+ if resp.Status == UpgradeFailedRollbackSuccess {
+ newResp.Event = "Rollback"
+ newResp.Action = api.ActionFailure
+ }
+
+ if resp.Status == UpgradeFailedRollbackFailed {
+ newResp.Event = "Rollback"
+ newResp.Action = api.ActionSuccess
+ }
+
+ msg := beehiveModel.NewMessage("").SetRoute(modules.CloudHubModuleName, modules.CloudHubModuleGroup).
+ SetResourceOperation(fmt.Sprintf("task/%s/node/%s", taskID, nodeID), taskType).FillBody(newResp)
+ 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 cf632345f..58bee37a7 100644
--- a/cloud/pkg/common/messagelayer/context.go
+++ b/cloud/pkg/common/messagelayer/context.go
@@ -104,22 +104,6 @@ func TaskManagerMessageLayer() MessageLayer {
}
}
-func NodeUpgradeJobControllerMessageLayer() MessageLayer {
- return &ContextMessageLayer{
- SendModuleName: modules.CloudHubModuleName,
- ReceiveModuleName: modules.NodeUpgradeJobControllerModuleName,
- ResponseModuleName: modules.CloudHubModuleName,
- }
-}
-
-func ImagePrePullControllerMessageLayer() MessageLayer {
- return &ContextMessageLayer{
- SendModuleName: modules.CloudHubModuleName,
- ReceiveModuleName: modules.ImagePrePullControllerModuleName,
- ResponseModuleName: modules.CloudHubModuleName,
- }
-}
-
func PolicyControllerMessageLayer() MessageLayer {
return &ContextMessageLayer{
SendModuleName: modules.CloudHubModuleName,
diff --git a/cloud/pkg/common/modules/modules.go b/cloud/pkg/common/modules/modules.go
index 98665ccab..12fa9927e 100644
--- a/cloud/pkg/common/modules/modules.go
+++ b/cloud/pkg/common/modules/modules.go
@@ -16,9 +16,6 @@ const (
NodeUpgradeJobControllerModuleName = "nodeupgradejobcontroller"
NodeUpgradeJobControllerModuleGroup = "nodeupgradejobcontroller"
- ImagePrePullControllerModuleName = "imageprepullcontroller"
- ImagePrePullControllerModuleGroup = "imageprepullcontroller"
-
TaskManagerModuleName = "taskmanager"
TaskManagerModuleGroup = "taskmanager"
diff --git a/cloud/pkg/imageprepullcontroller/config/config.go b/cloud/pkg/imageprepullcontroller/config/config.go
deleted file mode 100644
index 88572a842..000000000
--- a/cloud/pkg/imageprepullcontroller/config/config.go
+++ /dev/null
@@ -1,38 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package config
-
-import (
- "sync"
-
- "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
-)
-
-var Config Configure
-var once sync.Once
-
-type Configure struct {
- v1alpha1.ImagePrePullController
-}
-
-func InitConfigure(dc *v1alpha1.ImagePrePullController) {
- once.Do(func() {
- Config = Configure{
- ImagePrePullController: *dc,
- }
- })
-}
diff --git a/cloud/pkg/imageprepullcontroller/controller/downstream.go b/cloud/pkg/imageprepullcontroller/controller/downstream.go
deleted file mode 100644
index d82c58413..000000000
--- a/cloud/pkg/imageprepullcontroller/controller/downstream.go
+++ /dev/null
@@ -1,268 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package controller
-
-import (
- "fmt"
- "time"
-
- metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
- "k8s.io/apimachinery/pkg/util/wait"
- "k8s.io/apimachinery/pkg/watch"
- 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"
- "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/common/modules"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/util"
- "github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller/manager"
- commontypes "github.com/kubeedge/kubeedge/common/types"
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
- crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
- crdinformers "github.com/kubeedge/kubeedge/pkg/client/informers/externalversions"
-)
-
-type DownstreamController struct {
- kubeClient kubernetes.Interface
- informer k8sinformer.SharedInformerFactory
- crdClient crdClientset.Interface
- messageLayer messagelayer.MessageLayer
-
- imagePrePullJobManager *manager.ImagePrePullJobManager
-}
-
-// Start DownstreamController
-func (dc *DownstreamController) Start() error {
- klog.Info("Start ImagePrePullJob Downstream Controller")
- go dc.syncImagePrePullJob()
- return nil
-}
-
-// syncImagePrePullJob is used to get events from informer
-func (dc *DownstreamController) syncImagePrePullJob() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("stop sync ImagePrePullJob")
- return
- case e := <-dc.imagePrePullJobManager.Events():
- imagePrePull, ok := e.Object.(*v1alpha1.ImagePrePullJob)
- if !ok {
- klog.Warningf("object type: %T unsupported", e.Object)
- continue
- }
- switch e.Type {
- case watch.Added:
- dc.imagePrePullJobAdded(imagePrePull)
- case watch.Deleted:
- dc.imagePrePullJobDeleted(imagePrePull)
- case watch.Modified:
- dc.imagePrePullJobUpdate(imagePrePull)
- default:
- klog.Warningf("ImagePrePullJob event type: %s unsupported", e.Type)
- }
- }
- }
-}
-
-// imagePrePullJobAdded is used to process addition of new ImagePrePullJob in apiserver
-func (dc *DownstreamController) imagePrePullJobAdded(imagePrePull *v1alpha1.ImagePrePullJob) {
- klog.V(4).Infof("add ImagePrePullJob: %v", imagePrePull)
- // store in cache map
- dc.imagePrePullJobManager.ImagePrePullMap.Store(imagePrePull.Name, imagePrePull)
-
- // If imagePrePullJob is not initial state, we don't need to send message\
- if imagePrePull.Status.State != v1alpha1.PrePullInitialValue {
- klog.Errorf("The ImagePrePullJob %s is already running or completed, don't send message again", imagePrePull.Name)
- return
- }
-
- // get node list that need prepull images
- var nodesToPrePullImage []string
- if len(imagePrePull.Spec.ImagePrePullTemplate.NodeNames) != 0 {
- for _, node := range imagePrePull.Spec.ImagePrePullTemplate.NodeNames {
- nodeInfo, err := dc.informer.Core().V1().Nodes().Lister().Get(node)
- if err != nil {
- klog.Errorf("Failed to get node(%s) info: %v", node, err)
- continue
- }
-
- if validateNode(nodeInfo) {
- nodesToPrePullImage = append(nodesToPrePullImage, nodeInfo.Name)
- }
- }
- } else if imagePrePull.Spec.ImagePrePullTemplate.LabelSelector != nil {
- selector, err := metav1.LabelSelectorAsSelector(imagePrePull.Spec.ImagePrePullTemplate.LabelSelector)
- if err != nil {
- klog.Errorf("LabelSelector(%s) is not valid: %v", imagePrePull.Spec.ImagePrePullTemplate.LabelSelector, err)
- return
- }
-
- nodes, err := dc.informer.Core().V1().Nodes().Lister().List(selector)
- if err != nil {
- klog.Errorf("Failed to get nodes with label %s: %v", selector.String(), err)
- return
- }
-
- for _, node := range nodes {
- if validateNode(node) {
- nodesToPrePullImage = append(nodesToPrePullImage, node.Name)
- }
- }
- }
-
- // deduplicate: remove duplicate nodes to avoid repeating prepull images on the same node
- nodesToPrePullImage = util.RemoveDuplicateElement(nodesToPrePullImage)
-
- klog.Infof("Filtered finished, images will be prepulled on below nodes\n%v\n", nodesToPrePullImage)
-
- go func() {
- for _, node := range nodesToPrePullImage {
- dc.processPrePull(node, imagePrePull)
- }
- }()
-}
-
-// imagePrePullJobDeleted is used to process deleted ImagePrePullJob in apiServer
-func (dc *DownstreamController) imagePrePullJobDeleted(imagePrePull *v1alpha1.ImagePrePullJob) {
- // delete drom cache map
- dc.imagePrePullJobManager.ImagePrePullMap.Delete(imagePrePull.Name)
-}
-
-// imagePrePullJobUpdate is used to process update of ImagePrePullJob in apiServer
-// Now we don't allow update spec, so we only update the cache map in imagePrePullJobUpdate func.
-func (dc *DownstreamController) imagePrePullJobUpdate(imagePrePull *v1alpha1.ImagePrePullJob) {
- _, ok := dc.imagePrePullJobManager.ImagePrePullMap.Load(imagePrePull.Name)
- // store in cache map
- dc.imagePrePullJobManager.ImagePrePullMap.Store(imagePrePull.Name, imagePrePull)
- if !ok {
- klog.Infof("ImagePrePull Job %s not exist, and store it into first", imagePrePull.Name)
- // If ImagePrePullJob not present in ImagePrePull map means it is not modified and added.
- dc.imagePrePullJobAdded(imagePrePull)
- }
-}
-
-// processPrePull process prepull job added and send it to edge nodes.
-func (dc *DownstreamController) processPrePull(node string, imagePrePull *v1alpha1.ImagePrePullJob) {
- klog.V(4).Infof("begin to send imagePrePull message to %s", node)
-
- imagePrePullTemplateInfo := imagePrePull.Spec.ImagePrePullTemplate
- imagePrePullRequest := commontypes.ImagePrePullJobRequest{
- Images: imagePrePullTemplateInfo.Images,
- NodeName: node,
- Secret: imagePrePullTemplateInfo.ImageSecret,
- RetryTimes: imagePrePullTemplateInfo.RetryTimes,
- CheckItems: imagePrePullTemplateInfo.CheckItems,
- }
-
- // handle timeout: could not receive image prepull msg feedback from edge node
- // send prepull timeout response message to upstream
- go dc.handleImagePrePullJobTimeoutOnEachNode(node, imagePrePull.Name, imagePrePullTemplateInfo.TimeoutSecondsOnEachNode)
-
- // send prepull msg to edge node
- msg := model.NewMessage("")
- resource := buildPrePullResource(imagePrePull.Name, node)
- msg.BuildRouter(modules.ImagePrePullControllerModuleName, modules.ImagePrePullControllerModuleGroup, resource, ImagePrePull).
- FillBody(imagePrePullRequest)
-
- err := dc.messageLayer.Send(*msg)
- if err != nil {
- klog.Errorf("Failed to send prepull message %s due to error %v", msg.GetID(), err)
- return
- }
-
- // update imagePrePullJob status to prepulling
- status := &v1alpha1.ImagePrePullStatus{
- NodeName: node,
- State: v1alpha1.PrePulling,
- }
- err = patchImagePrePullStatus(dc.crdClient, imagePrePull, status)
- if err != nil {
- klog.Errorf("Failed to mark imagePrePullJob prepulling status: %v", err)
- }
-}
-
-// handleImagePrePullJobTimeoutOnEachNode is used to handle the situation that cloud don't receive prepull result
-// from edge node within the timeout period.
-// If so, the ImagePrePullJobState will update to timeout
-func (dc *DownstreamController) handleImagePrePullJobTimeoutOnEachNode(node, jobName string, timeoutSecondsEachNode *uint32) {
- var timeout uint32 = 360
- if timeoutSecondsEachNode != nil && *timeoutSecondsEachNode != 0 {
- timeout = *timeoutSecondsEachNode
- }
-
- receiveFeedback := false
-
- _ = wait.Poll(10*time.Second, time.Duration(timeout)*time.Second, func() (bool, error) {
- cacheValue, ok := dc.imagePrePullJobManager.ImagePrePullMap.Load(jobName)
- if !ok {
- receiveFeedback = true
- klog.Errorf("ImagePrePullJob %s is not exist", jobName)
- return false, fmt.Errorf("imagePrePullJob %s is not exist", jobName)
- }
- imagePrePullValue := cacheValue.(*v1alpha1.ImagePrePullJob)
- for _, statusValue := range imagePrePullValue.Status.Status {
- if statusValue.NodeName == node && (statusValue.State == v1alpha1.PrePullSuccessful || statusValue.State == v1alpha1.PrePullFailed) {
- receiveFeedback = true
- return true, nil
- }
- break
- }
- return false, nil
- })
-
- if receiveFeedback {
- return
- }
- klog.Errorf("TIMEOUT to receive image prepull %s response from edge node %s", jobName, node)
-
- // construct timeout image prepull response and send it to upstream controller
- responseResource := buildPrePullResource(jobName, node)
- resp := commontypes.ImagePrePullJobResponse{
- NodeName: node,
- State: v1alpha1.PrePullFailed,
- Reason: "timeout to receive response from edge",
- }
-
- respMsg := model.NewMessage("").
- BuildRouter(modules.ImagePrePullControllerModuleName, modules.ImagePrePullControllerModuleGroup, responseResource, ImagePrePull).
- FillBody(resp)
- beehiveContext.Send(modules.ImagePrePullControllerModuleName, *respMsg)
-}
-
-// NewDownstreamController new downstream controller to process downstream imageprepull msg to edge nodes.
-func NewDownstreamController(crdInformerFactory crdinformers.SharedInformerFactory) (*DownstreamController, error) {
- imagePrePullJobManager, err := manager.NewImagePrePullJobManager(crdInformerFactory.Operations().V1alpha1().ImagePrePullJobs().Informer())
- if err != nil {
- klog.Warningf("Create ImagePrePullJob manager failed with error: %s", err)
- return nil, err
- }
-
- dc := &DownstreamController{
- kubeClient: client.GetKubeClient(),
- informer: informers.GetInformersManager().GetKubeInformerFactory(),
- crdClient: client.GetCRDClient(),
- imagePrePullJobManager: imagePrePullJobManager,
- messageLayer: messagelayer.ImagePrePullControllerMessageLayer(),
- }
- return dc, nil
-}
diff --git a/cloud/pkg/imageprepullcontroller/controller/upstream.go b/cloud/pkg/imageprepullcontroller/controller/upstream.go
deleted file mode 100644
index 98fd28bea..000000000
--- a/cloud/pkg/imageprepullcontroller/controller/upstream.go
+++ /dev/null
@@ -1,147 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package controller
-
-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/imageprepullcontroller/config"
- "github.com/kubeedge/kubeedge/common/types"
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
- crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
-)
-
-// UpstreamController subscribe messages from edge and sync to k8s api server
-type UpstreamController struct {
- // downstream controller to update ImagePrePullJob status in cache
- dc *DownstreamController
-
- kubeClient kubernetes.Interface
- informer k8sinformer.SharedInformerFactory
- crdClient crdClientset.Interface
- messageLayer messagelayer.MessageLayer
- // message channel
- imagePrePullJobStatusChan chan model.Message
-}
-
-// Start imageprepull UpstreamController
-func (uc *UpstreamController) Start() error {
- klog.Info("Start ImagePrePullJob Upstream Controller")
-
- uc.imagePrePullJobStatusChan = make(chan model.Message, config.Config.Buffer.ImagePrePullJobStatus)
- go uc.dispatchMessage()
-
- for i := 0; i < int(config.Config.Load.ImagePrePullJobWorkers); i++ {
- go uc.updateImagePrePullJobStatus()
- }
- return nil
-}
-
-// updateImagePrePullJobStatus update imagePrePullJob status
-func (uc *UpstreamController) updateImagePrePullJobStatus() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("Stop update ImagePrePullJob status")
- return
- case msg := <-uc.imagePrePullJobStatusChan:
- klog.V(4).Infof("Message: %s, operation is: %s, and resource is: %s", msg.GetID(), msg.GetOperation(), msg.GetResource())
- nodeName, jobName, err := parsePrePullresource(msg.GetResource())
- if err != nil {
- klog.Errorf("message resource %s is not supported", msg.GetResource())
- continue
- }
-
- oldValue, ok := uc.dc.imagePrePullJobManager.ImagePrePullMap.Load(jobName)
- if !ok {
- klog.Errorf("ImagePrePullJob %s not exist", jobName)
- continue
- }
-
- imagePrePull, ok := oldValue.(*v1alpha1.ImagePrePullJob)
- if !ok {
- klog.Errorf("ImagePrePullJob info %T is not valid", oldValue)
- continue
- }
-
- data, err := msg.GetContentData()
- if err != nil {
- klog.Errorf("failed to get image prepull content from response msg, err: %v", err)
- continue
- }
- resp := &types.ImagePrePullJobResponse{}
- err = json.Unmarshal(data, resp)
- if err != nil {
- klog.Errorf("Failed to unmarshal image prepull response: %v", err)
- continue
- }
-
- status := &v1alpha1.ImagePrePullStatus{
- NodeName: nodeName,
- State: resp.State,
- Reason: resp.Reason,
- ImageStatus: resp.ImageStatus,
- }
- err = patchImagePrePullStatus(uc.crdClient, imagePrePull, status)
- if err != nil {
- klog.Errorf("Failed to patch ImagePrePullJob status, err: %v", err)
- }
- }
- }
-}
-
-// dispatchMessage receive the message from edge and write into imagePrePullJobStatusChan
-func (uc *UpstreamController) dispatchMessage() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("Stop dispatch ImagePrePullJob upstream message")
- return
- default:
- }
-
- msg, err := uc.messageLayer.Receive()
- if err != nil {
- klog.Warningf("Receive message failed, %v", err)
- continue
- }
-
- klog.V(4).Infof("ImagePrePullJob upstream controller receive msg %#v", msg)
- uc.imagePrePullJobStatusChan <- msg
- }
-}
-
-// 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.ImagePrePullControllerMessageLayer(),
- dc: dc,
- }
- return uc, nil
-}
diff --git a/cloud/pkg/imageprepullcontroller/controller/util.go b/cloud/pkg/imageprepullcontroller/controller/util.go
deleted file mode 100644
index 592164833..000000000
--- a/cloud/pkg/imageprepullcontroller/controller/util.go
+++ /dev/null
@@ -1,131 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package controller
-
-import (
- "context"
- "encoding/json"
- "fmt"
- "strings"
-
- jsonpatch "github.com/evanphx/json-patch"
- v1 "k8s.io/api/core/v1"
- metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
- apimachineryType "k8s.io/apimachinery/pkg/types"
- "k8s.io/klog/v2"
-
- "github.com/kubeedge/kubeedge/cloud/pkg/common/util"
- "github.com/kubeedge/kubeedge/common/constants"
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
- crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
-)
-
-const ImagePrePull = "prepull"
-
-func validateNode(node *v1.Node) bool {
- if !util.IsEdgeNode(node) {
- klog.Warningf("Node(%s) is not edge node", node.Name)
- return false
- }
-
- // if node is in NotReady state, cannot prepull images
- 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
-}
-
-// buildPrePullResource build prepull resource in msg send to edge node
-func buildPrePullResource(imagePrePullName, nodeName string) string {
- resource := fmt.Sprintf("%s/%s/%s/%s", "node", nodeName, ImagePrePull, imagePrePullName)
- return resource
-}
-
-func parsePrePullresource(resource string) (string, string, error) {
- var nodeName, jobName string
- sli := strings.Split(resource, constants.ResourceSep)
- if len(sli) != 4 {
- return nodeName, jobName, fmt.Errorf("the resource %s is not the standard type", resource)
- }
- return sli[1], sli[3], nil
-}
-
-func patchImagePrePullStatus(crdClient crdClientset.Interface, imagePrePull *v1alpha1.ImagePrePullJob, status *v1alpha1.ImagePrePullStatus) error {
- oldValue := imagePrePull.DeepCopy()
- newValue := updateNodeImagePrePullStatus(oldValue, status)
-
- var completeFlag int
- var failedFlag bool
- newValue.Status.State = v1alpha1.PrePulling
- for _, statusValue := range newValue.Status.Status {
- if statusValue.State == v1alpha1.PrePullFailed {
- failedFlag = true
- completeFlag++
- }
- if statusValue.State == v1alpha1.PrePullSuccessful {
- completeFlag++
- }
- }
- if completeFlag == len(newValue.Status.Status) {
- if failedFlag {
- newValue.Status.State = v1alpha1.PrePullFailed
- } else {
- newValue.Status.State = v1alpha1.PrePullSuccessful
- }
- }
-
- oldData, err := json.Marshal(oldValue)
- if err != nil {
- return fmt.Errorf("failed to marshal the old ImagePrePullJob(%s): %v", oldValue.Name, err)
- }
-
- newData, err := json.Marshal(newValue)
- if err != nil {
- return fmt.Errorf("failed to marshal the new ImagePrePullJob(%s): %v", newValue.Name, err)
- }
-
- patchBytes, err := jsonpatch.CreateMergePatch(oldData, newData)
- if err != nil {
- return fmt.Errorf("failed to create a merge patch: %v", err)
- }
-
- _, err = crdClient.OperationsV1alpha1().ImagePrePullJobs().Patch(context.TODO(), newValue.Name, apimachineryType.MergePatchType, patchBytes, metav1.PatchOptions{}, "status")
- if err != nil {
- return fmt.Errorf("failed to patch update ImagePrePullJob status: %v", err)
- }
-
- return nil
-}
-
-func updateNodeImagePrePullStatus(imagePrePull *v1alpha1.ImagePrePullJob, status *v1alpha1.ImagePrePullStatus) *v1alpha1.ImagePrePullJob {
- // return value imageprepull cannot populate the input parameter old
- newValue := imagePrePull.DeepCopy()
-
- for index, nodeStatus := range newValue.Status.Status {
- if nodeStatus.NodeName == status.NodeName {
- newValue.Status.Status[index] = *status
- return newValue
- }
- }
-
- newValue.Status.Status = append(newValue.Status.Status, *status)
- return newValue
-}
diff --git a/cloud/pkg/imageprepullcontroller/imageprepullcontroller.go b/cloud/pkg/imageprepullcontroller/imageprepullcontroller.go
deleted file mode 100644
index 4f772d258..000000000
--- a/cloud/pkg/imageprepullcontroller/imageprepullcontroller.go
+++ /dev/null
@@ -1,91 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package imageprepullcontroller
-
-import (
- "time"
-
- "k8s.io/klog/v2"
-
- "github.com/kubeedge/beehive/pkg/core"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/informers"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
- "github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller/config"
- "github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller/controller"
- "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
-)
-
-// ImagePrePullController is controller for processing prepull images on edge nodes
-type ImagePrePullController struct {
- downstream *controller.DownstreamController
- upstream *controller.UpstreamController
- enable bool
-}
-
-var _ core.Module = (*ImagePrePullController)(nil)
-
-func newImagePrePullController(enable bool) *ImagePrePullController {
- if !enable {
- return &ImagePrePullController{enable: enable}
- }
- downstream, err := controller.NewDownstreamController(informers.GetInformersManager().GetKubeEdgeInformerFactory())
- if err != nil {
- klog.Exitf("New ImagePrePull Controller downstream failed with error: %s", err)
- }
- upstream, err := controller.NewUpstreamController(downstream)
- if err != nil {
- klog.Exitf("New ImagePrePull Controller upstream failed with error: %s", err)
- }
- return &ImagePrePullController{
- downstream: downstream,
- upstream: upstream,
- enable: enable,
- }
-}
-
-func Register(dc *v1alpha1.ImagePrePullController) {
- config.InitConfigure(dc)
- core.Register(newImagePrePullController(dc.Enable))
-}
-
-// Name of controller
-func (uc *ImagePrePullController) Name() string {
- return modules.ImagePrePullControllerModuleName
-}
-
-// Group of controller
-func (uc *ImagePrePullController) Group() string {
- return modules.ImagePrePullControllerModuleGroup
-}
-
-// Enable indicates whether enable this module
-func (uc *ImagePrePullController) Enable() bool {
- return uc.enable
-}
-
-// Start controller
-func (uc *ImagePrePullController) Start() {
- if err := uc.downstream.Start(); err != nil {
- klog.Exitf("start ImagePrePullJob controller downstream failed with error: %s", err)
- }
- // wait for downstream controller to start and load ImagePrePullJob
- // TODO think about sync
- time.Sleep(1 * time.Second)
- if err := uc.upstream.Start(); err != nil {
- klog.Exitf("start ImagePrePullJob controller upstream failed with error: %s", err)
- }
-}
diff --git a/cloud/pkg/imageprepullcontroller/manager/common.go b/cloud/pkg/imageprepullcontroller/manager/common.go
deleted file mode 100644
index 9f371875a..000000000
--- a/cloud/pkg/imageprepullcontroller/manager/common.go
+++ /dev/null
@@ -1,62 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package manager
-
-import (
- "k8s.io/apimachinery/pkg/runtime"
- "k8s.io/apimachinery/pkg/watch"
- "k8s.io/klog/v2"
-)
-
-// Manager define the interface of a Manager, ImagePrePullJob Manager implement it
-type Manager interface {
- Events() chan watch.Event
-}
-
-// CommonResourceEventHandler can be used by ImagePrePullJob 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 ImagePrePullJob Manager
-func NewCommonResourceEventHandler(events chan watch.Event) *CommonResourceEventHandler {
- return &CommonResourceEventHandler{events: events}
-}
diff --git a/cloud/pkg/imageprepullcontroller/manager/imageprepull.go b/cloud/pkg/imageprepullcontroller/manager/imageprepull.go
deleted file mode 100644
index cdb2eaccc..000000000
--- a/cloud/pkg/imageprepullcontroller/manager/imageprepull.go
+++ /dev/null
@@ -1,52 +0,0 @@
-/*
-Copyright 2023 The KubeEdge Authors.
-
-Licensed under the Apache License, Version 2.0 (the "License");
-you may not use this file except in compliance with the License.
-You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
-*/
-
-package manager
-
-import (
- "sync"
-
- "k8s.io/apimachinery/pkg/watch"
- "k8s.io/client-go/tools/cache"
-
- "github.com/kubeedge/kubeedge/cloud/pkg/imageprepullcontroller/config"
-)
-
-// ImagePrePullJobManager is a manager watch ImagePrePullJob change event
-type ImagePrePullJobManager struct {
- // events from watch kubernetes api server
- events chan watch.Event
-
- // ImagePrePullMap, key is ImagePrePullJob.Name, value is *v1alpha1.ImagePrePullJob{}
- ImagePrePullMap sync.Map
-}
-
-// Events return a channel, can receive all ImagePrePullJob event
-func (dmm *ImagePrePullJobManager) Events() chan watch.Event {
- return dmm.events
-}
-
-// NewImagePrePullJobManager create ImagePrePullJobManager from config
-func NewImagePrePullJobManager(si cache.SharedIndexInformer) (*ImagePrePullJobManager, error) {
- events := make(chan watch.Event, config.Config.Buffer.ImagePrePullJobEvent)
- rh := NewCommonResourceEventHandler(events)
- _, err := si.AddEventHandler(rh)
- if err != nil {
- return nil, err
- }
-
- return &ImagePrePullJobManager{events: events}, nil
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/config/config.go b/cloud/pkg/nodeupgradejobcontroller/config/config.go
deleted file mode 100644
index b98f9ea89..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/config/config.go
+++ /dev/null
@@ -1,38 +0,0 @@
-/*
-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.NodeUpgradeJobController
-}
-
-func InitConfigure(dc *v1alpha1.NodeUpgradeJobController) {
- once.Do(func() {
- Config = Configure{
- NodeUpgradeJobController: *dc,
- }
- })
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go b/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go
deleted file mode 100644
index 6efee8012..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go
+++ /dev/null
@@ -1,447 +0,0 @@
-/*
-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 controller
-
-import (
- "context"
- "encoding/json"
- "fmt"
- "time"
-
- "github.com/google/uuid"
- v1 "k8s.io/api/core/v1"
- metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
- apimachineryType "k8s.io/apimachinery/pkg/types"
- "k8s.io/apimachinery/pkg/util/wait"
- "k8s.io/apimachinery/pkg/watch"
- 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"
- "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/common/modules"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/util"
- "github.com/kubeedge/kubeedge/cloud/pkg/nodeupgradejobcontroller/manager"
- "github.com/kubeedge/kubeedge/common/constants"
- commontypes "github.com/kubeedge/kubeedge/common/types"
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
- crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
- crdinformers "github.com/kubeedge/kubeedge/pkg/client/informers/externalversions"
-)
-
-type DownstreamController struct {
- kubeClient kubernetes.Interface
- informer k8sinformer.SharedInformerFactory
- crdClient crdClientset.Interface
- messageLayer messagelayer.MessageLayer
-
- nodeUpgradeJobManager *manager.NodeUpgradeJobManager
-}
-
-// Start DownstreamController
-func (dc *DownstreamController) Start() error {
- klog.Info("Start NodeUpgradeJob Downstream Controller")
-
- go dc.syncNodeUpgradeJob()
-
- return nil
-}
-
-// syncNodeUpgradeJob is used to get events from informer
-func (dc *DownstreamController) syncNodeUpgradeJob() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("stop sync NodeUpgradeJob")
- return
- case e := <-dc.nodeUpgradeJobManager.Events():
- upgrade, ok := e.Object.(*v1alpha1.NodeUpgradeJob)
- if !ok {
- klog.Warningf("object type: %T unsupported", e.Object)
- continue
- }
- switch e.Type {
- case watch.Added:
- dc.nodeUpgradeJobAdded(upgrade)
- case watch.Deleted:
- dc.nodeUpgradeJobDeleted(upgrade)
- case watch.Modified:
- dc.nodeUpgradeJobUpdated(upgrade)
- default:
- klog.Warningf("NodeUpgradeJob event type: %s unsupported", e.Type)
- }
- }
- }
-}
-
-func buildUpgradeResource(upgradeID, nodeID string) string {
- resource := fmt.Sprintf("%s%s%s%s%s%s%s", NodeUpgrade, constants.ResourceSep, upgradeID, constants.ResourceSep, "node", constants.ResourceSep, nodeID)
- return resource
-}
-
-// nodeUpgradeJobAdded is used to process addition of new NodeUpgradeJob in apiserver
-func (dc *DownstreamController) nodeUpgradeJobAdded(upgrade *v1alpha1.NodeUpgradeJob) {
- klog.V(4).Infof("add NodeUpgradeJob: %v", upgrade)
- // store in cache map
- dc.nodeUpgradeJobManager.UpgradeMap.Store(upgrade.Name, upgrade)
-
- // If all or partial edge nodes upgrade is upgrading or completed, we don't need to send upgrade message
- if isCompleted(upgrade) {
- klog.Errorf("The nodeUpgradeJob is already running or completed, don't send upgrade message again")
- return
- }
-
- // get node list that need upgrading
- var nodesToUpgrade []string
- if len(upgrade.Spec.NodeNames) != 0 {
- for _, node := range upgrade.Spec.NodeNames {
- nodeInfo, err := dc.informer.Core().V1().Nodes().Lister().Get(node)
- if err != nil {
- klog.Errorf("Failed to get node(%s) info: %v", node, err)
- continue
- }
-
- if needUpgrade(nodeInfo, upgrade.Spec.Version) {
- nodesToUpgrade = append(nodesToUpgrade, nodeInfo.Name)
- }
- }
- } else if upgrade.Spec.LabelSelector != nil {
- selector, err := metav1.LabelSelectorAsSelector(upgrade.Spec.LabelSelector)
- if err != nil {
- klog.Errorf("LabelSelector(%s) is not valid: %v", upgrade.Spec.LabelSelector, err)
- return
- }
-
- nodes, err := dc.informer.Core().V1().Nodes().Lister().List(selector)
- if err != nil {
- klog.Errorf("Failed to get nodes with label %s: %v", selector.String(), err)
- return
- }
-
- for _, node := range nodes {
- if needUpgrade(node, upgrade.Spec.Version) {
- nodesToUpgrade = append(nodesToUpgrade, node.Name)
- }
- }
- }
-
- // deduplicate: remove duplicate nodes to avoid repeating upgrade to the same node
- nodesToUpgrade = util.RemoveDuplicateElement(nodesToUpgrade)
-
- klog.Infof("Filtered finished, the below nodes are to upgrade\n%v\n", nodesToUpgrade)
-
- // upgrade most `UpgradeJob.Spec.Concurrency` nodes once a time
- nodesChan := make(chan string, upgrade.Spec.Concurrency)
-
- // select nodes to do upgrade operation
- go dc.selectConcurrentNodes(nodesChan, nodesToUpgrade, upgrade)
-
- go func() {
- for node := range nodesChan {
- dc.processUpgrade(node, upgrade)
- }
- }()
-}
-
-// processUpgrade do the upgrade operation on node
-func (dc *DownstreamController) processUpgrade(node string, upgrade *v1alpha1.NodeUpgradeJob) {
- klog.V(4).Infof("begin to upgrade node %s", node)
- // if users specify Image, we'll use upgrade Version as its image tag, even though Image contains tag.
- // if not, we'll use default image: kubeedge/installation-package:${Version}
- var repo string
- var err error
- repo = "kubeedge/installation-package"
- if upgrade.Spec.Image != "" {
- repo, err = GetImageRepo(upgrade.Spec.Image)
- if err != nil {
- klog.Errorf("Image format is not right: %v", err)
- return
- }
- }
- imageTag := upgrade.Spec.Version
- image := fmt.Sprintf("%s:%s", repo, imageTag)
-
- // send upgrade msg to edge node
- msg := model.NewMessage("")
-
- resource := buildUpgradeResource(upgrade.Name, node)
-
- upgradeReq := commontypes.NodeUpgradeJobRequest{
- UpgradeID: upgrade.Name,
- HistoryID: uuid.New().String(),
- UpgradeTool: upgrade.Spec.UpgradeTool,
- Version: upgrade.Spec.Version,
- Image: image,
- }
-
- msg.BuildRouter(modules.NodeUpgradeJobControllerModuleName, modules.NodeUpgradeJobControllerModuleGroup, resource, NodeUpgrade).
- FillBody(upgradeReq)
-
- err = dc.messageLayer.Send(*msg)
- if err != nil {
- klog.Errorf("Failed to send upgrade message %v due to error %v", msg.GetID(), err)
- return
- }
-
- // process time out: could not receive upgrade feedback from edge
- // send upgrade timeout response message to upstream
- go dc.handleNodeUpgradeJobTimeout(node, upgrade.Name, upgrade.Spec.Version, upgradeReq.HistoryID, upgrade.Spec.TimeoutSeconds)
-
- // mark Upgrade state upgrading
- status := &v1alpha1.UpgradeStatus{
- NodeName: node,
- State: v1alpha1.Upgrading,
- History: v1alpha1.History{
- HistoryID: upgradeReq.HistoryID,
- UpgradeTime: time.Now().Format(ISO8601UTC),
- },
- }
- err = patchNodeUpgradeJobStatus(dc.crdClient, upgrade, status)
- if err != nil {
- // not return, continue to mark node unschedulable
- klog.Errorf("Failed to mark Upgrade upgrading status: %v", err)
- }
-
- // mark edge node unschedulable
- // the effect is like running cmd: kubectl drain <node-to-drain> --ignore-daemonsets
- unscheduleNode := v1.Node{}
- unscheduleNode.Spec.Unschedulable = true
-
- // add a upgrade label
- unscheduleNode.Labels = map[string]string{NodeUpgradeJobStatusKey: NodeUpgradeJobStatusValue}
- byteNode, err := json.Marshal(unscheduleNode)
- if err != nil {
- klog.Warningf("marshal data failed: %v", err)
- return
- }
-
- _, err = dc.kubeClient.CoreV1().Nodes().Patch(context.Background(), node, apimachineryType.StrategicMergePatchType, byteNode, metav1.PatchOptions{})
- if err != nil {
- klog.Errorf("failed to drain node %s: %v", node, err)
- return
- }
-}
-
-// selectConcurrentNodes select the nodes to do upgrade operation, and put it into channel nodesChan
-func (dc *DownstreamController) selectConcurrentNodes(nodesChan chan string, allNodes []string, upgrade *v1alpha1.NodeUpgradeJob) {
- // the default concurrency is 1
- // this means that we will upgrade nodes one by one
- // only when the last one node upgrade finished, we'll continue to upgrade the next one node
- concurrency := 1
- if upgrade.Spec.Concurrency != 0 {
- concurrency = int(upgrade.Spec.Concurrency)
- }
-
- timeout := int(*upgrade.Spec.TimeoutSeconds) * (len(allNodes) + 1) / concurrency
-
- err := wait.Poll(10*time.Second, time.Duration(timeout)*time.Second, func() (bool, error) {
- upgradeJob, err := dc.crdClient.OperationsV1alpha1().NodeUpgradeJobs().Get(context.TODO(), upgrade.Name, metav1.GetOptions{})
- if err != nil {
- return false, nil
- }
-
- // if all the nodes upgrade is completed, close channel to inform that the upgrade is finished
- if upgradeJob.Status.State == v1alpha1.Completed {
- klog.Infof("all the nodes upgrade status are completed")
- close(nodesChan)
- return true, nil
- }
-
- // calculate the number of nodes in upgrading operation
- upgradingNum := 0
- for _, status := range upgradeJob.Status.Status {
- if status.State == v1alpha1.Upgrading {
- upgradingNum++
- }
- }
-
- // ensure the max number of upgrading nodes is Concurrency
- if upgradingNum < concurrency {
- for i := 0; i < concurrency-upgradingNum; i++ {
- nodesChan <- allNodes[i]
- }
-
- allNodes = allNodes[:int(concurrency)-upgradingNum]
- }
-
- if len(allNodes) == 0 {
- // all the nodes are to upgrade
- klog.Infof("The number of nodes need to be upgraded reach 0")
- close(nodesChan)
- return true, nil
- }
- return false, nil
- })
-
- if err != nil {
- klog.Errorf("failed to select all the related nodes to do upgrade operation: %v", err)
- close(nodesChan)
- }
-}
-
-func needUpgrade(node *v1.Node, upgradeVersion string) bool {
- if filterVersion(node.Status.NodeInfo.KubeletVersion, upgradeVersion) {
- klog.Warningf("Node(%s) version(%s) already on the expected version %s.", node.Name, node.Status.NodeInfo.KubeletVersion, upgradeVersion)
- return false
- }
-
- // we only care about edge nodes, so just remove not edge nodes
- if !util.IsEdgeNode(node) {
- klog.Warningf("Node(%s) is not edge node", node.Name)
- return false
- }
-
- // if node is in Upgrading state, don't need upgrade
- if _, ok := node.Labels[NodeUpgradeJobStatusKey]; ok {
- klog.Warningf("Node(%s) is in upgrade state", node.Name)
- return false
- }
-
- // if node is in NotReady state, don't need upgrade
- 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
-}
-
-// handleNodeUpgradeJobTimeout is used to handle the situation that cloud don't receive upgrade result from edge node
-// within the timeout period
-func (dc *DownstreamController) handleNodeUpgradeJobTimeout(node string, upgradeID string, upgradeVersion string, historyID string, timeoutSeconds *uint32) {
- // by default, if we don't receive upgrade response in 300s, we think it's timeout
- // if we have specified the timeout in Upgrade, we'll use it as the timeout time
- var timeout uint32 = 300
- if timeoutSeconds != nil && *timeoutSeconds != 0 {
- timeout = *timeoutSeconds
- }
-
- receiveFeedback := false
-
- // check whether edgecore report to the cloud about the upgrade result
- // we don't care about function Poll return error
- // if we don't receive Upgrade response, Poll function also return error: timed out waiting for the condition
- // we only care about variable: receiveFeedback
- _ = wait.Poll(10*time.Second, time.Duration(timeout)*time.Second, func() (bool, error) {
- v, ok := dc.nodeUpgradeJobManager.UpgradeMap.Load(upgradeID)
- if !ok {
- // we think it's receiveFeedback to avoid construct timeout response by ourselves
- receiveFeedback = true
- klog.Errorf("NodeUpgradeJob %v not exist", upgradeID)
- return false, fmt.Errorf("nodeUpgrade %v not exist", upgradeID)
- }
- upgradeValue := v.(*v1alpha1.NodeUpgradeJob)
- for index := range upgradeValue.Status.Status {
- if upgradeValue.Status.Status[index].NodeName == node {
- // if HistoryID matches and state is v1alpha1.Completed
- // it means we've received the specified Upgrade Operation
- if upgradeValue.Status.Status[index].History.HistoryID == historyID &&
- upgradeValue.Status.Status[index].State == v1alpha1.Completed {
- receiveFeedback = true
- return true, nil
- }
- break
- }
- }
- return false, nil
- })
-
- if receiveFeedback {
- // if already receive edge upgrade feedback, do nothing
- return
- }
-
- klog.Errorf("NOT receive node(%s) upgrade(%s) feedback response", node, upgradeID)
-
- // construct timeout upgrade response
- // and send it to upgrade controller upstream
- upgradeResource := buildUpgradeResource(upgradeID, node)
- resp := commontypes.NodeUpgradeJobResponse{
- UpgradeID: upgradeID,
- HistoryID: historyID,
- NodeName: node,
- FromVersion: "",
- ToVersion: upgradeVersion,
- Status: string(v1alpha1.UpgradeFailedRollbackSuccess),
- Reason: "timeout to get upgrade response from edge, maybe error due to cloud or edge",
- }
-
- updateMsg := model.NewMessage("").
- BuildRouter(modules.NodeUpgradeJobControllerModuleName, modules.NodeUpgradeJobControllerModuleGroup, upgradeResource, NodeUpgrade).
- FillBody(resp)
-
- // send upgrade resp message to upgrade controller upstream directly
- // let upgrade controller upstream to update upgrade status
- beehiveContext.Send(modules.NodeUpgradeJobControllerModuleName, *updateMsg)
-}
-
-// nodeUpgradeJobDeleted is used to process deleted NodeUpgradeJob in apiserver
-func (dc *DownstreamController) nodeUpgradeJobDeleted(upgrade *v1alpha1.NodeUpgradeJob) {
- // just need to delete from cache map
- dc.nodeUpgradeJobManager.UpgradeMap.Delete(upgrade.Name)
-}
-
-// upgradeAdded is used to process update of new NodeUpgradeJob in apiserver
-func (dc *DownstreamController) nodeUpgradeJobUpdated(upgrade *v1alpha1.NodeUpgradeJob) {
- oldValue, ok := dc.nodeUpgradeJobManager.UpgradeMap.Load(upgrade.Name)
- // store in cache map
- dc.nodeUpgradeJobManager.UpgradeMap.Store(upgrade.Name, upgrade)
- if !ok {
- klog.Infof("Upgrade %s not exist, and store it first", upgrade.Name)
- // If Upgrade not present in Upgrade map means it is not modified and added.
- dc.nodeUpgradeJobAdded(upgrade)
- return
- }
-
- old := oldValue.(*v1alpha1.NodeUpgradeJob)
- if !isUpgradeUpdated(upgrade, old) {
- klog.V(4).Infof("Upgrade %s no need to update", upgrade.Name)
- return
- }
-
- dc.nodeUpgradeJobAdded(upgrade)
-}
-
-// isUpgradeUpdated checks Upgrade is actually updated or not
-func isUpgradeUpdated(new *v1alpha1.NodeUpgradeJob, old *v1alpha1.NodeUpgradeJob) bool {
- // now we don't allow update spec fields
- // so always return false to avoid sending Upgrade msg to edge again when status fields changed
- return false
-}
-
-func NewDownstreamController(crdInformerFactory crdinformers.SharedInformerFactory) (*DownstreamController, error) {
- nodeUpgradeJobManager, err := manager.NewNodeUpgradeJobManager(crdInformerFactory.Operations().V1alpha1().NodeUpgradeJobs().Informer())
- if err != nil {
- klog.Warningf("Create NodeUpgradeJob manager failed with error: %s", err)
- return nil, err
- }
-
- dc := &DownstreamController{
- kubeClient: client.GetKubeClient(),
- informer: informers.GetInformersManager().GetKubeInformerFactory(),
- crdClient: client.GetCRDClient(),
- nodeUpgradeJobManager: nodeUpgradeJobManager,
- messageLayer: messagelayer.NodeUpgradeJobControllerMessageLayer(),
- }
- return dc, nil
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/controller/upstream.go b/cloud/pkg/nodeupgradejobcontroller/controller/upstream.go
deleted file mode 100644
index 6b3e53ffb..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/controller/upstream.go
+++ /dev/null
@@ -1,244 +0,0 @@
-/*
-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 controller
-
-import (
- "context"
- "encoding/json"
- "fmt"
- "strings"
-
- jsonpatch "github.com/evanphx/json-patch"
- metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
- apimachineryType "k8s.io/apimachinery/pkg/types"
- 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/nodeupgradejobcontroller/config"
- "github.com/kubeedge/kubeedge/common/types"
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
- crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
-)
-
-// 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
- nodeUpgradeJobStatusChan chan model.Message
-}
-
-// Start UpstreamController
-func (uc *UpstreamController) Start() error {
- klog.Info("Start NodeUpgradeJob Upstream Controller")
-
- uc.nodeUpgradeJobStatusChan = make(chan model.Message, config.Config.Buffer.UpdateNodeUpgradeJobStatus)
- go uc.dispatchMessage()
-
- for i := 0; i < int(config.Config.Load.NodeUpgradeJobWorkers); i++ {
- go uc.updateNodeUpgradeJobStatus()
- }
- return nil
-}
-
-// Start UpstreamController
-func (uc *UpstreamController) dispatchMessage() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("Stop dispatch NodeUpgradeJob upstream message")
- return
- default:
- }
-
- msg, err := uc.messageLayer.Receive()
- if err != nil {
- klog.Warningf("Receive message failed, %v", err)
- continue
- }
-
- klog.V(4).Infof("NodeUpgradeJob upstream controller receive msg %#v", msg)
-
- uc.nodeUpgradeJobStatusChan <- msg
- }
-}
-
-// updateNodeUpgradeJobStatus update NodeUpgradeJob status field
-func (uc *UpstreamController) updateNodeUpgradeJobStatus() {
- for {
- select {
- case <-beehiveContext.Done():
- klog.Info("Stop update NodeUpgradeJob status")
- return
- case msg := <-uc.nodeUpgradeJobStatusChan:
- 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 := getNodeName(msg.GetResource())
- upgradeID := getUpgradeID(msg.GetResource())
-
- oldValue, ok := uc.dc.nodeUpgradeJobManager.UpgradeMap.Load(upgradeID)
- if !ok {
- klog.Errorf("NodeUpgradeJob %s not exist", upgradeID)
- continue
- }
-
- upgrade, ok := oldValue.(*v1alpha1.NodeUpgradeJob)
- if !ok {
- klog.Errorf("NodeUpgradeJob info %T is not valid", oldValue)
- continue
- }
-
- data, err := msg.GetContentData()
- if err != nil {
- klog.Errorf("failed to get node upgrade content data: %v", err)
- continue
- }
- resp := &types.NodeUpgradeJobResponse{}
- err = json.Unmarshal(data, resp)
- if err != nil {
- klog.Errorf("Failed to unmarshal node upgrade response: %v", err)
- continue
- }
-
- status := &v1alpha1.UpgradeStatus{
- NodeName: nodeID,
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- HistoryID: resp.HistoryID,
- FromVersion: resp.FromVersion,
- ToVersion: resp.ToVersion,
- Result: v1alpha1.UpgradeResult(resp.Status),
- Reason: resp.Reason,
- },
- }
- err = patchNodeUpgradeJobStatus(uc.crdClient, upgrade, status)
- if err != nil {
- klog.Errorf("Failed to mark NodeUpgradeJob status to completed: %v", err)
- }
-
- // The below are to mark edge node schedulable
- // And to keep a successful record in node annotation only when upgrade is successful
- // like: nodeupgradejob.operations.kubeedge.io/history: "v1.9.0->v1.10.0->v1.11.1"
- nodeInfo, err := uc.informer.Core().V1().Nodes().Lister().Get(nodeID)
- if err != nil {
- klog.Errorf("Failed to get node info: %v", err)
- continue
- }
-
- // mark edge node schedulable
- // the effect is like running cmd: kubectl uncordon <node-to-drain>
- if nodeInfo.Labels != nil {
- if value, ok := nodeInfo.Labels[NodeUpgradeJobStatusKey]; ok {
- if value == NodeUpgradeJobStatusValue {
- nodeInfo.Spec.Unschedulable = false
- delete(nodeInfo.Labels, NodeUpgradeJobStatusKey)
- }
- }
- }
- // record upgrade logs in node annotation
- if v1alpha1.UpgradeResult(resp.Status) == v1alpha1.UpgradeSuccess {
- if nodeInfo.Annotations == nil {
- nodeInfo.Annotations = make(map[string]string)
- }
- nodeInfo.Annotations[NodeUpgradeHistoryKey] = mergeAnnotationUpgradeHistory(nodeInfo.Annotations[NodeUpgradeHistoryKey], resp.FromVersion, resp.ToVersion)
- }
- _, err = uc.kubeClient.CoreV1().Nodes().Update(context.Background(), nodeInfo, metav1.UpdateOptions{})
- if err != nil {
- // just log, and continue to process the next step
- klog.Errorf("Failed to mark node schedulable and add upgrade record: %v", err)
- }
- }
- }
-}
-
-// patchNodeUpgradeJobStatus call patch api to patch update NodeUpgradeJob status
-func patchNodeUpgradeJobStatus(crdClient crdClientset.Interface, upgrade *v1alpha1.NodeUpgradeJob, status *v1alpha1.UpgradeStatus) error {
- oldValue := upgrade.DeepCopy()
-
- newValue := UpdateNodeUpgradeJobStatus(oldValue, status)
-
- // after mark each node upgrade state, we also need to judge whether all edge node upgrade is completed
- // if all edge node is in completed state, we should set the total state to completed
- var completed int
- for _, v := range newValue.Status.Status {
- if v.State == v1alpha1.Completed {
- completed++
- }
- }
- if completed == len(newValue.Status.Status) {
- newValue.Status.State = v1alpha1.Completed
- } else {
- newValue.Status.State = v1alpha1.Upgrading
- }
-
- oldData, err := json.Marshal(oldValue)
- if err != nil {
- return fmt.Errorf("failed to marshal the old NodeUpgradeJob(%s): %v", oldValue.Name, err)
- }
-
- newData, err := json.Marshal(newValue)
- if err != nil {
- return fmt.Errorf("failed to marshal the new NodeUpgradeJob(%s): %v", newValue.Name, err)
- }
-
- patchBytes, err := jsonpatch.CreateMergePatch(oldData, newData)
- if err != nil {
- return fmt.Errorf("failed to create a merge patch: %v", err)
- }
-
- _, err = crdClient.OperationsV1alpha1().NodeUpgradeJobs().Patch(context.TODO(), newValue.Name, apimachineryType.MergePatchType, patchBytes, metav1.PatchOptions{}, "status")
- if err != nil {
- return fmt.Errorf("failed to patch update NodeUpgradeJob status: %v", err)
- }
-
- return nil
-}
-
-func getNodeName(resource string) string {
- // upgrade/${UpgradeID}/node/${NodeID}
- s := strings.Split(resource, "/")
- return s[3]
-}
-func getUpgradeID(resource string) string {
- // upgrade/${UpgradeID}/node/${NodeID}
- s := strings.Split(resource, "/")
- return s[1]
-}
-
-// 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.NodeUpgradeJobControllerMessageLayer(),
- dc: dc,
- }
- return uc, nil
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/controller/util.go b/cloud/pkg/nodeupgradejobcontroller/controller/util.go
deleted file mode 100644
index 8b80c4261..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/controller/util.go
+++ /dev/null
@@ -1,124 +0,0 @@
-/*
-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 controller
-
-import (
- "fmt"
- "strings"
- "time"
-
- "github.com/distribution/distribution/v3/reference"
-
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
-)
-
-const (
- NodeUpgradeJobStatusKey = "nodeupgradejob.operations.kubeedge.io/status"
- NodeUpgradeJobStatusValue = ""
- NodeUpgradeHistoryKey = "nodeupgradejob.operations.kubeedge.io/history"
-)
-
-const (
- NodeUpgrade = "upgrade"
-
- ISO8601UTC = "2006-01-02T15:04:05Z"
-)
-
-// 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
- index := strings.Index(version, "-kubeedge-")
- if index == -1 {
- return false
- }
-
- length := len("-kubeedge-")
-
- // filter nodes that already in the required version
- return version[index+length:] == expected
-}
-
-// isCompleted returns true only if some/all edge upgrade is upgrading or completed
-func isCompleted(upgrade *v1alpha1.NodeUpgradeJob) bool {
- // all edge node upgrade is upgrading or completed
- if upgrade.Status.State != v1alpha1.InitialValue {
- return true
- }
-
- // partial edge node upgrade is upgrading or completed
- for _, status := range upgrade.Status.Status {
- if status.State != v1alpha1.InitialValue {
- return true
- }
- }
-
- return false
-}
-
-// UpdateNodeUpgradeJobStatus updates the status
-// return the updated result
-func UpdateNodeUpgradeJobStatus(old *v1alpha1.NodeUpgradeJob, status *v1alpha1.UpgradeStatus) *v1alpha1.NodeUpgradeJob {
- // return value upgrade cannot populate the input parameter old
- upgrade := old.DeepCopy()
-
- for index := range upgrade.Status.Status {
- // If Node's Upgrade info exist, just overwrite
- if upgrade.Status.Status[index].NodeName == status.NodeName {
- // The input status no upgradeTime, we need set it with old value
- status.History.UpgradeTime = upgrade.Status.Status[index].History.UpgradeTime
- upgrade.Status.Status[index] = *status
- return upgrade
- }
- }
-
- // if Node's Upgrade info not exist, just append
- if status.History.UpgradeTime == "" {
- // If upgrade time is blank, set to the current time
- status.History.UpgradeTime = time.Now().Format(ISO8601UTC)
- }
- upgrade.Status.Status = append(upgrade.Status.Status, *status)
-
- return upgrade
-}
-
-// 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
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/controller/util_test.go b/cloud/pkg/nodeupgradejobcontroller/controller/util_test.go
deleted file mode 100644
index 499da62fa..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/controller/util_test.go
+++ /dev/null
@@ -1,223 +0,0 @@
-/*
-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 controller
-
-import (
- "reflect"
- "testing"
-
- "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
-)
-
-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 TestUpdateUpgradeStatus(t *testing.T) {
- upgrade := v1alpha1.NodeUpgradeJob{
- Status: v1alpha1.NodeUpgradeJobStatus{
- Status: []v1alpha1.UpgradeStatus{
- {
- NodeName: "edge-node",
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- Reason: "the first upgrade",
- UpgradeTime: "2023-09-22T17:33:00Z",
- },
- },
- },
- },
- }
- upgrade2 := upgrade.DeepCopy()
- upgrade2.Status.Status[0].History.Reason = "the second upgrade"
-
- upgrade3 := upgrade.DeepCopy()
- upgrade3.Status.Status = append(upgrade3.Status.Status, v1alpha1.UpgradeStatus{
- NodeName: "edge-node2",
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- Reason: "the first upgrade",
- UpgradeTime: "2023-09-22T17:35:00Z",
- },
- })
-
- tests := []struct {
- name string
- upgrade *v1alpha1.NodeUpgradeJob
- status *v1alpha1.UpgradeStatus
- expected *v1alpha1.NodeUpgradeJob
- }{
- {
- name: "case1: first add one node",
- upgrade: &v1alpha1.NodeUpgradeJob{},
- status: &v1alpha1.UpgradeStatus{
- NodeName: "edge-node",
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- Reason: "the first upgrade",
- UpgradeTime: "2023-09-22T17:33:00Z",
- },
- },
- expected: upgrade.DeepCopy(),
- },
- {
- name: "case2: add to one NOT exist node record",
- upgrade: upgrade.DeepCopy(),
- status: &v1alpha1.UpgradeStatus{
- NodeName: "edge-node2",
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- Reason: "the first upgrade",
- UpgradeTime: "2023-09-22T17:35:00Z",
- },
- },
- expected: upgrade3,
- },
- {
- name: "case3: add to one exist node record",
- upgrade: upgrade.DeepCopy(),
- status: &v1alpha1.UpgradeStatus{
- NodeName: "edge-node",
- State: v1alpha1.Completed,
- History: v1alpha1.History{
- Reason: "the second upgrade",
- },
- },
- expected: upgrade2,
- },
- }
- for _, test := range tests {
- t.Run(test.name, func(t *testing.T) {
- newValue := UpdateNodeUpgradeJobStatus(test.upgrade, test.status)
- if !reflect.DeepEqual(newValue, test.expected) {
- t.Errorf("Got = %v, Want = %v", newValue, 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)
- }
- })
- }
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/manager/common.go b/cloud/pkg/nodeupgradejobcontroller/manager/common.go
deleted file mode 100644
index 6e8008abf..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/manager/common.go
+++ /dev/null
@@ -1,62 +0,0 @@
-/*
-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/nodeupgradejobcontroller/manager/nodeupgrade.go b/cloud/pkg/nodeupgradejobcontroller/manager/nodeupgrade.go
deleted file mode 100644
index 7cdb78e1a..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/manager/nodeupgrade.go
+++ /dev/null
@@ -1,52 +0,0 @@
-/*
-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/nodeupgradejobcontroller/config"
-)
-
-// NodeUpgradeJobManager is a manager watch NodeUpgradeJob change event
-type NodeUpgradeJobManager struct {
- // events from watch kubernetes api server
- events chan watch.Event
-
- // UpgradeMap, key is NodeUpgradeJob.Name, value is *v1alpha1.NodeUpgradeJob{}
- UpgradeMap sync.Map
-}
-
-// Events return a channel, can receive all NodeUpgradeJob event
-func (dmm *NodeUpgradeJobManager) Events() chan watch.Event {
- return dmm.events
-}
-
-// NewNodeUpgradeJobManager create NodeUpgradeJobManager from config
-func NewNodeUpgradeJobManager(si cache.SharedIndexInformer) (*NodeUpgradeJobManager, error) {
- events := make(chan watch.Event, config.Config.Buffer.NodeUpgradeJobEvent)
- rh := NewCommonResourceEventHandler(events)
- _, err := si.AddEventHandler(rh)
- if err != nil {
- return nil, err
- }
-
- return &NodeUpgradeJobManager{events: events}, nil
-}
diff --git a/cloud/pkg/nodeupgradejobcontroller/nodeupgradejobcontroller.go b/cloud/pkg/nodeupgradejobcontroller/nodeupgradejobcontroller.go
deleted file mode 100644
index d9c4d8e4c..000000000
--- a/cloud/pkg/nodeupgradejobcontroller/nodeupgradejobcontroller.go
+++ /dev/null
@@ -1,91 +0,0 @@
-/*
-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 nodeupgradejobcontroller
-
-import (
- "time"
-
- "k8s.io/klog/v2"
-
- "github.com/kubeedge/beehive/pkg/core"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/informers"
- "github.com/kubeedge/kubeedge/cloud/pkg/common/modules"
- "github.com/kubeedge/kubeedge/cloud/pkg/nodeupgradejobcontroller/config"
- "github.com/kubeedge/kubeedge/cloud/pkg/nodeupgradejobcontroller/controller"
- "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/cloudcore/v1alpha1"
-)
-
-// NodeUpgradeJobController is controller for processing upgrading edge node from cloud
-type NodeUpgradeJobController struct {
- downstream *controller.DownstreamController
- upstream *controller.UpstreamController
- enable bool
-}
-
-var _ core.Module = (*NodeUpgradeJobController)(nil)
-
-func newNodeUpgradeJobController(enable bool) *NodeUpgradeJobController {
- if !enable {
- return &NodeUpgradeJobController{enable: enable}
- }
- downstream, err := controller.NewDownstreamController(informers.GetInformersManager().GetKubeEdgeInformerFactory())
- if err != nil {
- klog.Exitf("New NodeUpgradeJob Controller downstream failed with error: %s", err)
- }
- upstream, err := controller.NewUpstreamController(downstream)
- if err != nil {
- klog.Exitf("New NodeUpgradeJob Controller upstream failed with error: %s", err)
- }
- return &NodeUpgradeJobController{
- downstream: downstream,
- upstream: upstream,
- enable: enable,
- }
-}
-
-func Register(dc *v1alpha1.NodeUpgradeJobController) {
- config.InitConfigure(dc)
- core.Register(newNodeUpgradeJobController(dc.Enable))
-}
-
-// 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 uc.enable
-}
-
-// Start controller
-func (uc *NodeUpgradeJobController) Start() {
- if err := uc.downstream.Start(); err != nil {
- klog.Exitf("start NodeUpgradeJob controller 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 NodeUpgradeJob controller upstream failed with error: %s", err)
- }
-}
diff --git a/cloud/pkg/taskmanager/config/config.go b/cloud/pkg/taskmanager/config/config.go
index 1366a7904..1940f14f8 100644
--- a/cloud/pkg/taskmanager/config/config.go
+++ b/cloud/pkg/taskmanager/config/config.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
diff --git a/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_controller.go b/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_controller.go
new file mode 100644
index 000000000..31dc525aa
--- /dev/null
+++ b/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_controller.go
@@ -0,0 +1,338 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
+package imageprepullcontroller
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "os"
+ "strconv"
+ "sync"
+ "time"
+
+ jsonpatch "github.com/evanphx/json-patch"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ apimachineryType "k8s.io/apimachinery/pkg/types"
+ "k8s.io/apimachinery/pkg/util/wait"
+ "k8s.io/apimachinery/pkg/watch"
+ "k8s.io/klog/v2"
+
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ keclient "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/informers"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/controller"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/manager"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
+ 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 ImagePrePullController struct {
+ sync.Mutex
+ *controller.BaseController
+}
+
+var cache *manager.TaskCache
+
+func NewImagePrePullController(messageChan chan util.TaskMessage) (*ImagePrePullController, error) {
+ var err error
+ cache, err = manager.NewTaskCache(
+ informers.GetInformersManager().GetKubeEdgeInformerFactory().Operations().V1alpha1().ImagePrePullJobs().Informer())
+ if err != nil {
+ klog.Warningf("Create image pre pull controller failed with error: %s", err)
+ return nil, err
+ }
+ return &ImagePrePullController{
+ BaseController: &controller.BaseController{
+ Informer: informers.GetInformersManager().GetKubeInformerFactory(),
+ TaskManager: cache,
+ MessageChan: messageChan,
+ CrdClient: client.GetCRDClient(),
+ KubeClient: keclient.GetKubeClient(),
+ },
+ }, nil
+}
+
+func (ndc *ImagePrePullController) ReportNodeStatus(taskID, nodeID string, event fsm.Event) (api.State, error) {
+ nodeFSM := NewImagePrePullNodeFSM(taskID, nodeID)
+ err := nodeFSM.AllowTransit(event)
+ if err != nil {
+ return "", err
+ }
+ state, err := nodeFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ ndc.Lock()
+ defer ndc.Unlock()
+ err = nodeFSM.Transit(event)
+ if err != nil {
+ return "", err
+ }
+ checkStatusChanged(nodeFSM, state)
+ state, err = nodeFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ return state, nil
+}
+
+func checkStatusChanged(nodeFSM *fsm.FSM, state api.State) {
+ err := wait.Poll(100*time.Millisecond, time.Second, func() (bool, error) {
+ nowState, err := nodeFSM.CurrentState()
+ if err != nil {
+ return false, nil
+ }
+ if nowState == state {
+ return false, nil
+ }
+ return true, err
+ })
+ if err != nil {
+ klog.V(4).Infof("check status changed failed: %s", err.Error())
+ }
+}
+
+func (ndc *ImagePrePullController) ReportTaskStatus(taskID string, event fsm.Event) (api.State, error) {
+ taskFSM := NewImagePrePullTaskFSM(taskID)
+ state, err := taskFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ err = taskFSM.AllowTransit(event)
+ if err != nil {
+ return "", err
+ }
+ err = taskFSM.Transit(event)
+ if err != nil {
+ return "", err
+ }
+ checkStatusChanged(taskFSM, state)
+ return taskFSM.CurrentState()
+}
+
+func (ndc *ImagePrePullController) StageCompleted(taskID string, state api.State) bool {
+ taskFSM := NewImagePrePullTaskFSM(taskID)
+ return taskFSM.TaskStagCompleted(state)
+}
+
+func (ndc *ImagePrePullController) GetNodeStatus(name string) ([]v1alpha1.TaskStatus, error) {
+ imagePrePull, err := ndc.CrdClient.OperationsV1alpha1().ImagePrePullJobs().Get(context.TODO(), name, metav1.GetOptions{})
+ if err != nil {
+ return nil, err
+ }
+ statusList := make([]v1alpha1.TaskStatus, len(imagePrePull.Status.Status))
+ for i, status := range imagePrePull.Status.Status {
+ if status.TaskStatus == nil {
+ statusList[i] = v1alpha1.TaskStatus{}
+ continue
+ }
+ statusList[i] = *status.TaskStatus
+ }
+ return statusList, nil
+}
+
+func (ndc *ImagePrePullController) UpdateNodeStatus(name string, nodeStatus []v1alpha1.TaskStatus) error {
+ imagePrePull, err := ndc.CrdClient.OperationsV1alpha1().ImagePrePullJobs().Get(context.TODO(), name, metav1.GetOptions{})
+ if err != nil {
+ return err
+ }
+ status := imagePrePull.Status
+ statusList := make([]v1alpha1.ImagePrePullStatus, len(nodeStatus))
+ for i := 0; i < len(nodeStatus); i++ {
+ statusList[i].TaskStatus = &nodeStatus[i]
+ }
+ status.Status = statusList
+ err = patchStatus(imagePrePull, status, ndc.CrdClient)
+ if err != nil {
+ return err
+ }
+ return nil
+}
+
+func patchStatus(imagePrePullJob *v1alpha1.ImagePrePullJob, status v1alpha1.ImagePrePullJobStatus, crdClient crdClientset.Interface) error {
+ oldData, err := json.Marshal(imagePrePullJob)
+ if err != nil {
+ return fmt.Errorf("failed to marshal the old ImagePrePullJob(%s): %v", imagePrePullJob.Name, err)
+ }
+ imagePrePullJob.Status = status
+ newData, err := json.Marshal(imagePrePullJob)
+ if err != nil {
+ return fmt.Errorf("failed to marshal the new ImagePrePullJob(%s): %v", imagePrePullJob.Name, err)
+ }
+
+ patchBytes, err := jsonpatch.CreateMergePatch(oldData, newData)
+ if err != nil {
+ return fmt.Errorf("failed to create a merge patch: %v", err)
+ }
+
+ result, err := crdClient.OperationsV1alpha1().ImagePrePullJobs().Patch(context.TODO(), imagePrePullJob.Name, apimachineryType.MergePatchType, patchBytes, metav1.PatchOptions{}, "status")
+ if err != nil {
+ return fmt.Errorf("failed to patch update ImagePrePullJob status: %v", err)
+ }
+ klog.V(4).Info("patch update task status result: ", result)
+ return nil
+}
+
+func (ndc *ImagePrePullController) Start() error {
+ go ndc.startSync()
+ return nil
+}
+
+func (ndc *ImagePrePullController) startSync() {
+ imagePrePullList, err := ndc.CrdClient.OperationsV1alpha1().ImagePrePullJobs().List(context.TODO(), metav1.ListOptions{})
+ if err != nil {
+ klog.Errorf(err.Error())
+ os.Exit(2)
+ }
+ for _, imagePrePull := range imagePrePullList.Items {
+ if fsm.TaskFinish(imagePrePull.Status.State) {
+ continue
+ }
+ ndc.imagePrePullJobAdded(&imagePrePull)
+ }
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("stop sync ImagePrePullJob")
+ return
+ case e := <-ndc.TaskManager.Events():
+ prePull, ok := e.Object.(*v1alpha1.ImagePrePullJob)
+ if !ok {
+ klog.Warningf("object type: %T unsupported", e.Object)
+ continue
+ }
+ switch e.Type {
+ case watch.Added:
+ ndc.imagePrePullJobAdded(prePull)
+ case watch.Deleted:
+ ndc.imagePrePullJobDeleted(prePull)
+ case watch.Modified:
+ ndc.imagePrePullJobUpdated(prePull)
+ default:
+ klog.Warningf("ImagePrePullJob event type: %s unsupported", e.Type)
+ }
+ }
+ }
+}
+
+// imagePrePullJobAdded is used to process addition of new ImagePrePullJob in apiserver
+func (ndc *ImagePrePullController) imagePrePullJobAdded(imagePrePull *v1alpha1.ImagePrePullJob) {
+ klog.V(4).Infof("add ImagePrePullJob: %v", imagePrePull)
+ // store in cache map
+ ndc.TaskManager.CacheMap.Store(imagePrePull.Name, imagePrePull)
+
+ // If all or partial edge nodes image pull is pulling or completed, we don't need to send pull message
+ if fsm.TaskFinish(imagePrePull.Status.State) {
+ klog.Warning("The ImagePrePullJob is completed, don't send pull message again")
+ return
+ }
+
+ ndc.processPrePull(imagePrePull)
+}
+
+// processPrePull do the pre pull operation on node
+func (ndc *ImagePrePullController) processPrePull(imagePrePull *v1alpha1.ImagePrePullJob) {
+ imagePrePullTemplateInfo := imagePrePull.Spec.ImagePrePullTemplate
+ imagePrePullRequest := commontypes.ImagePrePullJobRequest{
+ Images: imagePrePullTemplateInfo.Images,
+ Secret: imagePrePullTemplateInfo.ImageSecret,
+ RetryTimes: imagePrePullTemplateInfo.RetryTimes,
+ CheckItems: imagePrePullTemplateInfo.CheckItems,
+ }
+ tolerate, err := strconv.ParseFloat(imagePrePull.Spec.ImagePrePullTemplate.FailureTolerate, 64)
+ if err != nil {
+ klog.Errorf("convert FailureTolerate to float64 failed: %v", err)
+ tolerate = 0.1
+ }
+
+ concurrency := imagePrePull.Spec.ImagePrePullTemplate.Concurrency
+ if concurrency <= 0 {
+ concurrency = 1
+ }
+ klog.V(4).Infof("deal task message: %v", imagePrePull)
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskPrePull,
+ CheckItem: imagePrePull.Spec.ImagePrePullTemplate.CheckItems,
+ Name: imagePrePull.Name,
+ TimeOutSeconds: imagePrePull.Spec.ImagePrePullTemplate.TimeoutSeconds,
+ Concurrency: concurrency,
+ FailureTolerate: tolerate,
+ NodeNames: imagePrePull.Spec.ImagePrePullTemplate.NodeNames,
+ LabelSelector: imagePrePull.Spec.ImagePrePullTemplate.LabelSelector,
+ Status: v1alpha1.TaskStatus{},
+ Msg: imagePrePullRequest,
+ }
+}
+
+// imagePrePullJobDeleted is used to process deleted ImagePrePullJob in apiserver
+func (ndc *ImagePrePullController) imagePrePullJobDeleted(imagePrePull *v1alpha1.ImagePrePullJob) {
+ // just need to delete from cache map
+ ndc.TaskManager.CacheMap.Delete(imagePrePull.Name)
+ klog.Errorf("image pre pull job %s delete", imagePrePull.Name)
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskPrePull,
+ Name: imagePrePull.Name,
+ ShutDown: true,
+ }
+}
+
+// imagePrePullJobUpdated is used to process update of new ImagePrePullJob in apiserver
+func (ndc *ImagePrePullController) imagePrePullJobUpdated(pullJob *v1alpha1.ImagePrePullJob) {
+ oldValue, ok := ndc.TaskManager.CacheMap.Load(pullJob.Name)
+ old := oldValue.(*v1alpha1.ImagePrePullJob)
+ if !ok {
+ klog.Infof("Update %s not exist, and store it first", pullJob.Name)
+ // If PrePull not present in PrePull map means it is not modified and added.
+ ndc.imagePrePullJobAdded(pullJob)
+ return
+ }
+
+ // store in cache map
+ ndc.TaskManager.CacheMap.Store(pullJob.Name, pullJob)
+
+ node := checkUpdateNode(old, pullJob)
+ if node == nil {
+ klog.Info("none node update")
+ return
+ }
+
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskPrePull,
+ Name: pullJob.Name,
+ Status: *node,
+ }
+}
+
+func checkUpdateNode(old, new *v1alpha1.ImagePrePullJob) *v1alpha1.TaskStatus {
+ if len(old.Status.Status) == 0 {
+ return nil
+ }
+ for i, updateNode := range new.Status.Status {
+ oldNode := old.Status.Status[i]
+ if !util.NodeUpdated(*oldNode.TaskStatus, *updateNode.TaskStatus) {
+ continue
+ }
+ return updateNode.TaskStatus
+ }
+ return nil
+}
diff --git a/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_task.go b/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_task.go
new file mode 100644
index 000000000..3778a0ed8
--- /dev/null
+++ b/cloud/pkg/taskmanager/imageprepullcontroller/image_prepull_task.go
@@ -0,0 +1,133 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
+package imageprepullcontroller
+
+import (
+ "encoding/json"
+ "fmt"
+ "time"
+
+ "k8s.io/klog/v2"
+
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ fsmapi "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ v1alpha12 "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/util/fsm"
+)
+
+func currentPrePullNodeState(id, nodeName string) (v1alpha12.State, error) {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return "", fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.ImagePrePullJob)
+ var state v1alpha12.State
+ for _, status := range task.Status.Status {
+ if status.NodeName == nodeName {
+ state = status.State
+ break
+ }
+ }
+ if state == "" {
+ state = v1alpha12.TaskInit
+ }
+ return state, nil
+}
+
+func updatePrePullNodeState(id, nodeName string, state v1alpha12.State, event fsm.Event) error {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.ImagePrePullJob)
+ newTask := task.DeepCopy()
+ status := newTask.Status.DeepCopy()
+ for i, nodeStatus := range status.Status {
+ if nodeStatus.NodeName == nodeName {
+ var imagesStatus []v1alpha1.ImageStatus
+ err := json.Unmarshal([]byte(event.ExternalMessage), &imagesStatus)
+ if err != nil {
+ klog.Warningf("Failed to unmarshal images status: %v", err)
+ }
+ status.Status[i] = v1alpha1.ImagePrePullStatus{
+ TaskStatus: &v1alpha1.TaskStatus{
+ NodeName: nodeName,
+ State: state,
+ Event: event.Type,
+ Action: event.Action,
+ Time: time.Now().Format(util.ISO8601UTC),
+ Reason: event.Msg,
+ },
+ ImageStatus: imagesStatus,
+ }
+ break
+ }
+ }
+ err := patchStatus(newTask, *status, client.GetCRDClient())
+ if err != nil {
+ return err
+ }
+ return nil
+}
+
+func NewImagePrePullNodeFSM(taskName, nodeName string) *fsm.FSM {
+ fsm := &fsm.FSM{}
+ return fsm.NodeName(nodeName).ID(taskName).Guard(fsmapi.PrePullRule).StageSequence(fsmapi.PrePullStageSequence).CurrentFunc(currentPrePullNodeState).UpdateFunc(updatePrePullNodeState)
+}
+
+func NewImagePrePullTaskFSM(taskName string) *fsm.FSM {
+ fsm := &fsm.FSM{}
+ return fsm.ID(taskName).Guard(fsmapi.PrePullRule).StageSequence(fsmapi.PrePullStageSequence).CurrentFunc(currentPrePullTaskState).UpdateFunc(updateUpgradeTaskState)
+}
+
+func currentPrePullTaskState(id, _ string) (v1alpha12.State, error) {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return "", fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.ImagePrePullJob)
+ state := task.Status.State
+ if state == "" {
+ state = v1alpha12.TaskInit
+ }
+ return state, nil
+}
+
+func updateUpgradeTaskState(id, _ string, state v1alpha12.State, event fsm.Event) error {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.ImagePrePullJob)
+ newTask := task.DeepCopy()
+ status := newTask.Status.DeepCopy()
+
+ status.Event = event.Type
+ status.Action = event.Action
+ status.Reason = event.Msg
+ status.State = state
+ status.Time = time.Now().Format(util.ISO8601UTC)
+
+ err := patchStatus(newTask, *status, client.GetCRDClient())
+
+ if err != nil {
+ return err
+ }
+ return nil
+}
diff --git a/cloud/pkg/taskmanager/manager/downstream.go b/cloud/pkg/taskmanager/manager/downstream.go
index dca8b5c9c..9e8e47263 100644
--- a/cloud/pkg/taskmanager/manager/downstream.go
+++ b/cloud/pkg/taskmanager/manager/downstream.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
@@ -31,7 +31,7 @@ type DownstreamController struct {
// Start DownstreamController
func (dc *DownstreamController) Start() error {
- klog.Info("Start NodeUpgradeJob Downstream Controller")
+ klog.Info("Start TaskManager Downstream Controller")
go dc.syncTask()
diff --git a/cloud/pkg/taskmanager/manager/executor.go b/cloud/pkg/taskmanager/manager/executor.go
index 83005c075..d3a7aec45 100644
--- a/cloud/pkg/taskmanager/manager/executor.go
+++ b/cloud/pkg/taskmanager/manager/executor.go
@@ -1,3 +1,19 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
package manager
import (
@@ -26,11 +42,16 @@ import (
"github.com/kubeedge/kubeedge/pkg/util/fsm"
)
+const TimeOutSecond = 300
+
type Executor struct {
- task util.TaskMessage
- statusChan chan *v1alpha1.TaskStatus
- nodes []v1alpha1.TaskStatus
- controller controller.Controller
+ task util.TaskMessage
+ statusChan chan *v1alpha1.TaskStatus
+ nodes []v1alpha1.TaskStatus
+ controller controller.Controller
+ maxFailedNodes float64
+ failedNodes map[string]bool
+ workers workers
}
func NewExecutorMachine(messageChan chan util.TaskMessage, downStreamChan chan model.Message) (*ExecutorMachine, error) {
@@ -71,7 +92,7 @@ func (em *ExecutorMachine) syncTask() {
}
err := GetExecutor(msg).HandleMessage(msg.Status)
if err != nil {
- klog.Errorf("Failed to handel upgrade message due to error %s", err.Error())
+ klog.Errorf("Failed to handel %s message due to error %s", msg.Type, err.Error())
break
}
}
@@ -97,7 +118,7 @@ func GetExecutor(msg util.TaskMessage) *Executor {
}
e, err := initExecutor(msg)
if err != nil {
- klog.Error("executor init failed, error: %s", err.Error())
+ klog.Errorf("executor init failed, error: %s", err.Error())
return nil
}
return e
@@ -141,7 +162,7 @@ func (e *Executor) initMessage(node v1alpha1.TaskStatus) *model.Message {
CheckItem: e.task.CheckItem,
}
}
- msg.BuildRouter(modules.TaskManagerModuleName, modules.TaskManagerModuleGroup, resource, util.TaskUpgrade).
+ msg.BuildRouter(modules.TaskManagerModuleName, modules.TaskManagerModuleGroup, resource, e.task.Type).
FillBody(taskReq)
return msg
}
@@ -201,10 +222,18 @@ func initExecutor(message util.TaskMessage) (*Executor, error) {
}
}
e := &Executor{
- task: message,
- statusChan: make(chan *v1alpha1.TaskStatus, 10),
- controller: controller,
- nodes: nodeStatus,
+ task: message,
+ statusChan: make(chan *v1alpha1.TaskStatus, 10),
+ nodes: nodeStatus,
+ controller: controller,
+ maxFailedNodes: float64(len(nodeStatus)) * (message.FailureTolerate),
+ failedNodes: map[string]bool{},
+ workers: workers{
+ number: int(message.Concurrency),
+ jobs: make(map[string]int),
+ shuttingDown: false,
+ Mutex: sync.Mutex{},
+ },
}
go e.start()
executorMachine.executors[fmt.Sprintf("%s::%s", message.Type, message.Name)] = e
@@ -212,78 +241,43 @@ func initExecutor(message util.TaskMessage) (*Executor, error) {
}
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)
+ index, err := e.initWorker(0)
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{}) {
+ if 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)
+ endNode, err = e.workers.endJob(status.NodeName)
if err != nil {
klog.Errorf(err.Error())
break
}
e.nodes[endNode] = *status
- err = dealCompletedNode(*status)
+ err = e.dealFailedNode(*status)
if err != nil {
klog.Warning(err.Error())
break
}
if index >= len(e.nodes) {
- if len(worker.jobs) != 0 {
+ if len(e.workers.jobs) != 0 {
break
}
var state api.State
- state, err = e.completedTaskStage(*status)
+ state, err = e.completedTaskStage()
if err != nil {
klog.Errorf(err.Error())
break
@@ -293,28 +287,54 @@ func (e *Executor) start() {
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
+ index = 0
}
- nextNode := e.nodes[index]
- err = worker.addJob(nextNode, index, e)
+ index, err = e.initWorker(index)
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{
+func (e *Executor) dealFailedNode(node v1alpha1.TaskStatus) error {
+ if node.State == api.TaskFailed {
+ e.failedNodes[node.NodeName] = true
+ }
+ if float64(len(e.failedNodes)) < e.maxFailedNodes {
+ return nil
+ }
+ e.workers.shuttingDown = true
+ if len(e.workers.jobs) > 0 {
+ klog.Warningf("wait for all workers(%d/%d) for task %s to finish running ", len(e.workers.jobs), e.workers.number, e.task.Name)
+ return nil
+ }
+
+ errMsg := fmt.Sprintf("the number of failed nodes is %d/%d, which exceeds the failure tolerance threshold.", len(e.failedNodes), len(e.nodes))
+ _, err := e.controller.ReportTaskStatus(e.task.Name, fsm.Event{
Type: node.Event,
+ Action: api.ActionFailure,
+ Msg: errMsg,
+ })
+ if err != nil {
+ return fmt.Errorf("%s, report status failed, %s", errMsg, err.Error())
+ }
+ return fmt.Errorf(errMsg)
+}
+
+func (e *Executor) completedTaskStage() (api.State, error) {
+ var event = e.nodes[0].Event
+ for _, node := range e.nodes {
+ if node.State != api.TaskFailed {
+ event = node.Event
+ break
+ }
+ }
+ state, err := e.controller.ReportTaskStatus(e.task.Name, fsm.Event{
+ Type: event,
Action: api.ActionSuccess,
})
if err != nil {
@@ -323,28 +343,22 @@ func (e *Executor) completedTaskStage(node v1alpha1.TaskStatus) (api.State, erro
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 {
+func (e *Executor) initWorker(index int) (int, error) {
+ for ; index < len(e.nodes); index++ {
+ node := e.nodes[index]
if e.controller.StageCompleted(e.task.Name, node.State) {
- err := dealCompletedNode(node)
+ err := e.dealFailedNode(node)
if err != nil {
return 0, err
}
continue
}
- err := worker.addJob(node, index, e)
+ err := e.workers.addJob(node, index, e)
if err != nil {
- klog.Info(err.Error())
- isEndNode = false
+ klog.V(4).Info(err.Error())
break
}
}
- if isEndNode {
- index++
- }
return index, nil
}
@@ -374,7 +388,11 @@ func (w *workers) addJob(node v1alpha1.TaskStatus, index int, e *Executor) error
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) {
+ timeoutSecond := *e.task.TimeOutSeconds
+ if timeoutSecond == 0 {
+ timeoutSecond = TimeOutSecond
+ }
+ err := wait.Poll(1*time.Second, time.Duration(timeoutSecond)*time.Second, func() (bool, error) {
if lastState != e.nodes[index].State || fsm.TaskFinish(e.nodes[index].State) {
return true, nil
}
@@ -383,9 +401,9 @@ func (e *Executor) handelTimeOutJob(index int) {
})
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()),
+ Type: api.EventTimeOut,
+ Action: api.ActionFailure,
+ Msg: fmt.Sprintf("node task %s execution timeout, %s", lastState, err.Error()),
})
if err != nil {
klog.Warningf(err.Error())
diff --git a/cloud/pkg/taskmanager/manager/upstream.go b/cloud/pkg/taskmanager/manager/upstream.go
index 12cad9f46..b94370cf3 100644
--- a/cloud/pkg/taskmanager/manager/upstream.go
+++ b/cloud/pkg/taskmanager/manager/upstream.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
@@ -31,7 +31,7 @@ import (
"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"
+ "github.com/kubeedge/kubeedge/common/types"
crdClientset "github.com/kubeedge/kubeedge/pkg/client/clientset/versioned"
"github.com/kubeedge/kubeedge/pkg/util/fsm"
)
@@ -110,16 +110,17 @@ func (uc *UpstreamController) updateTaskStatus() {
continue
}
- resp := v1alpha1.TaskStatus{}
+ resp := types.NodeTaskResponse{}
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,
+ Type: resp.Event,
+ Action: resp.Action,
+ Msg: resp.Reason,
+ ExternalMessage: resp.ExternalMessage,
}
_, err = c.ReportNodeStatus(taskID, nodeID, event)
diff --git a/cloud/pkg/taskmanager/nodeupgradecontroller/node_upgrade_controller.go b/cloud/pkg/taskmanager/nodeupgradecontroller/node_upgrade_controller.go
new file mode 100644
index 000000000..649b95cbb
--- /dev/null
+++ b/cloud/pkg/taskmanager/nodeupgradecontroller/node_upgrade_controller.go
@@ -0,0 +1,386 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
+package nodeupgradecontroller
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "os"
+ "strconv"
+ "strings"
+ "sync"
+ "time"
+
+ jsonpatch "github.com/evanphx/json-patch"
+ "github.com/google/uuid"
+ v1 "k8s.io/api/core/v1"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ apimachineryType "k8s.io/apimachinery/pkg/types"
+ "k8s.io/apimachinery/pkg/util/wait"
+ "k8s.io/apimachinery/pkg/watch"
+ "k8s.io/klog/v2"
+
+ beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ keclient "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/informers"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/controller"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util/manager"
+ commontypes "github.com/kubeedge/kubeedge/common/types"
+ 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"
+)
+
+const NodeUpgrade = "NodeUpgradeController"
+
+type NodeUpgradeController struct {
+ sync.Mutex
+ *controller.BaseController
+}
+
+var cache *manager.TaskCache
+
+func NewNodeUpgradeController(messageChan chan util.TaskMessage) (*NodeUpgradeController, error) {
+ var err error
+ cache, err = manager.NewTaskCache(
+ informers.GetInformersManager().GetKubeEdgeInformerFactory().Operations().V1alpha1().NodeUpgradeJobs().Informer())
+ if err != nil {
+ klog.Warningf("Create NodeUpgradeJob manager failed with error: %s", err)
+ return nil, err
+ }
+ return &NodeUpgradeController{
+ BaseController: &controller.BaseController{
+ Informer: informers.GetInformersManager().GetKubeInformerFactory(),
+ TaskManager: cache,
+ MessageChan: messageChan,
+ CrdClient: client.GetCRDClient(),
+ KubeClient: keclient.GetKubeClient(),
+ },
+ }, nil
+}
+
+func (ndc *NodeUpgradeController) ReportNodeStatus(taskID, nodeID string, event fsm.Event) (api.State, error) {
+ nodeFSM := NewUpgradeNodeFSM(taskID, nodeID)
+ err := nodeFSM.AllowTransit(event)
+ if err != nil {
+ return "", err
+ }
+ state, err := nodeFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ ndc.Lock()
+ defer ndc.Unlock()
+ err = nodeFSM.Transit(event)
+ if err != nil {
+ return "", err
+ }
+ checkStatusChanged(nodeFSM, state)
+ state, err = nodeFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ return state, nil
+}
+
+func checkStatusChanged(nodeFSM *fsm.FSM, state api.State) {
+ err := wait.Poll(100*time.Millisecond, time.Second, func() (bool, error) {
+ nowState, err := nodeFSM.CurrentState()
+ if err != nil {
+ return false, nil
+ }
+ if nowState == state {
+ return false, nil
+ }
+ return true, err
+ })
+ if err != nil {
+ klog.V(4).Infof("check status changed failed: %s", err.Error())
+ }
+}
+
+func (ndc *NodeUpgradeController) ReportTaskStatus(taskID string, event fsm.Event) (api.State, error) {
+ taskFSM := NewUpgradeTaskFSM(taskID)
+ state, err := taskFSM.CurrentState()
+ if err != nil {
+ return "", err
+ }
+ err = taskFSM.AllowTransit(event)
+ if err != nil {
+ return "", err
+ }
+ err = taskFSM.Transit(event)
+ if err != nil {
+ return "", err
+ }
+ checkStatusChanged(taskFSM, state)
+ return taskFSM.CurrentState()
+}
+
+func (ndc *NodeUpgradeController) ValidateNode(taskMessage util.TaskMessage) []v1.Node {
+ var validateNodes []v1.Node
+ nodes := ndc.BaseController.ValidateNode(taskMessage)
+ req, ok := taskMessage.Msg.(commontypes.NodeUpgradeJobRequest)
+ if !ok {
+ klog.Errorf("convert message to commontypes.NodeUpgradeJobRequest failed")
+ return nil
+ }
+ for _, node := range nodes {
+ if needUpgrade(node, req.Version) {
+ validateNodes = append(validateNodes, node)
+ }
+ }
+ return validateNodes
+}
+
+func (ndc *NodeUpgradeController) StageCompleted(taskID string, state api.State) bool {
+ taskFSM := NewUpgradeTaskFSM(taskID)
+ return taskFSM.TaskStagCompleted(state)
+}
+
+func (ndc *NodeUpgradeController) GetNodeStatus(name string) ([]v1alpha1.TaskStatus, error) {
+ nodeUpgrade, err := ndc.CrdClient.OperationsV1alpha1().NodeUpgradeJobs().Get(context.TODO(), name, metav1.GetOptions{})
+ if err != nil {
+ return nil, err
+ }
+ return nodeUpgrade.Status.Status, nil
+}
+
+func (ndc *NodeUpgradeController) GetNodeVersion(name string) (string, error) {
+ node, err := ndc.Informer.Core().V1().Nodes().Lister().Get(name)
+ if err != nil {
+ return "", err
+ }
+ strs := strings.Split(node.Status.NodeInfo.KubeletVersion, "-")
+ return strs[2], nil
+}
+
+func (ndc *NodeUpgradeController) UpdateNodeStatus(name string, nodeStatus []v1alpha1.TaskStatus) error {
+ nodeUpgrade, err := ndc.CrdClient.OperationsV1alpha1().NodeUpgradeJobs().Get(context.TODO(), name, metav1.GetOptions{})
+ if err != nil {
+ return err
+ }
+ status := nodeUpgrade.Status
+ status.Status = nodeStatus
+ err = patchStatus(nodeUpgrade, status, ndc.CrdClient)
+ if err != nil {
+ return err
+ }
+ return nil
+}
+
+func patchStatus(nodeUpgrade *v1alpha1.NodeUpgradeJob, status v1alpha1.NodeUpgradeJobStatus, crdClient crdClientset.Interface) error {
+ oldData, err := json.Marshal(nodeUpgrade)
+ if err != nil {
+ return fmt.Errorf("failed to marshal the old NodeUpgradeJob(%s): %v", nodeUpgrade.Name, err)
+ }
+ nodeUpgrade.Status = status
+ newData, err := json.Marshal(nodeUpgrade)
+ if err != nil {
+ return fmt.Errorf("failed to marshal the new NodeUpgradeJob(%s): %v", nodeUpgrade.Name, err)
+ }
+
+ patchBytes, err := jsonpatch.CreateMergePatch(oldData, newData)
+ if err != nil {
+ return fmt.Errorf("failed to create a merge patch: %v", err)
+ }
+
+ result, err := crdClient.OperationsV1alpha1().NodeUpgradeJobs().Patch(context.TODO(), nodeUpgrade.Name, apimachineryType.MergePatchType, patchBytes, metav1.PatchOptions{}, "status")
+ if err != nil {
+ return fmt.Errorf("failed to patch update NodeUpgradeJob status: %v", err)
+ }
+ klog.V(4).Info("patch upgrade task status result: ", result)
+ return nil
+}
+
+func (ndc *NodeUpgradeController) Start() error {
+ go ndc.startSync()
+ return nil
+}
+
+func (ndc *NodeUpgradeController) startSync() {
+ nodeUpgradeList, err := ndc.CrdClient.OperationsV1alpha1().NodeUpgradeJobs().List(context.TODO(), metav1.ListOptions{})
+ if err != nil {
+ klog.Errorf(err.Error())
+ os.Exit(2)
+ }
+ for _, nodeUpgrade := range nodeUpgradeList.Items {
+ if fsm.TaskFinish(nodeUpgrade.Status.State) {
+ continue
+ }
+ ndc.nodeUpgradeJobAdded(&nodeUpgrade)
+ }
+ for {
+ select {
+ case <-beehiveContext.Done():
+ klog.Info("stop sync NodeUpgradeJob")
+ return
+ case e := <-ndc.TaskManager.Events():
+ upgrade, ok := e.Object.(*v1alpha1.NodeUpgradeJob)
+ if !ok {
+ klog.Warningf("object type: %T unsupported", e.Object)
+ continue
+ }
+ switch e.Type {
+ case watch.Added:
+ ndc.nodeUpgradeJobAdded(upgrade)
+ case watch.Deleted:
+ ndc.nodeUpgradeJobDeleted(upgrade)
+ case watch.Modified:
+ ndc.nodeUpgradeJobUpdated(upgrade)
+ default:
+ klog.Warningf("NodeUpgradeJob event type: %s unsupported", e.Type)
+ }
+ }
+ }
+}
+
+// nodeUpgradeJobAdded is used to process addition of new NodeUpgradeJob in apiserver
+func (ndc *NodeUpgradeController) nodeUpgradeJobAdded(upgrade *v1alpha1.NodeUpgradeJob) {
+ klog.V(4).Infof("add NodeUpgradeJob: %v", upgrade)
+ // store in cache map
+ ndc.TaskManager.CacheMap.Store(upgrade.Name, upgrade)
+
+ // If all or partial edge nodes upgrade is upgrading or completed, we don't need to send upgrade message
+ if fsm.TaskFinish(upgrade.Status.State) {
+ klog.Warning("The nodeUpgradeJob is completed, don't send upgrade message again")
+ return
+ }
+
+ ndc.processUpgrade(upgrade)
+}
+
+// processUpgrade do the upgrade operation on node
+func (ndc *NodeUpgradeController) processUpgrade(upgrade *v1alpha1.NodeUpgradeJob) {
+ // if users specify Image, we'll use upgrade Version as its image tag, even though Image contains tag.
+ // if not, we'll use default image: kubeedge/installation-package:${Version}
+ var repo string
+ var err error
+ repo = "kubeedge/installation-package"
+ if upgrade.Spec.Image != "" {
+ repo, err = util.GetImageRepo(upgrade.Spec.Image)
+ if err != nil {
+ klog.Errorf("Image format is not right: %v", err)
+ return
+ }
+ }
+ imageTag := upgrade.Spec.Version
+ image := fmt.Sprintf("%s:%s", repo, imageTag)
+
+ upgradeReq := commontypes.NodeUpgradeJobRequest{
+ UpgradeID: upgrade.Name,
+ HistoryID: uuid.New().String(),
+ Version: upgrade.Spec.Version,
+ Image: image,
+ }
+
+ tolerate, err := strconv.ParseFloat(upgrade.Spec.FailureTolerate, 64)
+ if err != nil {
+ klog.Errorf("convert FailureTolerate to float64 failed: %v", err)
+ tolerate = 0.1
+ }
+
+ concurrency := upgrade.Spec.Concurrency
+ if concurrency <= 0 {
+ concurrency = 1
+ }
+ klog.V(4).Infof("deal task message: %v", upgrade)
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskUpgrade,
+ CheckItem: upgrade.Spec.CheckItems,
+ Name: upgrade.Name,
+ TimeOutSeconds: upgrade.Spec.TimeoutSeconds,
+ Concurrency: concurrency,
+ FailureTolerate: tolerate,
+ NodeNames: upgrade.Spec.NodeNames,
+ LabelSelector: upgrade.Spec.LabelSelector,
+ Status: v1alpha1.TaskStatus{},
+ Msg: upgradeReq,
+ }
+}
+
+func needUpgrade(node v1.Node, upgradeVersion string) bool {
+ if util.FilterVersion(node.Status.NodeInfo.KubeletVersion, upgradeVersion) {
+ klog.Warningf("Node(%s) version(%s) already on the expected version %s.", node.Name, node.Status.NodeInfo.KubeletVersion, upgradeVersion)
+ return false
+ }
+
+ // if node is in Upgrading state, don't need upgrade
+ if _, ok := node.Labels[util.NodeUpgradeJobStatusKey]; ok {
+ klog.Warningf("Node(%s) is in upgrade state", node.Name)
+ return false
+ }
+
+ return true
+}
+
+// nodeUpgradeJobDeleted is used to process deleted NodeUpgradeJob in apiserver
+func (ndc *NodeUpgradeController) nodeUpgradeJobDeleted(upgrade *v1alpha1.NodeUpgradeJob) {
+ // just need to delete from cache map
+ ndc.TaskManager.CacheMap.Delete(upgrade.Name)
+ klog.Errorf("upgrade job %s delete", upgrade.Name)
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskUpgrade,
+ Name: upgrade.Name,
+ ShutDown: true,
+ }
+}
+
+// upgradeAdded is used to process update of new NodeUpgradeJob in apiserver
+func (ndc *NodeUpgradeController) nodeUpgradeJobUpdated(upgrade *v1alpha1.NodeUpgradeJob) {
+ oldValue, ok := ndc.TaskManager.CacheMap.Load(upgrade.Name)
+ old := oldValue.(*v1alpha1.NodeUpgradeJob)
+ if !ok {
+ klog.Infof("Upgrade %s not exist, and store it first", upgrade.Name)
+ // If Upgrade not present in Upgrade map means it is not modified and added.
+ ndc.nodeUpgradeJobAdded(upgrade)
+ return
+ }
+
+ // store in cache map
+ ndc.TaskManager.CacheMap.Store(upgrade.Name, upgrade)
+
+ node := checkUpdateNode(old, upgrade)
+ if node == nil {
+ klog.Info("none node update")
+ return
+ }
+
+ ndc.MessageChan <- util.TaskMessage{
+ Type: util.TaskUpgrade,
+ Name: upgrade.Name,
+ Status: *node,
+ }
+}
+
+func checkUpdateNode(old, new *v1alpha1.NodeUpgradeJob) *v1alpha1.TaskStatus {
+ if len(old.Status.Status) == 0 {
+ return nil
+ }
+ for i, updateNode := range new.Status.Status {
+ oldNode := old.Status.Status[i]
+ if !util.NodeUpdated(oldNode, updateNode) {
+ continue
+ }
+ return &updateNode
+ }
+ return nil
+}
diff --git a/cloud/pkg/taskmanager/nodeupgradecontroller/upgrade_task.go b/cloud/pkg/taskmanager/nodeupgradecontroller/upgrade_task.go
new file mode 100644
index 000000000..6f28eb9ae
--- /dev/null
+++ b/cloud/pkg/taskmanager/nodeupgradecontroller/upgrade_task.go
@@ -0,0 +1,122 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
+package nodeupgradecontroller
+
+import (
+ "fmt"
+ "time"
+
+ "github.com/kubeedge/kubeedge/cloud/pkg/common/client"
+ "github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
+ fsmapi "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ v1alpha12 "github.com/kubeedge/kubeedge/pkg/apis/fsm/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
+ "github.com/kubeedge/kubeedge/pkg/util/fsm"
+)
+
+func currentUpgradeNodeState(id, nodeName string) (v1alpha12.State, error) {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return "", fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.NodeUpgradeJob)
+ var state v1alpha12.State
+ for _, status := range task.Status.Status {
+ if status.NodeName == nodeName {
+ state = status.State
+ break
+ }
+ }
+ if state == "" {
+ state = v1alpha12.TaskInit
+ }
+ return state, nil
+}
+
+func updateUpgradeNodeState(id, nodeName string, state v1alpha12.State, event fsm.Event) error {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.NodeUpgradeJob)
+ newTask := task.DeepCopy()
+ status := newTask.Status.DeepCopy()
+ for i, nodeStatus := range status.Status {
+ if nodeStatus.NodeName == nodeName {
+ status.Status[i] = v1alpha1.TaskStatus{
+ NodeName: nodeName,
+ State: state,
+ Event: event.Type,
+ Action: event.Action,
+ Time: time.Now().Format(util.ISO8601UTC),
+ Reason: event.Msg,
+ }
+ break
+ }
+ }
+ err := patchStatus(newTask, *status, client.GetCRDClient())
+ if err != nil {
+ return err
+ }
+ return nil
+}
+
+func NewUpgradeNodeFSM(taskName, nodeName string) *fsm.FSM {
+ fsm := &fsm.FSM{}
+ return fsm.NodeName(nodeName).ID(taskName).Guard(fsmapi.UpgradeRule).StageSequence(fsmapi.UpdateStageSequence).CurrentFunc(currentUpgradeNodeState).UpdateFunc(updateUpgradeNodeState)
+}
+
+func NewUpgradeTaskFSM(taskName string) *fsm.FSM {
+ fsm := &fsm.FSM{}
+ return fsm.ID(taskName).Guard(fsmapi.UpgradeRule).StageSequence(fsmapi.UpdateStageSequence).CurrentFunc(currentUpgradeTaskState).UpdateFunc(updateUpgradeTaskState)
+}
+
+func currentUpgradeTaskState(id, _ string) (v1alpha12.State, error) {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return "", fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.NodeUpgradeJob)
+ state := task.Status.State
+ if state == "" {
+ state = v1alpha12.TaskInit
+ }
+ return state, nil
+}
+
+func updateUpgradeTaskState(id, _ string, state v1alpha12.State, event fsm.Event) error {
+ v, ok := cache.CacheMap.Load(id)
+ if !ok {
+ return fmt.Errorf("can not find task %s", id)
+ }
+ task := v.(*v1alpha1.NodeUpgradeJob)
+ newTask := task.DeepCopy()
+ status := newTask.Status.DeepCopy()
+
+ status.Event = event.Type
+ status.Action = event.Action
+ status.Reason = event.Msg
+ status.State = state
+ status.Time = time.Now().Format(util.ISO8601UTC)
+
+ err := patchStatus(newTask, *status, client.GetCRDClient())
+
+ if err != nil {
+ return err
+ }
+ return nil
+}
diff --git a/cloud/pkg/taskmanager/task_manager.go b/cloud/pkg/taskmanager/task_manager.go
index 167d65d8f..425ecff64 100644
--- a/cloud/pkg/taskmanager/task_manager.go
+++ b/cloud/pkg/taskmanager/task_manager.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
@@ -25,6 +25,7 @@ import (
"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/imageprepullcontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/manager"
"github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/nodeupgradecontroller"
"github.com/kubeedge/kubeedge/cloud/pkg/taskmanager/util"
@@ -32,34 +33,6 @@ import (
"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
@@ -92,7 +65,13 @@ func newTaskManager(enable bool) *TaskManager {
if err != nil {
klog.Exitf("New upgrade node controller failed with error: %s", err)
}
+
+ imagePrePullController, err := imageprepullcontroller.NewImagePrePullController(taskMessage)
+ if err != nil {
+ klog.Exitf("New upgrade node controller failed with error: %s", err)
+ }
controller.Register(util.TaskUpgrade, upgradeNodeController)
+ controller.Register(util.TaskPrePull, imagePrePullController)
return &TaskManager{
downstream: downstream,
diff --git a/cloud/pkg/taskmanager/util/controller/controller.go b/cloud/pkg/taskmanager/util/controller/controller.go
index 55d43f92c..78259d3ab 100644
--- a/cloud/pkg/taskmanager/util/controller/controller.go
+++ b/cloud/pkg/taskmanager/util/controller/controller.go
@@ -1,3 +1,19 @@
+/*
+Copyright 2023 The KubeEdge Authors.
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+*/
+
package controller
import (
diff --git a/cloud/pkg/taskmanager/util/manager/common.go b/cloud/pkg/taskmanager/util/manager/common.go
index 6e8008abf..59200daa3 100644
--- a/cloud/pkg/taskmanager/util/manager/common.go
+++ b/cloud/pkg/taskmanager/util/manager/common.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
diff --git a/cloud/pkg/taskmanager/util/manager/task_cache.go b/cloud/pkg/taskmanager/util/manager/task_cache.go
index 98c793fb1..9614c7732 100644
--- a/cloud/pkg/taskmanager/util/manager/task_cache.go
+++ b/cloud/pkg/taskmanager/util/manager/task_cache.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
diff --git a/cloud/pkg/taskmanager/util/util.go b/cloud/pkg/taskmanager/util/util.go
index c68e8574a..bf3aea571 100644
--- a/cloud/pkg/taskmanager/util/util.go
+++ b/cloud/pkg/taskmanager/util/util.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
@@ -18,13 +18,13 @@ 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"
+ "k8s.io/klog/v2"
"github.com/kubeedge/kubeedge/common/constants"
"github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1"
@@ -37,10 +37,10 @@ const (
)
const (
- TaskUpgrade = "upgrade"
- TaskRollback = "rollback"
- TaskBackup = "backup"
- TaskDockerPrePull = "pre-pull"
+ TaskUpgrade = "upgrade"
+ TaskRollback = "rollback"
+ TaskBackup = "backup"
+ TaskPrePull = "prepull"
ISO8601UTC = "2006-01-02T15:04:05Z"
)
@@ -167,3 +167,15 @@ func VersionLess(version1, version2 string) (bool, error) {
}
return less, nil
}
+
+func NodeUpdated(old, new v1alpha1.TaskStatus) bool {
+ if old.NodeName != new.NodeName {
+ klog.V(4).Infof("old node %s and new node %s is not same", old.NodeName, new.NodeName)
+ return false
+ }
+ if old.State == new.State || new.State == "" {
+ klog.V(4).Infof("node %s state is not change", old.NodeName)
+ return false
+ }
+ return true
+}
diff --git a/cloud/pkg/taskmanager/util/util_test.go b/cloud/pkg/taskmanager/util/util_test.go
index ced6de9ef..9a13e2295 100644
--- a/cloud/pkg/taskmanager/util/util_test.go
+++ b/cloud/pkg/taskmanager/util/util_test.go
@@ -1,5 +1,5 @@
/*
-Copyright 2022 The KubeEdge Authors.
+Copyright 2023 The KubeEdge Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
@@ -38,13 +38,13 @@ func TestFilterVersion(t *testing.T) {
name: "not match expected version",
version: "v1.22.6-kubeedge-v1.10.0-beta.0.194+77ea462f402efb",
expected: "v1.10.0",
- expectResult: false,
+ expectResult: true,
},
{
name: "no right format version",
version: "v1.22.6",
expected: "v1.10.0",
- expectResult: false,
+ expectResult: true,
},
{
name: "match expected version",