diff options
| author | gy95 <1015105054@qq.com> | 2022-11-15 09:32:50 +0800 |
|---|---|---|
| committer | gy95 <1015105054@qq.com> | 2023-01-16 10:18:52 +0800 |
| commit | c6592f35a63a77c1875c6b820db75c7c0a1314dd (patch) | |
| tree | 0159f0cabf8df5af97fc35bca1637ea79430e282 | |
| parent | Merge pull request #4468 from wackxu/cleanapplica (diff) | |
| download | kubeedge-c6592f35a63a77c1875c6b820db75c7c0a1314dd.tar.gz | |
add concurrency to NodeUpgradeJob
Signed-off-by: gy95 <1015105054@qq.com>
6 files changed, 240 insertions, 47 deletions
diff --git a/build/crds/operations/operations_v1alpha1_nodeupgradejob.yaml b/build/crds/operations/operations_v1alpha1_nodeupgradejob.yaml index a3ea81624..130df65ac 100644 --- a/build/crds/operations/operations_v1alpha1_nodeupgradejob.yaml +++ b/build/crds/operations/operations_v1alpha1_nodeupgradejob.yaml @@ -36,6 +36,12 @@ spec: spec: description: Specification of the desired behavior of NodeUpgradeJob. properties: + concurrency: + description: Concurrency specifies the max number of edge nodes that + can be upgraded at the same time. The default Concurrency value + is 1. + format: int32 + type: integer image: description: 'Image specifies a container image name, the image contains: keadm and edgecore. keadm is used as upgradetool, to install the diff --git a/cloud/pkg/admissioncontroller/admission.go b/cloud/pkg/admissioncontroller/admission.go index 8fae052e7..db7f12416 100644 --- a/cloud/pkg/admissioncontroller/admission.go +++ b/cloud/pkg/admissioncontroller/admission.go @@ -36,6 +36,9 @@ const ( OfflineMigrationConfigName = "mutate-offlinemigration" OfflineMigrationWebhookName = "mutateofflinemigration.kubeedge.io" + MutatingAdmissionWebhookName = "kubeedge-mutating-webhook" + MutatingNodeUpgradeWebhookName = "mutatingnodeupgradejob.kubeedge.io" + AutonomyLabel = "app-offline.kubeedge.io=autonomy" ) @@ -101,6 +104,7 @@ func Run(opt *options.AdmissionOptions) error { http.HandleFunc("/ruleendpoints", serveRuleEndpoint) http.HandleFunc("/offlinemigration", serveOfflineMigration) http.HandleFunc("/nodeupgradejobs", serveNodeUpgradeJob) + http.HandleFunc("/mutating/nodeupgradejobs", serveMutatingNodeUpgradeJob) tlsConfig, err := configTLS(opt, restConfig) if err != nil { @@ -317,8 +321,46 @@ func (ac *AdmissionController) registerWebhooks(opt *options.AdmissionOptions, c }, } + // NodeUpgradeJob mutating webhook + nodeUpgradeJobWebhook := admissionregistrationv1.MutatingWebhook{ + Name: MutatingNodeUpgradeWebhookName, + Rules: []admissionregistrationv1.RuleWithOperations{{ + Operations: []admissionregistrationv1.OperationType{ + admissionregistrationv1.Create, + admissionregistrationv1.Update, + }, + + Rule: admissionregistrationv1.Rule{ + APIGroups: []string{"operations.kubeedge.io"}, + APIVersions: []string{"v1alpha1"}, + Resources: []string{"nodeupgradejobs"}, + }, + }}, + ClientConfig: admissionregistrationv1.WebhookClientConfig{ + Service: &admissionregistrationv1.ServiceReference{ + Namespace: opt.AdmissionServiceNamespace, + Name: opt.AdmissionServiceName, + Path: strPtr("/mutating/nodeupgradejobs"), + Port: &opt.Port, + }, + CABundle: cabundle, + }, + FailurePolicy: &ignorePolicy, + SideEffects: &noneSideEffect, + AdmissionReviewVersions: []string{"v1"}, + } + // mutatingWebhook contains all the kubeedge related Mutating webhook + mutatingWebhook := admissionregistrationv1.MutatingWebhookConfiguration{ + ObjectMeta: metav1.ObjectMeta{ + Name: MutatingAdmissionWebhookName, + }, + Webhooks: []admissionregistrationv1.MutatingWebhook{ + nodeUpgradeJobWebhook, + }, + } + return registerMutatingWebhook(ac.Client.AdmissionregistrationV1().MutatingWebhookConfigurations(), - []admissionregistrationv1.MutatingWebhookConfiguration{offlineMigrationWebhook}) + []admissionregistrationv1.MutatingWebhookConfiguration{offlineMigrationWebhook, mutatingWebhook}) } func (ac *AdmissionController) getRuleEndpoint(namespace, name string) (*v1.RuleEndpoint, error) { diff --git a/cloud/pkg/admissioncontroller/admit_nodeupgradejob.go b/cloud/pkg/admissioncontroller/admit_nodeupgradejob.go index b2d4bb10f..9881dfd9c 100644 --- a/cloud/pkg/admissioncontroller/admit_nodeupgradejob.go +++ b/cloud/pkg/admissioncontroller/admit_nodeupgradejob.go @@ -17,6 +17,7 @@ limitations under the License. package admissioncontroller import ( + "encoding/json" "errors" "fmt" "net/http" @@ -26,6 +27,7 @@ import ( "github.com/blang/semver" admissionv1 "k8s.io/api/admission/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/klog/v2" "github.com/kubeedge/kubeedge/pkg/apis/operations/v1alpha1" ) @@ -34,6 +36,10 @@ func serveNodeUpgradeJob(w http.ResponseWriter, r *http.Request) { serve(w, r, admitNodeUpgradeJob) } +func serveMutatingNodeUpgradeJob(w http.ResponseWriter, r *http.Request) { + serve(w, r, mutatingNodeUpgradeJob) +} + func admitNodeUpgradeJob(review admissionv1.AdmissionReview) *admissionv1.AdmissionResponse { switch review.Request.Operation { case admissionv1.Create: @@ -111,3 +117,60 @@ func admissionResponse(err error) *admissionv1.AdmissionResponse { Allowed: true, } } + +func mutatingNodeUpgradeJob(review admissionv1.AdmissionReview) *admissionv1.AdmissionResponse { + reviewResponse := admissionv1.AdmissionResponse{ + Allowed: true, + } + + var upgrade v1alpha1.NodeUpgradeJob + if err := json.Unmarshal(review.Request.Object.Raw, &upgrade); err != nil { + klog.Errorf("Could not unmarshal raw object: %v", err) + return toAdmissionResponse(err) + } + + payload := generateNodeUpgradeJobPatch(upgrade.Spec) + if len(payload) == 0 { + return &reviewResponse + } + + patch, err := json.Marshal(payload) + if err != nil { + return toAdmissionResponse(err) + } + + reviewResponse.Patch = patch + pt := admissionv1.PatchTypeJSONPatch + reviewResponse.PatchType = &pt + return &reviewResponse +} + +func generateNodeUpgradeJobPatch(spec v1alpha1.NodeUpgradeJobSpec) []patchValue { + patch := make([]patchValue, 0) + + // mutate .spec.concurrency to default value 1 if not specified + if spec.Concurrency == 0 { + patch = append(patch, patchValue{ + Op: "replace", + Path: "/spec/concurrency", + Value: 1, + }) + } + // mutate .spec.timeoutSeconds to default value 300 if not specified + if spec.TimeoutSeconds == nil { + var defaultTimeoutSeconds uint32 = 300 + patch = append(patch, patchValue{ + Op: "replace", + Path: "/spec/timeoutSeconds", + Value: &defaultTimeoutSeconds, + }) + } + + return patch +} + +type patchValue struct { + Op string `json:"op"` + Path string `json:"path"` + Value interface{} `json:"value,omitempty"` +} diff --git a/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go b/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go index e95e33d13..1e49f814b 100644 --- a/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go +++ b/cloud/pkg/nodeupgradejobcontroller/controller/downstream.go @@ -147,6 +147,22 @@ func (dc *DownstreamController) nodeUpgradeJobAdded(upgrade *v1alpha1.NodeUpgrad 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 @@ -162,65 +178,121 @@ func (dc *DownstreamController) nodeUpgradeJobAdded(upgrade *v1alpha1.NodeUpgrad imageTag := upgrade.Spec.Version image := fmt.Sprintf("%s:%s", repo, imageTag) - for _, node := range nodesToUpgrade { - // send upgrade msg to every edge node - msg := model.NewMessage("") + // send upgrade msg to edge node + msg := model.NewMessage("") - resource := buildUpgradeResource(upgrade.Name, node) + 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, - } + 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) - 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: cloud did 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().String(), + }, + } + 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.messageLayer.Send(*msg) + _, 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 { - klog.Errorf("Failed to send upgrade message %v due to error %v", msg.GetID(), err) - continue + return false, nil } - // process time out: cloud did 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().String(), - }, + // 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 } - 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) + + // calculate the number of nodes in upgrading operation + upgradingNum := 0 + for _, status := range upgradeJob.Status.Status { + if status.State == v1alpha1.Upgrading { + upgradingNum++ + } } - // 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) - continue + // 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] } - _, 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) - continue + 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) } } diff --git a/manifests/charts/cloudcore/crds/operations_v1alpha1_nodeupgradejob.yaml b/manifests/charts/cloudcore/crds/operations_v1alpha1_nodeupgradejob.yaml index a3ea81624..130df65ac 100644 --- a/manifests/charts/cloudcore/crds/operations_v1alpha1_nodeupgradejob.yaml +++ b/manifests/charts/cloudcore/crds/operations_v1alpha1_nodeupgradejob.yaml @@ -36,6 +36,12 @@ spec: spec: description: Specification of the desired behavior of NodeUpgradeJob. properties: + concurrency: + description: Concurrency specifies the max number of edge nodes that + can be upgraded at the same time. The default Concurrency value + is 1. + format: int32 + type: integer image: description: 'Image specifies a container image name, the image contains: keadm and edgecore. keadm is used as upgradetool, to install the diff --git a/pkg/apis/operations/v1alpha1/type.go b/pkg/apis/operations/v1alpha1/type.go index fb276a3f6..1b45d5860 100644 --- a/pkg/apis/operations/v1alpha1/type.go +++ b/pkg/apis/operations/v1alpha1/type.go @@ -87,6 +87,10 @@ type NodeUpgradeJobSpec struct { // The default image name is: kubeedge/installation-package. // +optional Image string `json:"image,omitempty"` + // Concurrency specifies the max number of edge nodes that can be upgraded at the same time. + // The default Concurrency value is 1. + // +optional + Concurrency int32 `json:"concurrency,omitempty"` } // UpgradeResult describe the result status of upgrade operation on edge nodes. |
