summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorgy95 <1015105054@qq.com>2022-11-15 09:32:50 +0800
committergy95 <1015105054@qq.com>2023-01-16 10:18:52 +0800
commitc6592f35a63a77c1875c6b820db75c7c0a1314dd (patch)
tree0159f0cabf8df5af97fc35bca1637ea79430e282
parentMerge pull request #4468 from wackxu/cleanapplica (diff)
downloadkubeedge-c6592f35a63a77c1875c6b820db75c7c0a1314dd.tar.gz
add concurrency to NodeUpgradeJob
Signed-off-by: gy95 <1015105054@qq.com>
-rw-r--r--build/crds/operations/operations_v1alpha1_nodeupgradejob.yaml6
-rw-r--r--cloud/pkg/admissioncontroller/admission.go44
-rw-r--r--cloud/pkg/admissioncontroller/admit_nodeupgradejob.go63
-rw-r--r--cloud/pkg/nodeupgradejobcontroller/controller/downstream.go164
-rw-r--r--manifests/charts/cloudcore/crds/operations_v1alpha1_nodeupgradejob.yaml6
-rw-r--r--pkg/apis/operations/v1alpha1/type.go4
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.