diff options
| author | Shelley-BaoYue <baoyue2@huawei.com> | 2023-06-19 19:35:00 +0800 |
|---|---|---|
| committer | Shelley-BaoYue <baoyue2@huawei.com> | 2023-06-20 15:53:21 +0800 |
| commit | 087e3c76e090f5ea2a7d047d2992b0a2df5f9e67 (patch) | |
| tree | b7bae4e6c4e17d4864d27ca7a576fa26f5d37bbb /build | |
| parent | Merge pull request #4813 from RyanZhaoXB/fix-rbac-configmap (diff) | |
| download | kubeedge-087e3c76e090f5ea2a7d047d2992b0a2df5f9e67.tar.gz | |
add node runner test
Signed-off-by: Shelley-BaoYue <baoyue2@huawei.com>
Diffstat (limited to 'build')
| -rw-r--r-- | build/conformance/e2e-runner/run.go | 291 | ||||
| -rw-r--r-- | build/conformance/kubernetes/kube_node_conformance_test.go | 21 | ||||
| -rw-r--r-- | build/conformance/node-e2e-runner/run.go | 159 | ||||
| -rw-r--r-- | build/conformance/nodeconformance.Dockerfile | 45 | ||||
| -rw-r--r-- | build/conformance/util/util.go | 266 |
5 files changed, 511 insertions, 271 deletions
diff --git a/build/conformance/e2e-runner/run.go b/build/conformance/e2e-runner/run.go index 1c18eac39..cdf742824 100644 --- a/build/conformance/e2e-runner/run.go +++ b/build/conformance/e2e-runner/run.go @@ -17,8 +17,6 @@ limitations under the License. package main import ( - "context" - "encoding/json" "fmt" "io" "log" @@ -28,19 +26,10 @@ import ( "path/filepath" "strings" "syscall" - "time" "github.com/pkg/errors" - "gopkg.in/yaml.v3" - v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/types" - "k8s.io/apimachinery/pkg/util/strategicpatch" - "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/client-go/kubernetes" - "k8s.io/client-go/tools/clientcmd" - "k8s.io/client-go/util/retry" + + "github.com/kubeedge/kubeedge/build/conformance/util" ) const ( @@ -61,8 +50,6 @@ const ( defaultReportPrefix = "conformance" defaultGinkgoBinary = "/usr/local/bin/ginkgo" defaultTestBinary = "/usr/local/bin/e2e.test" - - edgeNodeLabelKey = "node-role.kubernetes.io/edge" ) func main() { @@ -71,7 +58,7 @@ func main() { go func() { sig := <-c log.Printf("Received signal %v, exiting", sig) - err := afterRunConformance() + err := util.AfterRunConformance() if err != nil { log.Printf("failed to cleanup after conformance, err: %v\n", err) } @@ -83,19 +70,19 @@ func main() { } func RunE2E() error { - err := beforeRunConformance() + err := util.BeforeRunConformance() if err != nil { return fmt.Errorf("failed to prepare for run conformance, err: %v", err) } defer func() { - err := afterRunConformance() + err := util.AfterRunConformance() if err != nil { log.Printf("failed to cleanup after conformance, err: %v\n", err) } }() - resultsDir := getEnvWithDefault(resultsDirEnvKey, defaultResultsDir) + resultsDir := util.GetEnvWithDefault(resultsDirEnvKey, defaultResultsDir) // Print the output to stdout and a logfile which will be returned // as part of the results' tarball. @@ -112,7 +99,7 @@ func RunE2E() error { return err } - log.Printf("Running command:\n%v\n", cmdInfo(cmd)) + log.Printf("Running command:\n%v\n", util.CmdInfo(cmd)) err = cmd.Start() if err != nil { @@ -125,7 +112,7 @@ func RunE2E() error { func makeCmd(w io.Writer) (*exec.Cmd, error) { var ginkgoArgs []string - skipCommands, err := skipCommands() + skipCommands, err := util.SkipCommands() if err != nil { return nil, err } @@ -134,277 +121,39 @@ func makeCmd(w io.Writer) (*exec.Cmd, error) { ginkgoArgs = append(ginkgoArgs, "--skip="+skipped) - skipEnvValue := getEnvWithDefault(skipEnvKey, "") + skipEnvValue := util.GetEnvWithDefault(skipEnvKey, "") if len(skipEnvValue) > 0 { ginkgoArgs = append(ginkgoArgs, "--skip="+skipEnvValue) } - focusEnvValue := getEnvWithDefault(focusEnvKey, defaultFocus) + focusEnvValue := util.GetEnvWithDefault(focusEnvKey, defaultFocus) ginkgoArgs = append(ginkgoArgs, "--focus="+focusEnvValue) ginkgoArgs = append(ginkgoArgs, "--noColor=true") - if len(getEnvWithDefault(dryRunEnvKey, "")) > 0 { + if len(util.GetEnvWithDefault(dryRunEnvKey, "")) > 0 { ginkgoArgs = append(ginkgoArgs, "--dryRun=true") } extraArgs := []string{ - "--report-dir=" + getEnvWithDefault(resultsDirEnvKey, defaultResultsDir), - "--report-prefix=" + getEnvWithDefault(reportPrefixEnvKey, defaultReportPrefix), - "--kubeconfig=" + getEnvWithDefault(kubeConfigEnvKey, ""), - "--image-url=" + getEnvWithDefault(imageURL, "nginx"), - "--image-url=" + getEnvWithDefault(imageURL, "nginx"), - "--test-with-device=" + getEnvWithDefault(testWithDevice, "false"), + "--report-dir=" + util.GetEnvWithDefault(resultsDirEnvKey, defaultResultsDir), + "--report-prefix=" + util.GetEnvWithDefault(reportPrefixEnvKey, defaultReportPrefix), + "--kubeconfig=" + util.GetEnvWithDefault(kubeConfigEnvKey, ""), + "--image-url=" + util.GetEnvWithDefault(imageURL, "nginx"), + "--test-with-device=" + util.GetEnvWithDefault(testWithDevice, "false"), } - if len(getEnvWithDefault(extraArgsEnvKey, "")) > 0 { - extraArgs = append(extraArgs, strings.Split(getEnvWithDefault(extraArgsEnvKey, ""), ",")...) + if len(util.GetEnvWithDefault(extraArgsEnvKey, "")) > 0 { + extraArgs = append(extraArgs, strings.Split(util.GetEnvWithDefault(extraArgsEnvKey, ""), ",")...) } var args []string args = append(args, ginkgoArgs...) - args = append(args, getEnvWithDefault(testBinEnvKey, defaultTestBinary)) + args = append(args, util.GetEnvWithDefault(testBinEnvKey, defaultTestBinary)) args = append(args, "--") args = append(args, extraArgs...) - cmd := exec.Command(getEnvWithDefault(ginkgoEnvKey, defaultGinkgoBinary), args...) + cmd := exec.Command(util.GetEnvWithDefault(ginkgoEnvKey, defaultGinkgoBinary), args...) cmd.Stdout = w cmd.Stderr = w return cmd, nil } - -func getEnvWithDefault(envKey, defaultValue string) string { - value := os.Getenv(envKey) - if len(value) == 0 { - return defaultValue - } - return value -} - -type Tests struct { - TestName string `yaml:"testname"` - CodeName string `yaml:"codename"` - Description string `yaml:"description"` - Release string `yaml:"release"` - File string `yaml:"file"` -} - -func skipCommands() ([]string, error) { - tests, err := skipCases() - if err != nil { - return nil, err - } - - var skipCommands []string - for _, test := range tests { - skipCommands = append(skipCommands, test.CodeName) - } - - return skipCommands, nil -} - -func skipCases() ([]Tests, error) { - data, err := Read("/testdata/edge_skip_case.yaml") - if err != nil { - return nil, fmt.Errorf("read skip test case err: %v", err) - } - - var skipTests []Tests - - if err := yaml.Unmarshal(data, &skipTests); err != nil { - return nil, fmt.Errorf("unmarshal skip test case err: %v", err) - } - - return skipTests, err -} - -func Read(filePath string) ([]byte, error) { - data, err := os.ReadFile(filePath) - if os.IsNotExist(err) { - // Not an error (yet), some other provider may have the file. - return nil, nil - } - return data, err -} - -func cmdInfo(cmd *exec.Cmd) string { - return fmt.Sprintf( - `Command env: %v -Run from directory: %v -Executable path: %v -Args (comma-delimited): %v`, cmd.Env, cmd.Dir, cmd.Path, strings.Join(cmd.Args, ","), - ) -} - -// tempTaints is temporarily added to center node when run kubeEdge conformance -// to make sure that all the pod created by conformance to run on the edge node -var tempTaints = &v1.Taint{ - Key: "node.kubeedge.io/conformance", - Value: "remove-when-completed", - Effect: v1.TaintEffectNoSchedule, -} - -var updateTaintBackoff = wait.Backoff{ - Steps: 5, - Duration: 100 * time.Millisecond, - Jitter: 1.0, -} - -// beforeRunConformance do prepare work before run conformance -func beforeRunConformance() error { - kubeClient, err := getKubeClient() - if err != nil { - return err - } - - nodeList, err := kubeClient.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return err - } - - for _, node := range nodeList.Items { - if isEdgeNode(node) { - continue - } - - err = addConformanceTaintOnNode(kubeClient, &node) - if err != nil { - return err - } - } - - return nil -} - -// afterRunConformance do clean work after conformance done -func afterRunConformance() error { - kubeClient, err := getKubeClient() - if err != nil { - return err - } - - nodeList, err := kubeClient.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return err - } - - for _, node := range nodeList.Items { - if isEdgeNode(node) { - continue - } - - err = deleteConformanceTaintOnNode(kubeClient, &node) - if err != nil { - log.Printf("failed delete taint for node:%v\n", node.Name) - } - } - - return nil -} - -func addConformanceTaintOnNode(c kubernetes.Interface, node *v1.Node) error { - newNode, updated := addTaint(node, tempTaints) - if !updated { - return nil - } - - return retry.RetryOnConflict(updateTaintBackoff, func() error { - return patchNodeTaints(c, node, newNode) - }) -} - -func deleteConformanceTaintOnNode(c kubernetes.Interface, node *v1.Node) error { - newNode, updated := removeTaint(node, tempTaints) - if !updated { - return nil - } - - return retry.RetryOnConflict(updateTaintBackoff, func() error { - return patchNodeTaints(c, node, newNode) - }) -} - -func isEdgeNode(node v1.Node) bool { - if node.Labels == nil { - return false - } - - _, ok := node.Labels[edgeNodeLabelKey] - return ok -} - -func getKubeClient() (kubernetes.Interface, error) { - configPath := getEnvWithDefault(kubeConfigEnvKey, "") - kubeConfig, err := clientcmd.BuildConfigFromFlags("", configPath) - if err != nil { - return nil, err - } - - kubeConfig.ContentType = runtime.ContentTypeProtobuf - kubeClient := kubernetes.NewForConfigOrDie(kubeConfig) - return kubeClient, nil -} - -func addTaint(node *v1.Node, taint *v1.Taint) (*v1.Node, bool) { - newNode := node.DeepCopy() - nodeTaints := newNode.Spec.Taints - - var newTaints []v1.Taint - for i := range nodeTaints { - if taint.MatchTaint(&nodeTaints[i]) { - log.Printf("taint already exist for node:%v\n", node.Name) - return node, false - } - - newTaints = append(newTaints, nodeTaints[i]) - } - - newTaints = append(newTaints, *taint) - newNode.Spec.Taints = newTaints - - return newNode, true -} - -func removeTaint(node *v1.Node, taintToDelete *v1.Taint) (*v1.Node, bool) { - newNode := node.DeepCopy() - nodeTaints := newNode.Spec.Taints - if len(nodeTaints) == 0 { - return newNode, false - } - - var newTaints []v1.Taint - deleted := false - for i := range nodeTaints { - if taintToDelete.MatchTaint(&nodeTaints[i]) { - deleted = true - continue - } - newTaints = append(newTaints, nodeTaints[i]) - } - - newNode.Spec.Taints = newTaints - - return newNode, deleted -} - -func patchNodeTaints(c kubernetes.Interface, oldNode *v1.Node, newNode *v1.Node) error { - oldData, err := json.Marshal(oldNode) - if err != nil { - return fmt.Errorf("failed to marshal old node %#v for node %q: %v", oldNode, oldNode.Name, err) - } - - newTaints := newNode.Spec.Taints - newNodeClone := oldNode.DeepCopy() - newNodeClone.Spec.Taints = newTaints - newData, err := json.Marshal(newNodeClone) - if err != nil { - return fmt.Errorf("failed to marshal new node %#v for node %q: %v", newNodeClone, oldNode.Name, err) - } - - patchBytes, err := strategicpatch.CreateTwoWayMergePatch(oldData, newData, v1.Node{}) - if err != nil { - return fmt.Errorf("failed to create patch for node %q: %v", oldNode.Name, err) - } - - _, err = c.CoreV1().Nodes().Patch(context.TODO(), oldNode.Name, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}) - return err -} diff --git a/build/conformance/kubernetes/kube_node_conformance_test.go b/build/conformance/kubernetes/kube_node_conformance_test.go new file mode 100644 index 000000000..f514cac93 --- /dev/null +++ b/build/conformance/kubernetes/kube_node_conformance_test.go @@ -0,0 +1,21 @@ +/* +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 e2e + +import ( + _ "k8s.io/kubernetes/test/e2e/common/node" +) diff --git a/build/conformance/node-e2e-runner/run.go b/build/conformance/node-e2e-runner/run.go new file mode 100644 index 000000000..7af4cce01 --- /dev/null +++ b/build/conformance/node-e2e-runner/run.go @@ -0,0 +1,159 @@ +/* +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 main + +import ( + "fmt" + "io" + "log" + "os" + "os/exec" + "os/signal" + "path/filepath" + "strings" + "syscall" + + "github.com/pkg/errors" + + "github.com/kubeedge/kubeedge/build/conformance/util" +) + +const ( + dryRunEnvKey = "E2E_DRYRUN" + skipEnvKey = "E2E_SKIP" + focusEnvKey = "E2E_FOCUS" + ginkgoEnvKey = "GINKGO_BIN" + testBinEnvKey = "TEST_BIN" + resultsDirEnvKey = "RESULTS_DIR" + reportPrefixEnvKey = "REPORT_PREFIX" + imageURL = "IMAGE_URL" + testWithDevice = "TEST_WITH_DEVICE" + kubeConfigEnvKey = "KUBECONFIG" + logFileName = "e2e.log" + defaultFocus = "\\[NodeConformance\\]" + extraArgsEnvKey = "E2E_EXTRA_ARGS" + defaultResultsDir = "/tmp/results" + defaultReportPrefix = "conformance" + defaultGinkgoBinary = "/usr/local/bin/ginkgo" + defaultTestBinary = "/usr/local/bin/e2e.test" +) + +func main() { + c := make(chan os.Signal, 1) + signal.Notify(c, syscall.SIGTERM, syscall.SIGINT, syscall.SIGHUP) + go func() { + sig := <-c + log.Printf("Received signal %v, exiting", sig) + err := util.AfterRunConformance() + if err != nil { + log.Printf("failed to cleanup after conformance, err: %v\n", err) + } + }() + + if err := RunE2E(); err != nil { + log.Fatal(err) + } +} + +func RunE2E() error { + err := util.BeforeRunConformance() + if err != nil { + return fmt.Errorf("failed to prepare for run conformance, err: %v", err) + } + + defer func() { + err := util.AfterRunConformance() + if err != nil { + log.Printf("failed to cleanup after conformance, err: %v\n", err) + } + }() + + resultsDir := util.GetEnvWithDefault(resultsDirEnvKey, defaultResultsDir) + + // Print the output to stdout and a logfile which will be returned + // as part of the results' tarball. + logFilePath := filepath.Join(resultsDir, logFileName) + logFile, err := os.Create(logFilePath) + if err != nil { + return fmt.Errorf("failed to create log file %v, err: %v", logFilePath, err) + } + + mw := io.MultiWriter(os.Stdout, logFile) + + cmd, err := makeCmd(mw) + if err != nil { + return err + } + + log.Printf("Running command:\n%v\n", util.CmdInfo(cmd)) + + err = cmd.Start() + if err != nil { + return errors.Wrap(err, "starting command") + } + + return errors.Wrap(cmd.Wait(), "running command") +} + +func makeCmd(w io.Writer) (*exec.Cmd, error) { + var ginkgoArgs []string + + skipCommands, err := util.SkipCommands() + if err != nil { + return nil, err + } + + skipped := strings.Join(skipCommands, "|") + + ginkgoArgs = append(ginkgoArgs, "--skip="+skipped) + + skipEnvValue := util.GetEnvWithDefault(skipEnvKey, "") + if len(skipEnvValue) > 0 { + ginkgoArgs = append(ginkgoArgs, "--skip="+skipEnvValue) + } + + focusEnvValue := util.GetEnvWithDefault(focusEnvKey, defaultFocus) + ginkgoArgs = append(ginkgoArgs, "--focus="+focusEnvValue) + ginkgoArgs = append(ginkgoArgs, "--noColor=true") + + if len(util.GetEnvWithDefault(dryRunEnvKey, "")) > 0 { + ginkgoArgs = append(ginkgoArgs, "--dryRun=true") + } + + extraArgs := []string{ + "--report-dir=" + util.GetEnvWithDefault(resultsDirEnvKey, defaultResultsDir), + "--report-prefix=" + util.GetEnvWithDefault(reportPrefixEnvKey, defaultReportPrefix), + "--kubeconfig=" + util.GetEnvWithDefault(kubeConfigEnvKey, ""), + "--image-url=" + util.GetEnvWithDefault(imageURL, "nginx"), + "--test-with-device=" + util.GetEnvWithDefault(testWithDevice, "false"), + } + + if len(util.GetEnvWithDefault(extraArgsEnvKey, "")) > 0 { + extraArgs = append(extraArgs, strings.Split(util.GetEnvWithDefault(extraArgsEnvKey, ""), ",")...) + } + + var args []string + args = append(args, ginkgoArgs...) + args = append(args, util.GetEnvWithDefault(testBinEnvKey, defaultTestBinary)) + args = append(args, "--") + args = append(args, extraArgs...) + + cmd := exec.Command(util.GetEnvWithDefault(ginkgoEnvKey, defaultGinkgoBinary), args...) + cmd.Stdout = w + cmd.Stderr = w + return cmd, nil +} diff --git a/build/conformance/nodeconformance.Dockerfile b/build/conformance/nodeconformance.Dockerfile new file mode 100644 index 000000000..b7139d3c1 --- /dev/null +++ b/build/conformance/nodeconformance.Dockerfile @@ -0,0 +1,45 @@ +# 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. + +FROM golang:1.17.13-alpine3.16 AS builder + +ARG GO_LDFLAGS + +RUN go install github.com/onsi/ginkgo/ginkgo@v1.16.5 + +COPY . /go/src/github.com/kubeedge/kubeedge + +RUN cp /go/src/github.com/kubeedge/kubeedge/build/conformance/kubernetes/kube_node_conformance_test.go \ + /go/src/github.com/kubeedge/kubeedge/tests/e2e/ + +RUN cd /go/src/github.com/kubeedge/kubeedge && go mod vendor + +RUN CGO_ENABLED=0 GO111MODULE=off ginkgo build -ldflags "-w -s -extldflags -static" -r /go/src/github.com/kubeedge/kubeedge/tests/e2e + +RUN CGO_ENABLED=0 GO111MODULE=off go build -v -o /usr/local/bin/node-e2e-runner -ldflags "$GO_LDFLAGS -w -s" \ + /go/src/github.com/kubeedge/kubeedge/build/conformance/node-e2e-runner + +FROM alpine:3.13 + +COPY --from=builder /go/bin/ginkgo /usr/local/bin/ginkgo + +COPY --from=builder /usr/local/bin/node-e2e-runner /usr/local/bin/node-e2e-runner + +COPY --from=builder /go/src/github.com/kubeedge/kubeedge/tests/e2e/e2e.test /usr/local/bin/e2e.test + +COPY --from=builder /go/src/github.com/kubeedge/kubeedge/build/conformance/kubernetes/edge_skip_case.yaml /testdata/edge_skip_case.yaml + +RUN mkdir -p /tmp/results + +ENTRYPOINT ["node-e2e-runner"]
\ No newline at end of file diff --git a/build/conformance/util/util.go b/build/conformance/util/util.go new file mode 100644 index 000000000..20f8d555b --- /dev/null +++ b/build/conformance/util/util.go @@ -0,0 +1,266 @@ +package util + +import ( + "context" + "encoding/json" + "fmt" + "log" + "os" + "os/exec" + "strings" + "time" + + "gopkg.in/yaml.v3" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/strategicpatch" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/clientcmd" + "k8s.io/client-go/util/retry" +) + +const ( + kubeConfigEnvKey = "KUBECONFIG" + + edgeNodeLabelKey = "node-role.kubernetes.io/edge" +) + +type Tests struct { + TestName string `yaml:"testname"` + CodeName string `yaml:"codename"` + Description string `yaml:"description"` + Release string `yaml:"release"` + File string `yaml:"file"` +} + +func GetEnvWithDefault(envKey, defaultValue string) string { + value := os.Getenv(envKey) + if len(value) == 0 { + return defaultValue + } + return value +} + +func SkipCommands() ([]string, error) { + tests, err := skipCases() + if err != nil { + return nil, err + } + + var skipCommands []string + for _, test := range tests { + skipCommands = append(skipCommands, test.CodeName) + } + + return skipCommands, nil +} + +func skipCases() ([]Tests, error) { + data, err := Read("/testdata/edge_skip_case.yaml") + if err != nil { + return nil, fmt.Errorf("read skip test case err: %v", err) + } + + var skipTests []Tests + + if err := yaml.Unmarshal(data, &skipTests); err != nil { + return nil, fmt.Errorf("unmarshal skip test case err: %v", err) + } + + return skipTests, err +} + +func Read(filePath string) ([]byte, error) { + data, err := os.ReadFile(filePath) + if os.IsNotExist(err) { + // Not an error (yet), some other provider may have the file. + return nil, nil + } + return data, err +} + +func CmdInfo(cmd *exec.Cmd) string { + return fmt.Sprintf( + `Command env: %v +Run from directory: %v +Executable path: %v +Args (comma-delimited): %v`, cmd.Env, cmd.Dir, cmd.Path, strings.Join(cmd.Args, ","), + ) +} + +// tempTaints is temporarily added to center node when run kubeEdge conformance +// to make sure that all the pod created by conformance to run on the edge node +var tempTaints = &v1.Taint{ + Key: "node.kubeedge.io/conformance", + Value: "remove-when-completed", + Effect: v1.TaintEffectNoSchedule, +} + +var updateTaintBackoff = wait.Backoff{ + Steps: 5, + Duration: 100 * time.Millisecond, + Jitter: 1.0, +} + +// beforeRunConformance do prepare work before run conformance +func BeforeRunConformance() error { + kubeClient, err := getKubeClient() + if err != nil { + return err + } + + nodeList, err := kubeClient.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) + if err != nil { + return err + } + + for _, node := range nodeList.Items { + if isEdgeNode(node) { + continue + } + + err = addConformanceTaintOnNode(kubeClient, &node) + if err != nil { + return err + } + } + + return nil +} + +// AfterRunConformance do clean work after conformance done +func AfterRunConformance() error { + kubeClient, err := getKubeClient() + if err != nil { + return err + } + + nodeList, err := kubeClient.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) + if err != nil { + return err + } + + for _, node := range nodeList.Items { + if isEdgeNode(node) { + continue + } + + err = deleteConformanceTaintOnNode(kubeClient, &node) + if err != nil { + log.Printf("failed delete taint for node:%v\n", node.Name) + } + } + + return nil +} + +func addConformanceTaintOnNode(c kubernetes.Interface, node *v1.Node) error { + newNode, updated := addTaint(node, tempTaints) + if !updated { + return nil + } + + return retry.RetryOnConflict(updateTaintBackoff, func() error { + return patchNodeTaints(c, node, newNode) + }) +} + +func deleteConformanceTaintOnNode(c kubernetes.Interface, node *v1.Node) error { + newNode, updated := removeTaint(node, tempTaints) + if !updated { + return nil + } + + return retry.RetryOnConflict(updateTaintBackoff, func() error { + return patchNodeTaints(c, node, newNode) + }) +} + +func isEdgeNode(node v1.Node) bool { + if node.Labels == nil { + return false + } + + _, ok := node.Labels[edgeNodeLabelKey] + return ok +} + +func getKubeClient() (kubernetes.Interface, error) { + configPath := GetEnvWithDefault(kubeConfigEnvKey, "") + kubeConfig, err := clientcmd.BuildConfigFromFlags("", configPath) + if err != nil { + return nil, err + } + + kubeConfig.ContentType = runtime.ContentTypeProtobuf + kubeClient := kubernetes.NewForConfigOrDie(kubeConfig) + return kubeClient, nil +} + +func addTaint(node *v1.Node, taint *v1.Taint) (*v1.Node, bool) { + newNode := node.DeepCopy() + nodeTaints := newNode.Spec.Taints + + var newTaints []v1.Taint + for i := range nodeTaints { + if taint.MatchTaint(&nodeTaints[i]) { + log.Printf("taint already exist for node:%v\n", node.Name) + return node, false + } + + newTaints = append(newTaints, nodeTaints[i]) + } + + newTaints = append(newTaints, *taint) + newNode.Spec.Taints = newTaints + + return newNode, true +} + +func removeTaint(node *v1.Node, taintToDelete *v1.Taint) (*v1.Node, bool) { + newNode := node.DeepCopy() + nodeTaints := newNode.Spec.Taints + if len(nodeTaints) == 0 { + return newNode, false + } + + var newTaints []v1.Taint + deleted := false + for i := range nodeTaints { + if taintToDelete.MatchTaint(&nodeTaints[i]) { + deleted = true + continue + } + newTaints = append(newTaints, nodeTaints[i]) + } + + newNode.Spec.Taints = newTaints + + return newNode, deleted +} + +func patchNodeTaints(c kubernetes.Interface, oldNode *v1.Node, newNode *v1.Node) error { + oldData, err := json.Marshal(oldNode) + if err != nil { + return fmt.Errorf("failed to marshal old node %#v for node %q: %v", oldNode, oldNode.Name, err) + } + + newTaints := newNode.Spec.Taints + newNodeClone := oldNode.DeepCopy() + newNodeClone.Spec.Taints = newTaints + newData, err := json.Marshal(newNodeClone) + if err != nil { + return fmt.Errorf("failed to marshal new node %#v for node %q: %v", newNodeClone, oldNode.Name, err) + } + + patchBytes, err := strategicpatch.CreateTwoWayMergePatch(oldData, newData, v1.Node{}) + if err != nil { + return fmt.Errorf("failed to create patch for node %q: %v", oldNode.Name, err) + } + + _, err = c.CoreV1().Nodes().Patch(context.TODO(), oldNode.Name, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}) + return err +} |
