diff options
| author | zhengxinwei <zhengxinwei@huawei.com> | 2024-01-04 15:28:35 +0800 |
|---|---|---|
| committer | zhengxinwei-f <zhengxinwei@huawei.com> | 2024-01-16 20:18:43 +0800 |
| commit | d091c086a5c8620bf43c90a64835d55253227005 (patch) | |
| tree | faa19b8f9a64684541631e1580e9dc7d78a63f48 /cloud | |
| parent | Implement task manager to complete cloud edge task execution (diff) | |
| download | kubeedge-d091c086a5c8620bf43c90a64835d55253227005.tar.gz | |
support edge nodes upgrade and image pre pull
Signed-off-by: zhengxinwei-f <zhengxinwei@huawei.com>
Diffstat (limited to 'cloud')
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", |
